use std::collections::BTreeMap;
use std::io;
use std::sync::Arc;
use std::sync::atomic::{AtomicU64, Ordering};
#[derive(Debug, Clone, Default)]
pub struct Generation(Arc<AtomicU64>);
impl Generation {
#[must_use]
pub fn new() -> Self {
Self::default()
}
pub fn bump(&self) -> u64 {
self.0.fetch_add(1, Ordering::SeqCst) + 1
}
#[must_use]
pub fn current(&self) -> u64 {
self.0.load(Ordering::SeqCst)
}
}
#[derive(Debug, Clone)]
pub struct ReplayQueue<W: Ord + Clone, O: Clone> {
pending: BTreeMap<W, O>,
}
impl<W: Ord + Clone, O: Clone> Default for ReplayQueue<W, O> {
fn default() -> Self {
Self {
pending: BTreeMap::new(),
}
}
}
impl<W: Ord + Clone, O: Clone> ReplayQueue<W, O> {
#[must_use]
pub fn new() -> Self {
Self::default()
}
pub fn set(&mut self, worktree: W, overlay: O) {
self.pending.insert(worktree, overlay);
}
pub fn remove(&mut self, worktree: &W) {
self.pending.remove(worktree);
}
#[must_use]
pub fn len(&self) -> usize {
self.pending.len()
}
#[must_use]
pub fn is_empty(&self) -> bool {
self.pending.is_empty()
}
#[must_use]
pub fn snapshot_sorted(&self) -> Vec<(W, O)> {
self.pending
.iter()
.map(|(w, o)| (w.clone(), o.clone()))
.collect()
}
}
pub trait RecoverySink<W, O> {
fn reapply(&mut self, worktree: &W, overlay: &O) -> io::Result<()>;
}
#[derive(Debug)]
pub enum ReplayOutcome {
Completed { applied: usize },
Superseded { applied_before_abort: usize },
SinkError {
applied_before_error: usize,
error: io::Error,
},
}
pub fn replay<W, O>(
queue: &ReplayQueue<W, O>,
sink: &mut dyn RecoverySink<W, O>,
generation: &Generation,
expected_generation: u64,
) -> ReplayOutcome
where
W: Ord + Clone,
O: Clone,
{
let snapshot = queue.snapshot_sorted();
let mut applied = 0usize;
for (w, o) in &snapshot {
if generation.current() != expected_generation {
return ReplayOutcome::Superseded {
applied_before_abort: applied,
};
}
if let Err(error) = sink.reapply(w, o) {
return ReplayOutcome::SinkError {
applied_before_error: applied,
error,
};
}
applied += 1;
}
if generation.current() != expected_generation {
return ReplayOutcome::Superseded {
applied_before_abort: applied,
};
}
ReplayOutcome::Completed { applied }
}
#[cfg(test)]
mod tests {
use super::*;
use std::cell::RefCell;
struct MockSink<'g> {
applied: RefCell<Vec<(String, String)>>,
fail_at: Option<usize>,
bump_at: Option<(usize, &'g Generation)>,
}
impl<'g> MockSink<'g> {
fn new() -> Self {
Self {
applied: RefCell::new(Vec::new()),
fail_at: None,
bump_at: None,
}
}
fn failing_at(mut self, n: usize) -> Self {
self.fail_at = Some(n);
self
}
fn bumping_at(mut self, n: usize, g: &'g Generation) -> Self {
self.bump_at = Some((n, g));
self
}
}
impl RecoverySink<String, String> for MockSink<'_> {
fn reapply(&mut self, w: &String, o: &String) -> io::Result<()> {
let n = self.applied.borrow().len();
if let Some((at, g)) = self.bump_at {
if n == at {
g.bump(); }
}
if self.fail_at == Some(n) {
return Err(io::Error::other("sink boom"));
}
self.applied.borrow_mut().push((w.clone(), o.clone()));
Ok(())
}
}
fn q(pairs: &[(&str, &str)]) -> ReplayQueue<String, String> {
let mut q = ReplayQueue::new();
for (w, o) in pairs {
q.set((*w).to_string(), (*o).to_string());
}
q
}
#[test]
fn recovery_enqueue_is_idempotent_same_pair() {
let mut queue = ReplayQueue::new();
queue.set("A".to_string(), "ov1".to_string());
queue.set("A".to_string(), "ov1".to_string());
assert_eq!(queue.len(), 1);
}
#[test]
fn recovery_latest_wins_per_worktree_stale_overlay_never_replayed() {
let mut queue = ReplayQueue::new();
queue.set("A".to_string(), "old".to_string());
queue.set("A".to_string(), "new".to_string()); let snap = queue.snapshot_sorted();
assert_eq!(snap, vec![("A".to_string(), "new".to_string())]);
assert!(!snap.iter().any(|(_, o)| o == "old"));
}
#[test]
fn recovery_snapshot_is_deterministic_sorted_by_worktree() {
let queue = q(&[("C", "c"), ("A", "a"), ("B", "b")]);
let snap = queue.snapshot_sorted();
let order: Vec<&str> = snap.iter().map(|(w, _)| w.as_str()).collect();
assert_eq!(order, ["A", "B", "C"], "replay order must be reproducible");
}
#[test]
fn recovery_remove_drops_worktree_from_replay() {
let mut queue = q(&[("A", "a"), ("B", "b")]);
queue.remove(&"A".to_string());
assert_eq!(
queue.snapshot_sorted(),
vec![("B".to_string(), "b".to_string())]
);
}
#[test]
fn recovery_replay_applies_all_in_sorted_order() {
let queue = q(&[("C", "c"), ("A", "a"), ("B", "b")]);
let genr = Generation::new();
let g = genr.bump(); let mut sink = MockSink::new();
let out = replay(&queue, &mut sink, &genr, g);
match out {
ReplayOutcome::Completed { applied } => assert_eq!(applied, 3),
other => panic!("expected Completed, got {other:?}"),
}
let got: Vec<String> = sink
.applied
.borrow()
.iter()
.map(|(w, _)| w.clone())
.collect();
assert_eq!(got, ["A", "B", "C"], "applied in deterministic WT order");
}
#[test]
fn recovery_empty_queue_replay_is_completed_noop() {
let queue: ReplayQueue<String, String> = ReplayQueue::new();
let genr = Generation::new();
let g = genr.bump();
let mut sink = MockSink::new();
match replay(&queue, &mut sink, &genr, g) {
ReplayOutcome::Completed { applied } => assert_eq!(applied, 0),
other => panic!("expected empty Completed, got {other:?}"),
}
assert!(sink.applied.borrow().is_empty());
}
#[test]
fn recovery_replay_propagates_sink_error_without_false_completion() {
let queue = q(&[("A", "a"), ("B", "b"), ("C", "c")]);
let genr = Generation::new();
let g = genr.bump();
let mut sink = MockSink::new().failing_at(1); match replay(&queue, &mut sink, &genr, g) {
ReplayOutcome::SinkError {
applied_before_error,
..
} => assert_eq!(
applied_before_error, 1,
"A applied, B errored, C not reached"
),
other => panic!("expected SinkError, got {other:?}"),
}
assert_eq!(sink.applied.borrow().len(), 1);
}
#[test]
fn recovery_replay_superseded_by_newer_generation_aborts_cleanly() {
let queue = q(&[("A", "a"), ("B", "b"), ("C", "c"), ("D", "d")]);
let genr = Generation::new();
let g = genr.bump(); let mut sink = MockSink::new().bumping_at(2, &genr);
match replay(&queue, &mut sink, &genr, g) {
ReplayOutcome::Superseded {
applied_before_abort,
} => assert_eq!(
applied_before_abort, 3,
"A,B normal; respawn during C (in-flight, atomic) ⇒ C applied, \
D never issued — Superseded, never Completed"
),
other => panic!("expected Superseded, got {other:?}"),
}
let names: Vec<String> = sink
.applied
.borrow()
.iter()
.map(|(w, _)| w.clone())
.collect();
assert_eq!(names, ["A", "B", "C"], "D must never reach the doomed RA");
}
#[test]
fn recovery_final_fence_catches_respawn_on_last_overlay() {
let queue = q(&[("A", "a")]);
let genr = Generation::new();
let g = genr.bump();
let mut sink = MockSink::new().bumping_at(0, &genr);
match replay(&queue, &mut sink, &genr, g) {
ReplayOutcome::Superseded {
applied_before_abort,
} => assert_eq!(applied_before_abort, 1),
other => panic!("expected Superseded at final fence, got {other:?}"),
}
}
#[test]
fn recovery_generation_bumps_monotonically() {
let genr = Generation::new();
assert_eq!(genr.current(), 0);
assert_eq!(genr.bump(), 1);
assert_eq!(genr.bump(), 2);
assert_eq!(genr.current(), 2);
}
#[test]
fn recovery_replay_into_correct_generation_completes() {
let queue = q(&[("A", "a"), ("B", "b")]);
let genr = Generation::new();
genr.bump();
genr.bump(); let mut sink = MockSink::new();
match replay(&queue, &mut sink, &genr, 2) {
ReplayOutcome::Completed { applied } => assert_eq!(applied, 2),
other => panic!("expected Completed at stable gen 2, got {other:?}"),
}
}
}