use std::collections::HashMap;
use std::collections::VecDeque;
use std::io::Write;
use std::path::{Path, PathBuf};
use std::sync::{Arc, Mutex, MutexGuard, PoisonError};
use crate::error_capture::rotation::{self, RotationPolicy};
use crate::error_capture::types::CapturedError;
pub const DEFAULT_CAPTURE_CAPACITY: usize = 500;
const ERRORS_FILENAME: &str = "errors.jsonl";
#[derive(Clone)]
pub struct ErrorStore {
inner: Arc<Mutex<Inner>>,
capacity: usize,
}
struct Inner {
ring: VecDeque<CapturedError>,
file_path: Option<PathBuf>,
policy: RotationPolicy,
refused_disk_writes: u64,
corrupt_lines_skipped: u64,
}
impl ErrorStore {
#[must_use]
pub fn open(app_name: &str, capacity: usize) -> Self {
let file_path = match crate::resolve_data_dir(app_name) {
Ok(dir) => Some(dir.join(ERRORS_FILENAME)),
Err(e) => {
eprintln!("[bug-capture] cannot resolve data dir for {app_name}: {e}");
None
}
};
Self::with_path(file_path, capacity)
}
#[must_use]
pub fn with_path(file_path: Option<PathBuf>, capacity: usize) -> Self {
Self::with_path_and_rotation(file_path, capacity, RotationPolicy::default())
}
#[must_use]
pub fn with_path_and_rotation(
file_path: Option<PathBuf>,
capacity: usize,
policy: RotationPolicy,
) -> Self {
let capacity = capacity.max(1);
let (ring, corrupt_lines_skipped) = match file_path {
Some(ref path) => load_ring_from_disk(path, capacity, policy),
None => (VecDeque::with_capacity(capacity), 0),
};
let inner = Inner {
ring,
file_path,
policy,
refused_disk_writes: 0,
corrupt_lines_skipped,
};
Self {
inner: Arc::new(Mutex::new(inner)),
capacity,
}
}
fn lock(&self) -> MutexGuard<'_, Inner> {
self.inner.lock().unwrap_or_else(PoisonError::into_inner)
}
pub fn append(&self, record: CapturedError) {
let (file_path, policy) = {
let guard = self.lock();
(guard.file_path.clone(), guard.policy)
};
super::test_hook::fire(super::test_hook::Point::BeforeDiskWrite);
let mut refused = false;
if let Some(path) = file_path.as_deref() {
if let Err(e) = rotation::rotate_if_due(path, policy) {
refused = true;
eprintln!(
"[bug-capture] rotating {} failed ({e}); record kept in memory only",
path.display()
);
} else if let Err(e) = serialise_and_append(path, &record) {
eprintln!("[bug-capture] write to {}: {e}", path.display());
}
}
let mut guard = self.lock();
if refused {
guard.refused_disk_writes += 1;
}
guard.ring.push_back(record);
while guard.ring.len() > self.capacity {
guard.ring.pop_front();
}
}
#[must_use]
pub fn recent_errors(&self, n: usize) -> Vec<CapturedError> {
let guard = match self.inner.lock() {
Ok(g) => g,
Err(poisoned) => poisoned.into_inner(),
};
let skip = guard.ring.len().saturating_sub(n);
guard.ring.iter().skip(skip).cloned().collect()
}
#[must_use]
pub fn errors_by_fingerprint(&self) -> Vec<(CapturedError, usize)> {
let guard = match self.inner.lock() {
Ok(g) => g,
Err(poisoned) => poisoned.into_inner(),
};
let mut counts: HashMap<String, usize> = HashMap::new();
let mut latest: HashMap<String, CapturedError> = HashMap::new();
for rec in &guard.ring {
*counts.entry(rec.fingerprint.clone()).or_insert(0) += 1;
latest.insert(rec.fingerprint.clone(), rec.clone());
}
let mut result: Vec<(CapturedError, usize)> = latest
.into_iter()
.map(|(fp, rec)| (rec, *counts.get(&fp).unwrap_or(&1)))
.collect();
result.sort_by_key(|item| std::cmp::Reverse(item.1));
result
}
#[must_use]
pub fn len(&self) -> usize {
match self.inner.lock() {
Ok(g) => g.ring.len(),
Err(p) => p.into_inner().ring.len(),
}
}
#[must_use]
pub fn is_empty(&self) -> bool {
self.len() == 0
}
#[must_use]
pub fn refused_disk_writes(&self) -> u64 {
match self.inner.lock() {
Ok(g) => g.refused_disk_writes,
Err(p) => p.into_inner().refused_disk_writes,
}
}
#[must_use]
pub fn corrupt_lines_skipped(&self) -> u64 {
match self.inner.lock() {
Ok(g) => g.corrupt_lines_skipped,
Err(p) => p.into_inner().corrupt_lines_skipped,
}
}
#[must_use]
pub fn read_records(path: &Path, limit: usize) -> Vec<CapturedError> {
load_ring_from_disk(path, limit, RotationPolicy::default())
.0
.into_iter()
.collect()
}
}
fn load_ring_from_disk(
path: &Path,
capacity: usize,
policy: RotationPolicy,
) -> (VecDeque<CapturedError>, u64) {
let mut ring: VecDeque<CapturedError> = VecDeque::with_capacity(capacity);
let mut skipped_total = 0u64;
let mut seen = Vec::new();
for (i, file) in rotation::files_newest_first(path, policy)
.iter()
.enumerate()
{
if ring.len() >= capacity {
break;
}
if i > 0 {
super::test_hook::fire(super::test_hook::Point::BetweenStoreFiles);
}
let bytes = match read_unseen(file, policy.read_limit(), &mut seen) {
Ok(Some(b)) => b,
Ok(None) => continue,
Err(e) if e.kind() == std::io::ErrorKind::NotFound => continue,
Err(e) => {
eprintln!("[bug-capture] cannot read {}: {e}", file.display());
continue;
}
};
let mut skipped = 0u64;
let records: Vec<CapturedError> = bytes
.split(|b| *b == b'\n')
.filter(|line| !line.trim_ascii().is_empty())
.filter_map(|line| {
let parsed = serde_json::from_slice(line.trim_ascii()).ok();
skipped += u64::from(parsed.is_none());
parsed
})
.collect();
for rec in records.into_iter().rev() {
if ring.len() >= capacity {
break;
}
ring.push_front(rec);
}
if skipped > 0 {
eprintln!(
"[bug-capture] skipped {skipped} corrupt record(s) in {}",
file.display()
);
}
skipped_total += skipped;
}
(ring, skipped_total)
}
fn read_unseen(
file: &Path,
limit: u64,
seen: &mut Vec<(u64, u64)>,
) -> std::io::Result<Option<Vec<u8>>> {
let mut handle = std::fs::File::open(file)?;
#[cfg(unix)]
{
use std::os::unix::fs::MetadataExt as _;
let meta = handle.metadata()?;
let id = (meta.dev(), meta.ino());
if seen.contains(&id) {
return Ok(None);
}
seen.push(id);
}
#[cfg(not(unix))]
let _ = &seen;
rotation::read_tail_from(&mut handle, limit).map(|(buf, _)| Some(buf))
}
fn serialise_and_append(path: &Path, record: &CapturedError) -> std::io::Result<()> {
let mut file = std::fs::OpenOptions::new()
.append(true)
.create(true)
.open(path)?;
write_record_line(&mut file, record)
}
pub(super) fn write_record_line(
out: &mut impl Write,
record: &CapturedError,
) -> std::io::Result<()> {
let mut line = serde_json::to_vec(record)
.map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidData, e))?;
line.push(b'\n');
out.write_all(&line)
}
#[cfg(test)]
#[path = "store_lock_tests.rs"]
mod store_lock_tests;
#[cfg(test)]
mod tests {
use super::*;
use crate::error_capture::types::CapturedError;
fn make_record(msg: &str, fp: &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(),
}
}
#[test]
fn store_ring_bounded() {
let store = ErrorStore::with_path(None, 3);
assert!(store.is_empty());
for i in 0..5u32 {
store.append(make_record(&format!("err {i}"), &format!("fp{i}")));
}
assert_eq!(store.len(), 3);
let recent = store.recent_errors(10);
assert_eq!(recent.len(), 3);
assert_eq!(recent[0].message, "err 2");
assert_eq!(recent[2].message, "err 4");
}
#[test]
fn store_round_trip_write_read() {
let tmp_dir = {
let pid = std::process::id();
let nanos = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_nanos())
.unwrap_or(0);
std::env::temp_dir().join(format!("bugcap-test-{pid}-{nanos}"))
};
std::fs::create_dir_all(&tmp_dir).unwrap();
let file_path = tmp_dir.join(ERRORS_FILENAME);
{
let store = ErrorStore::with_path(Some(file_path.clone()), 10);
store.append(make_record("first error", "fp1"));
store.append(make_record("second error", "fp2"));
assert_eq!(store.len(), 2);
}
let store2 = ErrorStore::with_path(Some(file_path), 10);
let records = store2.recent_errors(10);
assert_eq!(records.len(), 2, "expected 2 records after reload");
assert_eq!(records[0].message, "first error");
assert_eq!(records[1].message, "second error");
}
#[test]
fn store_handles_missing_file_gracefully() {
let pid = std::process::id();
let nanos = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_nanos())
.unwrap_or(0);
let nonexistent = std::env::temp_dir().join(format!("bugcap-missing-{pid}-{nanos}.jsonl"));
let store = ErrorStore::with_path(Some(nonexistent), 10);
assert!(store.is_empty());
store.append(make_record("hello", "fp1"));
assert_eq!(store.len(), 1);
}
#[test]
fn store_corrupt_line_skipped() {
let tmp_dir = {
let pid = std::process::id();
let nanos = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_nanos())
.unwrap_or(0);
std::env::temp_dir().join(format!("bugcap-corrupt-{pid}-{nanos}"))
};
std::fs::create_dir_all(&tmp_dir).unwrap();
let file_path = tmp_dir.join(ERRORS_FILENAME);
{
let valid = serde_json::to_string(&make_record("valid first", "fp1")).unwrap();
let valid2 = serde_json::to_string(&make_record("valid second", "fp2")).unwrap();
let content = format!("{valid}\nnot-json-at-all\n{valid2}\n");
std::fs::write(&file_path, content).unwrap();
}
let store = ErrorStore::with_path(Some(file_path), 10);
assert_eq!(store.len(), 2, "corrupt line should be skipped");
assert_eq!(
store.corrupt_lines_skipped(),
1,
"#8028: the skip is counted"
);
let records = store.recent_errors(10);
assert_eq!(records[0].message, "valid first");
assert_eq!(records[1].message, "valid second");
}
#[test]
fn store_errors_by_fingerprint() {
let store = ErrorStore::with_path(None, 20);
store.append(make_record("err a", "fp1"));
store.append(make_record("err b", "fp2"));
store.append(make_record("err c", "fp1"));
let by_fp = store.errors_by_fingerprint();
assert_eq!(by_fp.len(), 2, "expected 2 unique fingerprints");
assert_eq!(by_fp[0].1, 2);
assert_eq!(by_fp[0].0.fingerprint, "fp1");
assert_eq!(by_fp[0].0.message, "err c");
assert_eq!(by_fp[1].1, 1);
}
#[test]
fn read_records_loads_file() {
let tmp_dir = {
let pid = std::process::id();
let nanos = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_nanos())
.unwrap_or(0);
std::env::temp_dir().join(format!("bugcap-readrec-{pid}-{nanos}"))
};
std::fs::create_dir_all(&tmp_dir).unwrap();
let file_path = tmp_dir.join(ERRORS_FILENAME);
let store = ErrorStore::with_path(Some(file_path.clone()), 10);
store.append(make_record("alpha error", "fp-a"));
store.append(make_record("beta error", "fp-b"));
let records = ErrorStore::read_records(&file_path, 10);
assert_eq!(records.len(), 2, "expected 2 records");
assert_eq!(records[0].message, "alpha error");
assert_eq!(records[1].message, "beta error");
}
#[test]
fn read_records_missing_file_is_empty() {
let nonexistent = std::env::temp_dir().join("bugcap-no-such-file-x99.jsonl");
let records = ErrorStore::read_records(&nonexistent, 50);
assert!(records.is_empty(), "missing file must yield empty vec");
}
}