From 387e3820d8e1e5fc29ac0ad705d38460369047f1 Mon Sep 17 00:00:00 2001 From: DenzelPenzel Date: Wed, 30 Sep 2026 10:56:24 +0100 Subject: [PATCH 1/5] statement-store: test replica outage and recovery in v2 DHT soak --- .../zombie_ci/statement_store/v2_dht_soak.rs | 596 +++++++++++++++++- 1 file changed, 584 insertions(+), 12 deletions(-) diff --git a/cumulus/zombienet/zombienet-sdk/tests/zombie_ci/statement_store/v2_dht_soak.rs b/cumulus/zombienet/zombienet-sdk/tests/zombie_ci/statement_store/v2_dht_soak.rs index b79db44362c2..1f1409e031e8 100644 --- a/cumulus/zombienet/zombienet-sdk/tests/zombie_ci/statement_store/v2_dht_soak.rs +++ b/cumulus/zombienet/zombienet-sdk/tests/zombie_ci/statement_store/v2_dht_soak.rs @@ -9,8 +9,13 @@ //! Every wave of load asserts that the subscriber received every ring statement, and that each //! probe and three ring samples are stored on their `replication_factor` XOR-closest nodes (ring //! ones also on the subscriber, whose subscription grants affinity) and nowhere else. The soak -//! ends with a scan of every node's log for statement errors. `STATEMENT_V2_SOAK_NODES` (default -//! 12) and `STATEMENT_V2_SOAK_SECS` (default 900) size the run. +//! also stops and restarts one full node between two successful ordinary waves. During the outage, +//! every cohort statement must reach an online subscriber and stay on its original online holders. +//! After restart, the node must recover every replica and subscription statement. +//! The soak ends with a scan of every node's log for statement errors, including the victim's +//! pre-outage log. `STATEMENT_V2_SOAK_NODES` (default 12) and `STATEMENT_V2_SOAK_SECS` (default +//! 900) size the run; the lifecycle and post-recovery wave are mandatory even with a zero duration. +//! The lifecycle requires a full statement mesh. //! //! Runs on demand only: a dispatch of .github/workflows/zombienet_statement-store-soak.yml, or //! locally — against the cluster given a kubeconfig: @@ -32,26 +37,26 @@ use super::common::{ launch_network_with_commands, submit_at_rate, submit_statement, subscribe_topic, subscribe_topic_filter, wait_for_first_block, Load, }; -use anyhow::anyhow; +use anyhow::{anyhow, ensure, Context}; use codec::Encode; +use futures::FutureExt; use log::{info, warn}; use sp_crypto_hashing::blake2_256; use sp_statement_store::{StatementEvent, SubmitResult, Topic, TopicFilter}; use std::{ collections::HashSet, + panic::AssertUnwindSafe, path::Path, time::{Duration, Instant}, }; use zombienet_orchestrator::network::node::LogLineCountOptions; use zombienet_sdk::{ subxt::{backend::rpc::RpcClient, ext::subxt_rpcs::rpc_params}, - LocalFileSystem, Network, NetworkNode, + AssetLocation, LocalFileSystem, Network, NetworkNode, }; const SOAK_SECS_ENV: &str = "STATEMENT_V2_SOAK_SECS"; const DEFAULT_SOAK_SECS: u64 = 900; -/// How many statement-store nodes to spawn: the 4 authoring collators plus soak-* full nodes for -/// the rest. const NODES_ENV: &str = "STATEMENT_V2_SOAK_NODES"; const DEFAULT_NODES: usize = 12; const CONNECTED_PEER_FLOOR: usize = 50; @@ -69,6 +74,18 @@ const DELIVERY_TIMEOUT_SECS: u64 = 120; const PLACEMENT_TIMEOUT_SECS: u64 = 90; const CONNECTED_PEERS_METRIC: &str = "substrate_sync_statement_v2dht_connected_peers"; const ELIGIBLE_PEERS_METRIC: &str = "substrate_sync_statement_v2dht_eligible_peers"; +const OUTAGE_BATCHES: usize = 3; +const OUTAGE_RATE: usize = 4; +const OUTAGE_DELIVERY_SECS: u64 = 30; +const LIFECYCLE_WAIT_SECS: u64 = 240; +const OUTAGE_TIMEOUT_SECS: u64 = 600; +const STABLE_PLACEMENT_SECS: u64 = 35; +const QUIET_METRICS: [&str; 4] = [ + "substrate_sync_initial_sync_peers_active", + "substrate_sync_initial_sync_in_flight_bytes", + "substrate_sync_propagation_in_flight_bytes", + "substrate_sync_pending_statement_validations", +]; fn env_parsed(name: &str) -> Result, anyhow::Error> where @@ -89,8 +106,6 @@ fn xor_distance(a: [u8; 32], b: [u8; 32]) -> [u8; 32] { distance } -/// Node indices by the XOR distance of their peer key to `topic`, closest first. Equal distances -/// would mean equal keys, so the order needs no tie-break. fn ranked_by_distance(peer_keys: &[[u8; 32]], topic: Topic) -> Vec { let mut order: Vec = (0..peer_keys.len()).collect(); order.sort_by_cached_key(|&idx| xor_distance(*topic, peer_keys[idx])); @@ -171,10 +186,12 @@ async fn collect_peer_keys(nodes: &[NodeHandle<'_>]) -> Result, an .map_err(|e| anyhow!("{}: cannot decode peer id {peer_id}: {e}", handle.name()))?; keys.push(blake2_256(&id)); } + ensure!(keys.iter().collect::>().len() == keys.len(), "duplicate peer keys"); Ok(keys) } -/// Reads a node's persistent store from a fresh subscription's replay. +/// Reads every active stored body, including transient records, through a completed replay. +/// `Any` contributes no explicit topics and therefore does not change the placement oracle. async fn store_snapshot(rpc: &RpcClient) -> Result>, anyhow::Error> { let mut subscription = subscribe_topic_filter(rpc, TopicFilter::Any).await?; let mut snapshot = HashSet::new(); @@ -193,6 +210,556 @@ async fn store_snapshot(rpc: &RpcClient) -> Result>, anyhow::Err } } +#[derive(Clone)] +struct Placement { + hash: String, + encoded_statement: Vec, + holders: Vec, +} + +fn placement_disagreement( + phase: &str, + expected: &[Placement], + snapshots: &[(usize, HashSet>)], + names: &[&str], +) -> Option { + for placement in expected { + let mut missing = Vec::new(); + let mut extra = Vec::new(); + for (idx, snapshot) in snapshots { + match (placement.holders.contains(idx), snapshot.contains(&placement.encoded_statement)) + { + (true, false) => missing.push(names[*idx]), + (false, true) => extra.push(names[*idx]), + _ => {}, + } + } + if !missing.is_empty() || !extra.is_empty() { + return Some(format!( + "{phase} {}: missing={missing:?}, extra={extra:?}", + placement.hash, + )); + } + } + None +} + +#[derive(Clone, Copy, PartialEq, Eq)] +enum OutageRole { + Replica, + Subscriber, + NonAffine, +} + +impl std::fmt::Display for OutageRole { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + f.write_str(match self { + Self::Replica => "replica", + Self::Subscriber => "subscriber", + Self::NonAffine => "non-affine", + }) + } +} + +struct OutageTopic { + role: OutageRole, + topic: Topic, + holders: Vec, + donor: usize, + witness: usize, +} + +fn outage_topics(peer_keys: &[[u8; 32]]) -> Result<(usize, Vec), anyhow::Error> { + // A fixed peer need not have every XOR rank. Choose the non-affine topic and victim together. + // Peer keys follow the soak's node order: authoring collators, then full nodes. + let (non_affine_topic, victim) = (0..100_000) + .find_map(|candidate| { + let topic = soak_topic(b"soak-outage", 0, candidate); + let victim = ranked_by_distance(peer_keys, topic)[REPLICATION_FACTOR]; + (victim >= AUTHORING_COLLATORS.len()).then_some((topic, victim)) + }) + .ok_or_else(|| anyhow!("no non-affine full-node victim in 100000 bounded candidates"))?; + + let mut topics = Vec::new(); + for role in [OutageRole::Replica, OutageRole::Subscriber, OutageRole::NonAffine] { + let (topic, order) = (0..100_000) + .find_map(|candidate| { + let topic = soak_topic(b"soak-outage", 0, candidate); + let order = ranked_by_distance(peer_keys, topic); + let suitable = match role { + OutageRole::Replica => { + topic != non_affine_topic && order[..REPLICATION_FACTOR].contains(&victim) + }, + OutageRole::Subscriber => { + topic != non_affine_topic && !order[..REPLICATION_FACTOR].contains(&victim) + }, + OutageRole::NonAffine => topic == non_affine_topic, + }; + (suitable && topics.iter().all(|t: &OutageTopic| t.topic != topic)) + .then_some((topic, order)) + }) + .ok_or_else(|| { + anyhow!("no {role} topic for node index {victim} in 100000 bounded candidates") + })?; + + let witness = *order[..REPLICATION_FACTOR] + .iter() + .find(|&&idx| idx != victim) + .expect("K > 1 provides an online replica; qed"); + + let donor = if role == OutageRole::NonAffine { + // A farther explicit donor makes the non-replica victim a routing target on reconnect. + order[REPLICATION_FACTOR + 1] + } else { + *order[..REPLICATION_FACTOR] + .iter() + .find(|&&idx| idx != victim && idx != witness) + .expect("K > 2 provides a donor distinct from the delivery witness; qed") + }; + let mut holders = order[..REPLICATION_FACTOR].to_vec(); + match role { + OutageRole::Subscriber => holders.push(victim), + OutageRole::NonAffine => holders.push(donor), + OutageRole::Replica => {}, + } + topics.push(OutageTopic { role, topic, holders, donor, witness }); + } + + Ok((victim, topics)) +} + +/// Require a full statement mesh, excluding the stopped node. +/// After stopping the node, N-2 open statement substreams on *every* survivor is the barrier; +/// eligibility stays N-1, so none of the oracle's original replica slots is replaced. +async fn topology_disagreement( + nodes: &[NodeHandle<'_>], + offline: Option, +) -> Result, anyhow::Error> { + let connected = nodes.len() - 1 - usize::from(offline.is_some()); + let observations = futures::future::try_join_all( + nodes.iter().enumerate().filter(|(idx, _)| Some(*idx) != offline).map( + |(_, handle)| async move { + let eligible = handle.node.reports(ELIGIBLE_PEERS_METRIC).await?; + let actual_connected = handle.node.reports(CONNECTED_PEERS_METRIC).await?; + let health: serde_json::Value = + handle.rpc.request("system_health", rpc_params![]).await?; + let syncing = health["isSyncing"] + .as_bool() + .ok_or_else(|| anyhow!("{}: missing system_health.isSyncing", handle.name()))?; + Ok::<_, anyhow::Error>( + (eligible != (nodes.len() - 1) as f64 || + actual_connected != connected as f64 || + syncing) + .then(|| { + format!( + "{}: eligible={eligible}, connected={actual_connected}, syncing={syncing}; \ + expected eligible={}, connected={connected}, syncing=false", + handle.name(), nodes.len() - 1, + ) + }), + ) + }, + ), + ) + .await?; + Ok(observations.into_iter().flatten().next()) +} + +async fn wait_lifecycle_topology( + phase: &str, + nodes: &[NodeHandle<'_>], + offline: Option, +) -> Result<(), anyhow::Error> { + let mut last = String::from("no observation"); + tokio::time::timeout(Duration::from_secs(LIFECYCLE_WAIT_SECS), async { + loop { + match topology_disagreement(nodes, offline).await? { + None => return Ok::<_, anyhow::Error>(()), + Some(problem) => last = problem, + } + tokio::time::sleep(Duration::from_secs(2)).await; + } + }) + .await + .map_err(|_| anyhow!("{phase}: topology deadline: {last}"))??; + info!("Lifecycle {phase}: full eligible topology, expected open substreams, no major sync"); + Ok(()) +} + +/// Require exact placement throughout a quiet interval longer than the maintenance period. +/// Store snapshots prove the retention outcome without polling each node's entire log. +async fn wait_cohort_placement( + phase: &str, + nodes: &[NodeHandle<'_>], + expected: &[Placement], + offline: Option, +) -> Result<(), anyhow::Error> { + let names: Vec<_> = nodes.iter().map(NodeHandle::name).collect(); + let active: Vec<_> = + nodes.iter().enumerate().filter(|(idx, _)| Some(*idx) != offline).collect(); + let mut last = String::from("no complete snapshot"); + tokio::time::timeout(Duration::from_secs(LIFECYCLE_WAIT_SECS), async { + let mut stable_since: Option = None; + loop { + let snapshots = + futures::future::try_join_all(active.iter().map(|&(idx, handle)| async move { + let snapshot = + tokio::time::timeout(Duration::from_secs(30), store_snapshot(&handle.rpc)) + .await + .with_context(|| { + format!("{phase}: {} snapshot deadline", handle.name()) + })? + .with_context(|| { + format!("{phase}: {} snapshot failed", handle.name()) + })?; + Ok::<_, anyhow::Error>((idx, snapshot)) + })) + .await?; + let mut problem = placement_disagreement(phase, expected, &snapshots, &names); + if problem.is_none() { + problem = topology_disagreement(nodes, offline).await?; + } + if problem.is_none() { + for &(_, handle) in &active { + for metric in QUIET_METRICS { + let value = handle.node.reports(metric).await?; + if value != 0.0 { + problem = Some(format!("{}: {metric}={value}", handle.name())); + } + } + } + } + if let Some(problem) = problem { + last = problem; + stable_since = None; + } else { + let since = stable_since.get_or_insert_with(Instant::now); + if since.elapsed() >= Duration::from_secs(STABLE_PLACEMENT_SECS) { + return Ok::<_, anyhow::Error>(()); + } + last = String::from("waiting for stable placement and quiet queues"); + } + tokio::time::sleep(Duration::from_secs(2)).await; + } + }) + .await + .map_err(|_| anyhow!("{phase}: placement deadline: {last}"))??; + info!( + "Lifecycle {phase}: all {} cohort hashes exactly placed, \ + stable for {STABLE_PLACEMENT_SECS}s with quiet queues", + expected.len(), + ); + Ok(()) +} + +fn record_cohort( + phase: &str, + topic: &OutageTopic, + encoded_statements: &[Vec], + cohort: &mut Vec, +) { + for encoded_statement in encoded_statements { + let hash = hex::encode(blake2_256(encoded_statement)); + info!("Lifecycle {phase} topic {}: hash={hash}", topic.role); + cohort.push(Placement { + hash, + encoded_statement: encoded_statement.clone(), + holders: topic.holders.clone(), + }); + } +} + +async fn admissions_by_reason(node: &NetworkNode) -> Result { + let mut parts = Vec::new(); + for reason in ["dht", "explicit", "both", "transient", "persistent"] { + let metric = + format!("substrate_sub_statement_store_submitted_statements{{reason=\"{reason}\"}}"); + let value = node.reports(metric).await?; + parts.push(format!("{reason}={value}")); + } + Ok(parts.join(" ")) +} + +async fn victim_unavailable(victim: &NodeHandle<'_>) -> Result<(), anyhow::Error> { + let probe = victim.rpc.request::("system_localPeerId", rpc_params![]); + match tokio::time::timeout(Duration::from_secs(3), probe).await { + Ok(Ok(peer)) => Err(anyhow!("stopped victim still answers RPC as {peer}")), + Ok(Err(error)) => { + info!("Lifecycle outage: expected victim RPC failure: {error}"); + Ok(()) + }, + Err(_) => { + info!("Lifecycle outage: expected victim RPC timeout (substreams separately checked)"); + Ok(()) + }, + } +} + +async fn archive_victim_log(victim: &NetworkNode, phase: &str) -> Result<(), anyhow::Error> { + let logs = tokio::time::timeout(Duration::from_secs(30), victim.logs()) + .await + .context("victim log deadline")??; + let log_dir = base_dir()?.join("logs"); + std::fs::create_dir_all(&log_dir)?; + let path = log_dir.join(format!("{}.{phase}.log", victim.name())); + std::fs::write(&path, &logs)?; + info!("Lifecycle {phase}: preserved victim log at {}", path.display()); + ensure!( + !logs.lines().any(|line| line.contains("ERROR") && line.contains("statement")), + "{}: statement errors in {}", + victim.name(), + path.display(), + ); + Ok(()) +} + +async fn run_replica_outage( + nodes: &mut [NodeHandle<'_>], + peer_keys: &[[u8; 32]], +) -> Result { + let (victim_idx, topics) = outage_topics(peer_keys)?; + let subscriber_topic = topics + .iter() + .find(|topic| topic.role == OutageRole::Subscriber) + .expect("outage topics include a subscriber; qed") + .topic; + let victim = nodes[victim_idx].node; + let mut cohort = Vec::new(); + let mut subscribed_statements = Vec::new(); + let mut load = Load::new(PARTICIPANTS); + let mut witnesses = Vec::new(); + let mut donor_subscriptions = Vec::new(); + + let (original_peer, old_subscriber_subscription) = + tokio::time::timeout(Duration::from_secs(LIFECYCLE_WAIT_SECS), async { + wait_lifecycle_topology("baseline-ready", nodes, None).await?; + let original_peer: String = + nodes[victim_idx].rpc.request("system_localPeerId", rpc_params![]).await?; + let mut subscriber_subscription = + subscribe_topic(&nodes[victim_idx].rpc, subscriber_topic).await?; + for topic in &topics { + info!( + "Lifecycle topic {}: topic={}, holders={:?}, donor={}, witness={}", + topic.role, + hex::encode(*topic.topic), + topic.holders.iter().map(|&idx| nodes[idx].name()).collect::>(), + nodes[topic.donor].name(), + nodes[topic.witness].name() + ); + witnesses.push(subscribe_topic(&nodes[topic.witness].rpc, topic.topic).await?); + if topic.role != OutageRole::Replica { + donor_subscriptions + .push(subscribe_topic(&nodes[topic.donor].rpc, topic.topic).await?); + } + let donor = &nodes[topic.donor]; + let blobs = submit_at_rate( + &mut load, + u64::MAX, + topic.topic, + 1, + 1, + &[(donor.name(), &donor.rpc)], + ) + .await?; + assert_statements_match( + witnesses.last_mut().expect("just subscribed; qed"), + &blobs, + OUTAGE_DELIVERY_SECS, + nodes[topic.witness].name(), + ) + .await?; + if topic.role == OutageRole::Subscriber { + assert_statements_match( + &mut subscriber_subscription, + &blobs, + OUTAGE_DELIVERY_SECS, + victim.name(), + ) + .await?; + subscribed_statements.extend(blobs.clone()); + } + record_cohort("baseline", topic, &blobs, &mut cohort); + } + wait_cohort_placement("baseline", nodes, &cohort, None).await?; + Ok::<_, anyhow::Error>((original_peer, subscriber_subscription)) + }) + .await + .context("lifecycle baseline deadline")??; + + // Both providers upload these assets as executables; k8s requires a /scripts command even + // when restoring. restart_with regenerates the original node arguments for each invocation. + let offline_program = "soak-replica-offline.sh"; + let online_program = "soak-replica-online.sh"; + let dir = base_dir()?; + let offline_path = dir.join(offline_program); + let online_path = dir.join(online_program); + std::fs::write(&offline_path, "#!/bin/sh\nexec sleep 3600\n")?; + std::fs::write(&online_path, "#!/bin/sh\nexec polkadot-parachain \"$@\"\n")?; + + // Keep this outside run_wave's reconnect-and-skip catch. Even helper assertions/panics or a + // timeout must attempt recovery before returning the original failure. Replacing the node + // closes its TCP connections; SIGSTOP alone does not guarantee a statement disconnect. + let outage = tokio::time::timeout( + Duration::from_secs(OUTAGE_TIMEOUT_SECS), + AssertUnwindSafe(async { + archive_victim_log(victim, "before-outage").await?; + info!("Lifecycle stop: {} peer={original_peer}", victim.name()); + victim + .restart_with( + vec![AssetLocation::FilePath(offline_path)], + Some(offline_program.into()), + None, + None, + ) + .await?; + + // Pause only the idle replacement so the SDK monitor treats the node as offline. + victim.pause().await?; + wait_lifecycle_topology("disconnected", nodes, Some(victim_idx)).await?; + victim_unavailable(&nodes[victim_idx]).await?; + + for batch in 0..OUTAGE_BATCHES { + for (idx, topic) in topics.iter().enumerate() { + // Mint fresh hashes only after the statement-substream disconnect barrier. + let started = Instant::now(); + let blobs = + tokio::time::timeout(Duration::from_secs(OUTAGE_DELIVERY_SECS), async { + let donor = &nodes[topic.donor]; + let blobs = submit_at_rate( + &mut load, + u64::MAX, + topic.topic, + OUTAGE_RATE, + 1, + &[(donor.name(), &donor.rpc)], + ) + .await?; + assert_statements_match( + &mut witnesses[idx], + &blobs, + OUTAGE_DELIVERY_SECS, + nodes[topic.witness].name(), + ) + .await?; + Ok::<_, anyhow::Error>(blobs) + }) + .await + .with_context(|| { + format!("outage topic {} batch {batch}: delivery deadline", topic.role) + })??; + info!( + "Lifecycle outage topic {} batch {batch}: all {} hashes delivered \ + to {} in {:.1}s while victim stopped", + topic.role, + blobs.len(), + nodes[topic.witness].name(), + started.elapsed().as_secs_f64(), + ); + if topic.role == OutageRole::Subscriber { + subscribed_statements.extend(blobs.clone()); + } + record_cohort("outage", topic, &blobs, &mut cohort); + } + } + // The replica topic keeps its original K-1 online holders, without a replacement. + wait_cohort_placement("outage", nodes, &cohort, Some(victim_idx)).await?; + victim_unavailable(&nodes[victim_idx]).await?; + Ok::<_, anyhow::Error>(()) + }) + .catch_unwind(), + ) + .await; + + let restart = tokio::time::timeout(Duration::from_secs(LIFECYCLE_WAIT_SECS), async { + victim + .restart_with( + vec![AssetLocation::FilePath(online_path)], + Some(online_program.into()), + None, + None, + ) + .await?; + victim.wait_until_is_up(120u64).await?; + loop { + match victim.rpc().await { + Ok(rpc) => break Ok::<_, anyhow::Error>(rpc), + Err(error) => warn!("Lifecycle RPC not ready after restart: {error:#}"), + } + tokio::time::sleep(Duration::from_secs(1)).await; + } + }) + .await + .context("victim restart/up deadline") + .and_then(|result| result); + if let Err(error) = &restart { + warn!("Lifecycle recovery failed: {error:#}; original node is not confirmed running"); + } + match outage { + Err(_) => return Err(anyhow!("replica outage exceeded {OUTAGE_TIMEOUT_SECS}s")), + Ok(Err(panic)) => std::panic::resume_unwind(panic), + Ok(Ok(result)) => result?, + } + nodes[victim_idx].rpc = restart?; + drop(old_subscriber_subscription); + + let recovered = tokio::time::timeout( + Duration::from_secs(OUTAGE_TIMEOUT_SECS), + AssertUnwindSafe(async { + let peer: String = + nodes[victim_idx].rpc.request("system_localPeerId", rpc_params![]).await?; + ensure!(peer == original_peer, "victim PeerId changed: {original_peer} -> {peer}"); + let eligible = victim.reports(ELIGIBLE_PEERS_METRIC).await?; + info!( + "Lifecycle restarted: unchanged PeerId {peer}, same base directory {}, \ + {eligible} eligible peers known at the first RPC", + victim.base_dir().display() + ); + let mut subscriber_subscription = + subscribe_topic(&nodes[victim_idx].rpc, subscriber_topic).await?; + assert_statements_match( + &mut subscriber_subscription, + &subscribed_statements, + DELIVERY_TIMEOUT_SECS, + victim.name(), + ) + .await?; + wait_lifecycle_topology("recovered-ready", nodes, None).await?; + // Defer non-affine recovery placement: a cold topology can grant permanent DHT + // retention. + let expected: Vec = cohort + .iter() + .filter(|placement| placement.holders.contains(&victim_idx)) + .cloned() + .collect(); + let placement = wait_cohort_placement("recovered", nodes, &expected, None).await; + let admissions = admissions_by_reason(victim).await?; + placement.with_context(|| { + format!("{} admissions by reason since restart: {admissions}", victim.name()) + })?; + info!( + "Lifecycle recovered: {} holds every required replica/subscription statement; \ + admissions by reason: {admissions}", + victim.name() + ); + Ok::<_, anyhow::Error>(()) + }) + .catch_unwind(), + ) + .await; + let after_logs = archive_victim_log(victim, "after-restart").await; + if let Err(error) = &after_logs { + warn!("Lifecycle post-restart log check failed: {error:#}"); + } + match recovered { + Err(_) => return Err(anyhow!("replica recovery exceeded {OUTAGE_TIMEOUT_SECS}s")), + Ok(Err(panic)) => std::panic::resume_unwind(panic), + Ok(Ok(result)) => result?, + } + after_logs?; + drop(donor_subscriptions); + Ok(cohort.len()) +} + struct WaveReport { ring_statements: usize, submit_time: Duration, @@ -358,6 +925,7 @@ async fn statement_store_v2_dht_soak() -> Result<(), anyhow::Error> { let mut wave: u64 = 0; let mut total_statements = 0usize; let mut reconnects = 0usize; + let mut outage_completed = false; loop { let report = match run_wave(wave, &nodes, &peer_keys, &mut load).await { Ok(report) => report, @@ -381,13 +949,17 @@ async fn statement_store_v2_dht_soak() -> Result<(), anyhow::Error> { report.verify_time.as_secs_f64(), ); wave += 1; - if soak_started.elapsed() >= Duration::from_secs(soak_secs) { + if !outage_completed { + // Mandatory exactly once after the first successful wave, outside the ordinary + // dropped-RPC skip path. Do not check elapsed time until a later wave succeeds. + total_statements += run_replica_outage(&mut nodes, &peer_keys).await?; + outage_completed = true; + } else if soak_started.elapsed() >= Duration::from_secs(soak_secs) { break; } } - // No statement errors on any node; the count-zero predicate only passes after a full log - // scan, so the timeout is minimal. + // No statement errors on any node let options = LogLineCountOptions::new(|n| n == 0, Duration::from_secs(1), false); for handle in &nodes { let result = handle From 5993516bea20aba515d99e1cdeb5fa1234ecf1e9 Mon Sep 17 00:00:00 2001 From: DenzelPenzel Date: Wed, 30 Sep 2026 11:04:27 +0100 Subject: [PATCH 2/5] statement-store: prdoc for the v2 soak replica outage step --- prdoc/pr_13378.prdoc | 17 +++++++++++++++++ 1 file changed, 17 insertions(+) create mode 100644 prdoc/pr_13378.prdoc diff --git a/prdoc/pr_13378.prdoc b/prdoc/pr_13378.prdoc new file mode 100644 index 000000000000..5f9ff8253c37 --- /dev/null +++ b/prdoc/pr_13378.prdoc @@ -0,0 +1,17 @@ +title: 'statement-store: v2 soak step, a replica goes offline and returns' + +doc: +- audience: Node Dev + description: |- + Adds a mandatory lifecycle step to the on-demand v2 DHT soak zombienet test. Between the + first and the second wave one full node is stopped by swapping its executable for an idle + one, load continues on a bounded cohort of topics for which the node is a DHT replica, an + explicit subscriber and the closest non-replica, and the node is restored with its database + and node key. During the outage every cohort statement must reach an online subscriber and + stay exactly on its original online holders, and the stopped node must keep its rank in every + peer's topology. After the restart the node must hold every replica and subscription + statement and its re-opened subscription must receive the whole subscription backlog. Its + admissions by retention reason are logged; the absence of the non-replica statements on the + returned node is not asserted yet because of a known startup race in the node. The lifecycle + needs a full statement mesh. +crates: [] From 4609d035687eeb9ddc8fc074b7c1fddd9501245b Mon Sep 17 00:00:00 2001 From: DenzelPenzel Date: Wed, 30 Sep 2026 13:16:19 +0100 Subject: [PATCH 3/5] DNM: build against the litep2p provider-store fix Patches litep2p to paritytech/litep2p#666 so this branch can run the soak past the 20-providers-per-key limit where #665 panics nodes out of the network. The KademliaEvent::PeersDiscovered arm is unrelated to the fix - the branch sits on litep2p master, which added that event in #611, and sc-network has to compile against it. --- .../zombienet_statement_store_soak_tests.yml | 3 ++ Cargo.lock | 45 +++++++++---------- Cargo.toml | 5 +++ .../client/network/src/litep2p/discovery.rs | 4 ++ 4 files changed, 34 insertions(+), 23 deletions(-) diff --git a/.github/zombienet-tests/zombienet_statement_store_soak_tests.yml b/.github/zombienet-tests/zombienet_statement_store_soak_tests.yml index d93d867dade9..9e9d16105d0a 100644 --- a/.github/zombienet-tests/zombienet_statement_store_soak_tests.yml +++ b/.github/zombienet-tests/zombienet_statement_store_soak_tests.yml @@ -5,6 +5,9 @@ # Sized down for paritytech/litep2p#665, which takes nodes down above 20 providers per key; # raise it again once the fix (paritytech/litep2p#666) ships. +# DO NOT MERGE. This branch patches in the fix for paritytech/litep2p#665; 40 nodes is above the +# 20-providers-per-key limit where that bug bites, so the run exercises the fix while staying +# small enough to finish quickly. - job-name: "zombienet-statement-store-soak-0001-v2-dht" test-filter: "zombie_ci::statement_store::v2_dht_soak::statement_store_v2_dht_soak" cumulus-image: "polkadot-parachain-debug" diff --git a/Cargo.lock b/Cargo.lock index 006a161ea559..1d6df5930284 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2499,7 +2499,7 @@ dependencies = [ "bitflags 2.9.4", "cexpr", "clang-sys", - "itertools 0.10.5", + "itertools 0.11.0", "lazy_static", "lazycell", "log", @@ -6850,7 +6850,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "33d852cb9b869c2a9b3df2f71a3074817f01e1844f839a144f5fcef059a4eb5d" dependencies = [ "libc", - "windows-sys 0.52.0", + "windows-sys 0.59.0", ] [[package]] @@ -9074,7 +9074,7 @@ dependencies = [ "libc", "percent-encoding", "pin-project-lite", - "socket2 0.5.9", + "socket2 0.6.5", "tokio", "tower-service", "tracing", @@ -10919,9 +10919,8 @@ dependencies = [ [[package]] name = "litep2p" -version = "0.15.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "5e588ab8f917ed65ebcf00ab52da6001e5724fbabac9ff6f06623ea7fa004cda" +version = "0.15.2" +source = "git+https://github.com/paritytech/litep2p?branch=denzelpenzel/kad-local-provider-tracking#6c4e2dff76b4a93ae96ffc2a19bbac6e0a7141c4" dependencies = [ "async-trait", "bs58", @@ -19276,7 +19275,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "f8650aabb6c35b860610e9cff5dc1af886c9e25073b7b1712a68972af4281302" dependencies = [ "bytes", - "heck 0.4.1", + "heck 0.5.0", "itertools 0.13.0", "log", "multimap", @@ -19296,8 +19295,8 @@ version = "0.14.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ac6c3320f9abac597dcbc668774ef006702672474aad53c6d596b62e487b40b1" dependencies = [ - "heck 0.4.1", - "itertools 0.10.5", + "heck 0.5.0", + "itertools 0.11.0", "log", "multimap", "once_cell", @@ -19330,7 +19329,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "81bddcdb20abf9501610992b6759a4c888aef7d1a7247ef75e2404275ac24af1" dependencies = [ "anyhow", - "itertools 0.10.5", + "itertools 0.11.0", "proc-macro2 1.0.95", "quote 1.0.40", "syn 2.0.98", @@ -19343,7 +19342,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "8a56d757972c98b346a9b766e3f02746cde6dd1cd1d1d563472929fdd74bec4d" dependencies = [ "anyhow", - "itertools 0.10.5", + "itertools 0.11.0", "proc-macro2 1.0.95", "quote 1.0.40", "syn 2.0.98", @@ -19356,7 +19355,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9120690fafc389a67ba3803df527d0ec9cbbc9cc45e4cc20b332996dfb672425" dependencies = [ "anyhow", - "itertools 0.10.5", + "itertools 0.11.0", "proc-macro2 1.0.95", "quote 1.0.40", "syn 2.0.98", @@ -19513,7 +19512,7 @@ dependencies = [ "quinn-udp 0.5.4", "rustc-hash 2.1.1", "rustls 0.23.31", - "socket2 0.5.9", + "socket2 0.6.5", "thiserror 2.0.18", "tokio", "tracing", @@ -19564,9 +19563,9 @@ checksum = "76150b617afc75e6e21ac5f39bc196e80b65415ae48d62dbef8e2519d040ce42" dependencies = [ "cfg_aliases 0.2.1", "libc", - "socket2 0.5.9", + "socket2 0.6.5", "tracing", - "windows-sys 0.52.0", + "windows-sys 0.59.0", ] [[package]] @@ -20812,7 +20811,7 @@ dependencies = [ "errno", "libc", "linux-raw-sys 0.4.14", - "windows-sys 0.52.0", + "windows-sys 0.59.0", ] [[package]] @@ -20825,7 +20824,7 @@ dependencies = [ "errno", "libc", "linux-raw-sys 0.9.4", - "windows-sys 0.52.0", + "windows-sys 0.59.0", ] [[package]] @@ -20928,7 +20927,7 @@ dependencies = [ "security-framework 3.5.1", "security-framework-sys", "webpki-root-certs 0.26.11", - "windows-sys 0.52.0", + "windows-sys 0.59.0", ] [[package]] @@ -20949,7 +20948,7 @@ dependencies = [ "security-framework 3.5.1", "security-framework-sys", "webpki-root-certs 1.0.3", - "windows-sys 0.52.0", + "windows-sys 0.59.0", ] [[package]] @@ -22005,7 +22004,7 @@ dependencies = [ "ip_network", "libp2p", "linked_hash_set", - "litep2p 0.15.1", + "litep2p 0.15.2", "log", "mockall", "multistream-select", @@ -22107,7 +22106,7 @@ dependencies = [ "bytes", "cid", "futures", - "litep2p 0.15.1", + "litep2p 0.15.2", "log", "proptest", "rand 0.8.5", @@ -22327,7 +22326,7 @@ dependencies = [ "ed25519-dalek", "libp2p-identity", "libp2p-kad", - "litep2p 0.15.1", + "litep2p 0.15.2", "log", "multiaddr", "multihash", @@ -28280,7 +28279,7 @@ dependencies = [ "fastrand", "once_cell", "rustix 0.38.42", - "windows-sys 0.52.0", + "windows-sys 0.59.0", ] [[package]] diff --git a/Cargo.toml b/Cargo.toml index c19bbf6e6b46..3065d9ce8bc6 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1525,6 +1525,11 @@ zombienet-orchestrator = { version = "0.4.18" } zombienet-sdk = { version = "0.4.18" } zstd = { version = "0.12.4", default-features = false } +# DO NOT MERGE. Points litep2p at the fix for paritytech/litep2p#665 so this soak can run past +# the 20-providers-per-key limit without nodes panicking out of the network. +[patch.crates-io] +litep2p = { git = "https://github.com/paritytech/litep2p", branch = "denzelpenzel/kad-local-provider-tracking" } + [profile.release] # Polkadot runtime requires unwinding. opt-level = 3 diff --git a/substrate/client/network/src/litep2p/discovery.rs b/substrate/client/network/src/litep2p/discovery.rs index 7c4f0c957c76..d42e8f55a6ee 100644 --- a/substrate/client/network/src/litep2p/discovery.rs +++ b/substrate/client/network/src/litep2p/discovery.rs @@ -721,6 +721,10 @@ impl Stream for Discovery { }, // We do not validate incoming providers. Poll::Ready(Some(KademliaEvent::IncomingProvider { .. })) => {}, + // DO NOT MERGE. Only reachable because this branch patches litep2p to master, which + // added this event in paritytech/litep2p#611. litep2p has already recorded the + // addresses by the time it reports them, so there is nothing to do here. + Poll::Ready(Some(KademliaEvent::PeersDiscovered { .. })) => {}, } match Pin::new(&mut this.identify_event_stream).poll_next(cx) { From 084fd94b84c279d5bc3faa8cb4cc7ee6280cc4fc Mon Sep 17 00:00:00 2001 From: DenzelPenzel Date: Wed, 30 Sep 2026 15:47:46 +0100 Subject: [PATCH 4/5] DNM: spawn the network in batches The orchestrator otherwise brings a whole level up at once, and the apiserver drops the exec connections partway through: two runs here died at 35 and 39 of 43 nodes with 'deadline has elapsed'. A 102-node run with this set spawned cleanly in ten minutes, so the size is not the problem - the burst is. --- .github/workflows/zombienet_statement-store-soak.yml | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/.github/workflows/zombienet_statement-store-soak.yml b/.github/workflows/zombienet_statement-store-soak.yml index 7af3d0521ccb..e7b863b10441 100644 --- a/.github/workflows/zombienet_statement-store-soak.yml +++ b/.github/workflows/zombienet_statement-store-soak.yml @@ -37,6 +37,11 @@ permissions: read-all env: GHA_CLUSTER_SERVER_ADDR: "https://kubernetes.default:443" KUBECONFIG: "/data/config" + # DO NOT MERGE. Without this the orchestrator brings a whole level up at once - the default + # works out to one exec per node, all concurrent - and the apiserver drops connections partway + # through: two runs here died at 35 and 39 of 43 nodes with "deadline has elapsed", while a + # 102-node run with this set spawned cleanly. + ZOMBIE_SPAWN_CONCURRENCY: "20" ZOMBIE_CLEANER_DISABLED: 1 jobs: From d7f24b6899777946d002c5a4466181b86aab3b35 Mon Sep 17 00:00:00 2001 From: DenzelPenzel Date: Wed, 30 Sep 2026 17:08:21 +0100 Subject: [PATCH 5/5] statement-store: clarify v2 soak naming and simplify lifecycle checks --- .../zombienet_statement-store-soak.yml | 1 + .../zombie_ci/statement_store/v2_dht_soak.rs | 459 +++++++++--------- 2 files changed, 239 insertions(+), 221 deletions(-) diff --git a/.github/workflows/zombienet_statement-store-soak.yml b/.github/workflows/zombienet_statement-store-soak.yml index 7af3d0521ccb..1466d5c4fde9 100644 --- a/.github/workflows/zombienet_statement-store-soak.yml +++ b/.github/workflows/zombienet_statement-store-soak.yml @@ -37,6 +37,7 @@ permissions: read-all env: GHA_CLUSTER_SERVER_ADDR: "https://kubernetes.default:443" KUBECONFIG: "/data/config" + ZOMBIE_SPAWN_CONCURRENCY: "20" ZOMBIE_CLEANER_DISABLED: 1 jobs: diff --git a/cumulus/zombienet/zombienet-sdk/tests/zombie_ci/statement_store/v2_dht_soak.rs b/cumulus/zombienet/zombienet-sdk/tests/zombie_ci/statement_store/v2_dht_soak.rs index 1f1409e031e8..b6cd27207595 100644 --- a/cumulus/zombienet/zombienet-sdk/tests/zombie_ci/statement_store/v2_dht_soak.rs +++ b/cumulus/zombienet/zombienet-sdk/tests/zombie_ci/statement_store/v2_dht_soak.rs @@ -12,8 +12,8 @@ //! also stops and restarts one full node between two successful ordinary waves. During the outage, //! every cohort statement must reach an online subscriber and stay on its original online holders. //! After restart, the node must recover every replica and subscription statement. -//! The soak ends with a scan of every node's log for statement errors, including the victim's -//! pre-outage log. `STATEMENT_V2_SOAK_NODES` (default 12) and `STATEMENT_V2_SOAK_SECS` (default +//! The soak ends with a scan of every node's log for statement errors. +//! `STATEMENT_V2_SOAK_NODES` (default 12) and `STATEMENT_V2_SOAK_SECS` (default //! 900) size the run; the lifecycle and post-recovery wave are mandatory even with a zero duration. //! The lifecycle requires a full statement mesh. //! @@ -75,8 +75,8 @@ const PLACEMENT_TIMEOUT_SECS: u64 = 90; const CONNECTED_PEERS_METRIC: &str = "substrate_sync_statement_v2dht_connected_peers"; const ELIGIBLE_PEERS_METRIC: &str = "substrate_sync_statement_v2dht_eligible_peers"; const OUTAGE_BATCHES: usize = 3; -const OUTAGE_RATE: usize = 4; -const OUTAGE_DELIVERY_SECS: u64 = 30; +const OUTAGE_RATE_PER_SECOND: usize = 4; +const OUTAGE_DELIVERY_TIMEOUT_SECS: u64 = 30; const LIFECYCLE_WAIT_SECS: u64 = 240; const OUTAGE_TIMEOUT_SECS: u64 = 600; const STABLE_PLACEMENT_SECS: u64 = 35; @@ -190,8 +190,6 @@ async fn collect_peer_keys(nodes: &[NodeHandle<'_>]) -> Result, an Ok(keys) } -/// Reads every active stored body, including transient records, through a completed replay. -/// `Any` contributes no explicit topics and therefore does not change the placement oracle. async fn store_snapshot(rpc: &RpcClient) -> Result>, anyhow::Error> { let mut subscription = subscribe_topic_filter(rpc, TopicFilter::Any).await?; let mut snapshot = HashSet::new(); @@ -211,32 +209,34 @@ async fn store_snapshot(rpc: &RpcClient) -> Result>, anyhow::Err } #[derive(Clone)] -struct Placement { +struct ExpectedPlacement { hash: String, encoded_statement: Vec, - holders: Vec, + holder_node_indices: Vec, } -fn placement_disagreement( +fn placement_mismatch( phase: &str, - expected: &[Placement], + expected_placements: &[ExpectedPlacement], snapshots: &[(usize, HashSet>)], - names: &[&str], + node_names: &[&str], ) -> Option { - for placement in expected { - let mut missing = Vec::new(); - let mut extra = Vec::new(); - for (idx, snapshot) in snapshots { - match (placement.holders.contains(idx), snapshot.contains(&placement.encoded_statement)) - { - (true, false) => missing.push(names[*idx]), - (false, true) => extra.push(names[*idx]), + for placement in expected_placements { + let mut missing_nodes = Vec::new(); + let mut extra_nodes = Vec::new(); + for (node_idx, snapshot) in snapshots { + match ( + placement.holder_node_indices.contains(node_idx), + snapshot.contains(&placement.encoded_statement), + ) { + (true, false) => missing_nodes.push(node_names[*node_idx]), + (false, true) => extra_nodes.push(node_names[*node_idx]), _ => {}, } } - if !missing.is_empty() || !extra.is_empty() { + if !missing_nodes.is_empty() || !extra_nodes.is_empty() { return Some(format!( - "{phase} {}: missing={missing:?}, extra={extra:?}", + "{phase} {}: missing={missing_nodes:?}, extra={extra_nodes:?}", placement.hash, )); } @@ -245,13 +245,13 @@ fn placement_disagreement( } #[derive(Clone, Copy, PartialEq, Eq)] -enum OutageRole { +enum OutageNodeRole { Replica, Subscriber, NonAffine, } -impl std::fmt::Display for OutageRole { +impl std::fmt::Display for OutageNodeRole { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { f.write_str(match self { Self::Replica => "replica", @@ -262,82 +262,95 @@ impl std::fmt::Display for OutageRole { } struct OutageTopic { - role: OutageRole, + outage_node_role: OutageNodeRole, topic: Topic, - holders: Vec, - donor: usize, - witness: usize, + holder_node_indices: Vec, + submitter_node_idx: usize, + receiver_node_idx: usize, } fn outage_topics(peer_keys: &[[u8; 32]]) -> Result<(usize, Vec), anyhow::Error> { - // A fixed peer need not have every XOR rank. Choose the non-affine topic and victim together. - // Peer keys follow the soak's node order: authoring collators, then full nodes. - let (non_affine_topic, victim) = (0..100_000) + // A fixed peer need not have every XOR rank. Choose the non-affine topic and outage node + // together. Peer keys follow the soak's node order: authoring collators, then full nodes. + let (non_affine_topic, outage_node_idx) = (0..100_000) .find_map(|candidate| { let topic = soak_topic(b"soak-outage", 0, candidate); - let victim = ranked_by_distance(peer_keys, topic)[REPLICATION_FACTOR]; - (victim >= AUTHORING_COLLATORS.len()).then_some((topic, victim)) + let outage_node_idx = ranked_by_distance(peer_keys, topic)[REPLICATION_FACTOR]; + (outage_node_idx >= AUTHORING_COLLATORS.len()).then_some((topic, outage_node_idx)) }) - .ok_or_else(|| anyhow!("no non-affine full-node victim in 100000 bounded candidates"))?; + .ok_or_else(|| anyhow!("no non-affine full outage node in 100000 bounded candidates"))?; let mut topics = Vec::new(); - for role in [OutageRole::Replica, OutageRole::Subscriber, OutageRole::NonAffine] { + for outage_node_role in [ + OutageNodeRole::Replica, + OutageNodeRole::Subscriber, + OutageNodeRole::NonAffine, + ] { let (topic, order) = (0..100_000) .find_map(|candidate| { let topic = soak_topic(b"soak-outage", 0, candidate); let order = ranked_by_distance(peer_keys, topic); - let suitable = match role { - OutageRole::Replica => { - topic != non_affine_topic && order[..REPLICATION_FACTOR].contains(&victim) + let suitable = match outage_node_role { + OutageNodeRole::Replica => { + topic != non_affine_topic && + order[..REPLICATION_FACTOR].contains(&outage_node_idx) }, - OutageRole::Subscriber => { - topic != non_affine_topic && !order[..REPLICATION_FACTOR].contains(&victim) + OutageNodeRole::Subscriber => { + topic != non_affine_topic && + !order[..REPLICATION_FACTOR].contains(&outage_node_idx) }, - OutageRole::NonAffine => topic == non_affine_topic, + OutageNodeRole::NonAffine => topic == non_affine_topic, }; (suitable && topics.iter().all(|t: &OutageTopic| t.topic != topic)) .then_some((topic, order)) }) .ok_or_else(|| { - anyhow!("no {role} topic for node index {victim} in 100000 bounded candidates") + anyhow!( + "no {outage_node_role} topic for node index {outage_node_idx} \ + in 100000 bounded candidates" + ) })?; - let witness = *order[..REPLICATION_FACTOR] + let receiver_node_idx = *order[..REPLICATION_FACTOR] .iter() - .find(|&&idx| idx != victim) + .find(|&&idx| idx != outage_node_idx) .expect("K > 1 provides an online replica; qed"); - let donor = if role == OutageRole::NonAffine { - // A farther explicit donor makes the non-replica victim a routing target on reconnect. + let submitter_node_idx = if outage_node_role == OutageNodeRole::NonAffine { + // A farther explicit submitter makes the outage node a routing target on reconnect. order[REPLICATION_FACTOR + 1] } else { *order[..REPLICATION_FACTOR] .iter() - .find(|&&idx| idx != victim && idx != witness) - .expect("K > 2 provides a donor distinct from the delivery witness; qed") + .find(|&&idx| idx != outage_node_idx && idx != receiver_node_idx) + .expect("K > 2 provides a submitter distinct from the delivery receiver; qed") }; - let mut holders = order[..REPLICATION_FACTOR].to_vec(); - match role { - OutageRole::Subscriber => holders.push(victim), - OutageRole::NonAffine => holders.push(donor), - OutageRole::Replica => {}, + let mut holder_node_indices = order[..REPLICATION_FACTOR].to_vec(); + match outage_node_role { + OutageNodeRole::Subscriber => holder_node_indices.push(outage_node_idx), + OutageNodeRole::NonAffine => holder_node_indices.push(submitter_node_idx), + OutageNodeRole::Replica => {}, } - topics.push(OutageTopic { role, topic, holders, donor, witness }); + topics.push(OutageTopic { + outage_node_role, + topic, + holder_node_indices, + submitter_node_idx, + receiver_node_idx, + }); } - Ok((victim, topics)) + Ok((outage_node_idx, topics)) } -/// Require a full statement mesh, excluding the stopped node. -/// After stopping the node, N-2 open statement substreams on *every* survivor is the barrier; -/// eligibility stays N-1, so none of the oracle's original replica slots is replaced. -async fn topology_disagreement( +async fn topology_mismatch( nodes: &[NodeHandle<'_>], - offline: Option, + offline_node_idx: Option, ) -> Result, anyhow::Error> { - let connected = nodes.len() - 1 - usize::from(offline.is_some()); + let expected_eligible = nodes.len() - 1; + let expected_connected = expected_eligible - usize::from(offline_node_idx.is_some()); let observations = futures::future::try_join_all( - nodes.iter().enumerate().filter(|(idx, _)| Some(*idx) != offline).map( + nodes.iter().enumerate().filter(|(idx, _)| Some(*idx) != offline_node_idx).map( |(_, handle)| async move { let eligible = handle.node.reports(ELIGIBLE_PEERS_METRIC).await?; let actual_connected = handle.node.reports(CONNECTED_PEERS_METRIC).await?; @@ -346,18 +359,18 @@ async fn topology_disagreement( let syncing = health["isSyncing"] .as_bool() .ok_or_else(|| anyhow!("{}: missing system_health.isSyncing", handle.name()))?; - Ok::<_, anyhow::Error>( - (eligible != (nodes.len() - 1) as f64 || - actual_connected != connected as f64 || - syncing) - .then(|| { - format!( - "{}: eligible={eligible}, connected={actual_connected}, syncing={syncing}; \ - expected eligible={}, connected={connected}, syncing=false", - handle.name(), nodes.len() - 1, - ) - }), - ) + if eligible != expected_eligible as f64 || + actual_connected != expected_connected as f64 || + syncing + { + return Ok(Some(format!( + "{}: eligible={eligible}, connected={actual_connected}, syncing={syncing}; \ + expected eligible={expected_eligible}, connected={expected_connected}, \ + syncing=false", + handle.name(), + ))); + } + Ok::<_, anyhow::Error>(None) }, ), ) @@ -365,15 +378,15 @@ async fn topology_disagreement( Ok(observations.into_iter().flatten().next()) } -async fn wait_lifecycle_topology( +async fn wait_for_topology( phase: &str, nodes: &[NodeHandle<'_>], - offline: Option, + offline_node_idx: Option, ) -> Result<(), anyhow::Error> { let mut last = String::from("no observation"); tokio::time::timeout(Duration::from_secs(LIFECYCLE_WAIT_SECS), async { loop { - match topology_disagreement(nodes, offline).await? { + match topology_mismatch(nodes, offline_node_idx).await? { None => return Ok::<_, anyhow::Error>(()), Some(problem) => last = problem, } @@ -386,23 +399,24 @@ async fn wait_lifecycle_topology( Ok(()) } -/// Require exact placement throughout a quiet interval longer than the maintenance period. -/// Store snapshots prove the retention outcome without polling each node's entire log. -async fn wait_cohort_placement( +async fn wait_for_stable_placement( phase: &str, nodes: &[NodeHandle<'_>], - expected: &[Placement], - offline: Option, + expected_placements: &[ExpectedPlacement], + offline_node_idx: Option, ) -> Result<(), anyhow::Error> { - let names: Vec<_> = nodes.iter().map(NodeHandle::name).collect(); - let active: Vec<_> = - nodes.iter().enumerate().filter(|(idx, _)| Some(*idx) != offline).collect(); + let node_names: Vec<_> = nodes.iter().map(NodeHandle::name).collect(); + let online_nodes: Vec<_> = nodes + .iter() + .enumerate() + .filter(|(idx, _)| Some(*idx) != offline_node_idx) + .collect(); let mut last = String::from("no complete snapshot"); tokio::time::timeout(Duration::from_secs(LIFECYCLE_WAIT_SECS), async { let mut stable_since: Option = None; loop { let snapshots = - futures::future::try_join_all(active.iter().map(|&(idx, handle)| async move { + futures::future::try_join_all(online_nodes.iter().map(|&(idx, handle)| async move { let snapshot = tokio::time::timeout(Duration::from_secs(30), store_snapshot(&handle.rpc)) .await @@ -415,12 +429,12 @@ async fn wait_cohort_placement( Ok::<_, anyhow::Error>((idx, snapshot)) })) .await?; - let mut problem = placement_disagreement(phase, expected, &snapshots, &names); + let mut problem = placement_mismatch(phase, expected_placements, &snapshots, &node_names); if problem.is_none() { - problem = topology_disagreement(nodes, offline).await?; + problem = topology_mismatch(nodes, offline_node_idx).await?; } if problem.is_none() { - for &(_, handle) in &active { + for &(_, handle) in &online_nodes { for metric in QUIET_METRICS { let value = handle.node.reports(metric).await?; if value != 0.0 { @@ -447,24 +461,24 @@ async fn wait_cohort_placement( info!( "Lifecycle {phase}: all {} cohort hashes exactly placed, \ stable for {STABLE_PLACEMENT_SECS}s with quiet queues", - expected.len(), + expected_placements.len(), ); Ok(()) } -fn record_cohort( +fn record_expected_placements( phase: &str, topic: &OutageTopic, encoded_statements: &[Vec], - cohort: &mut Vec, + expected_placements: &mut Vec, ) { for encoded_statement in encoded_statements { let hash = hex::encode(blake2_256(encoded_statement)); - info!("Lifecycle {phase} topic {}: hash={hash}", topic.role); - cohort.push(Placement { + info!("Lifecycle {phase} topic {}: hash={hash}", topic.outage_node_role); + expected_placements.push(ExpectedPlacement { hash, encoded_statement: encoded_statement.clone(), - holders: topic.holders.clone(), + holder_node_indices: topic.holder_node_indices.clone(), }); } } @@ -480,8 +494,8 @@ async fn admissions_by_reason(node: &NetworkNode) -> Result) -> Result<(), anyhow::Error> { - let probe = victim.rpc.request::("system_localPeerId", rpc_params![]); +async fn ensure_outage_node_unavailable(outage_node: &NodeHandle<'_>) -> Result<(), anyhow::Error> { + let probe = outage_node.rpc.request::("system_localPeerId", rpc_params![]); match tokio::time::timeout(Duration::from_secs(3), probe).await { Ok(Ok(peer)) => Err(anyhow!("stopped victim still answers RPC as {peer}")), Ok(Err(error)) => { @@ -495,93 +509,89 @@ async fn victim_unavailable(victim: &NodeHandle<'_>) -> Result<(), anyhow::Error } } -async fn archive_victim_log(victim: &NetworkNode, phase: &str) -> Result<(), anyhow::Error> { - let logs = tokio::time::timeout(Duration::from_secs(30), victim.logs()) - .await - .context("victim log deadline")??; - let log_dir = base_dir()?.join("logs"); - std::fs::create_dir_all(&log_dir)?; - let path = log_dir.join(format!("{}.{phase}.log", victim.name())); - std::fs::write(&path, &logs)?; - info!("Lifecycle {phase}: preserved victim log at {}", path.display()); - ensure!( - !logs.lines().any(|line| line.contains("ERROR") && line.contains("statement")), - "{}: statement errors in {}", - victim.name(), - path.display(), - ); - Ok(()) -} - +/// Full-node outage and recovery: +/// 1. Pick replica/subscriber/non-affine topics; check baseline delivery and placement. +/// 2. Stop the node; wait for its statement substreams to close. +/// 3. Submit fresh statements; check delivery and placement on original online holders. +/// 4. Restart even if the outage fails, panics or times out. +/// 5. Check unchanged PeerId, subscription backlog and replica/subscriber placement. async fn run_replica_outage( nodes: &mut [NodeHandle<'_>], peer_keys: &[[u8; 32]], ) -> Result { - let (victim_idx, topics) = outage_topics(peer_keys)?; + let (outage_node_idx, topics) = outage_topics(peer_keys)?; let subscriber_topic = topics .iter() - .find(|topic| topic.role == OutageRole::Subscriber) + .find(|topic| topic.outage_node_role == OutageNodeRole::Subscriber) .expect("outage topics include a subscriber; qed") .topic; - let victim = nodes[victim_idx].node; - let mut cohort = Vec::new(); - let mut subscribed_statements = Vec::new(); + let outage_node = nodes[outage_node_idx].node; + let mut expected_placements = Vec::new(); + let mut expected_subscription_backlog = Vec::new(); let mut load = Load::new(PARTICIPANTS); - let mut witnesses = Vec::new(); - let mut donor_subscriptions = Vec::new(); + let mut receiver_subscriptions = Vec::new(); + let mut submitter_subscriptions = Vec::new(); - let (original_peer, old_subscriber_subscription) = + let (original_peer_id, pre_outage_subscription) = tokio::time::timeout(Duration::from_secs(LIFECYCLE_WAIT_SECS), async { - wait_lifecycle_topology("baseline-ready", nodes, None).await?; - let original_peer: String = - nodes[victim_idx].rpc.request("system_localPeerId", rpc_params![]).await?; - let mut subscriber_subscription = - subscribe_topic(&nodes[victim_idx].rpc, subscriber_topic).await?; + wait_for_topology("baseline-ready", nodes, None).await?; + let original_peer_id: String = + nodes[outage_node_idx].rpc.request("system_localPeerId", rpc_params![]).await?; + let mut outage_node_subscription = + subscribe_topic(&nodes[outage_node_idx].rpc, subscriber_topic).await?; for topic in &topics { info!( - "Lifecycle topic {}: topic={}, holders={:?}, donor={}, witness={}", - topic.role, + "Lifecycle topic {}: topic={}, holders={:?}, submitter={}, receiver={}", + topic.outage_node_role, hex::encode(*topic.topic), - topic.holders.iter().map(|&idx| nodes[idx].name()).collect::>(), - nodes[topic.donor].name(), - nodes[topic.witness].name() + topic.holder_node_indices.iter().map(|&idx| nodes[idx].name()).collect::>(), + nodes[topic.submitter_node_idx].name(), + nodes[topic.receiver_node_idx].name() ); - witnesses.push(subscribe_topic(&nodes[topic.witness].rpc, topic.topic).await?); - if topic.role != OutageRole::Replica { - donor_subscriptions - .push(subscribe_topic(&nodes[topic.donor].rpc, topic.topic).await?); + receiver_subscriptions + .push(subscribe_topic(&nodes[topic.receiver_node_idx].rpc, topic.topic).await?); + if topic.outage_node_role != OutageNodeRole::Replica { + submitter_subscriptions + .push(subscribe_topic(&nodes[topic.submitter_node_idx].rpc, topic.topic).await?); } - let donor = &nodes[topic.donor]; - let blobs = submit_at_rate( + let submitter = &nodes[topic.submitter_node_idx]; + let encoded_statements = submit_at_rate( &mut load, u64::MAX, topic.topic, 1, 1, - &[(donor.name(), &donor.rpc)], + &[(submitter.name(), &submitter.rpc)], ) .await?; assert_statements_match( - witnesses.last_mut().expect("just subscribed; qed"), - &blobs, - OUTAGE_DELIVERY_SECS, - nodes[topic.witness].name(), + receiver_subscriptions.last_mut().expect("just subscribed; qed"), + &encoded_statements, + OUTAGE_DELIVERY_TIMEOUT_SECS, + nodes[topic.receiver_node_idx].name(), ) .await?; - if topic.role == OutageRole::Subscriber { + + if topic.outage_node_role == OutageNodeRole::Subscriber { assert_statements_match( - &mut subscriber_subscription, - &blobs, - OUTAGE_DELIVERY_SECS, - victim.name(), + &mut outage_node_subscription, + &encoded_statements, + OUTAGE_DELIVERY_TIMEOUT_SECS, + outage_node.name(), ) .await?; - subscribed_statements.extend(blobs.clone()); + expected_subscription_backlog.extend(encoded_statements.clone()); } - record_cohort("baseline", topic, &blobs, &mut cohort); + + record_expected_placements( + "baseline", + topic, + &encoded_statements, + &mut expected_placements, + ); } - wait_cohort_placement("baseline", nodes, &cohort, None).await?; - Ok::<_, anyhow::Error>((original_peer, subscriber_subscription)) + wait_for_stable_placement("baseline", nodes, &expected_placements, None).await?; + Ok::<_, anyhow::Error>((original_peer_id, outage_node_subscription)) }) .await .context("lifecycle baseline deadline")??; @@ -596,15 +606,11 @@ async fn run_replica_outage( std::fs::write(&offline_path, "#!/bin/sh\nexec sleep 3600\n")?; std::fs::write(&online_path, "#!/bin/sh\nexec polkadot-parachain \"$@\"\n")?; - // Keep this outside run_wave's reconnect-and-skip catch. Even helper assertions/panics or a - // timeout must attempt recovery before returning the original failure. Replacing the node - // closes its TCP connections; SIGSTOP alone does not guarantee a statement disconnect. - let outage = tokio::time::timeout( + let outage_result = tokio::time::timeout( Duration::from_secs(OUTAGE_TIMEOUT_SECS), AssertUnwindSafe(async { - archive_victim_log(victim, "before-outage").await?; - info!("Lifecycle stop: {} peer={original_peer}", victim.name()); - victim + info!("Lifecycle stop: {} peer={original_peer_id}", outage_node.name()); + outage_node .restart_with( vec![AssetLocation::FilePath(offline_path)], Some(offline_program.into()), @@ -614,64 +620,74 @@ async fn run_replica_outage( .await?; // Pause only the idle replacement so the SDK monitor treats the node as offline. - victim.pause().await?; - wait_lifecycle_topology("disconnected", nodes, Some(victim_idx)).await?; - victim_unavailable(&nodes[victim_idx]).await?; + outage_node.pause().await?; + wait_for_topology("disconnected", nodes, Some(outage_node_idx)).await?; + ensure_outage_node_unavailable(&nodes[outage_node_idx]).await?; + // Send data during the outage for batch in 0..OUTAGE_BATCHES { - for (idx, topic) in topics.iter().enumerate() { + for (topic_idx, topic) in topics.iter().enumerate() { // Mint fresh hashes only after the statement-substream disconnect barrier. let started = Instant::now(); - let blobs = - tokio::time::timeout(Duration::from_secs(OUTAGE_DELIVERY_SECS), async { - let donor = &nodes[topic.donor]; - let blobs = submit_at_rate( + let encoded_statements = + tokio::time::timeout(Duration::from_secs(OUTAGE_DELIVERY_TIMEOUT_SECS), async { + let submitter = &nodes[topic.submitter_node_idx]; + let encoded_statements = submit_at_rate( &mut load, u64::MAX, topic.topic, - OUTAGE_RATE, + OUTAGE_RATE_PER_SECOND, 1, - &[(donor.name(), &donor.rpc)], + &[(submitter.name(), &submitter.rpc)], ) .await?; assert_statements_match( - &mut witnesses[idx], - &blobs, - OUTAGE_DELIVERY_SECS, - nodes[topic.witness].name(), + &mut receiver_subscriptions[topic_idx], + &encoded_statements, + OUTAGE_DELIVERY_TIMEOUT_SECS, + nodes[topic.receiver_node_idx].name(), ) .await?; - Ok::<_, anyhow::Error>(blobs) + Ok::<_, anyhow::Error>(encoded_statements) }) .await .with_context(|| { - format!("outage topic {} batch {batch}: delivery deadline", topic.role) + format!( + "outage topic {} batch {batch}: delivery deadline", + topic.outage_node_role + ) })??; info!( "Lifecycle outage topic {} batch {batch}: all {} hashes delivered \ - to {} in {:.1}s while victim stopped", - topic.role, - blobs.len(), - nodes[topic.witness].name(), + to {} in {:.1}s while outage node stopped", + topic.outage_node_role, + encoded_statements.len(), + nodes[topic.receiver_node_idx].name(), started.elapsed().as_secs_f64(), ); - if topic.role == OutageRole::Subscriber { - subscribed_statements.extend(blobs.clone()); + if topic.outage_node_role == OutageNodeRole::Subscriber { + expected_subscription_backlog.extend(encoded_statements.clone()); } - record_cohort("outage", topic, &blobs, &mut cohort); + record_expected_placements( + "outage", + topic, + &encoded_statements, + &mut expected_placements, + ); } } // The replica topic keeps its original K-1 online holders, without a replacement. - wait_cohort_placement("outage", nodes, &cohort, Some(victim_idx)).await?; - victim_unavailable(&nodes[victim_idx]).await?; + wait_for_stable_placement("outage", nodes, &expected_placements, Some(outage_node_idx)) + .await?; + ensure_outage_node_unavailable(&nodes[outage_node_idx]).await?; Ok::<_, anyhow::Error>(()) }) .catch_unwind(), ) .await; - let restart = tokio::time::timeout(Duration::from_secs(LIFECYCLE_WAIT_SECS), async { - victim + let restart_rpc_result = tokio::time::timeout(Duration::from_secs(LIFECYCLE_WAIT_SECS), async { + outage_node .restart_with( vec![AssetLocation::FilePath(online_path)], Some(online_program.into()), @@ -679,9 +695,9 @@ async fn run_replica_outage( None, ) .await?; - victim.wait_until_is_up(120u64).await?; + outage_node.wait_until_is_up(120u64).await?; loop { - match victim.rpc().await { + match outage_node.rpc().await { Ok(rpc) => break Ok::<_, anyhow::Error>(rpc), Err(error) => warn!("Lifecycle RPC not ready after restart: {error:#}"), } @@ -689,75 +705,76 @@ async fn run_replica_outage( } }) .await - .context("victim restart/up deadline") + .context("outage node restart/up deadline") .and_then(|result| result); - if let Err(error) = &restart { + if let Err(error) = &restart_rpc_result { warn!("Lifecycle recovery failed: {error:#}; original node is not confirmed running"); } - match outage { + + match outage_result { Err(_) => return Err(anyhow!("replica outage exceeded {OUTAGE_TIMEOUT_SECS}s")), Ok(Err(panic)) => std::panic::resume_unwind(panic), Ok(Ok(result)) => result?, } - nodes[victim_idx].rpc = restart?; - drop(old_subscriber_subscription); - let recovered = tokio::time::timeout( + nodes[outage_node_idx].rpc = restart_rpc_result?; + drop(pre_outage_subscription); + + let recovery_result = tokio::time::timeout( Duration::from_secs(OUTAGE_TIMEOUT_SECS), AssertUnwindSafe(async { - let peer: String = - nodes[victim_idx].rpc.request("system_localPeerId", rpc_params![]).await?; - ensure!(peer == original_peer, "victim PeerId changed: {original_peer} -> {peer}"); - let eligible = victim.reports(ELIGIBLE_PEERS_METRIC).await?; + let restarted_peer_id: String = + nodes[outage_node_idx].rpc.request("system_localPeerId", rpc_params![]).await?; + ensure!( + restarted_peer_id == original_peer_id, + "outage node PeerId changed: {original_peer_id} -> {restarted_peer_id}" + ); + let eligible = outage_node.reports(ELIGIBLE_PEERS_METRIC).await?; info!( - "Lifecycle restarted: unchanged PeerId {peer}, same base directory {}, \ + "Lifecycle restarted: unchanged PeerId {restarted_peer_id}, same base directory {}, \ {eligible} eligible peers known at the first RPC", - victim.base_dir().display() + outage_node.base_dir().display() ); - let mut subscriber_subscription = - subscribe_topic(&nodes[victim_idx].rpc, subscriber_topic).await?; + let mut outage_node_subscription = + subscribe_topic(&nodes[outage_node_idx].rpc, subscriber_topic).await?; assert_statements_match( - &mut subscriber_subscription, - &subscribed_statements, + &mut outage_node_subscription, + &expected_subscription_backlog, DELIVERY_TIMEOUT_SECS, - victim.name(), + outage_node.name(), ) .await?; - wait_lifecycle_topology("recovered-ready", nodes, None).await?; // Defer non-affine recovery placement: a cold topology can grant permanent DHT // retention. - let expected: Vec = cohort + let recovery_placements: Vec = expected_placements .iter() - .filter(|placement| placement.holders.contains(&victim_idx)) + .filter(|placement| placement.holder_node_indices.contains(&outage_node_idx)) .cloned() .collect(); - let placement = wait_cohort_placement("recovered", nodes, &expected, None).await; - let admissions = admissions_by_reason(victim).await?; - placement.with_context(|| { - format!("{} admissions by reason since restart: {admissions}", victim.name()) + let placement_result = + wait_for_stable_placement("recovered", nodes, &recovery_placements, None).await; + let admissions = admissions_by_reason(outage_node).await?; + placement_result.with_context(|| { + format!("{} admissions by reason since restart: {admissions}", outage_node.name()) })?; info!( "Lifecycle recovered: {} holds every required replica/subscription statement; \ admissions by reason: {admissions}", - victim.name() + outage_node.name() ); Ok::<_, anyhow::Error>(()) }) .catch_unwind(), ) .await; - let after_logs = archive_victim_log(victim, "after-restart").await; - if let Err(error) = &after_logs { - warn!("Lifecycle post-restart log check failed: {error:#}"); - } - match recovered { + + match recovery_result { Err(_) => return Err(anyhow!("replica recovery exceeded {OUTAGE_TIMEOUT_SECS}s")), Ok(Err(panic)) => std::panic::resume_unwind(panic), Ok(Ok(result)) => result?, } - after_logs?; - drop(donor_subscriptions); - Ok(cohort.len()) + drop(submitter_subscriptions); + Ok(expected_placements.len()) } struct WaveReport {