commit 1d96d46f7c7512b1db287a38b7972ed289243222
parent 1d6a115a68a47a7f56df89eabea531978db876c8
Author: Joris Hartog <jorishartog@hotmail.com>
Date: Sat, 5 Sep 2026 13:04:06 +0200
Fix P2P convergence after recovery partitions
Diffstat:
4 files changed, 171 insertions(+), 5 deletions(-)
diff --git a/src/adapters/p2p/fetch.rs b/src/adapters/p2p/fetch.rs
@@ -238,8 +238,17 @@ pub(super) async fn validate_blocks_extension(
.iter()
.position(|block| block.hash == blocks[0].prev_hash)
.context("block page has no common ancestor with local chain")?;
+ #[cfg(feature = "e2e")]
+ for block in &blocks {
+ if !verify_vdf(&block.vdf_seed(), block.vdf_rounds, &block.vdf_output) {
+ anyhow::bail!("block VDF output is invalid");
+ }
+ }
candidate.blocks.truncate(ancestor + 1);
candidate.blocks.extend(blocks);
+ #[cfg(feature = "e2e")]
+ ledger.extend_from_preverified_snapshot_for_e2e(candidate)?;
+ #[cfg(not(feature = "e2e"))]
ledger.extend_from_snapshot_at(candidate, now_ms)?;
if let Some(target_height) = target_height {
on_progress(target_height);
diff --git a/src/app/receive.rs b/src/app/receive.rs
@@ -168,8 +168,13 @@ impl NodeCore {
anyhow::bail!("chain snapshot genesis does not match local chain");
}
let previous_height = self.ledger.height();
- if !replaces_setup_placeholder && ledger.height() <= previous_height {
- return Ok(false);
+ if !replaces_setup_placeholder {
+ if ledger.height() < previous_height {
+ return Ok(false);
+ }
+ if ledger.height() == previous_height && ledger.tip_hash() == self.ledger.tip_hash() {
+ return Ok(false);
+ }
}
self.ledger = ledger;
@@ -225,6 +230,30 @@ mod tests {
}
#[test]
+ fn verified_same_height_fork_replaces_local_ledger() {
+ let wallet = Wallet::from_seed("verified-same-height-fork");
+ let parent = funded_ledger(std::slice::from_ref(&wallet));
+ let mut local = parent.clone();
+ let burn = local.build_burn(&wallet, 1, 1).unwrap();
+ local.submit_transaction(burn.clone()).unwrap();
+ let local_block = local.mine_next_block(&wallet, 1).unwrap();
+ local.apply_locally_mined_block(local_block).unwrap();
+
+ let mut remote = parent;
+ remote.submit_transaction(burn).unwrap();
+ let remote_block = remote.mine_next_block(&wallet, 2).unwrap();
+ remote
+ .apply_locally_mined_block(remote_block.clone())
+ .unwrap();
+ assert_eq!(local.height(), remote.height());
+ assert_ne!(local.tip_hash(), remote.tip_hash());
+
+ let mut node = NodeCore::from_ledger(wallet, local, 0);
+ assert!(node.import_verified_ledger(remote).unwrap());
+ assert_eq!(node.ledger().tip_hash(), remote_block.hash);
+ }
+
+ #[test]
fn gossiped_mine_over_anchor_limit_is_silently_ignored() {
let wallets = (0..4)
.map(|index| {
diff --git a/src/domain/ledger_chain.rs b/src/domain/ledger_chain.rs
@@ -253,6 +253,7 @@ impl Ledger {
self.extend_from_snapshot_with_vdf_policy(snapshot, false, u64::MAX)
}
+ #[cfg_attr(feature = "e2e", allow(dead_code))]
pub(crate) fn extend_from_snapshot_at(
&mut self,
snapshot: ChainSnapshot,
diff --git a/tests/properties.rs b/tests/properties.rs
@@ -12,7 +12,7 @@ use iuna::{
adapters::{
chain_store::SqliteChainStore, p2p::GossipNetwork, stratum::StratumServer, wallet_store,
},
- app::{NodeCore, PeerBook, SharedNode, now_ms},
+ app::{NodeCore, PeerBook, SharedNode, SharedPeerBook, now_ms},
domain::{
FinalizerMode, Ledger, OBJECTIVE_FINALITY_ACTIVATION_HEIGHT, Transaction,
VDF_TARGET_BLOCK_MS, Wallet, configure_e2e_vdf_round_divisor_for_tests, run_vdf,
@@ -69,11 +69,12 @@ async fn release_soak_post_activation_auto_finalization_p2p_stratum_and_restarts
.collect::<Vec<_>>();
let peers = Arc::new(Mutex::new(PeerBook::from_addresses(peer_addresses)));
let network =
- GossipNetwork::start(node.clone(), peers, p2p_addrs[index], None, true).await?;
+ GossipNetwork::start(node.clone(), peers.clone(), p2p_addrs[index], None, true).await?;
nodes.push(SoakNode {
wallet: wallets[index].clone(),
burn_per_block,
node,
+ peers,
network,
store: stores[index].clone(),
});
@@ -124,6 +125,109 @@ async fn release_soak_post_activation_auto_finalization_p2p_stratum_and_restarts
Ok(())
}
+#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
+#[ignore = "accelerated post-activation P2P recovery partition; run with cargo test --release --features e2e --test properties -- --ignored"]
+async fn post_activation_p2p_partition_recovers_converges_and_resumes_tickets() -> Result<()> {
+ configure_e2e_vdf_round_divisor_for_tests(SOAK_VDF_ROUND_DIVISOR);
+ let (wallets, parent) = post_activation_fixture()?;
+ let wallets = &wallets[..4];
+ let p2p_addrs = reserve_loopback_addrs(wallets.len())?;
+ let store_dirs = (0..wallets.len())
+ .map(|_| tempdir())
+ .collect::<std::result::Result<Vec<_>, _>>()?;
+ let stores = store_dirs
+ .iter()
+ .map(|dir| SqliteChainStore::open(dir.path().join("chain.sqlite3")))
+ .collect::<Result<Vec<_>>>()?;
+ let mut nodes = Vec::new();
+
+ for (index, wallet) in wallets.iter().cloned().enumerate() {
+ let island = if index < 2 { 0..2 } else { 2..4 };
+ let peer_addresses = island
+ .filter(|peer_index| *peer_index != index)
+ .map(|peer_index| p2p_addrs[peer_index].to_string())
+ .collect::<Vec<_>>();
+ let mut core = NodeCore::from_ledger_with_burn_fee_and_enabled(
+ wallet.clone(),
+ parent.clone(),
+ true,
+ 2,
+ 1,
+ );
+ core.set_recovery_vdf_top_rank_percent(100);
+ let node = Arc::new(Mutex::new(core));
+ let peers = Arc::new(Mutex::new(PeerBook::from_addresses(peer_addresses)));
+ let network =
+ GossipNetwork::start(node.clone(), peers.clone(), p2p_addrs[index], None, true).await?;
+ nodes.push(SoakNode {
+ wallet,
+ burn_per_block: 2,
+ node,
+ peers,
+ network,
+ store: stores[index].clone(),
+ });
+ }
+ sleep(Duration::from_secs(2)).await;
+
+ let recovery_height = parent.height() + 1;
+ let recovery_timestamp = parent.recovery_block_min_timestamp();
+ let left_block = produce_recovery_block(&nodes[0], recovery_timestamp).await?;
+ let right_block = produce_recovery_block(&nodes[2], recovery_timestamp).await?;
+ assert_ne!(left_block.hash, right_block.hash);
+ wait_for_convergence(&nodes[..2], recovery_height, Duration::from_secs(8)).await?;
+ wait_for_convergence(&nodes[2..], recovery_height, Duration::from_secs(8)).await?;
+ assert_eq!(nodes[0].node.lock().await.chain_tip_hash(), left_block.hash);
+ assert_eq!(
+ nodes[2].node.lock().await.chain_tip_hash(),
+ right_block.hash
+ );
+
+ nodes[0]
+ .peers
+ .lock()
+ .await
+ .add_peer(p2p_addrs[2].to_string());
+ nodes[2]
+ .peers
+ .lock()
+ .await
+ .add_peer(p2p_addrs[0].to_string());
+ wait_for_convergence(&nodes, recovery_height, Duration::from_secs(20)).await?;
+
+ let converged_recovery_hash = nodes[0].node.lock().await.chain_tip_hash();
+ for node in &nodes {
+ let core = node.node.lock().await;
+ assert_eq!(core.chain_tip_hash(), converged_recovery_hash);
+ assert_eq!(
+ core.chain()
+ .last()
+ .context("P2P recovery chain has no tip")?
+ .finalizer_mode,
+ FinalizerMode::Recovery
+ );
+ }
+
+ restart_node_core(&nodes[3]).await?;
+ wait_for_convergence(&nodes, recovery_height, Duration::from_secs(8)).await?;
+ finalize_one_block(&nodes, recovery_height + 1).await?;
+ wait_for_convergence(&nodes, recovery_height + 1, Duration::from_secs(8)).await?;
+ for node in &nodes {
+ let core = node.node.lock().await;
+ let tip = core
+ .chain()
+ .last()
+ .context("continued P2P chain has no tip")?;
+ assert_eq!(tip.finalizer_mode, FinalizerMode::Ticket);
+ assert_eq!(tip.finalizer_rank, 0);
+ assert_eq!(
+ core.ledger().objective_finality_checkpoint(),
+ Some((recovery_height, converged_recovery_hash.as_str()))
+ );
+ }
+ Ok(())
+}
+
#[test]
#[ignore = "accelerated post-activation recovery soak; run with cargo test --release --features e2e --test properties -- --ignored"]
fn post_activation_recovery_candidates_converge_and_ticket_finalization_resumes() -> Result<()> {
@@ -229,6 +333,7 @@ struct SoakNode {
wallet: Wallet,
burn_per_block: u64,
node: SharedNode,
+ peers: SharedPeerBook,
network: GossipNetwork,
store: SqliteChainStore,
}
@@ -443,6 +548,25 @@ async fn complete_if_ready(
Ok(Some(block))
}
+async fn produce_recovery_block(node: &SoakNode, timestamp_ms: u64) -> Result<iuna::domain::Block> {
+ let work = {
+ let mut core = node.node.lock().await;
+ let _ = core.prepare_automatic_finalization(timestamp_ms);
+ core.wallet_view_ledger()?
+ .prepare_recovery_block(node.wallet.address(), timestamp_ms)?
+ };
+ let vdf_output = run_vdf(work.vdf_seed(), work.vdf_rounds());
+ let block =
+ node.node
+ .lock()
+ .await
+ .complete_prepared_block_at(work, vdf_output, timestamp_ms)?;
+ node.network
+ .broadcast(node.node.lock().await.drain_outbox())
+ .await?;
+ Ok(block)
+}
+
async fn wait_for_convergence(nodes: &[SoakNode], height: u64, duration: Duration) -> Result<()> {
let deadline = tokio::time::Instant::now() + duration;
loop {
@@ -458,7 +582,10 @@ async fn wait_for_convergence(nodes: &[SoakNode], height: u64, duration: Duratio
return Ok(());
}
if tokio::time::Instant::now() >= deadline {
- bail!("nodes did not converge at height {height}: {tips:?}");
+ bail!(
+ "nodes did not converge at height {height}: {tips:?}\n{}",
+ soak_diagnostics(nodes, height).await
+ );
}
sleep(Duration::from_millis(250)).await;
}