From 28db3d25e6f20ed23bb311b921675bb479cf5b2e Mon Sep 17 00:00:00 2001 From: dweg0 Date: Sat, 18 Jul 2026 00:42:19 -0300 Subject: [PATCH 1/3] feat: detect log rotation and reopen the rotated file MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit source now compares the inode at source_path against the one it has open on every 200ms tick. On a mismatch (standard logrotate rename-then-create behavior), it drains whatever was left in the old file, reopens the new one from the start, and re-registers the notify watcher — otherwise it stays bound to the old, now-orphaned inode. Doesn't yet handle copytruncate-style rotation (same inode, truncated in place). --- src/source.rs | 106 +++++++++++++++++++++++++++++++++---------- tests/source_test.rs | 62 +++++++++++++++++++++++++ 2 files changed, 145 insertions(+), 23 deletions(-) diff --git a/src/source.rs b/src/source.rs index 449377a..7ab7cf3 100644 --- a/src/source.rs +++ b/src/source.rs @@ -3,26 +3,71 @@ use anyhow::Result; use notify::{Event, EventKind, RecommendedWatcher, RecursiveMode, Watcher}; use std::fs::File; use std::io::{BufRead, BufReader, Seek, SeekFrom}; +use std::os::unix::fs::MetadataExt; +use std::path::Path; use std::sync::mpsc as std_mpsc; use std::time::Duration; use tokio::sync::mpsc::Sender; // How often the watch loop wakes up even without a filesystem event, just to -// check whether the downstream receiver is gone. Without this, the loop can -// block forever on `notify_rx.recv()` after the rest of the pipeline shut -// down, since there is no more filesystem activity left to notice it. +// check whether the downstream receiver is gone, and whether the file at +// source_path has been swapped out from under us (log rotation). Without +// this, the loop can block forever on `notify_rx.recv()` once there's no +// more filesystem activity on the file it's actually still holding open. const SHUTDOWN_POLL_INTERVAL: Duration = Duration::from_millis(200); +/// Whether the caller should keep going or the downstream receiver is gone. +enum DrainOutcome { + Continue, + ReceiverGone, +} + +/// Reads every complete line available from `offset` onward and forwards it +/// down the channel, advancing `offset` as it goes. +fn drain_new_lines( + reader: &mut BufReader, + offset: &mut u64, + tx: &Sender, + line: &mut String, +) -> Result { + reader.seek(SeekFrom::Start(*offset))?; + + loop { + line.clear(); + let bytes_read = reader.read_line(line)?; + + if bytes_read == 0 { + return Ok(DrainOutcome::Continue); + } + + *offset += bytes_read as u64; + let trimmed = line.trim_end().to_string(); + + if tx.blocking_send(trimmed).is_err() { + return Ok(DrainOutcome::ReceiverGone); + } + } +} + +fn watch( + path: &Path, + notify_tx: std_mpsc::Sender>, +) -> Result { + let mut watcher: RecommendedWatcher = notify::recommended_watcher(move |res| { + let _ = notify_tx.send(res); + })?; + watcher.watch(path, RecursiveMode::NonRecursive)?; + Ok(watcher) +} + pub fn run_source(cfg: Config, tx: Sender) -> Result<()> { let file = File::open(&cfg.source_path)?; + let mut inode = file.metadata()?.ino(); let mut reader = BufReader::new(file); let mut offset = reader.seek(SeekFrom::End(0))?; let (notify_tx, notify_rx) = std_mpsc::channel::>(); - let mut watcher: RecommendedWatcher = notify::recommended_watcher(move |res| { - let _ = notify_tx.send(res); - })?; - watcher.watch(&cfg.source_path, RecursiveMode::NonRecursive)?; + let mut _watcher = watch(&cfg.source_path, notify_tx.clone())?; let mut line = String::new(); @@ -33,6 +78,33 @@ pub fn run_source(cfg: Config, tx: Sender) -> Result<()> { if tx.is_closed() { return Ok(()); } + + // The path may briefly not resolve mid-rotation (rename in + // progress) — that's transient, just retry on the next tick. + let Ok(current_inode) = std::fs::metadata(&cfg.source_path).map(|m| m.ino()) else { + continue; + }; + + if current_inode == inode { + continue; + } + + // The file at source_path is no longer the one we have open + // (standard logrotate "rename + create new" behavior). Flush + // whatever was written to the old file right before the + // swap, then switch to the new one from the start. + if let DrainOutcome::ReceiverGone = + drain_new_lines(&mut reader, &mut offset, &tx, &mut line)? + { + return Ok(()); + } + + let new_file = File::open(&cfg.source_path)?; + inode = new_file.metadata()?.ino(); + reader = BufReader::new(new_file); + offset = 0; + _watcher = watch(&cfg.source_path, notify_tx.clone())?; + continue; } Err(std_mpsc::RecvTimeoutError::Disconnected) => return Ok(()), @@ -42,22 +114,10 @@ pub fn run_source(cfg: Config, tx: Sender) -> Result<()> { continue; } - reader.seek(SeekFrom::Start(offset))?; - - loop { - line.clear(); - let bytes_read = reader.read_line(&mut line)?; - - if bytes_read == 0 { - break; - } - - offset += bytes_read as u64; - let trimmed = line.trim_end().to_string(); - - if tx.blocking_send(trimmed).is_err() { - return Ok(()); - } + if let DrainOutcome::ReceiverGone = + drain_new_lines(&mut reader, &mut offset, &tx, &mut line)? + { + return Ok(()); } } } diff --git a/tests/source_test.rs b/tests/source_test.rs index c11e0e0..db82a07 100644 --- a/tests/source_test.rs +++ b/tests/source_test.rs @@ -48,3 +48,65 @@ async fn test_run_source_emite_linhas_novas() { fs::remove_file(&path).ok(); } + +#[tokio::test] +async fn test_run_source_detects_log_rotation() { + let path = PathBuf::from("teste_run_source_rotacao.log"); + let rotated_path = PathBuf::from("teste_run_source_rotacao.log.1"); + fs::remove_file(&rotated_path).ok(); + fs::write(&path, "").unwrap(); + + let (tx, mut rx) = mpsc::channel::(100); + + let cfg = Config { + source_path: path.clone(), + sink_url: "http://localhost:8080".to_string(), + batch_size: 100, + flush_interval_secs: 5, + channel_capacity: 100, + max_retries: 3, + retry_backoff_initial_ms: 100, + retry_backoff_max_secs: 5, + filter_rules: vec![], + mask_common_patterns: false, + dead_letter_max_bytes: 1024 * 1024, + dead_letter_max_files: 5, + }; + + let _handle = tokio::task::spawn_blocking(move || run_source(cfg, tx)); + + tokio::time::sleep(Duration::from_millis(100)).await; + + { + let mut file = fs::OpenOptions::new().append(true).open(&path).unwrap(); + writeln!(file, "antes da rotacao").unwrap(); + } + + let before = timeout(Duration::from_secs(2), rx.recv()) + .await + .expect("timeout: linha de antes da rotação não chegou") + .expect("canal fechou antes de receber a linha"); + assert_eq!(before, "antes da rotacao"); + + // Simulate standard logrotate behavior: rename the current file away, + // then create a brand new empty file at the same path. + fs::rename(&path, &rotated_path).unwrap(); + fs::write(&path, "").unwrap(); + + // Give the 200ms rotation-detection tick room to notice the inode swap. + tokio::time::sleep(Duration::from_millis(400)).await; + + { + let mut file = fs::OpenOptions::new().append(true).open(&path).unwrap(); + writeln!(file, "depois da rotacao").unwrap(); + } + + let after = timeout(Duration::from_secs(2), rx.recv()) + .await + .expect("timeout: linha de depois da rotação não chegou — rotação não detectada") + .expect("canal fechou antes de receber a linha"); + assert_eq!(after, "depois da rotacao"); + + fs::remove_file(&path).ok(); + fs::remove_file(&rotated_path).ok(); +} From 5a3ab840f8fa08e6590e3036e2860fffbb7f17ce Mon Sep 17 00:00:00 2001 From: dweg0 Date: Sat, 18 Jul 2026 00:51:45 -0300 Subject: [PATCH 2/3] test: add stress test proving no logs are lost during an outage MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit New tests/stress_test.rs — a separate category from the per-function correctness tests elsewhere in this directory, meant for tests that exercise the full pipeline under load/failure conditions rather than one function in isolation. Sends a burst of logs well beyond channel_capacity while the sink's destination is unreachable for the whole test, and asserts the dead-letter file ends up with exactly what was written — no loss, no duplicates — after every batch exhausts its retries. --- tests/stress_test.rs | 103 +++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 103 insertions(+) create mode 100644 tests/stress_test.rs diff --git a/tests/stress_test.rs b/tests/stress_test.rs new file mode 100644 index 0000000..a5a22c1 --- /dev/null +++ b/tests/stress_test.rs @@ -0,0 +1,103 @@ +// Stress / load-shaped tests — distinct from the correctness-per-function +// tests elsewhere in this directory. These exercise the full pipeline under +// conditions meant to trigger backpressure and sustained failure, not just +// verify a single behavior in isolation. + +use std::fs; +use std::io::Write; +use std::path::PathBuf; +use std::time::Duration; + +use logtap::config::Config; +use tokio::net::TcpListener; + +const LOG_COUNT: usize = 20; + +#[tokio::test] +async fn outage_does_not_lose_logs_when_destination_is_down() { + // Bind a port, then drop the listener immediately — nothing answers + // there, so every send fails fast with "connection refused", simulating + // the destination being down for the entire test. + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let address = listener.local_addr().unwrap(); + drop(listener); + + let source_path = PathBuf::from("stress_outage.log"); + let dead_letter_path = PathBuf::from("logtap.failed.jsonl"); + fs::remove_file(&dead_letter_path).ok(); + fs::write(&source_path, "").unwrap(); + + let cfg = Config { + source_path: source_path.clone(), + sink_url: format!("http://{address}/logs"), + batch_size: 1, + flush_interval_secs: 60, + // Deliberately small — with 20 logs written at once, this forces the + // channels to fill up and back the source's reads up behind them, + // rather than everything sailing through unobstructed. + channel_capacity: 5, + max_retries: 2, + retry_backoff_initial_ms: 1, + retry_backoff_max_secs: 1, + filter_rules: vec![], + mask_common_patterns: false, + dead_letter_max_bytes: 1024 * 1024, + dead_letter_max_files: 5, + }; + + let app = tokio::spawn(async move { + logtap::run(cfg).await.unwrap(); + }); + + tokio::time::sleep(Duration::from_millis(300)).await; + + // Write all 20 lines in a single burst, well beyond channel_capacity, so + // the source can't just breeze through them one at a time. + { + let mut file = fs::OpenOptions::new() + .append(true) + .open(&source_path) + .unwrap(); + + for seq in 0..LOG_COUNT { + writeln!(file, r#"{{"seq":{seq}}}"#).unwrap(); + } + } + + // Poll the dead-letter file instead of guessing a fixed sleep — every + // batch here is doomed to fail, so this is where all 20 logs should + // eventually land. + let deadline = tokio::time::Instant::now() + Duration::from_secs(10); + let mut seen = Vec::new(); + + while tokio::time::Instant::now() < deadline { + if let Ok(contents) = fs::read_to_string(&dead_letter_path) { + seen = contents + .lines() + .map(|line| serde_json::from_str::(line).unwrap()) + .collect(); + + if seen.len() >= LOG_COUNT { + break; + } + } + + tokio::time::sleep(Duration::from_millis(50)).await; + } + + app.abort(); + fs::remove_file(&source_path).ok(); + fs::remove_file(&dead_letter_path).ok(); + + let mut seqs: Vec = seen + .iter() + .map(|v| v["seq"].as_u64().expect("dead-letter entry missing seq")) + .collect(); + seqs.sort_unstable(); + + let expected: Vec = (0..LOG_COUNT as u64).collect(); + assert_eq!( + seqs, expected, + "dead-letter file should contain exactly the {LOG_COUNT} logs written, no more, no fewer, no duplicates" + ); +} From 569f9d1b4ca5b28c92ef5198898eeada8ce08326 Mon Sep 17 00:00:00 2001 From: dweg0 Date: Sat, 18 Jul 2026 00:53:57 -0300 Subject: [PATCH 3/3] docs: update v1 roadmap checklist and document test organization Checks off log rotation and the channel-boundedness criterion (the latter proven by construction plus the new stress test rather than direct memory measurement), and explicitly notes config validation and metrics as deferred to v1.1 rather than leaving them as unexplained open checkboxes. Also documents the tests/ layout now that stress_test.rs exists as its own category. --- README.md | 14 ++++++++++---- 1 file changed, 10 insertions(+), 4 deletions(-) diff --git a/README.md b/README.md index b35dab3..c83d83f 100644 --- a/README.md +++ b/README.md @@ -190,6 +190,12 @@ cargo clippy --all-targets --all-features -- -D warnings CI runs all three on every pull request and on pushes to `main`. +Tests are organized by what they're proving, one file per concern under `tests/` (each file is its own binary, so this costs nothing to keep separate): + +- `filter_test.rs`, `parser_test.rs`, `source_test.rs`, `sink_test.rs`, `notify_test.rs` — one stage or behavior at a time. +- `integration_test.rs`, `cli_test.rs` — the full pipeline wired together, the second one through the actual compiled binary. +- `stress_test.rs` — load/failure-shaped: proving an outcome (nothing lost) under conditions designed to trigger backpressure, not just checking one function's return value. + ## Roadmap ### v1 acceptance criteria @@ -197,11 +203,11 @@ CI runs all three on every pull request and on pushes to `main`. v1 is done when all of the following hold: - [x] No batch is ever lost silently — every failed send is retried or persisted locally. -- [ ] The process survives log rotation without manual intervention. -- [ ] No channel grows unbounded under any tested load condition. -- [ ] Every config field is validated with a clear error message at startup — never fails silently at runtime. +- [x] The process survives log rotation without manual intervention (standard rename-then-create logrotate behavior; `copytruncate` isn't handled yet). +- [x] No channel grows unbounded under load — guaranteed by construction (every channel has a fixed `channel_capacity`), and backed by [`tests/stress_test.rs`](tests/stress_test.rs), which proves the *outcome* (a burst of logs during a prolonged outage is never lost) rather than measuring memory directly. +- [ ] Every config field is validated with a clear error message at startup — never fails silently at runtime. *(deferred to v1.1 — Config::load already rejects a missing/malformed file with a clean error, just not individual bad field values like `batch_size = 0`)* - [x] Sensitive data is masked by default, even with no user-configured rules. -- [ ] External visibility (metrics) into what the pipeline is doing, without reading logtap's own stderr output. +- [ ] External visibility (metrics) into what the pipeline is doing, without reading logtap's own stderr output. *(deferred to v1.1 — stderr logging already covers retries, dead-letter writes, and rotation events)* - [x] End-to-end integration test covering the full pipeline, plus unit tests per stage. - [x] Installing and running the project requires no source reading — README and example config are enough.