use serde_json::json;
use super::testkit::*;
use crate::debug::FaultKind;
async fn assert_dedupes_everywhere(
nodes: &[&TestNode],
key: &str,
text: &str,
emb: &serde_json::Value,
expected_rid: &str,
) {
let (rid, _) = keyed_write_until_accepted(nodes, key, text, emb).await;
assert_eq!(
rid, expected_rid,
"key {key} answered a DIFFERENT rid after chaos — double-write"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn kill_leader_under_keyed_load_never_double_writes() {
let _serial = serial_guard().await;
let (nodes, _tmps) = spawn_cluster(3, ClusterSpec::default()).await;
let mut live: Vec<TestNode> = nodes;
let emb = embedding(0.3);
let all: Vec<&TestNode> = live.iter().collect();
let mut rids = Vec::new();
for i in 0..5 {
let (rid, _) = keyed_write_until_accepted(
&all,
&format!("chaos-k{i}"),
&format!("chaos replicated memory {i}"),
&emb,
)
.await;
rids.push(rid);
}
let leader_idx = wait_leader(&live.iter().collect::<Vec<_>>()).await;
let dead = live.remove(leader_idx).kill().await;
let survivors: Vec<&TestNode> = live.iter().collect();
for i in 5..10 {
let (rid, _) = keyed_write_until_accepted(
&survivors,
&format!("chaos-k{i}"),
&format!("chaos replicated memory {i}"),
&emb,
)
.await;
rids.push(rid);
}
for i in 0..10 {
assert_dedupes_everywhere(
&survivors,
&format!("chaos-k{i}"),
&format!("chaos replicated memory {i}"),
&emb,
&rids[i],
)
.await;
}
let revived = dead.restart().await;
wait_for_recall(&revived, &emb, &rids[9]).await;
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn partition_and_heal_fences_stale_leader() {
let _serial = serial_guard().await;
let (nodes, _tmps) = spawn_cluster(3, ClusterSpec::default()).await;
let emb = embedding(1.1);
let all: Vec<&TestNode> = nodes.iter().collect();
let leader_idx = wait_leader(&all).await;
let leader_id = nodes[leader_idx].node_id as u32;
let others: Vec<u32> = nodes
.iter()
.filter(|n| n.node_id as u32 != leader_id)
.map(|n| n.node_id as u32)
.collect();
for n in &nodes {
n.state.fault_registry.inject(
FaultKind::Partition {
side_a: vec![leader_id],
side_b: others.clone(),
},
None,
);
}
let doomed = reqwest::Client::builder()
.timeout(std::time::Duration::from_secs(2))
.build()
.unwrap()
.post(format!("{}/v1/remember", nodes[leader_idx].base))
.bearer_auth(&nodes[leader_idx].token)
.json(&json!({
"text": "doomed partition write",
"embedding": emb,
"idempotency_key": "chaos-doomed",
}))
.send()
.await;
assert!(
doomed.is_err() || !doomed.unwrap().status().is_success(),
"a partitioned leader must not ack a write"
);
let majority: Vec<&TestNode> = nodes
.iter()
.filter(|n| n.node_id as u32 != leader_id)
.collect();
let (rid_k, _) =
keyed_write_until_accepted(&majority, "chaos-part-k", "canonical partition write", &emb)
.await;
for n in &nodes {
n.state.fault_registry.clear();
}
assert_dedupes_everywhere(
&all,
"chaos-part-k",
"canonical partition write",
&emb,
&rid_k,
)
.await;
wait_for_recall(&nodes[leader_idx], &emb, &rid_k).await;
let (rid_d1, _) =
keyed_write_until_accepted(&all, "chaos-doomed", "doomed partition write", &emb).await;
assert_dedupes_everywhere(
&all,
"chaos-doomed",
"doomed partition write",
&emb,
&rid_d1,
)
.await;
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn torn_state_boots_quarantined_then_rejoins_via_grant() {
let _serial = serial_guard().await;
let (nodes, _tmps) = spawn_cluster(3, ClusterSpec::default()).await;
let emb = embedding(2.2);
let all: Vec<&TestNode> = nodes.iter().collect();
let (_rid_a, _) =
keyed_write_until_accepted(&all, "chaos-torn-a", "pre-tear write", &emb).await;
let mut live = nodes;
let victim_idx = live
.iter()
.position(|n| n.node_id == 3 && !n.handle.is_leader())
.unwrap_or_else(|| {
live.iter()
.position(|n| !n.handle.is_leader())
.expect("some follower")
});
let victim_id = live[victim_idx].node_id as u32;
let dead = live.remove(victim_idx).kill().await;
std::fs::write(dead.data_dir.join("yrp.state"), b"CORRUPT GARBAGE BYTES").unwrap();
let survivors: Vec<&TestNode> = live.iter().collect();
let (rid_b, _) =
keyed_write_until_accepted(&survivors, "chaos-torn-b", "post-tear write", &emb).await;
let survivor_ids: Vec<u32> = survivors.iter().map(|n| n.node_id as u32).collect();
for n in &survivors {
n.state.fault_registry.inject(
FaultKind::Partition {
side_a: vec![victim_id],
side_b: survivor_ids.clone(),
},
None,
);
}
let revived = dead.restart().await;
let client = reqwest::Client::new();
for _ in 0..2 {
tokio::time::sleep(std::time::Duration::from_millis(2500)).await;
let health: serde_json::Value = client
.get(format!("{}/v1/health", revived.base))
.send()
.await
.unwrap()
.json()
.await
.unwrap();
assert_eq!(
health["status"],
json!("quarantined"),
"torn state must surface (and hold) as quarantine: {health}"
);
}
for n in &survivors {
n.state.fault_registry.clear();
}
let deadline = tokio::time::Instant::now() + std::time::Duration::from_secs(30);
while revived.handle.quarantine_reasons().is_some() {
assert!(
tokio::time::Instant::now() < deadline,
"quarantined node never rejoined"
);
tokio::time::sleep(std::time::Duration::from_millis(200)).await;
}
let preserved = std::fs::read_dir(&revived.data_dir)
.unwrap()
.filter_map(|e| e.ok())
.any(|e| {
e.file_name()
.to_string_lossy()
.starts_with("yrp.preserved-")
});
assert!(preserved, "old state must be preserved before resync");
wait_for_recall(&revived, &emb, &rid_b).await;
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn straggler_beyond_gc_catches_up_and_compacted_claims_survive() {
let _serial = serial_guard().await;
let (nodes, _tmps) = spawn_cluster(
3,
ClusterSpec {
compact_after_entries: 8,
leader_retain_entries: 2,
},
)
.await;
let emb = embedding(3.3);
let all: Vec<&TestNode> = nodes.iter().collect();
let (rid_first, _) =
keyed_write_until_accepted(&all, "chaos-gc-first", "the earliest keyed write", &emb).await;
let mut live = nodes;
let victim_idx = live
.iter()
.position(|n| !n.handle.is_leader())
.expect("some follower");
let dead = live.remove(victim_idx).kill().await;
let survivors: Vec<&TestNode> = live.iter().collect();
let mut last_rid = String::new();
for i in 0..30 {
let (rid, _) = keyed_write_until_accepted(
&survivors,
&format!("chaos-gc-{i}"),
&format!("bulk write {i}"),
&emb,
)
.await;
last_rid = rid;
}
let revived = dead.restart().await;
wait_for_recall(&revived, &emb, &last_rid).await;
let with_revived: Vec<&TestNode> = survivors
.iter()
.copied()
.chain(std::iter::once(&revived))
.collect();
assert_dedupes_everywhere(
&with_revived,
"chaos-gc-first",
"the earliest keyed write",
&emb,
&rid_first,
)
.await;
}