Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 7 additions & 0 deletions src/config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,9 @@ pub struct Config {

#[serde(default = "default_dead_letter_max_files")]
pub dead_letter_max_files: u64,

#[serde(default = "default_dead_letter_path")]
pub dead_letter_path: String,
}

fn default_batch_size() -> usize {
Expand Down Expand Up @@ -69,6 +72,10 @@ fn default_dead_letter_max_files() -> u64 {
5
}

fn default_dead_letter_path() -> String {
"logtap.failed.json".to_string()
}

impl Config {
pub fn load(path: &str) -> Result<Self> {
let text = std::fs::read_to_string(path)
Expand Down
32 changes: 16 additions & 16 deletions src/sink.rs
Original file line number Diff line number Diff line change
Expand Up @@ -8,8 +8,6 @@ use std::time::Duration;
use tokio::sync::mpsc::Receiver;
use tokio::time::{interval, sleep};

const DEAD_LETTER_PATH: &str = "logtap.failed.jsonl";

pub async fn run_sink(cfg: Config, mut rx: Receiver<LogLine>) {
let client = reqwest::Client::new();
let mut batch: Vec<LogLine> = Vec::with_capacity(cfg.batch_size);
Expand Down Expand Up @@ -67,9 +65,10 @@ async fn flush_with_retry(client: &reqwest::Client, cfg: &Config, batch: &mut Ve

if attempt >= cfg.max_retries {
eprintln!(
"logtap: abandoning batch of {} items after {} attempts — writing to {DEAD_LETTER_PATH}",
"logtap: abandoning batch of {} items after {} attempts — writing to {}",
batch.len(),
cfg.max_retries
cfg.max_retries,
cfg.dead_letter_path
);

write_dead_letter(cfg, batch);
Expand Down Expand Up @@ -98,12 +97,13 @@ fn write_dead_letter(cfg: &Config, batch: &[LogLine]) {
let mut file = match OpenOptions::new()
.create(true)
.append(true)
.open(DEAD_LETTER_PATH)
.open(&cfg.dead_letter_path)
{
Ok(file) => file,
Err(err) => {
eprintln!(
"logtap: could not open {DEAD_LETTER_PATH} ({err}) — {} item(s) lost for good",
"logtap: could not open {} ({err}) — {} item(s) lost for good",
cfg.dead_letter_path,
batch.len()
);
return;
Expand All @@ -112,14 +112,14 @@ fn write_dead_letter(cfg: &Config, batch: &[LogLine]) {

for log in batch {
if let Err(err) = writeln!(file, "{log}") {
eprintln!("logtap: failed writing to {DEAD_LETTER_PATH}: {err}");
eprintln!("logtap: failed writing to {}: {err}", cfg.dead_letter_path);
return;
}
}
}

fn rotated_dead_letter_path(n: u64) -> PathBuf {
PathBuf::from(format!("{DEAD_LETTER_PATH}.{n}"))
fn rotated_dead_letter_path(cfg: &Config, n: u64) -> PathBuf {
PathBuf::from(format!("{}.{n}", cfg.dead_letter_path))
}

// Caps how big logtap.failed.jsonl is allowed to get. Same idea as
Expand All @@ -129,7 +129,7 @@ fn rotated_dead_letter_path(n: u64) -> PathBuf {
// make room. That eviction is real, permanent data loss, so it's always
// logged loudly rather than happening quietly.
fn rotate_dead_letter_if_full(cfg: &Config) {
let current_size = match fs::metadata(DEAD_LETTER_PATH) {
let current_size = match fs::metadata(&cfg.dead_letter_path) {
Ok(meta) => meta.len(),
Err(_) => return, // no file yet — nothing to rotate
};
Expand All @@ -139,7 +139,7 @@ fn rotate_dead_letter_if_full(cfg: &Config) {
}

if cfg.dead_letter_max_files == 0 {
if let Err(err) = fs::remove_file(DEAD_LETTER_PATH) {
if let Err(err) = fs::remove_file(&cfg.dead_letter_path) {
eprintln!("logtap: failed to reset full dead-letter file: {err}");
} else {
eprintln!(
Expand All @@ -149,7 +149,7 @@ fn rotate_dead_letter_if_full(cfg: &Config) {
return;
}

let oldest = rotated_dead_letter_path(cfg.dead_letter_max_files);
let oldest = rotated_dead_letter_path(cfg, cfg.dead_letter_max_files);
if oldest.exists() {
eprintln!(
"logtap: dead-letter rotation limit ({} files) reached — discarding oldest file {}",
Expand All @@ -159,9 +159,9 @@ fn rotate_dead_letter_if_full(cfg: &Config) {
}

for n in (1..cfg.dead_letter_max_files).rev() {
let from = rotated_dead_letter_path(n);
let from = rotated_dead_letter_path(cfg, n);
if from.exists() {
let to = rotated_dead_letter_path(n + 1);
let to = rotated_dead_letter_path(cfg, n + 1);
if let Err(err) = fs::rename(&from, &to) {
eprintln!(
"logtap: failed to rotate {} -> {}: {err}",
Expand All @@ -172,7 +172,7 @@ fn rotate_dead_letter_if_full(cfg: &Config) {
}
}

if let Err(err) = fs::rename(DEAD_LETTER_PATH, rotated_dead_letter_path(1)) {
eprintln!("logtap: failed to rotate {DEAD_LETTER_PATH}: {err}");
if let Err(err) = fs::rename(&cfg.dead_letter_path, rotated_dead_letter_path(cfg, 1)) {
eprintln!("logtap: failed to rotate {}: {err}", cfg.dead_letter_path);
}
}
1 change: 1 addition & 0 deletions tests/integration_test.rs
Original file line number Diff line number Diff line change
Expand Up @@ -82,6 +82,7 @@ async fn integration_source_parser_filter_sink_sends_log_to_http_server() {
mask_common_patterns: false,
dead_letter_max_bytes: 1024 * 1024,
dead_letter_max_files: 5,
dead_letter_path: "logtap.failed.jsonl".to_string()
};

let app = tokio::spawn(async move {
Expand Down
3 changes: 3 additions & 0 deletions tests/sink_test.rs
Original file line number Diff line number Diff line change
Expand Up @@ -86,6 +86,7 @@ async fn sink_posts_logs_when_batch_size_is_reached() {
mask_common_patterns: false,
dead_letter_max_bytes: 0,
dead_letter_max_files: 0,
dead_letter_path: "logtap.failed.jsonl".to_string()
};

let (tx, rx) = tokio::sync::mpsc::channel::<LogLine>(cfg.channel_capacity);
Expand Down Expand Up @@ -144,6 +145,7 @@ async fn sink_writes_batch_to_dead_letter_file_after_exhausting_retries() {
mask_common_patterns: false,
dead_letter_max_bytes: 1024 * 1024,
dead_letter_max_files: 5,
dead_letter_path: "logtap.failed.jsonl".to_string()
};

let (tx, rx) = tokio::sync::mpsc::channel::<LogLine>(cfg.channel_capacity);
Expand Down Expand Up @@ -205,6 +207,7 @@ async fn sink_rotates_dead_letter_file_once_it_exceeds_max_bytes() {
// so the *second* failure is guaranteed to trigger a rotation.
dead_letter_max_bytes: 10,
dead_letter_max_files: 1,
dead_letter_path: "logtap.failed.jsonl".to_string()
};

let (tx, rx) = tokio::sync::mpsc::channel::<LogLine>(cfg.channel_capacity);
Expand Down
2 changes: 2 additions & 0 deletions tests/source_test.rs
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,7 @@ async fn test_run_source_emite_linhas_novas() {
mask_common_patterns: false,
dead_letter_max_bytes: 1024 * 1024,
dead_letter_max_files: 5,
dead_letter_path: "logtap.failed.jsonl".to_string()
};

let _handle = tokio::task::spawn_blocking(move || run_source(cfg, tx));
Expand Down Expand Up @@ -71,6 +72,7 @@ async fn test_run_source_detects_log_rotation() {
mask_common_patterns: false,
dead_letter_max_bytes: 1024 * 1024,
dead_letter_max_files: 5,
dead_letter_path: "logtap.failed.jsonl".to_string()
};

let _handle = tokio::task::spawn_blocking(move || run_source(cfg, tx));
Expand Down
1 change: 1 addition & 0 deletions tests/stress_test.rs
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,7 @@ async fn outage_does_not_lose_logs_when_destination_is_down() {
mask_common_patterns: false,
dead_letter_max_bytes: 1024 * 1024,
dead_letter_max_files: 5,
dead_letter_path: "logtap.failed.jsonl".to_string()
};

let app = tokio::spawn(async move {
Expand Down
Loading