use std::net::IpAddr;
use std::sync::Arc;
use std::time::Duration;
use crate::persistence::{PersistedState, Persistence};
use crate::replica::version_hash;
use crate::{FileSnapshot, ReplicatedMap};
use super::ephemeral_config;
#[tokio::test]
async fn persistence_roundtrip_recovers_entries_and_tombstones() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("snapshot.bin");
let store = ReplicatedMap::<i32, i32>::new(ephemeral_config())
.await
.expect("bind failed")
.with_persistence(Arc::new(FileSnapshot::new(&path)));
store.insert(1, 11); store.insert(2, 22);
store.remove(&2); let expected = store.fingerprint(..);
store.snapshot();
let restarted = ReplicatedMap::<i32, i32>::new(ephemeral_config())
.await
.expect("bind failed")
.with_persistence(Arc::new(FileSnapshot::new(&path)));
assert_eq!(restarted.get(&1).as_deref(), Some(&11));
assert!(restarted.get(&2).is_none(), "tombstone was not recovered");
assert_eq!(
restarted.fingerprint(..),
expected,
"recovered state must hash identically (timestamps preserved)"
);
assert!(restarted.tombstones.remove(&2).is_some());
}
struct FlakyLoad {
kind: std::io::ErrorKind,
failures_remaining: std::sync::atomic::AtomicU32,
}
impl<K: Send + Sync + 'static, V: Send + Sync + 'static> Persistence<K, V> for FlakyLoad {
fn load(&self) -> std::io::Result<Option<PersistedState<K, V>>> {
use std::sync::atomic::Ordering;
if self
.failures_remaining
.fetch_update(Ordering::AcqRel, Ordering::Acquire, |n| {
(n > 0).then(|| n - 1)
})
.is_ok()
{
return Err(std::io::Error::new(
self.kind,
"simulated transient failure",
));
}
Ok(None)
}
fn save(&self, _state: &PersistedState<K, V>) -> std::io::Result<()> {
Ok(())
}
}
#[test]
fn backoff_delay_doubles_from_the_base() {
assert_eq!(
super::super::persistence::backoff_delay(1),
super::super::persistence::LOAD_RETRY_BASE_DELAY
);
assert_eq!(
super::super::persistence::backoff_delay(2),
super::super::persistence::LOAD_RETRY_BASE_DELAY * 2
);
assert_eq!(
super::super::persistence::backoff_delay(3),
super::super::persistence::LOAD_RETRY_BASE_DELAY * 4
);
assert_eq!(
super::super::persistence::backoff_delay(4),
super::super::persistence::LOAD_RETRY_BASE_DELAY * 8
);
}
#[tokio::test]
async fn transient_load_failure_is_retried_not_fatal() {
let backend = Arc::new(FlakyLoad {
kind: std::io::ErrorKind::PermissionDenied,
failures_remaining: std::sync::atomic::AtomicU32::new(
super::super::persistence::LOAD_RETRY_ATTEMPTS - 1,
),
});
let _store = ReplicatedMap::<i32, i32>::new(ephemeral_config())
.await
.expect("bind failed")
.with_persistence(backend);
}
#[tokio::test]
#[should_panic(expected = "failed to load persisted state after")]
async fn load_failure_beyond_retry_budget_still_panics() {
let backend = Arc::new(FlakyLoad {
kind: std::io::ErrorKind::PermissionDenied,
failures_remaining: std::sync::atomic::AtomicU32::new(
super::super::persistence::LOAD_RETRY_ATTEMPTS,
),
});
let _store = ReplicatedMap::<i32, i32>::new(ephemeral_config())
.await
.expect("bind failed")
.with_persistence(backend);
}
#[tokio::test]
#[should_panic(expected = "persisted state is corrupt or from an incompatible format")]
async fn invalid_data_panics_without_retrying() {
let backend = Arc::new(FlakyLoad {
kind: std::io::ErrorKind::InvalidData,
failures_remaining: std::sync::atomic::AtomicU32::new(1),
});
let start = std::time::Instant::now();
let _store = ReplicatedMap::<i32, i32>::new(ephemeral_config())
.await
.expect("bind failed")
.with_persistence(backend);
assert!(start.elapsed() < super::super::persistence::LOAD_RETRY_BASE_DELAY);
}
#[tokio::test]
async fn snapshot_across_multiple_chunks_recovers_every_entry() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("snapshot.bin");
let n = super::super::persistence::SNAPSHOT_CHUNK_SIZE * 2 + 17;
let store = ReplicatedMap::<i32, i32>::new(ephemeral_config())
.await
.expect("bind failed")
.with_persistence(Arc::new(FileSnapshot::new(&path)));
for k in 0..n as i32 {
store.just_insert(k, k * 2);
}
let expected = store.fingerprint(..);
store.snapshot();
let restarted = ReplicatedMap::<i32, i32>::new(ephemeral_config())
.await
.expect("bind failed")
.with_persistence(Arc::new(FileSnapshot::new(&path)));
assert_eq!(
restarted.fingerprint(..),
expected,
"chunked snapshot must recover every entry across chunk boundaries"
);
for k in 0..n as i32 {
assert_eq!(restarted.get(&k).as_deref(), Some(&(k * 2)));
}
}
#[tokio::test]
async fn restart_preserves_membership_and_acks() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("snapshot.bin");
let peer: IpAddr = "127.0.0.99".parse().unwrap();
let store = ReplicatedMap::<i32, i32>::new(ephemeral_config())
.await
.expect("bind failed")
.with_persistence(Arc::new(FileSnapshot::new(&path)));
store.engine.members.write().insert(peer);
store.insert(5, 55);
store.remove(&5); store
.engine
.tombstone_acks
.write()
.entry(5)
.or_default()
.insert(peer, 123);
store.snapshot();
let restarted = ReplicatedMap::<i32, i32>::new(ephemeral_config())
.await
.expect("bind failed")
.with_persistence(Arc::new(FileSnapshot::new(&path)));
assert!(
restarted.engine.members.read().contains(&peer),
"membership set was not restored"
);
assert_eq!(
restarted
.engine
.tombstone_acks
.read()
.get(&5)
.and_then(|acks| acks.get(&peer)),
Some(&123),
"tombstone acknowledgments were not restored"
);
}
#[tokio::test]
async fn restart_keeps_tombstone_gc_gated() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("snapshot.bin");
let peer: IpAddr = "127.0.0.98".parse().unwrap();
let store = ReplicatedMap::<i32, i32>::new(ephemeral_config())
.await
.expect("bind failed")
.with_persistence(Arc::new(FileSnapshot::new(&path)));
store.engine.members.write().insert(peer);
store.insert(1, 11);
store.remove(&1);
let version = store.engine.map.read().get(&1).map(version_hash).unwrap();
assert!(
!store.engine.is_tombstone_stable(&1, version),
"precondition: tombstone is gated before restart"
);
store.snapshot();
let fresh = ReplicatedMap::<i32, i32>::new(ephemeral_config())
.await
.expect("bind failed");
fresh.insert(1, 11);
fresh.remove(&1);
let fresh_version = fresh.engine.map.read().get(&1).map(version_hash).unwrap();
assert!(
fresh.engine.is_tombstone_stable(&1, fresh_version),
"a fresh restart with no membership would (wrongly) GC the tombstone — the hazard this guards against"
);
let restarted = ReplicatedMap::<i32, i32>::new(ephemeral_config())
.await
.expect("bind failed")
.with_persistence(Arc::new(FileSnapshot::new(&path)));
assert!(restarted.get(&1).is_none(), "tombstone was not recovered");
let version = restarted
.engine
.map
.read()
.get(&1)
.map(version_hash)
.unwrap();
assert!(
!restarted.engine.is_tombstone_stable(&1, version),
"restart dropped causal-stability state: tombstone would be GC'd and could resurrect"
);
}
#[tokio::test]
async fn restart_clock_advanced_past_persisted_max_stamp() {
use crate::clock::{Hlc, LogicalCounter, ManualClock, NodeId, PhysicalTime};
use crate::persistence::{InMemoryPersistence, PersistedState};
let clock = Arc::new(ManualClock::new(NodeId::new(1)));
let persisted_stamp = crate::clock::Timestamp::new(
Hlc::new(PhysicalTime::from_millis(100), LogicalCounter::new(0)),
NodeId::new(1),
);
let backend = Arc::new(InMemoryPersistence::<i32, i32>::new());
backend
.save(&PersistedState::from(vec![(
42,
crate::entry::Entry::present(persisted_stamp, 999),
)]))
.unwrap();
let store = ReplicatedMap::<i32, i32>::new_with_clock(
ephemeral_config().with_node_id(NodeId::new(1)),
clock,
)
.await
.expect("bind failed")
.with_persistence(backend);
store.insert(99, 1);
let minted_stamp = store
.engine
.map
.read()
.get(&99)
.map(|entry| entry.stamp)
.expect("key 99 must be present after insert");
assert!(
minted_stamp > persisted_stamp,
"post-restart write timestamp {minted_stamp:?} is not strictly greater than the \
persisted max {persisted_stamp:?}; the clock was not advanced on load"
);
}
#[tokio::test]
async fn restart_insert_beats_persisted_tombstone() {
use crate::clock::{Hlc, LogicalCounter, ManualClock, NodeId, PhysicalTime};
use crate::persistence::{InMemoryPersistence, PersistedState};
let clock = Arc::new(ManualClock::new(NodeId::new(2)));
let tombstone_stamp = crate::clock::Timestamp::new(
Hlc::new(PhysicalTime::from_millis(200), LogicalCounter::new(0)),
NodeId::new(2),
);
let backend = Arc::new(InMemoryPersistence::<i32, i32>::new());
backend
.save(&PersistedState::from(vec![(
7,
crate::entry::Entry::tombstone(tombstone_stamp), )]))
.unwrap();
let store = ReplicatedMap::<i32, i32>::new_with_clock(
ephemeral_config().with_node_id(NodeId::new(2)),
clock,
)
.await
.expect("bind failed")
.with_persistence(backend);
assert!(
store.get(&7).is_none(),
"tombstone was not recovered: expected key 7 to be absent after loading"
);
store.insert(7, 42);
let minted_stamp = store
.engine
.map
.read()
.get(&7)
.map(|entry| entry.stamp)
.expect("key 7 must be present after insert");
assert!(
minted_stamp > tombstone_stamp,
"post-restart insert timestamp {minted_stamp:?} is not strictly greater than the \
persisted tombstone stamp {tombstone_stamp:?}; the clock was not advanced on load, \
so a peer reconciling with this node could resurrect the tombstone via LWW"
);
}
#[tokio::test]
async fn snapshot_periodically_actually_persists() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("periodic.bin");
let store = ReplicatedMap::<i32, i32>::new(ephemeral_config())
.await
.expect("bind failed")
.with_persistence(Arc::new(FileSnapshot::new(&path)));
store.just_insert(1, 10);
let _ = tokio::time::timeout(
super::super::persistence::SNAPSHOT_INTERVAL + Duration::from_secs(1),
store.snapshot_periodically(),
)
.await;
let restarted = ReplicatedMap::<i32, i32>::new(ephemeral_config())
.await
.expect("bind failed")
.with_persistence(Arc::new(FileSnapshot::new(&path)));
assert_eq!(
restarted.get(&1).as_deref(),
Some(&10),
"no snapshot was written after a full SNAPSHOT_INTERVAL"
);
}