use std::collections::HashMap;
use std::path::{Path, PathBuf};
use std::time::{Duration, Instant};
use tracing::{info, warn};
const BASE_DELAY: Duration = Duration::from_secs(1);
const MAX_DELAY: Duration = Duration::from_secs(300);
const QUARANTINE_AFTER: u32 = 10;
const WARN_AFTER: u32 = 5;
pub const STUCK_DIR: &str = "stuck";
#[derive(Debug, Clone, Copy)]
struct Attempt {
failures: u32,
next: Instant,
}
#[derive(Debug, PartialEq, Eq)]
pub enum AfterFailure {
Retry,
Quarantine,
}
#[derive(Debug, Default)]
pub struct RetryLedger {
attempts: HashMap<PathBuf, Attempt>,
}
impl RetryLedger {
pub fn new() -> Self {
Self::default()
}
pub fn is_due(&self, path: &Path, now: Instant) -> bool {
self.attempts.get(path).is_none_or(|a| now >= a.next)
}
pub fn record_failure(&mut self, path: &Path, now: Instant) -> AfterFailure {
let entry = self.attempts.entry(path.to_path_buf()).or_insert(Attempt {
failures: 0,
next: now,
});
entry.failures += 1;
let shift = entry.failures.saturating_sub(1).min(16);
entry.next = now + (BASE_DELAY * (1u32 << shift)).min(MAX_DELAY);
if entry.failures >= QUARANTINE_AFTER {
AfterFailure::Quarantine
} else {
AfterFailure::Retry
}
}
pub fn failures(&self, path: &Path) -> u32 {
self.attempts.get(path).map_or(0, |a| a.failures)
}
pub fn should_warn(&self, path: &Path) -> bool {
self.failures(path) >= WARN_AFTER
}
pub fn record_success(&mut self, path: &Path) {
self.attempts.remove(path);
}
pub fn retain_present(&mut self, present: &[PathBuf]) {
let present: std::collections::HashSet<&Path> =
present.iter().map(PathBuf::as_path).collect();
self.attempts.retain(|p, _| present.contains(p.as_path()));
}
#[cfg(test)]
fn tracked(&self) -> usize {
self.attempts.len()
}
}
pub fn quarantine(path: &Path, label: &str) -> bool {
let Some(dir) = path.parent() else {
warn!(path = %path.display(), "{label}: cannot quarantine a file with no parent");
return false;
};
let stuck = dir.join(STUCK_DIR);
if let Err(e) = std::fs::create_dir_all(&stuck) {
warn!(error = %e, dir = %stuck.display(), "{label}: cannot create quarantine dir");
return false;
}
let Some(name) = path.file_name() else {
return false;
};
let Some(dest) = free_name(&stuck, name) else {
warn!(
dir = %stuck.display(),
"{label}: no free name in the quarantine dir; leaving the file in place",
);
return false;
};
match std::fs::rename(path, &dest) {
Ok(()) => {
info!(
from = %path.display(),
to = %dest.display(),
"{label}: file could not be published after {QUARANTINE_AFTER} attempts — quarantined",
);
true
}
Err(e) => {
warn!(
error = %e,
path = %path.display(),
"{label}: quarantine move failed; leaving the file in place",
);
false
}
}
}
fn free_name(dir: &Path, name: &std::ffi::OsStr) -> Option<PathBuf> {
let plain = dir.join(name);
if !plain.exists() {
return Some(plain);
}
(1..1000).find_map(|n| {
let mut alt = name.to_os_string();
alt.push(format!(".{n}"));
let p = dir.join(alt);
(!p.exists()).then_some(p)
})
}
#[cfg(test)]
mod tests {
use super::*;
fn p(name: &str) -> PathBuf {
PathBuf::from(name)
}
#[test]
fn an_unseen_file_is_due_immediately() {
let l = RetryLedger::new();
assert!(l.is_due(&p("a.json"), Instant::now()));
}
#[test]
fn one_failure_delays_by_the_drain_interval_not_more() {
let mut l = RetryLedger::new();
let t = Instant::now();
assert_eq!(l.record_failure(&p("a.json"), t), AfterFailure::Retry);
assert!(!l.is_due(&p("a.json"), t));
assert!(l.is_due(&p("a.json"), t + BASE_DELAY));
}
#[test]
fn the_delay_doubles_and_then_stops_growing() {
let mut l = RetryLedger::new();
let mut t = Instant::now();
let mut seen = Vec::new();
for _ in 0..12 {
l.record_failure(&p("a.json"), t);
let next = l.attempts[&p("a.json")].next;
seen.push(next - t);
t = next;
}
assert_eq!(seen[0], Duration::from_secs(1));
assert_eq!(seen[1], Duration::from_secs(2));
assert_eq!(seen[2], Duration::from_secs(4));
assert!(seen.iter().all(|d| *d <= MAX_DELAY));
assert_eq!(*seen.last().unwrap(), MAX_DELAY);
}
#[test]
fn a_long_stuck_file_does_not_overflow_the_shift() {
let mut l = RetryLedger::new();
let t = Instant::now();
for _ in 0..200 {
l.record_failure(&p("a.json"), t);
}
assert_eq!(l.attempts[&p("a.json")].next - t, MAX_DELAY);
}
#[test]
fn the_ceiling_is_never_actually_waited_on_at_the_current_threshold() {
let mut l = RetryLedger::new();
let mut t = Instant::now();
let mut waited = Duration::ZERO;
for n in 1..=QUARANTINE_AFTER {
let outcome = l.record_failure(&p("a.json"), t);
let delay = l.attempts[&p("a.json")].next - t;
if n < QUARANTINE_AFTER {
assert_eq!(outcome, AfterFailure::Retry);
assert!(delay < MAX_DELAY, "delay {delay:?} at failure {n}");
waited += delay;
} else {
assert_eq!(outcome, AfterFailure::Quarantine);
assert_eq!(delay, MAX_DELAY, "the ceiling first binds here…");
}
t = l.attempts[&p("a.json")].next;
}
assert_eq!(waited, Duration::from_secs(511)); }
#[test]
fn it_gives_up_after_the_quarantine_threshold() {
let mut l = RetryLedger::new();
let t = Instant::now();
for i in 1..QUARANTINE_AFTER {
assert_eq!(
l.record_failure(&p("a.json"), t),
AfterFailure::Retry,
"failure {i} should still retry"
);
}
assert_eq!(l.record_failure(&p("a.json"), t), AfterFailure::Quarantine);
}
#[test]
fn warning_starts_late_enough_to_survive_a_broker_restart() {
let mut l = RetryLedger::new();
let t = Instant::now();
l.record_failure(&p("a.json"), t);
assert!(!l.should_warn(&p("a.json")));
for _ in 1..WARN_AFTER {
l.record_failure(&p("a.json"), t);
}
assert!(l.should_warn(&p("a.json")));
}
#[test]
fn success_clears_the_history() {
let mut l = RetryLedger::new();
let t = Instant::now();
l.record_failure(&p("a.json"), t);
l.record_failure(&p("a.json"), t);
l.record_success(&p("a.json"));
assert_eq!(l.failures(&p("a.json")), 0);
assert!(l.is_due(&p("a.json"), t));
}
#[test]
fn files_that_left_the_directory_stop_being_tracked() {
let mut l = RetryLedger::new();
let t = Instant::now();
l.record_failure(&p("a.json"), t);
l.record_failure(&p("b.json"), t);
assert_eq!(l.tracked(), 2);
l.retain_present(&[p("b.json")]);
assert_eq!(l.tracked(), 1);
assert_eq!(l.failures(&p("a.json")), 0);
}
#[test]
fn backoff_is_per_file_not_global() {
let mut l = RetryLedger::new();
let t = Instant::now();
for _ in 0..5 {
l.record_failure(&p("stuck.json"), t);
}
assert!(!l.is_due(&p("stuck.json"), t));
assert!(l.is_due(&p("fresh.json"), t));
}
fn tmp(tag: &str) -> PathBuf {
let d = std::env::temp_dir().join(format!("kanade-outbox-{}-{tag}", std::process::id()));
let _ = std::fs::remove_dir_all(&d);
std::fs::create_dir_all(&d).unwrap();
d
}
#[test]
fn a_failed_quarantine_says_so() {
let dir = tmp("failed");
std::fs::write(dir.join(STUCK_DIR), b"not a directory").unwrap();
let f = dir.join("a.json");
std::fs::write(&f, b"{}").unwrap();
assert!(!quarantine(&f, "outbox"));
assert!(
f.exists(),
"the file must stay put so its backoff still applies"
);
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn quarantine_does_not_overwrite_an_earlier_payload() {
let dir = tmp("collide");
let f = dir.join("a.json");
std::fs::write(&f, b"first").unwrap();
assert!(quarantine(&f, "outbox"));
std::fs::write(&f, b"second").unwrap();
assert!(quarantine(&f, "outbox"));
let stuck = dir.join(STUCK_DIR);
assert_eq!(std::fs::read(stuck.join("a.json")).unwrap(), b"first");
assert_eq!(std::fs::read(stuck.join("a.json.1")).unwrap(), b"second");
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn quarantine_moves_the_file_and_keeps_its_bytes() {
let dir = std::env::temp_dir().join(format!("kanade-outbox-test-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
std::fs::create_dir_all(&dir).unwrap();
let f = dir.join("a.json");
std::fs::write(&f, b"{\"payload\":1}").unwrap();
assert!(quarantine(&f, "outbox"));
assert!(!f.exists(), "the drain loop must stop seeing it");
let moved = dir.join(STUCK_DIR).join("a.json");
assert_eq!(
std::fs::read(&moved).unwrap(),
b"{\"payload\":1}",
"the payload is preserved for diagnosis, not deleted"
);
let _ = std::fs::remove_dir_all(&dir);
}
}