use std::path::Path;
use std::sync::{Arc, Barrier};
use crate::error_capture::rotation::{self, RotationPolicy};
use crate::error_capture::store::ErrorStore;
use crate::error_capture::types::CapturedError;
fn record(msg: &str) -> CapturedError {
CapturedError {
timestamp_secs: 1_000_000,
crate_target: "test_crate".to_string(),
crate_version: "0.1.0".to_string(),
message: msg.to_string(),
fields: String::new(),
file: Some("src/lib.rs".to_string()),
line: Some(10),
os: "linux".to_string(),
arch: "x86_64".to_string(),
fingerprint: "fp".to_string(),
}
}
fn raw_lines(path: &Path) -> Vec<Vec<u8>> {
let bytes = std::fs::read(path).unwrap_or_default();
bytes
.split(|b| *b == b'\n')
.filter(|l| !l.is_empty())
.map(<[u8]>::to_vec)
.collect()
}
#[test]
fn concurrent_writers_produce_only_well_formed_lines() {
const WRITERS: usize = 8;
const PER_WRITER: usize = 200;
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("errors.jsonl");
let barrier = Arc::new(Barrier::new(WRITERS));
let handles: Vec<_> = (0..WRITERS)
.map(|w| {
let path = path.clone();
let barrier = Arc::clone(&barrier);
std::thread::spawn(move || {
let store = ErrorStore::with_path(Some(path), 10);
let body = "x".repeat(512);
barrier.wait();
for i in 0..PER_WRITER {
store.append(record(&format!("w{w}-{i}-{body}")));
}
})
})
.collect();
for h in handles {
h.join().expect("writer thread");
}
let lines = raw_lines(&path);
let malformed = lines
.iter()
.filter(|l| serde_json::from_slice::<CapturedError>(l).is_err())
.count();
assert_eq!(malformed, 0, "interleaved writes produced malformed lines");
assert_eq!(
lines.len(),
WRITERS * PER_WRITER,
"every record is one line"
);
}
#[test]
fn malformed_line_does_not_break_reading() {
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("errors.jsonl");
let mut bytes = serde_json::to_vec(&record("before")).expect("json");
bytes.extend_from_slice(b"\n{\"message\":\"torn \xE2\x82\n");
bytes.extend_from_slice(&serde_json::to_vec(&record("after")).expect("json"));
bytes.push(b'\n');
std::fs::write(&path, bytes).expect("seed");
let records = ErrorStore::read_records(&path, 10);
let messages: Vec<&str> = records.iter().map(|r| r.message.as_str()).collect();
assert_eq!(messages, ["before", "after"]);
}
fn line_len(rec: &CapturedError) -> u64 {
serde_json::to_vec(rec).expect("json").len() as u64 + 1
}
fn store_bytes(dir: &Path) -> u64 {
std::fs::read_dir(dir)
.expect("read_dir")
.filter_map(Result::ok)
.filter(|e| {
let name = e.file_name().to_string_lossy().into_owned();
name.starts_with("errors.jsonl") && !name.ends_with(".lock")
})
.map(|e| e.metadata().expect("metadata").len())
.sum()
}
#[test]
fn writing_past_the_cap_rotates_and_bounds_disk_use() {
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("errors.jsonl");
let policy = RotationPolicy {
max_bytes: 4096,
keep: 2,
};
let store = ErrorStore::with_path_and_rotation(Some(path.clone()), 10, policy);
let rec_len = line_len(&record(&format!("m{:04}-{}", 0, "y".repeat(400))));
for i in 0..200 {
store.append(record(&format!("m{i:04}-{}", "y".repeat(400))));
}
let written = 200 * rec_len;
let bound = 3 * (policy.max_bytes + rec_len);
let on_disk = store_bytes(dir.path());
assert!(
on_disk <= bound,
"on-disk {on_disk} exceeds bound {bound} after writing {written}"
);
assert!(rotation::rotated_path(&path, 2).exists(), "rotation ran");
assert!(!rotation::rotated_path(&path, 3).exists(), "keep=2 holds");
for file in rotation::files_newest_first(&path, policy) {
for line in raw_lines(&file) {
assert!(serde_json::from_slice::<CapturedError>(&line).is_ok());
}
}
let reopened = ErrorStore::with_path_and_rotation(Some(path), 12, policy);
let recent = reopened.recent_errors(12);
let live_count = raw_lines(&dir.path().join("errors.jsonl")).len();
assert!(
live_count < 12,
"live holds {live_count}; the read must span files"
);
assert_eq!(recent.len(), 12, "reader spans the rotated files");
assert!(recent[11].message.starts_with("m0199-"), "newest last");
assert!(recent[0].message.starts_with("m0188-"), "oldest first");
}
#[test]
fn oversized_legacy_file_is_compacted_on_rotation() {
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("errors.jsonl");
let policy = RotationPolicy {
max_bytes: 4096,
keep: 2,
};
let mut legacy = Vec::new();
for i in 0..500 {
legacy.extend(serde_json::to_vec(&record(&format!("old{i:03}"))).expect("json"));
legacy.push(b'\n');
}
assert!(legacy.len() as u64 > 10 * policy.max_bytes);
std::fs::write(&path, &legacy).expect("seed legacy file");
let store = ErrorStore::with_path_and_rotation(Some(path.clone()), 5, policy);
store.append(record("fresh"));
let first = rotation::rotated_path(&path, 1);
let first_len = std::fs::metadata(&first).expect(".1 exists").len();
assert!(first_len <= policy.max_bytes, ".1 holds {first_len} bytes");
let tail = raw_lines(&first);
assert!(!tail.is_empty());
for line in &tail {
assert!(serde_json::from_slice::<CapturedError>(line).is_ok());
}
let last: CapturedError = serde_json::from_slice(&tail[tail.len() - 1]).expect("json");
assert_eq!(last.message, "old499", "the tail keeps the newest records");
let live: Vec<CapturedError> = raw_lines(&path)
.iter()
.map(|l| serde_json::from_slice(l).expect("json"))
.collect();
assert_eq!(live.len(), 1);
assert_eq!(live[0].message, "fresh");
assert!(
!dir.path().join("errors.jsonl.1.tmp").exists(),
"temp renamed away"
);
}
#[test]
fn a_rotation_failure_refuses_the_disk_write_and_counts_it() {
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("errors.jsonl");
let blocker = rotation::rotated_path(&path, 1);
std::fs::create_dir(&blocker).expect("blocker dir");
std::fs::write(blocker.join("occupant"), b"x").expect("occupant");
let policy = RotationPolicy {
max_bytes: 2048,
keep: 1,
};
let store = ErrorStore::with_path_and_rotation(Some(path.clone()), 100, policy);
let rec_len = line_len(&record(&format!("r{:02}-{}", 0, "z".repeat(300))));
for i in 0..40 {
store.append(record(&format!("r{i:02}-{}", "z".repeat(300))));
}
let live_len = std::fs::metadata(&path).expect("live").len();
assert!(
live_len < policy.max_bytes + rec_len,
"live file grew to {live_len} past the cap"
);
assert!(store.refused_disk_writes() > 0, "refusals are counted");
assert_eq!(store.len(), 40, "logging continues in memory");
std::fs::remove_dir_all(&blocker).expect("clear blocker");
store.append(record("after-recovery"));
assert!(blocker.is_file(), "rotation resumed");
let live = raw_lines(&path);
assert_eq!(live.len(), 1, "disk writes resumed on a fresh live file");
}
#[derive(Default)]
struct CallRecorder {
calls: Vec<Vec<u8>>,
}
impl std::io::Write for CallRecorder {
fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
self.calls.extend([buf.to_vec()]);
Ok(buf.len())
}
fn flush(&mut self) -> std::io::Result<()> {
Ok(())
}
}
#[test]
fn a_record_reaches_the_file_in_one_write_call() {
let mut out = CallRecorder::default();
crate::error_capture::store::write_record_line(&mut out, &record("one call")).expect("write");
assert_eq!(out.calls.len(), 1, "record split across write calls");
let line = &out.calls[0];
assert_eq!(line.last(), Some(&b'\n'));
let rec: CapturedError = serde_json::from_slice(&line[..line.len() - 1]).expect("json");
assert_eq!(rec.message, "one call");
}