use std::cell::Cell;
use std::error::Error;
use std::sync::Arc;
use std::time::Duration;
use super::expiry_race_tests::{Arms, RaceClock, actor, arm, fire, stamp};
use super::handle::BatchItem;
use super::{GroupOutcome, GroupWrite, ShardActor};
use crate::store::{MemoryStore, NodeStore, StoreError};
use crate::sync::ballot::Ballot;
use crate::sync::topology::SyncNodeId;
use crate::tree::{Hash, Node};
use crate::ttl::filter::{Visibility, visible_value_at};
use crate::wal::{DurableWal, FsyncPolicy};
#[derive(Debug)]
struct SyncFailStore {
inner: MemoryStore,
fail_sync: Cell<bool>,
}
impl NodeStore for SyncFailStore {
type Error = StoreError;
fn get(&self, hash: &Hash) -> Result<Option<Arc<Node>>, Self::Error> {
Ok(self.inner.get(hash))
}
fn put(&mut self, node: &Node) -> Result<Hash, Self::Error> {
Ok(self.inner.put(node))
}
fn sync_dirty_dirs(&self) -> Result<(), Self::Error> {
if self.fail_sync.get() {
Err(StoreError::Io(std::io::Error::other(
"injected merge barrier failure",
)))
} else {
Ok(())
}
}
}
#[test]
fn merge_adopt_failure_preserves_buffer_index_and_oracle() -> Result<(), Box<dyn Error>> {
let clock = Arc::new(RaceClock::at(1_000));
let source_dir = tempfile::tempdir()?;
let source_wal = DurableWal::new(
source_dir.path().join("source.wal"),
FsyncPolicy::CommitOnly,
)?;
let mut source = ShardActor::new_with_clock(source_wal, clock.clone());
let mut source_store = MemoryStore::new();
source.apply_durable(
b"source",
None,
b"v".to_vec(),
Some(Duration::from_nanos(5)),
stamp(1, 0),
&mut source_store,
)?;
let export = source.export_reachable(0, &source_store)?;
let target_dir = tempfile::tempdir()?;
let target_wal = DurableWal::new(
target_dir.path().join("target.wal"),
FsyncPolicy::CommitOnly,
)?;
let mut target = ShardActor::new_with_clock(target_wal, clock);
let mut store = SyncFailStore {
inner: MemoryStore::new(),
fail_sync: Cell::new(false),
};
target.apply_durable(
b"old",
None,
b"v".to_vec(),
Some(Duration::from_nanos(20)),
stamp(1, 1),
&mut store,
)?;
target.put_with_ttl(b"buffered", b"v", Some(Duration::from_nanos(30)), &store)?;
let mut arms = Arms::default();
let generation = arm(&mut target, &mut arms)?;
let before = target.expiry_metrics();
store.fail_sync.set(true);
assert!(target.merge_adopt(&[export], &mut store).is_err());
assert!(matches!(
target.buffer().get(b"buffered"),
crate::wal::LookupResult::BufferedValue(_)
));
assert!(target.get_raw(b"buffered", &store)?.is_some());
assert_eq!(target.current_expiry_generation(), Some(generation));
assert_eq!(target.expiry_metrics(), before);
target.arm_expiry(&mut arms)?;
assert_eq!(
arms.0.len(),
1,
"failed adoption must not speculatively re-arm"
);
Ok(())
}
#[test]
fn durable_and_batch_rejections_move_zero_index_or_arm_counters() -> Result<(), Box<dyn Error>> {
let clock = Arc::new(RaceClock::at(100));
let (_dir, mut actor, mut store) = actor(clock)?;
let mut arms = Arms::default();
actor.put_with_ttl(b"existing", b"old", Some(Duration::from_nanos(50)), &store)?;
arm(&mut actor, &mut arms)?;
actor.record_promise(Ballot::new(9, SyncNodeId::new("new-owner")))?;
let before_fence = actor.expiry_metrics();
assert!(
actor
.apply_durable(
b"fenced",
None,
b"v".to_vec(),
Some(Duration::from_nanos(1)),
stamp(1, 0),
&mut store
)
.is_err()
);
assert_eq!(actor.expiry_metrics(), before_fence);
let before_cas = actor.expiry_metrics();
assert!(
actor
.apply_durable(
b"existing",
Some(Hash::of(b"wrong")),
b"new".to_vec(),
Some(Duration::from_nanos(1)),
stamp(9, 1),
&mut store
)
.is_err()
);
assert_eq!(actor.expiry_metrics(), before_cas);
let items: Vec<BatchItem> = vec![
(
b"would-pass".to_vec(),
None,
b"v".to_vec(),
Some(Duration::from_nanos(1)),
),
(
b"existing".to_vec(),
None,
b"bad".to_vec(),
Some(Duration::from_nanos(2)),
),
];
let before_batch = actor.expiry_metrics();
assert!(
actor
.apply_durable_batch(items, stamp(9, 2), &mut store)
.is_err()
);
assert_eq!(actor.expiry_metrics(), before_batch);
assert!(actor.get_raw(b"would-pass", &store)?.is_none());
Ok(())
}
#[test]
fn grouped_partial_rejection_indexes_survivors_in_order() -> Result<(), Box<dyn Error>> {
let base = crate::branch::current_timestamp();
let clock = Arc::new(RaceClock::at(base));
let (_dir, mut actor, mut store) = actor(clock.clone())?;
actor.apply_durable(
b"rejected",
None,
b"seed".to_vec(),
Some(Duration::from_secs(300)),
stamp(1, 0),
&mut store,
)?;
let before = actor.expiry_metrics();
let outcomes = actor.apply_group(
vec![
GroupWrite::ApplyValue {
key: b"first".to_vec(),
expected: None,
value: b"v".to_vec(),
ttl: Some(Duration::from_secs(100)),
stamp: stamp(2, 0),
},
GroupWrite::ApplyValue {
key: b"rejected".to_vec(),
expected: None,
value: b"bad".to_vec(),
ttl: Some(Duration::from_secs(1)),
stamp: stamp(2, 1),
},
GroupWrite::ApplyValue {
key: b"second".to_vec(),
expected: None,
value: b"v".to_vec(),
ttl: Some(Duration::from_secs(200)),
stamp: stamp(2, 2),
},
],
&mut store,
);
assert!(matches!(
outcomes.as_slice(),
[
GroupOutcome::Committed,
GroupOutcome::Rejected(_),
GroupOutcome::Committed
]
));
assert_eq!(
actor.expiry_metrics().index_mutations - before.index_mutations,
2
);
let mut arms = Arms::default();
let first = arm(&mut actor, &mut arms)?;
clock.set(base.saturating_add(100_000_000_000));
assert!(fire(&mut actor, first, &store)?);
let second = arm(&mut actor, &mut arms)?;
assert_eq!(arms.0[1].0, Duration::from_secs(100));
clock.set(base.saturating_add(200_000_000_000));
assert!(fire(&mut actor, second, &store)?);
assert_eq!(actor.get(b"rejected", &store)?, Some(b"seed".to_vec()));
Ok(())
}
#[test]
fn overdue_delivery_is_immediate_once_and_backward_jump_allows_resurrection()
-> Result<(), Box<dyn Error>> {
let clock = Arc::new(RaceClock::at(100));
let (_dir, mut actor, store) = actor(clock.clone())?;
let mut arms = Arms::default();
actor.put_with_ttl(b"k", b"v", Some(Duration::ZERO), &store)?;
let generation = arm(&mut actor, &mut arms)?;
assert_eq!(arms.0, vec![(Duration::ZERO, generation)]);
assert!(fire(&mut actor, generation, &store)?);
actor.arm_expiry(&mut arms)?;
assert_eq!(arms.0.len(), 1);
let encoded = crate::ttl::entry::encode_optional_ttl_at(
b"v".to_vec(),
Some(Duration::from_nanos(10)),
100,
)?;
assert_eq!(visible_value_at(&encoded, 110)?, Visibility::Expired);
clock.set(90);
assert_eq!(
visible_value_at(&encoded, 90)?,
Visibility::Live(b"v".to_vec())
);
Ok(())
}
#[derive(Debug)]
struct AlwaysFailStore;
impl NodeStore for AlwaysFailStore {
type Error = StoreError;
fn get(&self, _hash: &Hash) -> Result<Option<Arc<Node>>, Self::Error> {
Ok(None)
}
fn put(&mut self, _node: &Node) -> Result<Hash, Self::Error> {
Err(StoreError::Io(std::io::Error::other(
"injected group commit failure",
)))
}
}
#[test]
fn grouped_commit_failure_rolls_back_every_expiry_effect() -> Result<(), Box<dyn Error>> {
let clock = Arc::new(RaceClock::at(1_000));
let dir = tempfile::tempdir()?;
let wal = DurableWal::new(
dir.path().join("group-failure.wal"),
FsyncPolicy::CommitOnly,
)?;
let mut actor = ShardActor::new_with_clock(wal, clock);
let before = actor.expiry_metrics();
let outcomes = actor.apply_group(
vec![
GroupWrite::ApplyValue {
key: b"a".to_vec(),
expected: None,
value: b"v".to_vec(),
ttl: Some(Duration::from_nanos(1)),
stamp: stamp(1, 0),
},
GroupWrite::ApplyValue {
key: b"b".to_vec(),
expected: None,
value: b"v".to_vec(),
ttl: Some(Duration::from_nanos(2)),
stamp: stamp(1, 1),
},
],
&mut AlwaysFailStore,
);
assert!(
outcomes
.iter()
.all(|outcome| matches!(outcome, GroupOutcome::CommitFailed(_)))
);
assert_eq!(actor.expiry_metrics(), before);
assert!(actor.buffer().is_empty());
let mut arms = Arms::default();
actor.arm_expiry(&mut arms)?;
assert!(arms.0.is_empty());
Ok(())
}