use crate::gossip::wire::{decode_delta, encode_delta};
use crate::gossip::PubSubManager;
use crate::identity::AgentId;
use crate::kv::store::AccessPolicy;
use crate::kv::{KvStore, KvStoreDelta, Result};
use saorsa_gossip_types::PeerId;
use serde::{Deserialize, Serialize};
use std::path::{Path, PathBuf};
use std::sync::Arc;
use tokio::sync::RwLock;
const STATE_SYNC_TOPIC_SUFFIX: &str = "/state-sync";
const STATE_REQUEST_RETRY_SECS: [u64; 4] = [1, 5, 15, 30];
const STATE_REQUEST_TAIL_START_SECS: u64 = 30;
const STATE_REQUEST_TAIL_CAP_SECS: u64 = 300;
fn state_request_delays() -> impl Iterator<Item = u64> {
let tail = std::iter::successors(Some(STATE_REQUEST_TAIL_START_SECS), |d| {
Some(d.saturating_mul(2).min(STATE_REQUEST_TAIL_CAP_SECS))
});
STATE_REQUEST_RETRY_SECS.into_iter().chain(tail)
}
const STATE_RESPONSE_COOLDOWN_SECS: u64 = 15;
fn jittered_secs(secs: u64) -> std::time::Duration {
let factor = 0.8 + rand::random::<f64>() * 0.4;
std::time::Duration::from_secs_f64(secs as f64 * factor)
}
#[derive(Debug, Serialize, Deserialize)]
enum KvSyncMessage {
StateRequest { requester: PeerId },
OwnerAnnounce {
owner: AgentId,
policy: AccessPolicy,
policy_version: u64,
},
StateServed {
responder: PeerId,
empty: bool,
checkpoint_seq: Option<u64>,
},
StateServedV2 {
responder: PeerId,
digest: [u8; 32],
entry_count: u32,
},
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
struct ServedState {
digest: [u8; 32],
entry_count: u32,
}
#[derive(Debug, Default, Clone)]
struct ServedEvidence {
saw_nonempty: bool,
saw_owner_empty: bool,
max_checkpoint_seq: u64,
digests: std::collections::HashMap<PeerId, ServedState>,
}
struct BootstrapGuard(std::sync::Arc<std::sync::atomic::AtomicBool>);
impl Drop for BootstrapGuard {
fn drop(&mut self) {
self.0.store(false, std::sync::atomic::Ordering::Relaxed);
}
}
fn bootstrap_converged(
ev: &ServedEvidence,
is_empty: bool,
highest_checkpoint_seq: u64,
local_digest: [u8; 32],
) -> bool {
if ev.max_checkpoint_seq > 0 {
return highest_checkpoint_seq >= ev.max_checkpoint_seq;
}
if !ev.digests.is_empty() {
let any_nonempty = ev.digests.values().any(|d| d.entry_count > 0);
return ev
.digests
.values()
.any(|d| d.digest == local_digest && (d.entry_count > 0 || !any_nonempty));
}
if ev.saw_nonempty {
!is_empty
} else {
ev.saw_owner_empty
}
}
pub struct KvStoreSync {
store: Arc<RwLock<KvStore>>,
pubsub: Arc<PubSubManager>,
topic: String,
local_peer_id: PeerId,
local_agent_id: Option<AgentId>,
persist: std::sync::Mutex<Option<Arc<PersistCtx>>>,
stopped: Arc<std::sync::atomic::AtomicBool>,
cancel: tokio_util::sync::CancellationToken,
}
impl Drop for KvStoreSync {
fn drop(&mut self) {
self.cancel.cancel();
}
}
struct PersistCtx {
path: PathBuf,
gate: tokio::sync::Mutex<Option<u64>>,
degraded: std::sync::atomic::AtomicBool,
}
impl KvStoreSync {
pub fn new(
store: KvStore,
pubsub: Arc<PubSubManager>,
topic: String,
local_peer_id: PeerId,
local_agent_id: Option<AgentId>,
) -> Result<Self> {
let store = Arc::new(RwLock::new(store));
Ok(Self {
store,
pubsub,
topic,
local_peer_id,
local_agent_id,
persist: std::sync::Mutex::new(None),
stopped: Arc::new(std::sync::atomic::AtomicBool::new(false)),
cancel: tokio_util::sync::CancellationToken::new(),
})
}
pub fn set_persist_path(&self, path: PathBuf) {
if let Ok(mut guard) = self.persist.lock() {
*guard = Some(Arc::new(PersistCtx {
path,
gate: tokio::sync::Mutex::new(None),
degraded: std::sync::atomic::AtomicBool::new(false),
}));
}
}
fn persist_ctx(&self) -> Option<Arc<PersistCtx>> {
self.persist.lock().ok().and_then(|g| g.clone())
}
pub async fn persist(&self) -> Result<()> {
match self.persist_ctx() {
Some(ctx) => persist_snapshot(&self.store, &ctx).await,
None => Ok(()),
}
}
pub fn durability_degraded(&self) -> bool {
self.persist_ctx()
.is_some_and(|c| c.degraded.load(std::sync::atomic::Ordering::Relaxed))
}
pub async fn ensure_durable(&self) -> Result<()> {
match self.persist_ctx() {
Some(ctx) if ctx.degraded.load(std::sync::atomic::Ordering::Relaxed) => {
persist_snapshot(&self.store, &ctx).await
}
_ => Ok(()),
}
}
fn state_sync_topic(&self) -> String {
format!("{}{}", self.topic, STATE_SYNC_TOPIC_SUFFIX)
}
pub async fn start(&self) -> Result<()> {
self.start_with_spawner(|fut| {
tokio::spawn(fut);
})
.await
}
pub async fn start_with_spawner<S>(&self, spawn: S) -> Result<()>
where
S: Fn(std::pin::Pin<Box<dyn std::future::Future<Output = ()> + Send + 'static>>)
+ Send
+ Sync,
{
let mut sub = self.pubsub.subscribe(self.topic.clone()).await;
let store = Arc::clone(&self.store);
let bootstrap_needed = {
let s = store.read().await;
let local_is_owner =
self.local_agent_id.is_some() && s.owner() == self.local_agent_id.as_ref();
!local_is_owner || s.is_empty()
};
let main_topic = self.topic.clone();
let persist_ctx = self.persist_ctx();
let loop_persist_ctx = persist_ctx.clone();
let listener_cancel = self.cancel.clone();
let served_evidence = Arc::new(std::sync::Mutex::new(ServedEvidence::default()));
let bootstrap_active = Arc::new(std::sync::atomic::AtomicBool::new(bootstrap_needed));
let listener_served = Arc::clone(&served_evidence);
let listener_bootstrap_active = Arc::clone(&bootstrap_active);
spawn(Box::pin(async move {
loop {
let msg = tokio::select! {
() = listener_cancel.cancelled() => return,
msg = sub.recv() => msg,
};
let Some(msg) = msg else {
listener_cancel.cancel();
return;
};
if msg.topic != main_topic {
continue;
}
let decoded = decode_delta::<KvStoreDelta>(&msg.payload);
match decoded {
Ok((peer_id, delta)) => {
let merged = {
let mut s = store.write().await;
let writer = msg.sender.as_ref();
match s.merge_delta(&delta, peer_id, writer) {
Ok(()) => {
if listener_bootstrap_active
.load(std::sync::atomic::Ordering::Relaxed)
{
let declared = listener_served
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.digests
.get(&peer_id)
.copied();
if let Some(declared) = declared {
let sole_author =
matches!(s.policy(), AccessPolicy::Signed)
&& writer.is_some()
&& s.owner() == writer;
if sole_author
&& delta.added.len()
== declared.entry_count as usize
&& delta
.served_digest(s.id())
.is_some_and(|dg| dg == declared.digest)
{
let pruned = s.prune_to_served_set(&delta);
if pruned > 0 {
tracing::info!(
"pruned {pruned} stale key(s) after \
digest-verified full serve for store {}",
s.id()
);
}
}
}
}
true
}
Err(e) => {
tracing::warn!("Failed to merge KvStore delta: {e}");
false
}
}
};
if merged {
if let Some(ctx) = loop_persist_ctx.as_ref() {
let _ = persist_snapshot(&store, ctx).await;
}
}
}
Err(e) => {
tracing::warn!("Failed to deserialize KvStore delta: {e}");
}
}
}
}));
let mut sync_sub = self.pubsub.subscribe(self.state_sync_topic()).await;
let responder_store = Arc::clone(&self.store);
let responder_persist_ctx = persist_ctx.clone();
let responder_pubsub = Arc::clone(&self.pubsub);
let responder_topic = self.topic.clone();
let sync_topic = self.state_sync_topic();
let local_peer_id = self.local_peer_id;
let local_agent_id = self.local_agent_id;
let responder_served = Arc::clone(&served_evidence);
let responder_cancel = self.cancel.clone();
spawn(Box::pin(async move {
let mut last_full_response: Option<tokio::time::Instant> = None;
loop {
let msg = tokio::select! {
() = responder_cancel.cancelled() => return,
msg = sync_sub.recv() => msg,
};
let Some(msg) = msg else {
responder_cancel.cancel();
return;
};
if msg.topic != sync_topic {
continue;
}
let Ok(sync_msg) = bincode::deserialize::<KvSyncMessage>(&msg.payload) else {
continue;
};
match sync_msg {
KvSyncMessage::StateRequest { requester } => {
if requester == local_peer_id {
continue;
}
let announce = {
let s = responder_store.read().await;
match (local_agent_id, s.owner()) {
(Some(me), Some(owner)) if me == *owner => {
Some(KvSyncMessage::OwnerAnnounce {
owner: me,
policy: s.policy().clone(),
policy_version: s.policy_version(),
})
}
_ => None,
}
};
if let Some(announce) = announce {
match bincode::serialize(&announce) {
Ok(serialized) => {
if let Err(e) = responder_pubsub
.publish(sync_topic.clone(), bytes::Bytes::from(serialized))
.await
{
tracing::warn!(
"KvStore owner-announce publish failed: {e}"
);
}
}
Err(e) => {
tracing::warn!("KvStore owner-announce serialize failed: {e}");
}
}
}
let cooled_down = last_full_response.is_some_and(|t| {
t.elapsed()
< std::time::Duration::from_secs(STATE_RESPONSE_COOLDOWN_SECS)
});
let (full, is_empty, is_owner, checkpoint_seq, has_payload, served) = {
let s = responder_store.read().await;
let is_owner =
local_agent_id.is_some() && s.owner() == local_agent_id.as_ref();
let cp =
(s.highest_checkpoint_seq > 0).then_some(s.highest_checkpoint_seq);
let has_payload = !s.is_empty() || s.latest_checkpoint.is_some();
let full = (has_payload && !cooled_down).then(|| s.full_delta());
let served = (s.served_digest(), s.checkpoint_pairs().len() as u32);
(full, s.is_empty(), is_owner, cp, has_payload, served)
};
let mut markers: Vec<KvSyncMessage> = Vec::new();
if let Some(full) = full {
if let Ok(serialized) = encode_delta(local_peer_id, &full) {
if let Err(e) = responder_pubsub
.publish(
responder_topic.clone(),
bytes::Bytes::from(serialized),
)
.await
{
tracing::warn!("KvStore state-response publish failed: {e}");
} else {
last_full_response = Some(tokio::time::Instant::now());
markers.push(KvSyncMessage::StateServed {
responder: local_peer_id,
empty: is_empty,
checkpoint_seq,
});
markers.push(KvSyncMessage::StateServedV2 {
responder: local_peer_id,
digest: served.0,
entry_count: served.1,
});
}
}
} else if !has_payload {
markers.push(KvSyncMessage::StateServedV2 {
responder: local_peer_id,
digest: served.0,
entry_count: 0,
});
if is_owner {
markers.push(KvSyncMessage::StateServed {
responder: local_peer_id,
empty: true,
checkpoint_seq: None,
});
}
}
for marker in markers {
match bincode::serialize(&marker) {
Ok(serialized) => {
if let Err(e) = responder_pubsub
.publish(sync_topic.clone(), bytes::Bytes::from(serialized))
.await
{
tracing::warn!(
"KvStore state-served marker publish failed: {e}"
);
}
}
Err(e) => {
tracing::warn!(
"KvStore state-served marker serialize failed: {e}"
);
}
}
}
}
KvSyncMessage::StateServed {
responder,
empty,
checkpoint_seq,
} => {
if responder == local_peer_id {
continue; }
let owner_verified = {
let anchored = responder_store.read().await.owner().copied();
anchored.is_some() && msg.sender.as_ref() == anchored.as_ref()
};
let mut ev = responder_served
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
if empty {
if owner_verified {
ev.saw_owner_empty = true;
}
} else {
ev.saw_nonempty = true;
}
if let Some(seq) = checkpoint_seq.filter(|_| owner_verified) {
ev.max_checkpoint_seq = ev.max_checkpoint_seq.max(seq);
}
}
KvSyncMessage::StateServedV2 {
responder,
digest,
entry_count,
} => {
if responder == local_peer_id {
continue; }
responder_served
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.digests
.insert(
responder,
ServedState {
digest,
entry_count,
},
);
}
KvSyncMessage::OwnerAnnounce {
owner,
policy,
policy_version,
} => {
let Some(sender) = msg.sender else {
tracing::warn!(
"ignoring unsigned KvStore ownership announcement on {}",
msg.topic
);
continue;
};
if local_agent_id.is_some_and(|me| me == sender) {
continue; }
let learned = {
let mut s = responder_store.write().await;
match s.learn_ownership(owner, policy, policy_version, &sender) {
Ok(()) => {
tracing::info!(
"KvStore {} processed owner announce from {} (policy {}, version {})",
s.id(),
hex::encode(owner.as_bytes()),
s.policy(),
s.policy_version()
);
true
}
Err(e) => {
tracing::warn!(
"rejected KvStore ownership announcement from {}: {e}",
hex::encode(sender.as_bytes())
);
false
}
}
};
if learned {
if let Some(ctx) = responder_persist_ctx.as_ref() {
let _ = persist_snapshot(&responder_store, ctx).await;
}
}
}
}
}
}));
if bootstrap_needed {
let requester_pubsub = Arc::clone(&self.pubsub);
let sync_topic = self.state_sync_topic();
let requester_store = Arc::downgrade(&self.store);
let stopped = Arc::clone(&self.stopped);
let requester_cancel = self.cancel.clone();
let requester_served = Arc::clone(&served_evidence);
let requester_bootstrap_active = Arc::clone(&bootstrap_active);
spawn(Box::pin(async move {
let _guard = BootstrapGuard(requester_bootstrap_active);
for (attempt, delay_secs) in state_request_delays().enumerate() {
tokio::select! {
() = requester_cancel.cancelled() => return,
() = tokio::time::sleep(jittered_secs(delay_secs)) => {}
}
if stopped.load(std::sync::atomic::Ordering::Relaxed) {
return; }
if attempt >= STATE_REQUEST_RETRY_SECS.len() {
let Some(store) = requester_store.upgrade() else {
return; };
let ev = requester_served
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone();
let (is_empty, cp_hwm, local_digest) = {
let s = store.read().await;
(s.is_empty(), s.highest_checkpoint_seq, s.served_digest())
};
if bootstrap_converged(&ev, is_empty, cp_hwm, local_digest) {
return; }
}
let request = KvSyncMessage::StateRequest {
requester: local_peer_id,
};
let Ok(serialized) = bincode::serialize(&request) else {
return;
};
if let Err(e) = requester_pubsub
.publish(sync_topic.clone(), bytes::Bytes::from(serialized))
.await
{
tracing::debug!("KvStore state-request publish failed: {e}");
}
}
}));
}
Ok(())
}
pub fn silence_bootstrap(&self) {
self.stopped
.store(true, std::sync::atomic::Ordering::Relaxed);
}
pub fn cancel_sync(&self) {
self.cancel.cancel();
}
pub async fn stop(&self) -> Result<()> {
self.cancel_sync();
self.pubsub.unsubscribe(&self.topic).await;
self.pubsub.unsubscribe(&self.state_sync_topic()).await;
Ok(())
}
pub async fn publish_delta(&self, local_peer_id: PeerId, delta: KvStoreDelta) -> Result<()> {
let serialized = encode_delta(local_peer_id, &delta)
.map_err(|e| crate::kv::KvError::Gossip(format!("serialize delta failed: {e}")))?;
self.pubsub
.publish(self.topic.clone(), bytes::Bytes::from(serialized))
.await
.map_err(|e| crate::kv::KvError::Gossip(format!("publish delta failed: {e}")))?;
Ok(())
}
pub async fn read(&self) -> tokio::sync::RwLockReadGuard<'_, KvStore> {
self.store.read().await
}
pub async fn write(&self) -> tokio::sync::RwLockWriteGuard<'_, KvStore> {
self.store.write().await
}
#[must_use]
pub fn topic(&self) -> &str {
&self.topic
}
}
static SNAPSHOT_TMP_COUNTER: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
const SNAPSHOT_MAGIC: &[u8; 8] = b"X0XKVS1\0";
#[derive(Deserialize)]
struct SnapshotBody {
store: KvStore,
seq_counter: u64,
}
#[derive(Serialize)]
struct SnapshotBodyRef<'a> {
store: &'a KvStore,
seq_counter: u64,
}
fn encode_snapshot(store: &KvStore) -> Result<Vec<u8>> {
let body = SnapshotBodyRef {
store,
seq_counter: store.seq_counter_value(),
};
let mut out = Vec::with_capacity(256);
out.extend_from_slice(SNAPSHOT_MAGIC);
out.extend_from_slice(&bincode::serialize(&body)?);
Ok(out)
}
async fn persist_snapshot(store: &Arc<RwLock<KvStore>>, ctx: &PersistCtx) -> Result<()> {
let result = async {
let mut last = ctx.gate.lock().await;
let (version, bytes) = {
let s = store.read().await;
(s.current_version(), encode_snapshot(&s)?)
};
if last.is_some_and(|l| l >= version) {
return Ok(());
}
write_snapshot_atomic(&ctx.path, &bytes)?;
*last = Some(version);
Ok(())
}
.await;
ctx.degraded
.store(result.is_err(), std::sync::atomic::Ordering::Relaxed);
if let Err(e) = &result {
tracing::error!(
"kv snapshot persist failed for {}: {e} — store is durability-degraded; \
local writes are refused until a snapshot succeeds",
ctx.path.display()
);
}
result
}
fn write_snapshot_atomic(path: &Path, bytes: &[u8]) -> std::io::Result<()> {
if let Some(parent) = path.parent() {
std::fs::create_dir_all(parent)?;
}
let n = SNAPSHOT_TMP_COUNTER.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
let tmp = path.with_extension(format!("tmp.{}.{n}", std::process::id()));
{
use std::io::Write;
let mut f = std::fs::File::create(&tmp)?;
f.write_all(bytes)?;
f.sync_all()?;
}
if let Err(e) = std::fs::rename(&tmp, path) {
let _ = std::fs::remove_file(&tmp);
return Err(e);
}
#[cfg(unix)]
if let Some(parent) = path.parent() {
std::fs::File::open(parent)?.sync_all()?;
}
Ok(())
}
pub fn load_snapshot(path: &Path) -> Result<Option<KvStore>> {
let bytes = match std::fs::read(path) {
Ok(b) => b,
Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(None),
Err(e) => return Err(e.into()),
};
let Some(body_bytes) = bytes.strip_prefix(SNAPSHOT_MAGIC.as_slice()) else {
return Err(std::io::Error::other(
"unrecognized kv snapshot format (missing v1 magic) — corrupt or foreign file; \
refusing to start with amnesia",
)
.into());
};
let body: SnapshotBody = bincode::deserialize(body_bytes)?;
let store = body.store;
store.restore_seq_counter(body.seq_counter.max(store.current_version()));
Ok(Some(store))
}
#[cfg(test)]
mod tests {
use super::*;
use crate::identity::AgentId;
use crate::kv::store::AccessPolicy;
use crate::kv::{KvEntry, KvStoreId};
use crate::network::{NetworkConfig, NetworkNode};
use std::time::Duration;
fn agent(n: u8) -> AgentId {
AgentId([n; 32])
}
fn peer(n: u8) -> PeerId {
PeerId::new([n; 32])
}
fn store_id(n: u8) -> KvStoreId {
KvStoreId::new([n; 32])
}
#[test]
fn snapshot_roundtrip_missing_and_corrupt() {
let dir = tempfile::tempdir().expect("tmpdir");
let path = dir.path().join("kv").join("snap.bin");
assert!(
matches!(load_snapshot(&path), Ok(None)),
"missing snapshot is a clean first run"
);
let mut store = KvStore::new(
store_id(7),
"log".to_string(),
agent(1),
AccessPolicy::AppendOnly,
);
store
.put(
"k1".to_string(),
b"v1".to_vec(),
"text/plain".to_string(),
peer(1),
)
.expect("put");
store.highest_checkpoint_seq = 5;
let _ = store.next_seq();
let _ = store.next_seq();
let counter_before = store.seq_counter_value();
let bytes = encode_snapshot(&store).expect("encode");
write_snapshot_atomic(&path, &bytes).expect("atomic write");
let restored = load_snapshot(&path)
.expect("load ok")
.expect("snapshot present");
assert_eq!(*restored.policy(), AccessPolicy::AppendOnly);
assert_eq!(
restored.get("k1").map(|e| e.value.clone()),
Some(b"v1".to_vec())
);
assert_eq!(restored.highest_checkpoint_seq, 5);
assert!(
restored.next_seq() > counter_before,
"restored seq counter must exceed every pre-restart seq"
);
std::fs::write(&path, bincode::serialize(&store).expect("serialize")).expect("write bare");
assert!(
load_snapshot(&path).is_err(),
"missing-magic snapshot must be an error (fail closed)"
);
std::fs::write(&path, b"not a snapshot").expect("corrupt");
assert!(
load_snapshot(&path).is_err(),
"corrupt snapshot must be an error (fail closed), not a silent fresh start"
);
let mut evil = SNAPSHOT_MAGIC.to_vec();
evil.extend_from_slice(b"\x01\x02\x03");
std::fs::write(&path, evil).expect("write garbage body");
assert!(
load_snapshot(&path).is_err(),
"garbage body must be an error (fail closed)"
);
}
async fn make_node() -> Arc<NetworkNode> {
Arc::new(
NetworkNode::new(NetworkConfig::default(), None, None)
.await
.expect("network node"),
)
}
async fn make_sync(topic: &str, policy: AccessPolicy) -> KvStoreSync {
let node = make_node().await;
let pubsub = Arc::new(PubSubManager::new(node, None).expect("pubsub"));
let store = KvStore::new(store_id(1), "Test".to_string(), agent(1), policy);
KvStoreSync::new(store, pubsub, topic.to_string(), peer(1), Some(agent(1)))
.expect("kv sync")
}
async fn make_sync_with_pubsub(
topic: &str,
policy: AccessPolicy,
) -> (KvStoreSync, Arc<PubSubManager>) {
let node = make_node().await;
let pubsub = Arc::new(PubSubManager::new(node, None).expect("pubsub"));
let store = KvStore::new(store_id(1), "Test".to_string(), agent(1), policy);
let sync = KvStoreSync::new(
store,
Arc::clone(&pubsub),
topic.to_string(),
peer(1),
Some(agent(1)),
)
.expect("kv sync");
(sync, pubsub)
}
#[tokio::test]
async fn test_kv_store_sync_creation() {
let owner = agent(1);
let store = KvStore::new(store_id(1), "Test".to_string(), owner, AccessPolicy::Signed);
let _store_for_sync = store;
}
#[tokio::test]
async fn test_apply_delta_directly() {
let owner = agent(1);
let writer = agent(2);
let p2 = peer(2);
let mut store = KvStore::new(
store_id(1),
"Test".to_string(),
owner,
AccessPolicy::Allowlisted,
);
store.allow_writer(writer, &owner).expect("allow");
let store_arc = Arc::new(RwLock::new(store));
let entry = KvEntry::new(
"newkey".to_string(),
b"value".to_vec(),
"text/plain".to_string(),
);
let mut delta = KvStoreDelta::new(1);
delta.added.insert("newkey".to_string(), (entry, (p2, 1)));
{
let mut s = store_arc.write().await;
s.merge_delta(&delta, p2, Some(&writer)).expect("merge");
}
{
let s = store_arc.read().await;
assert!(s.get("newkey").is_some());
}
}
#[tokio::test]
async fn test_concurrent_reads() {
let owner = agent(1);
let store = KvStore::new(store_id(1), "Test".to_string(), owner, AccessPolicy::Signed);
let store_arc = Arc::new(RwLock::new(store));
let s1 = store_arc.read().await;
let s2 = store_arc.read().await;
assert_eq!(s1.name(), "Test");
assert_eq!(s2.name(), "Test");
}
#[tokio::test]
async fn new_sets_topic_and_yields_accessible_guards() {
let sync = make_sync("store/A", AccessPolicy::Signed).await;
assert_eq!(sync.topic(), "store/A");
{
let s = sync.read().await;
assert_eq!(s.name(), "Test");
assert!(s.is_empty());
}
let owner = agent(1);
let entry = KvEntry::new(
"owner-key".to_string(),
b"v".to_vec(),
"text/plain".to_string(),
);
let mut delta = KvStoreDelta::new(1);
delta
.added
.insert("owner-key".to_string(), (entry, (peer(1), 1)));
{
let mut s = sync.write().await;
s.merge_delta(&delta, peer(1), Some(&owner))
.expect("owner merge");
}
let s = sync.read().await;
assert!(s.get("owner-key").is_some(), "owner write must be visible");
}
#[tokio::test]
async fn state_sync_topic_appends_side_channel_suffix() {
let sync = make_sync("store/B", AccessPolicy::Signed).await;
assert_eq!(sync.state_sync_topic(), "store/B/state-sync");
let sync2 = make_sync("store/B/nested", AccessPolicy::Signed).await;
assert_eq!(sync2.state_sync_topic(), "store/B/nested/state-sync");
}
#[tokio::test]
async fn publish_delta_delivers_encoded_pair_to_subscriber() {
let (sync, pubsub) = make_sync_with_pubsub("store/C", AccessPolicy::Signed).await;
let mut sub = pubsub.subscribe("store/C".to_string()).await;
let sender = peer(7);
let entry = KvEntry::new(
"remote".to_string(),
b"payload".to_vec(),
"application/octet-stream".to_string(),
);
let mut delta = KvStoreDelta::new(9);
delta
.added
.insert("remote".to_string(), (entry, (sender, 3)));
sync.publish_delta(sender, delta)
.await
.expect("publish_delta");
let msg = tokio::time::timeout(Duration::from_secs(2), sub.recv())
.await
.expect("timed out waiting for published delta")
.expect("subscriber stream closed");
let (observed_sender, observed_delta) =
decode_delta::<KvStoreDelta>(&msg.payload).expect("wire decode");
assert_eq!(observed_sender, sender);
assert_eq!(observed_delta.version, 9);
assert!(observed_delta.added.contains_key("remote"));
assert_eq!(msg.topic, "store/C");
let reencoded = encode_delta(sender, &observed_delta).expect("re-encode");
let (s2, d2) = decode_delta::<KvStoreDelta>(&reencoded).expect("re-decode");
assert_eq!(s2, sender);
assert_eq!(d2.version, 9);
}
#[tokio::test]
async fn start_with_spawner_subscribes_and_returns_ok() {
let sync = make_sync("store/D", AccessPolicy::Signed).await;
sync.start_with_spawner(|_fut| {
})
.await
.expect("start_with_spawner");
}
#[tokio::test]
async fn start_default_spawner_merges_remote_delta() {
let sync = make_sync(
"store/E",
AccessPolicy::Encrypted {
group_id: vec![1, 2, 3],
},
)
.await;
sync.start().await.expect("start");
tokio::time::sleep(Duration::from_millis(100)).await;
let entry = KvEntry::new(
"merged-key".to_string(),
b"hello".to_vec(),
"text/plain".to_string(),
);
let mut delta = KvStoreDelta::new(1);
delta
.added
.insert("merged-key".to_string(), (entry, (peer(2), 1)));
sync.publish_delta(peer(2), delta).await.expect("publish");
let landed = tokio::time::timeout(Duration::from_secs(2), async {
loop {
let present = {
let s = sync.read().await;
s.get("merged-key").is_some()
};
if present {
return;
}
tokio::time::sleep(Duration::from_millis(25)).await;
}
})
.await;
assert!(
landed.is_ok(),
"remote delta was not merged by start() loop"
);
}
#[tokio::test]
async fn stop_returns_ok_and_is_idempotent() {
let sync = make_sync("store/F", AccessPolicy::Signed).await;
sync.stop().await.expect("first stop");
sync.stop().await.expect("second stop (idempotent)");
}
#[tokio::test]
async fn stop_and_silence_arm_their_kill_switches() {
let sync = make_sync("store/stopflag", AccessPolicy::Signed).await;
assert!(
!sync.cancel.is_cancelled() && !sync.stopped.load(std::sync::atomic::Ordering::Relaxed),
"both switches must start disarmed"
);
sync.silence_bootstrap();
assert!(
sync.stopped.load(std::sync::atomic::Ordering::Relaxed),
"silence_bootstrap() arms the requester-only flag"
);
assert!(
!sync.cancel.is_cancelled(),
"silence_bootstrap() must NOT cancel the listener/responder loops"
);
sync.stop().await.expect("stop");
assert!(
sync.cancel.is_cancelled(),
"stop() must cancel ALL background loops via the token"
);
}
#[tokio::test]
async fn dropping_the_sync_cancels_all_loops() {
let sync = make_sync("store/dropcancel", AccessPolicy::Signed).await;
let token = sync.cancel.clone();
assert!(!token.is_cancelled());
drop(sync);
assert!(
token.is_cancelled(),
"dropping the last sync reference must cancel every loop"
);
}
#[test]
fn state_request_schedule_never_terminates_while_unconverged() {
let front: Vec<u64> = state_request_delays().take(4).collect();
assert_eq!(front, STATE_REQUEST_RETRY_SECS, "front burst unchanged");
let tail: Vec<u64> = state_request_delays().skip(4).take(8).collect();
assert_eq!(
tail,
[30, 60, 120, 240, 300, 300, 300, 300],
"tail doubles to the cap, then holds it"
);
assert_eq!(
state_request_delays().nth(10_000),
Some(STATE_REQUEST_TAIL_CAP_SECS),
"the schedule is infinite — convergence, not the schedule, \
is what ends the requester"
);
}
#[test]
fn bootstrap_convergence_rule() {
const D1: [u8; 32] = [1u8; 32];
const D2: [u8; 32] = [2u8; 32];
const D_EMPTY: [u8; 32] = [9u8; 32];
let none = ServedEvidence::default();
assert!(!bootstrap_converged(&none, true, 0, D1));
assert!(!bootstrap_converged(&none, false, 9, D1));
let cp = ServedEvidence {
saw_nonempty: true,
saw_owner_empty: false,
max_checkpoint_seq: 5,
digests: std::collections::HashMap::new(),
};
assert!(
!bootstrap_converged(&cp, false, 4, D1),
"behind the served HWM"
);
assert!(bootstrap_converged(&cp, false, 5, D1));
assert!(bootstrap_converged(&cp, false, 7, D1));
let nonempty = ServedEvidence {
saw_nonempty: true,
saw_owner_empty: false,
max_checkpoint_seq: 0,
digests: std::collections::HashMap::new(),
};
assert!(!bootstrap_converged(&nonempty, true, 0, D1));
assert!(bootstrap_converged(&nonempty, false, 0, D1));
let owner_empty = ServedEvidence {
saw_nonempty: false,
saw_owner_empty: true,
max_checkpoint_seq: 0,
digests: std::collections::HashMap::new(),
};
assert!(bootstrap_converged(&owner_empty, true, 0, D1));
let mixed = ServedEvidence {
saw_nonempty: true,
saw_owner_empty: true,
max_checkpoint_seq: 0,
digests: std::collections::HashMap::new(),
};
assert!(
!bootstrap_converged(&mixed, true, 0, D1),
"divergent holders: the data-bearing claim must win, keep asking"
);
let digest_ev = |decls: &[(&[u8; 32], u32)]| ServedEvidence {
saw_nonempty: false,
saw_owner_empty: false,
max_checkpoint_seq: 0,
digests: decls
.iter()
.enumerate()
.map(|(i, (d, c))| {
(
peer(i as u8 + 10),
ServedState {
digest: **d,
entry_count: *c,
},
)
})
.collect(),
};
let ev = digest_ev(&[(&D1, 3)]);
assert!(bootstrap_converged(&ev, false, 0, D1));
assert!(!bootstrap_converged(&ev, false, 0, D2));
let mut ev_v1_too = digest_ev(&[(&D1, 3)]);
ev_v1_too.saw_nonempty = true;
assert!(
!bootstrap_converged(&ev_v1_too, false, 0, D2),
"v2 declarations outrank weak v1 evidence"
);
let ev_empty = digest_ev(&[(&D_EMPTY, 0)]);
assert!(bootstrap_converged(&ev_empty, true, 0, D_EMPTY));
let ev_divergent = digest_ev(&[(&D_EMPTY, 0), (&D1, 2)]);
assert!(
!bootstrap_converged(&ev_divergent, true, 0, D_EMPTY),
"a data-bearing declaration outranks the empty one — keep asking"
);
assert!(bootstrap_converged(&ev_divergent, false, 0, D1));
}
#[tokio::test(start_paused = true)]
async fn restored_replica_keeps_asking_until_a_holder_serves() {
let node = make_node().await;
let kp = crate::identity::AgentKeypair::generate().expect("keypair");
let owner_id = kp.agent_id();
let ctx = Arc::new(crate::gossip::SigningContext::from_keypair(&kp));
let pubsub = Arc::new(PubSubManager::new(node, Some(ctx)).expect("pubsub"));
let topic = "kv-238-restored-late-owner";
let mut replica = KvStore::new_replica(
store_id(1),
String::new(),
Some(owner_id),
crate::kv::store::AnchorChannel::Persistence,
);
replica
.put(
"k_old".to_string(),
b"v_old".to_vec(),
"text/plain".to_string(),
peer(2),
)
.expect("seed restored key");
let joiner = KvStoreSync::new(
replica,
Arc::clone(&pubsub),
topic.to_string(),
peer(2),
Some(agent(2)),
)
.expect("joiner sync");
joiner.start().await.expect("start joiner");
tokio::time::sleep(Duration::from_secs(90)).await;
let mut owned = KvStore::new(
store_id(1),
"log".to_string(),
owner_id,
AccessPolicy::Signed,
);
owned
.put(
"k_old".to_string(),
b"v_old".to_vec(),
"text/plain".to_string(),
peer(1),
)
.expect("owner k_old");
owned
.put(
"k_new".to_string(),
b"v_new".to_vec(),
"text/plain".to_string(),
peer(1),
)
.expect("owner k_new");
let owner_sync = KvStoreSync::new(
owned,
Arc::clone(&pubsub),
topic.to_string(),
peer(1),
Some(owner_id),
)
.expect("owner sync");
owner_sync.start().await.expect("start owner");
let mut recovered = false;
for _ in 0..200 {
tokio::time::sleep(Duration::from_secs(5)).await;
if joiner.read().await.get("k_new").is_some() {
recovered = true;
break;
}
}
assert!(
recovered,
"a restored non-empty replica must keep requesting until a \
holder serves it — non-emptiness alone is not convergence"
);
}
#[tokio::test(start_paused = true)]
async fn checkpointed_empty_owner_deletes_stale_replica_state() {
let node = make_node().await;
let kp = crate::identity::AgentKeypair::generate().expect("keypair");
let owner_id = kp.agent_id();
let ctx = Arc::new(crate::gossip::SigningContext::from_keypair(&kp));
let pubsub = Arc::new(PubSubManager::new(node, Some(ctx)).expect("pubsub"));
let topic = "kv-238-deleted-to-empty";
let mut owned = KvStore::new(
store_id(1),
"log".to_string(),
owner_id,
AccessPolicy::Signed,
);
let (pub_bytes, sec_bytes) = kp.to_bytes();
let public_key =
ant_quic::MlDsaPublicKey::from_bytes(&pub_bytes).expect("public key bytes");
let secret_key =
ant_quic::MlDsaSecretKey::from_bytes(&sec_bytes).expect("secret key bytes");
let pairs = owned.checkpoint_pairs();
let root = crate::kv::store::content_root(owned.id(), owned.name(), &pairs);
let cp = crate::kv::store::make_owner_checkpoint(crate::kv::store::OwnerCheckpointParams {
topic,
store_id: &store_id(1),
secret_key: &secret_key,
public_key: &public_key,
policy: &AccessPolicy::Signed,
policy_version: owned.policy_version(),
checkpoint_seq: 3,
content_root: root,
timestamp: 1,
})
.expect("sign empty checkpoint");
owned.latest_checkpoint = Some(cp);
owned.highest_checkpoint_seq = 3;
let mut replica = KvStore::new_replica(
store_id(1),
String::new(),
Some(owner_id),
crate::kv::store::AnchorChannel::Persistence,
);
replica
.put(
"k_stale".to_string(),
b"obsolete".to_vec(),
"text/plain".to_string(),
peer(2),
)
.expect("seed stale key");
let joiner = KvStoreSync::new(
replica,
Arc::clone(&pubsub),
topic.to_string(),
peer(2),
Some(agent(2)),
)
.expect("joiner sync");
joiner.start().await.expect("start joiner");
let owner_sync = KvStoreSync::new(
owned,
Arc::clone(&pubsub),
topic.to_string(),
peer(1),
Some(owner_id),
)
.expect("owner sync");
owner_sync.silence_bootstrap();
owner_sync.start().await.expect("start owner");
let mut cleaned = false;
for _ in 0..400 {
tokio::time::sleep(Duration::from_secs(2)).await;
let s = joiner.read().await;
if s.get("k_stale").is_none() && s.highest_checkpoint_seq == 3 {
cleaned = true;
break;
}
}
assert!(
cleaned,
"the checkpointed empty owner's response must full-replace the \
stale replica (key removed, checkpoint HWM adopted)"
);
assert!(
joiner.read().await.is_empty(),
"replica converges to the owner's (empty) state"
);
}
#[tokio::test(start_paused = true)]
async fn restored_non_empty_replica_still_requests_missed_state() {
let node = make_node().await;
let kp = crate::identity::AgentKeypair::generate().expect("keypair");
let owner_id = kp.agent_id();
let ctx = Arc::new(crate::gossip::SigningContext::from_keypair(&kp));
let pubsub = Arc::new(PubSubManager::new(node, Some(ctx)).expect("pubsub"));
let topic = "kv-238-missed-delta";
let mut owned = KvStore::new(
store_id(1),
"log".to_string(),
owner_id,
AccessPolicy::Signed,
);
owned
.put(
"k_old".to_string(),
b"v_old".to_vec(),
"text/plain".to_string(),
peer(1),
)
.expect("owner put k_old");
owned
.put(
"k_new".to_string(),
b"v_new".to_vec(),
"text/plain".to_string(),
peer(1),
)
.expect("owner put k_new");
let owner_sync = KvStoreSync::new(
owned,
Arc::clone(&pubsub),
topic.to_string(),
peer(1),
Some(owner_id),
)
.expect("owner sync");
owner_sync.start().await.expect("start owner");
let mut replica = KvStore::new_replica(
store_id(1),
String::new(),
Some(owner_id),
crate::kv::store::AnchorChannel::Persistence,
);
replica
.put(
"k_old".to_string(),
b"v_old".to_vec(),
"text/plain".to_string(),
peer(2),
)
.expect("seed restored key");
let joiner = KvStoreSync::new(
replica,
Arc::clone(&pubsub),
topic.to_string(),
peer(2),
Some(agent(2)),
)
.expect("joiner sync");
joiner.start().await.expect("start joiner");
let mut recovered = false;
for _ in 0..30 {
tokio::time::sleep(Duration::from_secs(2)).await;
if joiner.read().await.get("k_new").is_some() {
recovered = true;
break;
}
}
assert!(
recovered,
"a restored non-empty non-owner replica must still request \
state and recover deltas it missed while offline"
);
}
#[tokio::test(start_paused = true)]
async fn requester_tail_recovers_when_owner_returns_after_front_schedule() {
let node = make_node().await;
let kp = crate::identity::AgentKeypair::generate().expect("keypair");
let owner_id = kp.agent_id();
let ctx = Arc::new(crate::gossip::SigningContext::from_keypair(&kp));
let pubsub = Arc::new(PubSubManager::new(node, Some(ctx)).expect("pubsub"));
let topic = "kv-238-zombie";
let replica = KvStore::new_replica(
store_id(1),
String::new(),
Some(owner_id),
crate::kv::store::AnchorChannel::RestParam,
);
let joiner = KvStoreSync::new(
replica,
Arc::clone(&pubsub),
topic.to_string(),
peer(2),
Some(agent(2)),
)
.expect("joiner sync");
joiner.start().await.expect("start joiner");
tokio::time::sleep(Duration::from_secs(60)).await;
assert!(
joiner.read().await.is_empty(),
"nobody was online to answer the front schedule"
);
assert_eq!(
*joiner.read().await.policy(),
AccessPolicy::Signed,
"replica still reports its construction-default policy while \
the owner is away (the transient misreport under test)"
);
let mut owned = KvStore::new(
store_id(1),
"log".to_string(),
owner_id,
AccessPolicy::AppendOnly,
);
owned
.put(
"k1".to_string(),
b"v1".to_vec(),
"text/plain".to_string(),
peer(1),
)
.expect("owner put");
let owner_sync = KvStoreSync::new(
owned,
Arc::clone(&pubsub),
topic.to_string(),
peer(1),
Some(owner_id),
)
.expect("owner sync");
owner_sync.start().await.expect("start owner");
let mut converged = false;
for _ in 0..200 {
tokio::time::sleep(Duration::from_secs(5)).await;
if !joiner.read().await.is_empty() {
converged = true;
break;
}
}
assert!(
converged,
"the tail requester must recover the store once the owner \
returns (zombie subscription, issue #238)"
);
assert_eq!(
joiner.read().await.get("k1").map(|e| e.value.clone()),
Some(b"v1".to_vec()),
"the owner's key must arrive via the state response"
);
let mut policy_ok = false;
for _ in 0..60 {
if *joiner.read().await.policy() == AccessPolicy::AppendOnly {
policy_ok = true;
break;
}
tokio::time::sleep(Duration::from_secs(1)).await;
}
assert!(
policy_ok,
"the owner announce must refresh the replica policy \
(transient `signed` misreport, issue #238)"
);
}
async fn drain_state_requests(probe: &mut crate::gossip::Subscription, from: PeerId) -> usize {
let mut n = 0;
while let Ok(Some(msg)) = tokio::time::timeout(Duration::from_millis(1), probe.recv()).await
{
if let Ok(KvSyncMessage::StateRequest { requester }) =
bincode::deserialize::<KvSyncMessage>(&msg.payload)
{
if requester == from {
n += 1;
}
}
}
n
}
#[test]
fn served_digest_is_deterministic_content_bound_and_name_independent() {
let owner = agent(1);
let mut a = KvStore::new(
store_id(1),
"alpha".to_string(),
owner,
AccessPolicy::Signed,
);
let mut b = KvStore::new(store_id(1), "beta".to_string(), owner, AccessPolicy::Signed);
assert_eq!(a.served_digest(), b.served_digest());
let c = KvStore::new(
store_id(2),
"alpha".to_string(),
owner,
AccessPolicy::Signed,
);
assert_ne!(
a.served_digest(),
c.served_digest(),
"the store id binds the digest (cross-store replay defense)"
);
let entry = KvEntry::new("k".to_string(), b"v".to_vec(), "text/plain".to_string());
let mut d1 = KvStoreDelta::new(1);
d1.added.insert("k".to_string(), (entry, (peer(1), 1)));
a.merge_delta(&d1, peer(1), Some(&owner)).expect("merge a");
b.merge_delta(&d1, peer(1), Some(&owner)).expect("merge b");
assert_eq!(a.served_digest(), b.served_digest());
assert_eq!(
d1.served_digest(&store_id(1)),
None,
"incremental shape must not impersonate a full serve"
);
d1.name_update = Some(a.name_register().clone());
assert_eq!(d1.served_digest(&store_id(1)), Some(a.served_digest()));
let entry2 = KvEntry::new("k2".to_string(), b"w".to_vec(), "text/plain".to_string());
let mut d2 = KvStoreDelta::new(2);
d2.added.insert("k2".to_string(), (entry2, (peer(1), 2)));
b.merge_delta(&d2, peer(1), Some(&owner)).expect("merge b2");
assert_ne!(a.served_digest(), b.served_digest());
}
#[tokio::test(start_paused = true)]
async fn lost_full_delta_with_surviving_marker_keeps_requester_asking() {
let node = make_node().await;
let kp = crate::identity::AgentKeypair::generate().expect("keypair");
let owner_id = kp.agent_id();
let ctx = Arc::new(crate::gossip::SigningContext::from_keypair(&kp));
let pubsub = Arc::new(PubSubManager::new(node, Some(ctx)).expect("pubsub"));
let topic = "kv-240-lost-broadcast";
let side = format!("{topic}{STATE_SYNC_TOPIC_SUFFIX}");
let mut owned = KvStore::new(
store_id(1),
"log".to_string(),
owner_id,
AccessPolicy::Signed,
);
owned
.put(
"k1".to_string(),
b"v1".to_vec(),
"text/plain".to_string(),
peer(1),
)
.expect("k1");
owned
.put(
"k2".to_string(),
b"v2".to_vec(),
"text/plain".to_string(),
peer(1),
)
.expect("k2");
let mut replica = KvStore::new_replica(
store_id(1),
String::new(),
Some(owner_id),
crate::kv::store::AnchorChannel::RestParam,
);
let k2_only = {
let full = owned.full_delta();
let (key, (entry, tag)) = full
.added
.iter()
.find(|(k, _)| k.as_str() == "k2")
.expect("k2 in full delta");
let mut d = KvStoreDelta::new(1);
d.added.insert(key.clone(), (entry.clone(), *tag));
d
};
replica
.merge_delta(&k2_only, peer(1), Some(&owner_id))
.expect("seed k2");
let joiner = KvStoreSync::new(
replica,
Arc::clone(&pubsub),
topic.to_string(),
peer(2),
Some(agent(2)),
)
.expect("joiner sync");
joiner.start().await.expect("start joiner");
let mut probe = pubsub.subscribe(side.clone()).await;
let marker = KvSyncMessage::StateServedV2 {
responder: peer(1),
digest: owned.served_digest(),
entry_count: 2,
};
let bytes = bincode::serialize(&marker).expect("serialize marker");
pubsub
.publish(side.clone(), bytes::Bytes::from(bytes))
.await
.expect("publish marker");
let mut requests = 0;
for _ in 0..10 {
tokio::time::sleep(Duration::from_secs(30)).await;
requests += drain_state_requests(&mut probe, peer(2)).await;
}
assert!(
requests > 0,
"a requester whose full delta was lost must keep asking (digest mismatch)"
);
assert!(
joiner.read().await.get("k1").is_none(),
"the lost broadcast never arrived"
);
let owner_sync = KvStoreSync::new(
owned,
Arc::clone(&pubsub),
topic.to_string(),
peer(1),
Some(owner_id),
)
.expect("owner sync");
owner_sync.start().await.expect("start owner");
let mut recovered = false;
for _ in 0..200 {
tokio::time::sleep(Duration::from_secs(5)).await;
if joiner.read().await.get("k1").is_some() {
recovered = true;
break;
}
}
assert!(
recovered,
"the still-alive requester must recover the lost key"
);
tokio::time::sleep(Duration::from_secs(160)).await;
drain_state_requests(&mut probe, peer(2)).await;
let mut relapse = 0;
for _ in 0..20 {
tokio::time::sleep(Duration::from_secs(30)).await;
relapse += drain_state_requests(&mut probe, peer(2)).await;
}
assert_eq!(relapse, 0, "a digest-matched requester must fall silent");
}
#[tokio::test(start_paused = true)]
async fn digest_verified_full_serve_prunes_stale_keys() {
let node = make_node().await;
let kp = crate::identity::AgentKeypair::generate().expect("keypair");
let owner_id = kp.agent_id();
let ctx = Arc::new(crate::gossip::SigningContext::from_keypair(&kp));
let pubsub = Arc::new(PubSubManager::new(node, Some(ctx)).expect("pubsub"));
let topic = "kv-240-prune-stale";
let side = format!("{topic}{STATE_SYNC_TOPIC_SUFFIX}");
let mut owned = KvStore::new(
store_id(1),
"log".to_string(),
owner_id,
AccessPolicy::Signed,
);
owned
.put(
"k_live".to_string(),
b"v".to_vec(),
"text/plain".to_string(),
peer(1),
)
.expect("owner put");
let mut replica = KvStore::new_replica(
store_id(1),
String::new(),
Some(owner_id),
crate::kv::store::AnchorChannel::Persistence,
);
let live_only = {
let full = owned.full_delta();
let (key, (entry, tag)) = full
.added
.iter()
.find(|(k, _)| k.as_str() == "k_live")
.expect("k_live in full delta");
let mut d = KvStoreDelta::new(1);
d.added.insert(key.clone(), (entry.clone(), *tag));
d
};
replica
.merge_delta(&live_only, peer(1), Some(&owner_id))
.expect("seed k_live");
replica
.put(
"k_stale".to_string(),
b"obsolete".to_vec(),
"text/plain".to_string(),
peer(2),
)
.expect("seed stale key");
let joiner = KvStoreSync::new(
replica,
Arc::clone(&pubsub),
topic.to_string(),
peer(2),
Some(agent(2)),
)
.expect("joiner sync");
joiner.start().await.expect("start joiner");
let mut probe = pubsub.subscribe(side.clone()).await;
let owner_sync = KvStoreSync::new(
owned,
Arc::clone(&pubsub),
topic.to_string(),
peer(1),
Some(owner_id),
)
.expect("owner sync");
owner_sync.start().await.expect("start owner");
let mut pruned = false;
for _ in 0..200 {
tokio::time::sleep(Duration::from_secs(5)).await;
let s = joiner.read().await;
if s.get("k_stale").is_none() && s.get("k_live").is_some() {
pruned = true;
break;
}
}
assert!(
pruned,
"the digest-verified full serve must prune the stale key \
(checkpoint-less deletion cold-sync)"
);
tokio::time::sleep(Duration::from_secs(160)).await;
drain_state_requests(&mut probe, peer(2)).await;
let mut relapse = 0;
for _ in 0..20 {
tokio::time::sleep(Duration::from_secs(30)).await;
relapse += drain_state_requests(&mut probe, peer(2)).await;
}
assert_eq!(relapse, 0, "the requester must stop once the digests match");
}
#[tokio::test(start_paused = true)]
async fn empty_holder_v2_marker_terminates_empty_requester() {
let node = make_node().await;
let kp = crate::identity::AgentKeypair::generate().expect("keypair");
let owner_id = kp.agent_id();
let ctx = Arc::new(crate::gossip::SigningContext::from_keypair(&kp));
let pubsub = Arc::new(PubSubManager::new(node, Some(ctx)).expect("pubsub"));
let topic = "kv-240-empty-silence";
let side = format!("{topic}{STATE_SYNC_TOPIC_SUFFIX}");
let joiner = KvStoreSync::new(
KvStore::new_replica(
store_id(1),
String::new(),
Some(owner_id),
crate::kv::store::AnchorChannel::RestParam,
),
Arc::clone(&pubsub),
topic.to_string(),
peer(2),
Some(agent(2)),
)
.expect("joiner sync");
joiner.start().await.expect("start joiner");
let holder = KvStoreSync::new(
KvStore::new_replica(
store_id(1),
String::new(),
Some(owner_id),
crate::kv::store::AnchorChannel::RestParam,
),
Arc::clone(&pubsub),
topic.to_string(),
peer(1),
Some(agent(3)),
)
.expect("holder sync");
holder.start().await.expect("start holder");
let mut probe = pubsub.subscribe(side.clone()).await;
tokio::time::sleep(Duration::from_secs(160)).await;
drain_state_requests(&mut probe, peer(2)).await;
let mut late = 0;
for _ in 0..20 {
tokio::time::sleep(Duration::from_secs(30)).await;
late += drain_state_requests(&mut probe, peer(2)).await;
}
assert_eq!(
late, 0,
"an empty holder's verifiable digest must terminate the empty \
requester's tail (genuinely-empty stores converge silently)"
);
}
#[tokio::test(start_paused = true)]
async fn v1_marker_from_old_peer_still_converges() {
let node = make_node().await;
let kp = crate::identity::AgentKeypair::generate().expect("keypair");
let owner_id = kp.agent_id();
let ctx = Arc::new(crate::gossip::SigningContext::from_keypair(&kp));
let pubsub = Arc::new(PubSubManager::new(node, Some(ctx)).expect("pubsub"));
let topic = "kv-240-v1-compat";
let side = format!("{topic}{STATE_SYNC_TOPIC_SUFFIX}");
let joiner = KvStoreSync::new(
KvStore::new_replica(
store_id(1),
String::new(),
Some(owner_id),
crate::kv::store::AnchorChannel::RestParam,
),
Arc::clone(&pubsub),
topic.to_string(),
peer(2),
Some(agent(2)),
)
.expect("joiner sync");
joiner.start().await.expect("start joiner");
let mut probe = pubsub.subscribe(side.clone()).await;
tokio::time::sleep(Duration::from_secs(20)).await;
let mut owned = KvStore::new(
store_id(1),
"log".to_string(),
owner_id,
AccessPolicy::Signed,
);
owned
.put(
"k1".to_string(),
b"v1".to_vec(),
"text/plain".to_string(),
peer(1),
)
.expect("k1");
let full = owned.full_delta();
let encoded = encode_delta(peer(1), &full).expect("encode full");
pubsub
.publish(topic.to_string(), bytes::Bytes::from(encoded))
.await
.expect("publish full delta");
let marker = KvSyncMessage::StateServed {
responder: peer(1),
empty: false,
checkpoint_seq: None,
};
let marker_bytes = bincode::serialize(&marker).expect("serialize v1 marker");
pubsub
.publish(side.clone(), bytes::Bytes::from(marker_bytes))
.await
.expect("publish v1 marker");
let mut recovered = false;
for _ in 0..60 {
tokio::time::sleep(Duration::from_secs(5)).await;
if joiner.read().await.get("k1").is_some() {
recovered = true;
break;
}
}
assert!(recovered, "the old peer's full delta must merge");
tokio::time::sleep(Duration::from_secs(160)).await;
drain_state_requests(&mut probe, peer(2)).await;
let mut late = 0;
for _ in 0..20 {
tokio::time::sleep(Duration::from_secs(30)).await;
late += drain_state_requests(&mut probe, peer(2)).await;
}
assert_eq!(
late, 0,
"v1 evidence from an old peer must still converge the requester"
);
}
#[tokio::test(start_paused = true)]
async fn tampered_digest_is_rejected_and_does_not_wedge_recovery() {
let node = make_node().await;
let kp = crate::identity::AgentKeypair::generate().expect("keypair");
let owner_id = kp.agent_id();
let ctx = Arc::new(crate::gossip::SigningContext::from_keypair(&kp));
let pubsub = Arc::new(PubSubManager::new(node, Some(ctx)).expect("pubsub"));
let topic = "kv-240-tampered";
let side = format!("{topic}{STATE_SYNC_TOPIC_SUFFIX}");
let joiner = KvStoreSync::new(
KvStore::new_replica(
store_id(1),
String::new(),
Some(owner_id),
crate::kv::store::AnchorChannel::RestParam,
),
Arc::clone(&pubsub),
topic.to_string(),
peer(2),
Some(agent(2)),
)
.expect("joiner sync");
joiner.start().await.expect("start joiner");
let mut probe = pubsub.subscribe(side.clone()).await;
let marker = KvSyncMessage::StateServedV2 {
responder: peer(1),
digest: [0xAB; 32],
entry_count: 2,
};
let bytes = bincode::serialize(&marker).expect("serialize marker");
pubsub
.publish(side.clone(), bytes::Bytes::from(bytes))
.await
.expect("publish tampered marker");
let mut requests = 0;
for _ in 0..8 {
tokio::time::sleep(Duration::from_secs(30)).await;
requests += drain_state_requests(&mut probe, peer(2)).await;
}
assert!(
requests > 0,
"a tampered digest must not converge the requester"
);
assert!(
joiner.read().await.is_empty(),
"no state can have been adopted from a forged declaration"
);
let mut owned = KvStore::new(
store_id(1),
"log".to_string(),
owner_id,
AccessPolicy::Signed,
);
owned
.put(
"k1".to_string(),
b"v1".to_vec(),
"text/plain".to_string(),
peer(1),
)
.expect("k1");
let owner_sync = KvStoreSync::new(
owned,
Arc::clone(&pubsub),
topic.to_string(),
peer(1),
Some(owner_id),
)
.expect("owner sync");
owner_sync.start().await.expect("start owner");
let mut recovered = false;
for _ in 0..200 {
tokio::time::sleep(Duration::from_secs(5)).await;
if joiner.read().await.get("k1").is_some() {
recovered = true;
break;
}
}
assert!(
recovered,
"a tampered marker must not wedge recovery from a genuine holder"
);
}
#[tokio::test(start_paused = true)]
async fn recreated_key_converges_after_owner_delete_and_recreate() {
let node = make_node().await;
let kp = crate::identity::AgentKeypair::generate().expect("keypair");
let owner_id = kp.agent_id();
let ctx = Arc::new(crate::gossip::SigningContext::from_keypair(&kp));
let pubsub = Arc::new(PubSubManager::new(node, Some(ctx)).expect("pubsub"));
let topic = "kv-240-recreate";
let side = format!("{topic}{STATE_SYNC_TOPIC_SUFFIX}");
let mut owned = KvStore::new(
store_id(1),
"log".to_string(),
owner_id,
AccessPolicy::Signed,
);
owned
.put(
"k".to_string(),
b"v1".to_vec(),
"text/plain".to_string(),
peer(1),
)
.expect("original put");
let mut aged = owned.get("k").expect("original entry").clone();
aged.created_at -= 10_000;
aged.updated_at -= 10_000;
owned.remove("k").expect("owner delete");
owned
.put(
"k".to_string(),
b"v2".to_vec(),
"text/plain".to_string(),
peer(1),
)
.expect("owner recreate");
let mut replica = KvStore::new_replica(
store_id(1),
String::new(),
Some(owner_id),
crate::kv::store::AnchorChannel::Persistence,
);
let mut seed = KvStoreDelta::new(1);
seed.added.insert("k".to_string(), (aged, (peer(1), 1)));
replica
.merge_delta(&seed, peer(1), Some(&owner_id))
.expect("seed original entry");
let joiner = KvStoreSync::new(
replica,
Arc::clone(&pubsub),
topic.to_string(),
peer(2),
Some(agent(2)),
)
.expect("joiner sync");
joiner.start().await.expect("start joiner");
let mut probe = pubsub.subscribe(side.clone()).await;
let owner_sync = KvStoreSync::new(
owned,
Arc::clone(&pubsub),
topic.to_string(),
peer(1),
Some(owner_id),
)
.expect("owner sync");
owner_sync.start().await.expect("start owner");
let mut recovered = false;
for _ in 0..200 {
tokio::time::sleep(Duration::from_secs(5)).await;
if joiner.read().await.get("k").map(|e| e.value.clone()) == Some(b"v2".to_vec()) {
recovered = true;
break;
}
}
assert!(recovered, "the re-created entry must merge");
tokio::time::sleep(Duration::from_secs(160)).await;
drain_state_requests(&mut probe, peer(2)).await;
let mut relapse = 0;
for _ in 0..20 {
tokio::time::sleep(Duration::from_secs(30)).await;
relapse += drain_state_requests(&mut probe, peer(2)).await;
}
assert_eq!(
relapse, 0,
"a delete+recreate must not wedge the requester on created_at"
);
}
#[tokio::test(start_paused = true)]
async fn allowlisted_writer_key_survives_owner_only_serve() {
let node = make_node().await;
let kp = crate::identity::AgentKeypair::generate().expect("keypair");
let owner_id = kp.agent_id();
let ctx = Arc::new(crate::gossip::SigningContext::from_keypair(&kp));
let pubsub = Arc::new(PubSubManager::new(node, Some(ctx)).expect("pubsub"));
let topic = "kv-240-allowlisted";
let writer = agent(7);
let mut owned = KvStore::new(
store_id(1),
"log".to_string(),
owner_id,
AccessPolicy::Allowlisted,
);
owned.allow_writer(writer, &owner_id).expect("allow writer");
owned
.put(
"k_owner".to_string(),
b"v".to_vec(),
"text/plain".to_string(),
peer(1),
)
.expect("owner put");
let mut replica = KvStore::new_replica(
store_id(1),
String::new(),
Some(owner_id),
crate::kv::store::AnchorChannel::Persistence,
);
replica
.learn_ownership(
owner_id,
AccessPolicy::Allowlisted,
owned.policy_version(),
&owner_id,
)
.expect("learn policy");
replica
.allow_writer(writer, &owner_id)
.expect("learn allowlist");
let k_owner_seed = {
let full = owned.full_delta();
let (key, (entry, tag)) = full
.added
.iter()
.find(|(k, _)| k.as_str() == "k_owner")
.expect("k_owner in full delta");
let mut aged = entry.clone();
aged.created_at -= 10_000;
aged.updated_at -= 10_000;
let mut d = KvStoreDelta::new(1);
d.added.insert(key.clone(), (aged, *tag));
d
};
replica
.merge_delta(&k_owner_seed, peer(1), Some(&owner_id))
.expect("seed k_owner");
let writer_entry = KvEntry::new(
"k_writer".to_string(),
b"w".to_vec(),
"text/plain".to_string(),
);
let mut writer_delta = KvStoreDelta::new(2);
writer_delta
.added
.insert("k_writer".to_string(), (writer_entry, (peer(7), 1)));
replica
.merge_delta(&writer_delta, peer(7), Some(&writer))
.expect("seed k_writer");
let owner_updated_at = owned.get("k_owner").expect("owner entry").updated_at;
let joiner = KvStoreSync::new(
replica,
Arc::clone(&pubsub),
topic.to_string(),
peer(2),
Some(agent(2)),
)
.expect("joiner sync");
joiner.start().await.expect("start joiner");
let owner_sync = KvStoreSync::new(
owned,
Arc::clone(&pubsub),
topic.to_string(),
peer(1),
Some(owner_id),
)
.expect("owner sync");
owner_sync.start().await.expect("start owner");
let mut serve_landed = false;
for _ in 0..60 {
tokio::time::sleep(Duration::from_secs(5)).await;
let s = joiner.read().await;
assert!(
s.get("k_writer").is_some(),
"an owner-only serve must NOT prune an allowlisted writer's key"
);
if s.get("k_owner").map(|e| e.updated_at) == Some(owner_updated_at) {
serve_landed = true;
}
}
assert!(
serve_landed,
"the owner's verified serve must have merged (aged entry refreshed)"
);
}
#[test]
fn pruned_key_is_accepted_when_re_served_with_fresh_tags() {
let owner = agent(1);
let mut holder = KvStore::new(store_id(1), "log".to_string(), owner, AccessPolicy::Signed);
holder
.put(
"k_live".to_string(),
b"v".to_vec(),
"text/plain".to_string(),
peer(1),
)
.expect("put live");
holder
.put(
"k_doomed".to_string(),
b"x".to_vec(),
"text/plain".to_string(),
peer(1),
)
.expect("put doomed");
let mut replica = KvStore::new_replica(
store_id(1),
String::new(),
Some(owner),
crate::kv::store::AnchorChannel::RestParam,
);
let s1 = holder.full_delta();
replica
.merge_delta(&s1, peer(1), Some(&owner))
.expect("serve 1");
assert!(replica.get("k_doomed").is_some());
holder.remove("k_doomed").expect("delete");
let s2 = holder.full_delta();
assert_eq!(
s2.served_digest(&store_id(1)),
Some(holder.served_digest()),
"the serve must carry the holder's declared digest"
);
replica
.merge_delta(&s2, peer(1), Some(&owner))
.expect("serve 2");
assert_eq!(replica.prune_to_served_set(&s2), 1);
assert!(replica.get("k_doomed").is_none());
holder
.put(
"k_doomed".to_string(),
b"y".to_vec(),
"text/plain".to_string(),
peer(1),
)
.expect("re-add");
let s3 = holder.full_delta();
replica
.merge_delta(&s3, peer(1), Some(&owner))
.expect("serve 3");
let entry = replica
.get("k_doomed")
.expect("a re-served key must be accepted after a prune");
assert_eq!(entry.value, b"y");
}
}