iuna

iuna

iuna - experimental mainnet-candidate protocol
git clone https://getiuna.org/git/iuna.git
Log | Files | Refs | README | LICENSE

properties.rs (23377B)


      1 #![cfg(feature = "e2e")]
      2 
      3 use std::{
      4     net::{Ipv4Addr, SocketAddr, TcpListener as StdTcpListener},
      5     path::Path,
      6     sync::Arc,
      7     time::Duration,
      8 };
      9 
     10 use anyhow::{Context, Result, bail};
     11 use iuna::{
     12     adapters::{
     13         chain_store::SqliteChainStore, p2p::GossipNetwork, stratum::StratumServer, wallet_store,
     14     },
     15     app::{NodeCore, PeerBook, SharedNode, SharedPeerBook, now_ms},
     16     domain::{
     17         FinalizerMode, Ledger, OBJECTIVE_FINALITY_ACTIVATION_HEIGHT, Transaction,
     18         VDF_TARGET_BLOCK_MS, Wallet, configure_e2e_vdf_round_divisor_for_tests, run_vdf,
     19         verify_vdf,
     20     },
     21 };
     22 use serde_json::{Value, json};
     23 use tempfile::tempdir;
     24 use tokio::{
     25     io::{AsyncBufReadExt, AsyncWriteExt, BufReader},
     26     net::TcpStream,
     27     sync::Mutex,
     28     time::{sleep, timeout},
     29 };
     30 
     31 const SOAK_BLOCKS: u64 = 12;
     32 const SOAK_VDF_ROUND_DIVISOR: u64 = 100;
     33 const BURN_COLLECTION_MS: u64 = VDF_TARGET_BLOCK_MS / 20 + 1;
     34 const SOAK_START_HEIGHT: u64 = OBJECTIVE_FINALITY_ACTIVATION_HEIGHT + 1;
     35 const FIXTURE_SERVICES: [&str; 6] = ["bootstrap", "node2", "node3", "node4", "node5", "node6"];
     36 
     37 #[tokio::test(flavor = "multi_thread", worker_threads = 4)]
     38 #[ignore = "long-running post-activation soak; run with cargo test --release --features e2e --test properties -- --ignored"]
     39 async fn release_soak_post_activation_auto_finalization_p2p_stratum_and_restarts() -> Result<()> {
     40     configure_e2e_vdf_round_divisor_for_tests(SOAK_VDF_ROUND_DIVISOR);
     41     let (wallets, genesis) = post_activation_fixture()?;
     42     let p2p_addrs = reserve_loopback_addrs(wallets.len())?;
     43     let stratum_addr = reserve_loopback_addrs(1)?.remove(0);
     44     let store_dirs = (0..wallets.len())
     45         .map(|_| tempdir())
     46         .collect::<std::result::Result<Vec<_>, _>>()?;
     47     let stores = store_dirs
     48         .iter()
     49         .map(|dir| SqliteChainStore::open(dir.path().join("chain.sqlite3")))
     50         .collect::<Result<Vec<_>>>()?;
     51     let mut nodes = Vec::new();
     52 
     53     for (index, wallet) in wallets.iter().cloned().enumerate() {
     54         let burn_per_block = if index == 0 { 2 } else { 0 };
     55         let mut core = NodeCore::from_ledger_with_burn_fee_and_enabled(
     56             wallet,
     57             genesis.clone(),
     58             true,
     59             burn_per_block,
     60             1,
     61         );
     62         core.set_recovery_vdf_top_rank_percent(0);
     63         let node = Arc::new(Mutex::new(core));
     64         let peer_addresses = p2p_addrs
     65             .iter()
     66             .enumerate()
     67             .filter(|(peer_index, _)| *peer_index != index)
     68             .map(|(_, addr)| addr.to_string())
     69             .collect::<Vec<_>>();
     70         let peers = Arc::new(Mutex::new(PeerBook::from_addresses(peer_addresses)));
     71         let network =
     72             GossipNetwork::start(node.clone(), peers.clone(), p2p_addrs[index], None, true).await?;
     73         nodes.push(SoakNode {
     74             wallet: wallets[index].clone(),
     75             burn_per_block,
     76             node,
     77             peers,
     78             network,
     79             store: stores[index].clone(),
     80         });
     81     }
     82 
     83     let _stratum = StratumServer::start(
     84         nodes[0].node.clone(),
     85         nodes[0].network.clone(),
     86         stratum_addr,
     87     )
     88     .await?;
     89     let stratum_worker = nodes[1]
     90         .node
     91         .lock()
     92         .await
     93         .wallet_receive_address()
     94         .context("release-soak wallet must have a checksummed receive address")?;
     95     assert_stratum_serves_work(stratum_addr, &stratum_worker).await?;
     96     sleep(Duration::from_secs(2)).await;
     97 
     98     for target_height in (SOAK_START_HEIGHT + 1)..=(SOAK_START_HEIGHT + SOAK_BLOCKS) {
     99         finalize_one_block(&nodes, target_height, None).await?;
    100         wait_for_convergence(&nodes, target_height, Duration::from_secs(8)).await?;
    101 
    102         if target_height % 3 == 0 {
    103             restart_node_core(&nodes[1]).await?;
    104             wait_for_convergence(&nodes, target_height, Duration::from_secs(8)).await?;
    105         }
    106         if target_height % 4 == 0 {
    107             restart_node_core(&nodes[2]).await?;
    108             wait_for_convergence(&nodes, target_height, Duration::from_secs(8)).await?;
    109         }
    110     }
    111 
    112     let final_tip = nodes[0].node.lock().await.chain_tip_hash();
    113     for node in &nodes {
    114         let core = node.node.lock().await;
    115         assert_eq!(core.chain_tip_hash(), final_tip);
    116         assert!(core.chain_height() > OBJECTIVE_FINALITY_ACTIVATION_HEIGHT);
    117         assert!(
    118             core.status()
    119                 .chain
    120                 .finalized_height
    121                 .is_some_and(|height| height >= OBJECTIVE_FINALITY_ACTIVATION_HEIGHT)
    122         );
    123     }
    124     assert!(configured_automatic_burn_was_included(&nodes[0]).await);
    125     Ok(())
    126 }
    127 
    128 #[tokio::test(flavor = "multi_thread", worker_threads = 4)]
    129 #[ignore = "accelerated post-activation P2P recovery partition; run with cargo test --release --features e2e --test properties -- --ignored"]
    130 async fn post_activation_p2p_partition_recovers_converges_and_resumes_tickets() -> Result<()> {
    131     configure_e2e_vdf_round_divisor_for_tests(SOAK_VDF_ROUND_DIVISOR);
    132     let (wallets, parent) = post_activation_fixture()?;
    133     let wallets = &wallets[..4];
    134     let p2p_addrs = reserve_loopback_addrs(wallets.len())?;
    135     let store_dirs = (0..wallets.len())
    136         .map(|_| tempdir())
    137         .collect::<std::result::Result<Vec<_>, _>>()?;
    138     let stores = store_dirs
    139         .iter()
    140         .map(|dir| SqliteChainStore::open(dir.path().join("chain.sqlite3")))
    141         .collect::<Result<Vec<_>>>()?;
    142     let mut nodes = Vec::new();
    143 
    144     for (index, wallet) in wallets.iter().cloned().enumerate() {
    145         let island = if index < 2 { 0..2 } else { 2..4 };
    146         let peer_addresses = island
    147             .filter(|peer_index| *peer_index != index)
    148             .map(|peer_index| p2p_addrs[peer_index].to_string())
    149             .collect::<Vec<_>>();
    150         let mut core = NodeCore::from_ledger_with_burn_fee_and_enabled(
    151             wallet.clone(),
    152             parent.clone(),
    153             true,
    154             2,
    155             1,
    156         );
    157         core.set_recovery_vdf_top_rank_percent(100);
    158         let node = Arc::new(Mutex::new(core));
    159         let peers = Arc::new(Mutex::new(PeerBook::from_addresses(peer_addresses)));
    160         let network =
    161             GossipNetwork::start(node.clone(), peers.clone(), p2p_addrs[index], None, true).await?;
    162         nodes.push(SoakNode {
    163             wallet,
    164             burn_per_block: 2,
    165             node,
    166             peers,
    167             network,
    168             store: stores[index].clone(),
    169         });
    170     }
    171     sleep(Duration::from_secs(2)).await;
    172 
    173     let recovery_height = parent.height() + 1;
    174     let recovery_timestamp = parent.recovery_block_min_timestamp();
    175     let left_block = produce_recovery_block(&nodes[0], recovery_timestamp).await?;
    176     let right_block = produce_recovery_block(&nodes[2], recovery_timestamp).await?;
    177     assert_ne!(left_block.hash, right_block.hash);
    178     wait_for_convergence(&nodes[..2], recovery_height, Duration::from_secs(8)).await?;
    179     wait_for_convergence(&nodes[2..], recovery_height, Duration::from_secs(8)).await?;
    180     assert_eq!(nodes[0].node.lock().await.chain_tip_hash(), left_block.hash);
    181     assert_eq!(
    182         nodes[2].node.lock().await.chain_tip_hash(),
    183         right_block.hash
    184     );
    185 
    186     nodes[0]
    187         .peers
    188         .lock()
    189         .await
    190         .add_peer(p2p_addrs[2].to_string());
    191     nodes[2]
    192         .peers
    193         .lock()
    194         .await
    195         .add_peer(p2p_addrs[0].to_string());
    196     wait_for_convergence(&nodes, recovery_height, Duration::from_secs(20)).await?;
    197 
    198     let converged_recovery_hash = nodes[0].node.lock().await.chain_tip_hash();
    199     for node in &nodes {
    200         let core = node.node.lock().await;
    201         assert_eq!(core.chain_tip_hash(), converged_recovery_hash);
    202         assert_eq!(
    203             core.chain()
    204                 .last()
    205                 .context("P2P recovery chain has no tip")?
    206                 .finalizer_mode,
    207             FinalizerMode::Recovery
    208         );
    209     }
    210 
    211     restart_node_core(&nodes[3]).await?;
    212     wait_for_convergence(&nodes, recovery_height, Duration::from_secs(8)).await?;
    213     // The committed fixture has historical timestamps, so wall-clock time is
    214     // already beyond the next recovery deadline. Exercise ticket resumption
    215     // immediately after the recovery block instead.
    216     finalize_one_block(
    217         &nodes,
    218         recovery_height + 1,
    219         Some(recovery_timestamp.saturating_add(1)),
    220     )
    221     .await?;
    222     wait_for_convergence(&nodes, recovery_height + 1, Duration::from_secs(8)).await?;
    223     for node in &nodes {
    224         let core = node.node.lock().await;
    225         let tip = core
    226             .chain()
    227             .last()
    228             .context("continued P2P chain has no tip")?;
    229         assert_eq!(tip.finalizer_mode, FinalizerMode::Ticket);
    230         assert_eq!(tip.finalizer_rank, 0);
    231         assert_eq!(
    232             core.ledger().objective_finality_checkpoint(),
    233             Some((recovery_height, converged_recovery_hash.as_str()))
    234         );
    235     }
    236     Ok(())
    237 }
    238 
    239 #[test]
    240 #[ignore = "accelerated post-activation recovery soak; run with cargo test --release --features e2e --test properties -- --ignored"]
    241 fn post_activation_recovery_candidates_converge_and_ticket_finalization_resumes() -> Result<()> {
    242     configure_e2e_vdf_round_divisor_for_tests(SOAK_VDF_ROUND_DIVISOR);
    243     let (wallets, parent) = post_activation_fixture()?;
    244     let recovery_height = parent.height() + 1;
    245     let recovery_timestamp = parent.recovery_block_min_timestamp();
    246 
    247     let mut branches = Vec::new();
    248     for wallet in &wallets {
    249         let mut candidate = NodeCore::from_ledger_with_burn_fee_and_enabled(
    250             wallet.clone(),
    251             parent.clone(),
    252             true,
    253             2,
    254             1,
    255         );
    256         candidate.set_recovery_vdf_top_rank_percent(100);
    257         let _ = candidate.prepare_automatic_finalization(recovery_timestamp);
    258         let Ok(mut branch) = candidate.wallet_view_ledger() else {
    259             continue;
    260         };
    261         let block = branch.mine_recovery_block(wallet, recovery_timestamp)?;
    262         assert_eq!(block.height, recovery_height);
    263         assert_eq!(block.finalizer_mode, FinalizerMode::Recovery);
    264         assert!(verify_vdf(
    265             &block.vdf_seed(),
    266             block.vdf_rounds,
    267             &block.vdf_output
    268         ));
    269         branch.apply_locally_mined_block(block.clone())?;
    270         branches.push((branch, block));
    271         if branches.len() == 2 {
    272             break;
    273         }
    274     }
    275     if branches.len() < 2 {
    276         bail!("post-activation fixture needs two funded recovery candidates");
    277     }
    278 
    279     let (mut left, left_block) = branches.remove(0);
    280     let (mut right, right_block) = branches.remove(0);
    281     assert_ne!(left_block.hash, right_block.hash);
    282     let left_snapshot = left.snapshot();
    283     let right_snapshot = right.snapshot();
    284 
    285     left.extend_from_preverified_snapshot_for_e2e(right_snapshot)?;
    286     right.extend_from_preverified_snapshot_for_e2e(left_snapshot)?;
    287     assert_eq!(left.height(), recovery_height);
    288     assert_eq!(right.height(), recovery_height);
    289     assert_eq!(left.tip_hash(), right.tip_hash());
    290     assert_eq!(
    291         left.chain()
    292             .last()
    293             .context("converged recovery chain has no tip")?
    294             .finalizer_mode,
    295         FinalizerMode::Recovery
    296     );
    297 
    298     let recovery_hash = left.tip_hash().to_string();
    299     let leader = left
    300         .expected_leader_for_next_block()
    301         .context("ticket finalization did not resume after recovery")?;
    302     let leader_wallet = wallets
    303         .iter()
    304         .find(|wallet| wallet.address() == leader)
    305         .context("selected post-recovery leader is absent from the fixture")?;
    306     let mut leader_node =
    307         NodeCore::from_ledger_with_burn_fee_and_enabled(leader_wallet.clone(), left, true, 2, 1);
    308     leader_node.set_recovery_vdf_top_rank_percent(100);
    309     let ticket_timestamp = recovery_timestamp.saturating_add(1);
    310     let _ = leader_node.prepare_automatic_finalization(ticket_timestamp);
    311     let mut continued = leader_node.wallet_view_ledger()?;
    312     let burn_bundles = wallets.iter().try_fold(Vec::new(), |mut bundles, wallet| {
    313         bundles.extend(continued.build_burn_bundles(wallet)?);
    314         Ok::<_, anyhow::Error>(bundles)
    315     })?;
    316     let prepared = continued.prepare_next_block_with_burn_bundles(
    317         leader_wallet.address(),
    318         ticket_timestamp,
    319         burn_bundles,
    320     )?;
    321     let output = run_vdf(prepared.vdf_seed(), prepared.vdf_rounds());
    322     let ticket_block = prepared.finish(leader_wallet, output);
    323     assert_eq!(ticket_block.finalizer_mode, FinalizerMode::Ticket);
    324     assert_eq!(ticket_block.finalizer_rank, 0);
    325     assert!(verify_vdf(
    326         &ticket_block.vdf_seed(),
    327         ticket_block.vdf_rounds,
    328         &ticket_block.vdf_output
    329     ));
    330     continued.apply_locally_mined_block(ticket_block)?;
    331 
    332     assert_eq!(continued.height(), recovery_height + 1);
    333     assert_eq!(
    334         continued.objective_finality_checkpoint(),
    335         Some((recovery_height, recovery_hash.as_str()))
    336     );
    337     Ok(())
    338 }
    339 
    340 struct SoakNode {
    341     wallet: Wallet,
    342     burn_per_block: u64,
    343     node: SharedNode,
    344     peers: SharedPeerBook,
    345     network: GossipNetwork,
    346     store: SqliteChainStore,
    347 }
    348 
    349 fn post_activation_fixture() -> Result<(Vec<Wallet>, Ledger)> {
    350     if !cfg!(feature = "e2e") {
    351         bail!("post-activation soak requires --features e2e");
    352     }
    353 
    354     let fixture =
    355         Path::new(env!("CARGO_MANIFEST_DIR")).join("e2e/snapshots/first-objective-checkpoint");
    356     let wallets = FIXTURE_SERVICES
    357         .iter()
    358         .map(|service| {
    359             wallet_store::load_with_password(
    360                 &fixture.join(service).join("wallet.json"),
    361                 "testtesttest",
    362             )
    363         })
    364         .collect::<Result<Vec<_>>>()?;
    365     // SQLite may create journals while opening a database; never open the committed fixture in place.
    366     let chain_copy_dir = tempdir()?;
    367     let chain_copy = chain_copy_dir.path().join("chain.sqlite3");
    368     std::fs::copy(fixture.join("bootstrap/chain.sqlite3"), &chain_copy)?;
    369     let snapshot = SqliteChainStore::open(chain_copy)?
    370         .load()?
    371         .context("post-activation checkpoint has no chain snapshot")?;
    372     let ledger = Ledger::from_persisted_snapshot(snapshot)?;
    373     if ledger.height() != SOAK_START_HEIGHT {
    374         bail!(
    375             "post-activation checkpoint height is {}, expected {SOAK_START_HEIGHT}",
    376             ledger.height()
    377         );
    378     }
    379     Ok((wallets, ledger))
    380 }
    381 
    382 fn reserve_loopback_addrs(count: usize) -> Result<Vec<SocketAddr>> {
    383     let mut listeners = Vec::new();
    384     let mut addrs = Vec::new();
    385     for _ in 0..count {
    386         let listener = StdTcpListener::bind((Ipv4Addr::LOCALHOST, 0))?;
    387         addrs.push(listener.local_addr()?);
    388         listeners.push(listener);
    389     }
    390     drop(listeners);
    391     Ok(addrs)
    392 }
    393 
    394 async fn assert_stratum_serves_work(addr: SocketAddr, worker: &str) -> Result<()> {
    395     let stream = timeout(Duration::from_secs(5), TcpStream::connect(addr)).await??;
    396     let (read, mut write) = stream.into_split();
    397     let mut lines = BufReader::new(read).lines();
    398 
    399     write
    400         .write_all(
    401             json_line(json!({"id": 1, "method": "mining.subscribe", "params": []}))?.as_bytes(),
    402         )
    403         .await?;
    404     write
    405         .write_all(
    406             json_line(json!({"id": 2, "method": "mining.authorize", "params": [worker, "x"]}))?
    407                 .as_bytes(),
    408         )
    409         .await?;
    410 
    411     let mut authorized = false;
    412     let mut notified = false;
    413     for _ in 0..4 {
    414         let line = timeout(Duration::from_secs(5), lines.next_line())
    415             .await??
    416             .context("stratum server closed before sending work")?;
    417         let value: Value = serde_json::from_str(&line)?;
    418         authorized |=
    419             value.get("id") == Some(&json!(2)) && value.get("result") == Some(&json!(true));
    420         notified |= value.get("method") == Some(&json!("mining.notify"));
    421         if authorized && notified {
    422             return Ok(());
    423         }
    424     }
    425     bail!("stratum did not authorize and send mining.notify")
    426 }
    427 
    428 fn json_line(value: Value) -> Result<String> {
    429     Ok(format!("{}\n", serde_json::to_string(&value)?))
    430 }
    431 
    432 async fn finalize_one_block(
    433     nodes: &[SoakNode],
    434     target_height: u64,
    435     fixed_timestamp_ms: Option<u64>,
    436 ) -> Result<()> {
    437     let timestamp_ms = || fixed_timestamp_ms.unwrap_or_else(now_ms);
    438     let start = timestamp_ms().saturating_sub(BURN_COLLECTION_MS + 1);
    439     prepare_and_broadcast(nodes, start).await?;
    440     sleep(Duration::from_millis(250)).await;
    441     // Post-activation finalizers must see the complete committee quorum before preparing VDF work.
    442     for _ in 0..3 {
    443         prepare_and_broadcast(nodes, timestamp_ms()).await?;
    444         sleep(Duration::from_millis(100)).await;
    445     }
    446 
    447     let deadline = tokio::time::Instant::now() + Duration::from_secs(8);
    448     loop {
    449         for node in nodes {
    450             let block = complete_if_ready(node, timestamp_ms()).await?;
    451             node.network
    452                 .broadcast(node.node.lock().await.drain_outbox())
    453                 .await?;
    454             if let Some(block) = block {
    455                 assert_eq!(block.height, target_height);
    456                 return Ok(());
    457             }
    458         }
    459         if tokio::time::Instant::now() >= deadline {
    460             bail!(
    461                 "no node finalized block {target_height}:\n{}",
    462                 soak_diagnostics(nodes, target_height).await
    463             );
    464         }
    465         sleep(Duration::from_millis(100)).await;
    466     }
    467 }
    468 
    469 async fn soak_diagnostics(nodes: &[SoakNode], target_height: u64) -> String {
    470     let mut lines = Vec::new();
    471     for (index, node) in nodes.iter().enumerate() {
    472         let core = node.node.lock().await;
    473         let status = core.status();
    474         let rank = core
    475             .burn_leader_ranks_for_block(target_height)
    476             .ok()
    477             .and_then(|ranks| {
    478                 ranks
    479                     .into_iter()
    480                     .find(|rank| rank.owner == core.wallet_address())
    481                     .map(|rank| rank.rank)
    482             });
    483         let pending = core.pending_transactions();
    484         let pending_burns = pending.iter().filter(|tx| tx.is_burn()).count();
    485         let pending_burn_details = pending
    486             .iter()
    487             .filter(|tx| tx.is_burn())
    488             .map(|tx| {
    489                 let Transaction::Burn { anchor, .. } = tx else {
    490                     unreachable!();
    491                 };
    492                 let eligible_anchor = core.chain()[core.chain().len() - 1].prev_hash.as_str();
    493                 format!(
    494                     "sender={}, amount={}, anchor={:?}, eligible={}",
    495                     tx.sender(),
    496                     tx.amount(),
    497                     anchor,
    498                     anchor.as_deref() == Some(eligible_anchor)
    499                 )
    500             })
    501             .collect::<Vec<_>>();
    502         let wallet_view_pending = core
    503             .wallet_view_ledger()
    504             .map(|ledger| ledger.pending().len())
    505             .unwrap_or_default();
    506         let metrics = node.network.metrics();
    507         lines.push(format!(
    508             "node {index}: height={}, tip={}, wallet_rank={rank:?}, leader={:?}, \
    509              last_finalization={:?}, pending={} (burns={pending_burns}, \
    510              wallet_view={wallet_view_pending}, details={pending_burn_details:?}), \
    511              burn_bundles_received={}, control_received={}, rejected_blocks={}, \
    512              session_failures={}, last_session_failure={:?}, last_chain_error={:?}, \
    513              sync_progress={:?}",
    514             status.chain.height,
    515             status.chain.tip_hash,
    516             status.mining.current_leader,
    517             status.mining.last_auto_finalization_status,
    518             pending.len(),
    519             metrics.burn_bundles_received,
    520             metrics.control_envelopes_received,
    521             metrics.rejected_blocks,
    522             metrics.session_failures,
    523             metrics.last_session_failure,
    524             metrics.last_chain_payload_error,
    525             node.network.sync_progress(),
    526         ));
    527     }
    528     lines.join("\n")
    529 }
    530 
    531 async fn prepare_and_broadcast(nodes: &[SoakNode], timestamp_ms: u64) -> Result<()> {
    532     for node in nodes {
    533         {
    534             let mut core = node.node.lock().await;
    535             let _ = core.prepare_automatic_finalization(timestamp_ms);
    536         }
    537         node.network
    538             .broadcast(node.node.lock().await.drain_outbox())
    539             .await?;
    540     }
    541     Ok(())
    542 }
    543 
    544 async fn complete_if_ready(
    545     node: &SoakNode,
    546     timestamp_ms: u64,
    547 ) -> Result<Option<iuna::domain::Block>> {
    548     let work = {
    549         let mut core = node.node.lock().await;
    550         core.prepare_automatic_finalization(timestamp_ms).work
    551     };
    552     let Some(work) = work else {
    553         return Ok(None);
    554     };
    555     let vdf_output = run_vdf(work.vdf_seed(), work.vdf_rounds());
    556     let block =
    557         node.node
    558             .lock()
    559             .await
    560             .complete_prepared_block_at(work, vdf_output, timestamp_ms)?;
    561     Ok(Some(block))
    562 }
    563 
    564 async fn produce_recovery_block(node: &SoakNode, timestamp_ms: u64) -> Result<iuna::domain::Block> {
    565     let work = {
    566         let mut core = node.node.lock().await;
    567         let _ = core.prepare_automatic_finalization(timestamp_ms);
    568         core.wallet_view_ledger()?
    569             .prepare_recovery_block(node.wallet.address(), timestamp_ms)?
    570     };
    571     let vdf_output = run_vdf(work.vdf_seed(), work.vdf_rounds());
    572     let block =
    573         node.node
    574             .lock()
    575             .await
    576             .complete_prepared_block_at(work, vdf_output, timestamp_ms)?;
    577     node.network
    578         .broadcast(node.node.lock().await.drain_outbox())
    579         .await?;
    580     Ok(block)
    581 }
    582 
    583 async fn wait_for_convergence(nodes: &[SoakNode], height: u64, duration: Duration) -> Result<()> {
    584     let deadline = tokio::time::Instant::now() + duration;
    585     loop {
    586         let mut tips = Vec::new();
    587         for node in nodes {
    588             let core = node.node.lock().await;
    589             tips.push((core.chain_height(), core.chain_tip_hash()));
    590         }
    591         if tips
    592             .iter()
    593             .all(|(node_height, tip)| *node_height >= height && tip == &tips[0].1)
    594         {
    595             return Ok(());
    596         }
    597         if tokio::time::Instant::now() >= deadline {
    598             bail!(
    599                 "nodes did not converge at height {height}: {tips:?}\n{}",
    600                 soak_diagnostics(nodes, height).await
    601             );
    602         }
    603         sleep(Duration::from_millis(250)).await;
    604     }
    605 }
    606 
    607 async fn configured_automatic_burn_was_included(node: &SoakNode) -> bool {
    608     node.node
    609         .lock()
    610         .await
    611         .chain()
    612         .iter()
    613         .filter(|block| block.height > SOAK_START_HEIGHT)
    614         .flat_map(|block| &block.transactions)
    615         .any(|tx| tx.is_burn() && tx.sender() == node.wallet.address() && tx.amount() == 2)
    616 }
    617 
    618 async fn restart_node_core(node: &SoakNode) -> Result<()> {
    619     let snapshot = node.node.lock().await.chain_snapshot();
    620     node.store.save(&snapshot)?;
    621     let restored = Ledger::from_persisted_snapshot(
    622         node.store
    623             .load()?
    624             .context("persisted snapshot should exist after save")?,
    625     )?;
    626     let mut restored_node = NodeCore::from_ledger_with_burn_fee_and_enabled(
    627         node.wallet.clone(),
    628         restored,
    629         true,
    630         node.burn_per_block,
    631         1,
    632     );
    633     restored_node.set_recovery_vdf_top_rank_percent(0);
    634     *node.node.lock().await = restored_node;
    635     Ok(())
    636 }