iuna

iuna

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

commit e327e8555908eaeeec857ab5f3422a4f602162cb
parent 0ab272120415c0ecff82d00f7ad7e4fd66547459
Author: Joris Hartog <jorishartog@hotmail.com>
Date:   Sun, 23 Aug 2026 23:47:12 +0200

Page P2P chain synchronization

Diffstat:
Mdocker-compose.yml | 3+--
Msrc/adapters/http/actions.rs | 2+-
Msrc/adapters/p2p.rs | 7+++----
Msrc/adapters/p2p/fetch.rs | 145++++++++++++++++++++++++++++++++++++++++++++++++++++++-------------------------
Msrc/adapters/p2p/handshake.rs | 12++++++------
Msrc/adapters/p2p/line_codec.rs | 52+++++++++++++++++++++++-----------------------------
Msrc/adapters/p2p/peer_addr.rs | 2+-
Msrc/adapters/p2p/peer_status.rs | 24++++++++++++------------
Msrc/adapters/p2p/process.rs | 74++++++++++++++++++++++++++++++++++++++++++++++++--------------------------
Msrc/adapters/p2p/sync.rs | 82++++++++++++++++++++++++++++++++++++++++----------------------------------------
Msrc/adapters/p2p/tests.rs | 87+++++++++++++++++++++++++++++++++++++++++++++++++++++++------------------------
Msrc/adapters/p2p/writer.rs | 78++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Msrc/app.rs | 5+++--
Msrc/app/automatic_mining.rs | 2+-
Msrc/app/gossip.rs | 24++++++++++++++++++++++--
Msrc/app/receive.rs | 5+++--
Msrc/app/types.rs | 20+++++++++++++++++---
Msrc/domain/ledger_queries.rs | 38++++++++++++++++++++++++++++++++++++++
18 files changed, 459 insertions(+), 203 deletions(-)

diff --git a/docker-compose.yml b/docker-compose.yml @@ -23,8 +23,7 @@ services: IUNA_SETUP_COMPLETE: "true" IUNA_LOCAL_TESTNET: "true" IUNA_AUTOMATIC_BURN_ENABLED: "true" - IUNA_POW_MINING_ENABLED: "true" - IUNA_POW_MINING_WORKERS: "1" + IUNA_POW_MINING_ENABLED: "false" entrypoint: ["/bin/sh", "-c"] command: - if [ -s /data/chain.sqlite3 ]; then exec iuna --data-dir /data --http 0.0.0.0:18661 --p2p 0.0.0.0:9444 --p2p-announce 172.28.0.10:9444 --debug; else exec iuna --genesis --data-dir /data --http 0.0.0.0:18661 --p2p 0.0.0.0:9444 --p2p-announce 172.28.0.10:9444 --debug; fi diff --git a/src/adapters/http/actions.rs b/src/adapters/http/actions.rs @@ -327,7 +327,7 @@ pub(super) async fn reset_local_chain(state: &HttpState, confirmation: &str) -> clear_ui_data(&state.ui_data_store).await?; state .gossip - .broadcast(vec![GossipEnvelope::ChainSnapshotRequest]) + .broadcast(vec![GossipEnvelope::ChainBootstrapRequest]) .await?; Ok(()) } diff --git a/src/adapters/p2p.rs b/src/adapters/p2p.rs @@ -29,8 +29,7 @@ mod test_support; mod writer; pub use fetch::{fetch_peer_height, fetch_snapshot, fetch_snapshot_with_announcement}; use fetch::{ - network_adjusted_time_ms, validate_blocks_extension, validate_snapshot_extension, - verify_block_vdf, + network_adjusted_time_ms, validate_blocks_extension, validate_chain_bootstrap, verify_block_vdf, }; #[cfg(test)] use handshake::verify_advertised_peer_node_id; @@ -60,13 +59,13 @@ use sync::{ apply_peer_list, envelopes_for_peer, maybe_request_catchup, push_catchup_to_peer, write_peer_exchange, }; -use writer::{write_envelope, write_payload}; +use writer::{byte_bounded_block_page, write_envelope, write_payload}; const MAX_BLOCK_BATCH: usize = 128; const MAX_OBJECT_REQUESTS: usize = 128; const MAX_INVENTORY_ITEMS: usize = 512; +const MAX_BLOCK_LOCATOR_HASHES: usize = 64; const MAX_PEER_LIST: usize = 128; -const MAX_SNAPSHOT_BLOCKS: usize = 10_000; const MAX_GOSSIP_LINE_BYTES: usize = 8 * 1024 * 1024; const MAX_INBOUND_SESSIONS: usize = 64; const MAX_INBOUND_SESSIONS_PER_IP: usize = 8; diff --git a/src/adapters/p2p/fetch.rs b/src/adapters/p2p/fetch.rs @@ -2,19 +2,18 @@ use std::net::SocketAddr; use anyhow::{Context, Result}; use tokio::{ - io::AsyncWriteExt, net::{TcpStream, tcp::OwnedReadHalf}, time::timeout, }; use crate::{ - app::{GossipEnvelope, NETWORK_ID, PROTOCOL_VERSION, now_ms}, + app::{ChainBootstrap, GossipEnvelope, NETWORK_ID, PROTOCOL_VERSION, now_ms}, domain::{Block, ChainSnapshot, Ledger, verify_vdf}, }; use super::{ - GossipNetwork, JOIN_RESPONSE_TIMEOUT, LimitedLineReader, MAX_JOIN_RESPONSE_ENVELOPES, - PeerStatus, parse_envelope, + GossipNetwork, JOIN_RESPONSE_TIMEOUT, LimitedLineReader, MAX_BLOCK_BATCH, + MAX_JOIN_RESPONSE_ENVELOPES, PeerStatus, parse_envelope, write_envelope, }; pub async fn fetch_snapshot(peer: &str) -> Result<ChainSnapshot> { @@ -100,65 +99,109 @@ pub async fn fetch_snapshot_with_announcement( other => anyhow::bail!("join peer {peer} sent {other:?} instead of peer status"), } - let line = serde_json::to_string(&GossipEnvelope::ChainSnapshotRequest)?; - writer.write_all(line.as_bytes()).await?; - writer.write_all(b"\n").await?; - let snapshot = read_join_snapshot_response(peer, &mut reader).await?; + write_envelope(&mut writer, &GossipEnvelope::ChainBootstrapRequest).await?; + let bootstrap = read_join_bootstrap_response(peer, &mut reader).await?; + + let mut snapshot = ChainSnapshot { + genesis_allocations: bootstrap.genesis_allocations, + vdf_rounds: bootstrap.vdf_rounds, + launch_profile: bootstrap.launch_profile, + blocks: vec![bootstrap.genesis_block], + }; + while snapshot.blocks.last().map_or(0, |block| block.height) < bootstrap.height { + let from_height = snapshot.blocks.last().map_or(0, |block| block.height) + 1; + let remaining = bootstrap.height - from_height + 1; + write_envelope( + &mut writer, + &GossipEnvelope::BlockRangeRequest { + from_height, + limit: remaining.min(MAX_BLOCK_BATCH as u64) as usize, + }, + ) + .await?; + let blocks = read_join_blocks_response(peer, &mut reader).await?; + if blocks.is_empty() { + anyhow::bail!("join peer {peer} returned an empty block page at height {from_height}"); + } + if blocks[0].height != from_height { + anyhow::bail!("join peer {peer} returned a non-contiguous block page"); + } + snapshot.blocks.extend(blocks); + } + if snapshot.blocks.last().map(|block| &block.hash) != Some(&bootstrap.tip_hash) { + anyhow::bail!("join peer {peer} changed tips while serving block pages"); + } Ok(snapshot) } -async fn read_join_snapshot_response( +async fn read_join_bootstrap_response( peer: &str, reader: &mut LimitedLineReader<OwnedReadHalf>, -) -> Result<ChainSnapshot> { +) -> Result<ChainBootstrap> { for _ in 0..MAX_JOIN_RESPONSE_ENVELOPES { - let line = timeout(JOIN_RESPONSE_TIMEOUT, reader.read_line()) - .await - .with_context(|| format!("join peer {peer} timed out waiting for a chain snapshot"))?? - .with_context(|| format!("join peer {peer} closed before sending a chain snapshot"))?; - match join_snapshot_response(peer, parse_envelope(&line)?)? { - Some(snapshot) => return Ok(snapshot), - None => continue, + let envelope = read_join_envelope(peer, reader, "chain bootstrap").await?; + match envelope { + GossipEnvelope::ChainBootstrap(bootstrap) => return Ok(bootstrap), + envelope if is_join_control_envelope(&envelope) => continue, + other => anyhow::bail!("join peer {peer} sent {other:?} instead of chain bootstrap"), } } + anyhow::bail!("join peer {peer} sent too many control envelopes while joining") +} - anyhow::bail!("join peer {peer} sent too many non-snapshot envelopes while joining") +async fn read_join_blocks_response( + peer: &str, + reader: &mut LimitedLineReader<OwnedReadHalf>, +) -> Result<Vec<Block>> { + for _ in 0..MAX_JOIN_RESPONSE_ENVELOPES { + let envelope = read_join_envelope(peer, reader, "block page").await?; + match envelope { + GossipEnvelope::Blocks { blocks } => return Ok(blocks), + envelope if is_join_control_envelope(&envelope) => continue, + other => anyhow::bail!("join peer {peer} sent {other:?} instead of a block page"), + } + } + anyhow::bail!("join peer {peer} sent too many control envelopes while joining") } -pub(super) fn join_snapshot_response( +async fn read_join_envelope( peer: &str, - envelope: GossipEnvelope, -) -> Result<Option<ChainSnapshot>> { - match envelope { - GossipEnvelope::ChainSnapshot(snapshot) => Ok(Some(snapshot)), + reader: &mut LimitedLineReader<OwnedReadHalf>, + expected: &str, +) -> Result<GossipEnvelope> { + let line = timeout(JOIN_RESPONSE_TIMEOUT, reader.read_line()) + .await + .with_context(|| format!("join peer {peer} timed out waiting for {expected}"))?? + .with_context(|| format!("join peer {peer} closed before sending {expected}"))?; + parse_envelope(&line) +} + +fn is_join_control_envelope(envelope: &GossipEnvelope) -> bool { + matches!( + envelope, GossipEnvelope::Hello(_) - | GossipEnvelope::PeerStatus { .. } - | GossipEnvelope::PeerList { .. } - | GossipEnvelope::PeerVerificationChallenge { .. } - | GossipEnvelope::PeerVerificationResponse { .. } - | GossipEnvelope::Inventory { .. } => Ok(None), - other => anyhow::bail!("join peer {peer} sent {other:?} instead of a chain snapshot"), - } + | GossipEnvelope::PeerStatus { .. } + | GossipEnvelope::PeerList { .. } + | GossipEnvelope::PeerVerificationChallenge { .. } + | GossipEnvelope::PeerVerificationResponse { .. } + | GossipEnvelope::Inventory { .. } + ) } -pub(super) async fn validate_snapshot_extension( - mut ledger: Ledger, - snapshot: ChainSnapshot, +pub(super) async fn validate_chain_bootstrap( + bootstrap: ChainBootstrap, now_ms: u64, ) -> Result<Ledger> { - if ledger.is_setup_placeholder() { - return tokio::task::spawn_blocking(move || Ledger::from_snapshot_at(snapshot, now_ms)) - .await - .context("chain snapshot adoption worker failed")?; - } - - tokio::task::spawn_blocking(move || { - ledger.extend_from_snapshot_at(snapshot, now_ms)?; - Ok(ledger) - }) - .await - .context("chain snapshot extension worker failed")? + let snapshot = ChainSnapshot { + genesis_allocations: bootstrap.genesis_allocations, + vdf_rounds: bootstrap.vdf_rounds, + launch_profile: bootstrap.launch_profile, + blocks: vec![bootstrap.genesis_block], + }; + tokio::task::spawn_blocking(move || Ledger::from_snapshot_at(snapshot, now_ms)) + .await + .context("chain bootstrap adoption worker failed")? } pub(super) async fn validate_blocks_extension( @@ -171,6 +214,18 @@ pub(super) async fn validate_blocks_extension( } tokio::task::spawn_blocking(move || { + if blocks[0].prev_hash != ledger.tip_hash() { + let mut candidate = ledger.snapshot(); + let ancestor = candidate + .blocks + .iter() + .position(|block| block.hash == blocks[0].prev_hash) + .context("block page has no common ancestor with local chain")?; + candidate.blocks.truncate(ancestor + 1); + candidate.blocks.extend(blocks); + ledger.extend_from_snapshot_at(candidate, now_ms)?; + return Ok(ledger); + } for block in blocks { ledger.apply_block_at(block, now_ms)?; } diff --git a/src/adapters/p2p/handshake.rs b/src/adapters/p2p/handshake.rs @@ -136,8 +136,8 @@ async fn process_hello_inner( let genesis_mismatch = hello.genesis_hash != local_genesis; let remote_is_setup_placeholder = hello.height == 0 && hello.genesis_hash == setup_placeholder_genesis_hash(); - let request_snapshot = genesis_mismatch && local_accepts_remote_genesis; - let push_snapshot = genesis_mismatch && remote_is_setup_placeholder; + let request_bootstrap = genesis_mismatch && local_accepts_remote_genesis; + let push_bootstrap = genesis_mismatch && remote_is_setup_placeholder; if genesis_mismatch && !local_accepts_remote_genesis && !remote_is_setup_placeholder { anyhow::bail!( "wrong genesis {}; expected {local_genesis}", @@ -180,14 +180,14 @@ async fn process_hello_inner( &PeerStatus::with_time(hello.height, hello.tip_hash.clone(), hello.time_ms), ) .await; - if request_snapshot { - Ok(PeerStatus::with_snapshot_request( + if request_bootstrap { + Ok(PeerStatus::with_bootstrap_request( hello.height, hello.tip_hash, hello.time_ms, )) - } else if push_snapshot { - Ok(PeerStatus::with_snapshot_push( + } else if push_bootstrap { + Ok(PeerStatus::with_bootstrap_push( hello.height, hello.tip_hash, hello.time_ms, diff --git a/src/adapters/p2p/line_codec.rs b/src/adapters/p2p/line_codec.rs @@ -10,8 +10,8 @@ use crate::{ }; use super::{ - GossipNetwork, MAX_BLOCK_BATCH, MAX_GOSSIP_LINE_BYTES, MAX_INVENTORY_ITEMS, - MAX_OBJECT_REQUESTS, MAX_PEER_LIST, MAX_SNAPSHOT_BLOCKS, metrics::P2pMetricsCounters, + GossipNetwork, MAX_BLOCK_BATCH, MAX_BLOCK_LOCATOR_HASHES, MAX_GOSSIP_LINE_BYTES, + MAX_INVENTORY_ITEMS, MAX_OBJECT_REQUESTS, MAX_PEER_LIST, metrics::P2pMetricsCounters, }; pub(super) struct LimitedLineReader<R> { @@ -133,10 +133,11 @@ pub(super) fn record_received_envelope_kind( } GossipEnvelope::Block(_) | GossipEnvelope::Blocks { .. } - | GossipEnvelope::ChainSnapshot(_) => { + | GossipEnvelope::ChainBootstrap(_) => { P2pMetricsCounters::inc(&metrics.data_envelopes_received); } - GossipEnvelope::ChainSnapshotRequest + GossipEnvelope::ChainBootstrapRequest + | GossipEnvelope::BlockLocatorRequest { .. } | GossipEnvelope::BlockRangeRequest { .. } | GossipEnvelope::BlockRequest { .. } | GossipEnvelope::BurnBundleRequest { .. } @@ -163,6 +164,10 @@ pub(super) fn validate_envelope_limits(envelope: &GossipEnvelope) -> Result<()> GossipEnvelope::BlockRangeRequest { limit, .. } => { ensure_len("block range request", *limit, MAX_BLOCK_BATCH)?; } + GossipEnvelope::BlockLocatorRequest { locator, limit } => { + ensure_len("block locator", locator.len(), MAX_BLOCK_LOCATOR_HASHES)?; + ensure_len("block locator request", *limit, MAX_BLOCK_BATCH)?; + } GossipEnvelope::BlockRequest { hashes } => { ensure_len("block request", hashes.len(), MAX_OBJECT_REQUESTS)?; } @@ -185,14 +190,12 @@ pub(super) fn validate_envelope_limits(envelope: &GossipEnvelope) -> Result<()> GossipEnvelope::Blocks { blocks } => { ensure_len("block batch", blocks.len(), MAX_BLOCK_BATCH)?; } - GossipEnvelope::ChainSnapshot(snapshot) => { - ensure_len("chain snapshot", snapshot.blocks.len(), MAX_SNAPSHOT_BLOCKS)?; - } GossipEnvelope::PeerList { peers } => { ensure_len("peer list", peers.len(), MAX_PEER_LIST)?; } GossipEnvelope::Hello(_) - | GossipEnvelope::ChainSnapshotRequest + | GossipEnvelope::ChainBootstrapRequest + | GossipEnvelope::ChainBootstrap(_) | GossipEnvelope::PeerStatus { .. } | GossipEnvelope::Transaction(_) | GossipEnvelope::BurnBundle(_) @@ -219,14 +222,14 @@ mod tests { adapters::p2p::metrics::P2pMetricsCounters, app::{BlockInventory, GossipEnvelope, TRANSACTION_BATCH_LIMIT}, domain::{ - BURN_COMMITTEE_SIZE, Block, BurnBundle, BurnBundleSection, ChainSnapshot, - FinalizerMode, LaunchProfile, OutPoint, Transaction, TxInput, TxOutput, + BURN_COMMITTEE_SIZE, Block, BurnBundle, BurnBundleSection, FinalizerMode, OutPoint, + Transaction, TxInput, TxOutput, }, }; use super::{ - LimitedLineReader, MAX_BLOCK_BATCH, MAX_GOSSIP_LINE_BYTES, MAX_INVENTORY_ITEMS, - MAX_OBJECT_REQUESTS, MAX_PEER_LIST, MAX_SNAPSHOT_BLOCKS, parse_envelope, + LimitedLineReader, MAX_BLOCK_BATCH, MAX_BLOCK_LOCATOR_HASHES, MAX_GOSSIP_LINE_BYTES, + MAX_INVENTORY_ITEMS, MAX_OBJECT_REQUESTS, MAX_PEER_LIST, parse_envelope, record_received_envelope_kind, validate_envelope_limits, }; @@ -279,17 +282,6 @@ mod tests { } } - fn dummy_snapshot(blocks: usize) -> ChainSnapshot { - ChainSnapshot { - genesis_allocations: Default::default(), - vdf_rounds: 1, - launch_profile: LaunchProfile::default(), - blocks: (0..blocks) - .map(|height| dummy_block(height as u64)) - .collect(), - } - } - #[test] fn metrics_count_transaction_and_burn_bundle_batches() { let metrics = P2pMetricsCounters::default(); @@ -432,15 +424,17 @@ mod tests { .is_err() ); assert!( - validate_envelope_limits(&GossipEnvelope::ChainSnapshot(dummy_snapshot( - MAX_SNAPSHOT_BLOCKS - ))) + validate_envelope_limits(&GossipEnvelope::BlockLocatorRequest { + locator: vec!["0".repeat(64); MAX_BLOCK_LOCATOR_HASHES], + limit: MAX_BLOCK_BATCH, + }) .is_ok() ); assert!( - validate_envelope_limits(&GossipEnvelope::ChainSnapshot(dummy_snapshot( - MAX_SNAPSHOT_BLOCKS + 1 - ))) + validate_envelope_limits(&GossipEnvelope::BlockLocatorRequest { + locator: vec!["0".repeat(64); MAX_BLOCK_LOCATOR_HASHES + 1], + limit: MAX_BLOCK_BATCH, + }) .is_err() ); assert!( diff --git a/src/adapters/p2p/peer_addr.rs b/src/adapters/p2p/peer_addr.rs @@ -12,7 +12,7 @@ pub(super) fn next_reconnect_delay(current: Duration, max_delay: Duration) -> Du (current * 2).min(max_delay) } -pub(super) fn peer_needs_snapshot(peer_height: u64, envelopes: &[GossipEnvelope]) -> bool { +pub(super) fn peer_has_block_gap(peer_height: u64, envelopes: &[GossipEnvelope]) -> bool { envelopes .iter() .filter_map(|envelope| match envelope { diff --git a/src/adapters/p2p/peer_status.rs b/src/adapters/p2p/peer_status.rs @@ -5,8 +5,8 @@ pub(super) struct PeerStatus { pub(super) height: u64, pub(super) tip_hash: String, pub(super) time_ms: u64, - pub(super) request_snapshot: bool, - pub(super) push_snapshot: bool, + pub(super) request_bootstrap: bool, + pub(super) push_bootstrap: bool, } impl PeerStatus { @@ -19,8 +19,8 @@ impl PeerStatus { height, tip_hash, time_ms, - request_snapshot: false, - push_snapshot: false, + request_bootstrap: false, + push_bootstrap: false, } } @@ -29,28 +29,28 @@ impl PeerStatus { height, tip_hash, time_ms, - request_snapshot: false, - push_snapshot: false, + request_bootstrap: false, + push_bootstrap: false, } } - pub(super) fn with_snapshot_request(height: u64, tip_hash: String, time_ms: u64) -> Self { + pub(super) fn with_bootstrap_request(height: u64, tip_hash: String, time_ms: u64) -> Self { Self { height, tip_hash, time_ms, - request_snapshot: true, - push_snapshot: false, + request_bootstrap: true, + push_bootstrap: false, } } - pub(super) fn with_snapshot_push(height: u64, tip_hash: String, time_ms: u64) -> Self { + pub(super) fn with_bootstrap_push(height: u64, tip_hash: String, time_ms: u64) -> Self { Self { height, tip_hash, time_ms, - request_snapshot: false, - push_snapshot: true, + request_bootstrap: false, + push_bootstrap: true, } } } diff --git a/src/adapters/p2p/process.rs b/src/adapters/p2p/process.rs @@ -11,7 +11,7 @@ use crate::{ use super::{ GossipNetwork, MAX_BLOCK_BATCH, P2pMetricsCounters, apply_peer_list, forget_stale_self_peer, is_possible_fork_error, normalize_advertised_peer, peer_verification_response, process_hello, - validate_blocks_extension, validate_snapshot_extension, verify_block_vdf, write_envelope, + validate_blocks_extension, validate_chain_bootstrap, verify_block_vdf, write_envelope, write_payload, }; @@ -40,9 +40,19 @@ pub(super) async fn process_envelope( GossipEnvelope::Hello(hello) => { let _ = process_hello(network, remote_addr, known_peer, hello).await?; } - GossipEnvelope::ChainSnapshotRequest => { - let snapshot = network.inner.node.lock().await.chain_snapshot(); - write_envelope(writer, &GossipEnvelope::ChainSnapshot(snapshot)).await?; + GossipEnvelope::ChainBootstrapRequest => { + let bootstrap = network.inner.node.lock().await.chain_bootstrap(); + write_envelope(writer, &GossipEnvelope::ChainBootstrap(bootstrap)).await?; + } + GossipEnvelope::BlockLocatorRequest { locator, limit } => { + let blocks = network + .inner + .node + .lock() + .await + .blocks_after_locator(&locator, limit.min(MAX_BLOCK_BATCH)); + let blocks = super::byte_bounded_block_page(blocks); + write_envelope(writer, &GossipEnvelope::Blocks { blocks }).await?; } GossipEnvelope::BlockRangeRequest { from_height, limit } => { let blocks = network @@ -51,11 +61,13 @@ pub(super) async fn process_envelope( .lock() .await .blocks_from(from_height, limit.min(MAX_BLOCK_BATCH)); + let blocks = super::byte_bounded_block_page(blocks); write_envelope(writer, &GossipEnvelope::Blocks { blocks }).await?; } GossipEnvelope::BlockRequest { hashes } => { let blocks = network.inner.node.lock().await.blocks_by_hash(&hashes); if !blocks.is_empty() { + let blocks = super::byte_bounded_block_page(blocks); write_envelope(writer, &GossipEnvelope::Blocks { blocks }).await?; } } @@ -76,8 +88,8 @@ pub(super) async fn process_envelope( } else if node_id.is_some() && debug_logging_enabled() { eprintln!("p2p peer announcement for {peer} ignored until hello verification"); } - let snapshot = network.inner.node.lock().await.chain_snapshot(); - write_envelope(writer, &GossipEnvelope::ChainSnapshot(snapshot)).await?; + let bootstrap = network.inner.node.lock().await.chain_bootstrap(); + write_envelope(writer, &GossipEnvelope::ChainBootstrap(bootstrap)).await?; } GossipEnvelope::PeerVerificationChallenge { address, nonce } => { if let Some(response) = peer_verification_response(network, &address, &nonce) { @@ -134,7 +146,7 @@ pub(super) async fn process_envelope( }, Err(error) => Err(error), }; - let request_snapshot = result.as_ref().err().is_some_and(is_possible_fork_error); + let request_locator = result.as_ref().err().is_some_and(is_possible_fork_error); record_rejected_chain_payload( network, &network.inner.metrics.rejected_blocks, @@ -142,8 +154,8 @@ pub(super) async fn process_envelope( &result, ); record_inbound_result(network, known_peer, remote_addr, result).await; - if request_snapshot { - write_envelope(writer, &GossipEnvelope::ChainSnapshotRequest).await?; + if request_locator { + request_fork_blocks(network, writer).await?; } network.forward_outbox().await; } @@ -161,7 +173,7 @@ pub(super) async fn process_envelope( .map(|_| ()), Err(error) => Err(error), }; - let request_snapshot = result.as_ref().err().is_some_and(is_possible_fork_error); + let request_locator = result.as_ref().err().is_some_and(is_possible_fork_error); record_rejected_chain_payload( network, &network.inner.metrics.rejected_block_batches, @@ -169,29 +181,27 @@ pub(super) async fn process_envelope( &result, ); record_inbound_result(network, known_peer, remote_addr, result).await; - if request_snapshot { - write_envelope(writer, &GossipEnvelope::ChainSnapshotRequest).await?; + if request_locator { + request_fork_blocks(network, writer).await?; } network.forward_outbox().await; } - GossipEnvelope::ChainSnapshot(snapshot) => { + GossipEnvelope::ChainBootstrap(bootstrap) => { let adjusted_time_ms = super::network_adjusted_time_ms(network).await; - let local_ledger = network.inner.node.lock().await.clone_ledger(); - let result = - match validate_snapshot_extension(local_ledger, snapshot, adjusted_time_ms).await { - Ok(ledger) => network - .inner - .node - .lock() - .await - .import_verified_ledger(ledger) - .map(|_| ()), - Err(error) => Err(error), - }; + let result = match validate_chain_bootstrap(bootstrap, adjusted_time_ms).await { + Ok(ledger) => network + .inner + .node + .lock() + .await + .import_verified_ledger(ledger) + .map(|_| ()), + Err(error) => Err(error), + }; record_rejected_chain_payload( network, &network.inner.metrics.rejected_snapshots, - "snapshot", + "chain bootstrap", &result, ); record_inbound_result(network, known_peer, remote_addr, result).await; @@ -206,6 +216,18 @@ pub(super) async fn process_envelope( Ok(()) } +async fn request_fork_blocks(network: &GossipNetwork, writer: &mut OwnedWriteHalf) -> Result<()> { + let locator = network.inner.node.lock().await.block_locator(); + write_envelope( + writer, + &GossipEnvelope::BlockLocatorRequest { + locator, + limit: MAX_BLOCK_BATCH, + }, + ) + .await +} + async fn process_transactions( network: &GossipNetwork, remote_addr: SocketAddr, diff --git a/src/adapters/p2p/sync.rs b/src/adapters/p2p/sync.rs @@ -5,10 +5,13 @@ use tokio::net::tcp::OwnedWriteHalf; use super::metrics::P2pMetricsCounters; use super::peer_addr::{ - is_self_peer_address_for, normalize_advertised_peer, peer_list_address_is_discoverable, - peer_needs_snapshot, + is_self_peer_address_for, normalize_advertised_peer, peer_has_block_gap, + peer_list_address_is_discoverable, +}; +use super::{ + GossipNetwork, MAX_BLOCK_BATCH, PeerStatus, byte_bounded_block_page, write_envelope, + write_payload, }; -use super::{GossipNetwork, MAX_BLOCK_BATCH, PeerStatus, write_envelope, write_payload}; use crate::app::{GossipEnvelope, SharedNode, debug_logging_enabled}; pub(super) async fn maybe_request_catchup( @@ -21,8 +24,8 @@ pub(super) async fn maybe_request_catchup( let status = node.ledger().status(); (status.height, status.tip_hash) }; - if peer_status.request_snapshot { - write_envelope(writer, &GossipEnvelope::ChainSnapshotRequest).await?; + if peer_status.request_bootstrap { + write_envelope(writer, &GossipEnvelope::ChainBootstrapRequest).await?; } else if peer_status.height > local_height { write_envelope( writer, @@ -33,7 +36,15 @@ pub(super) async fn maybe_request_catchup( ) .await?; } else if peer_status.height == local_height && peer_status.tip_hash != local_tip_hash { - write_envelope(writer, &GossipEnvelope::ChainSnapshotRequest).await?; + let locator = network.inner.node.lock().await.block_locator(); + write_envelope( + writer, + &GossipEnvelope::BlockLocatorRequest { + locator, + limit: MAX_BLOCK_BATCH, + }, + ) + .await?; } Ok(()) } @@ -52,10 +63,9 @@ pub(super) async fn push_catchup_to_peer( GossipEnvelope::Blocks { blocks } => blocks .last() .map(|block| PeerStatus::new(block.height, block.hash.clone())), - GossipEnvelope::ChainSnapshot(snapshot) => snapshot - .blocks - .last() - .map(|block| PeerStatus::new(block.height, block.hash.clone())), + GossipEnvelope::ChainBootstrap(bootstrap) => { + Some(PeerStatus::new(0, bootstrap.genesis_block.hash.clone())) + } _ => None, }); write_payload(writer, &payload).await?; @@ -71,30 +81,22 @@ pub(super) async fn catchup_payload_for_peer( if node.ledger().is_setup_placeholder() { return Vec::new(); } - let mempool = node.mempool_gossip(); - if peer_status.push_snapshot { - let mut payload = vec![GossipEnvelope::ChainSnapshot(node.chain_snapshot())]; - payload.extend(mempool); - return payload; + if peer_status.push_bootstrap { + return vec![GossipEnvelope::ChainBootstrap(node.chain_bootstrap())]; } if peer_status.height < local_status.height { - let blocks = node.blocks_from(peer_status.height + 1, MAX_BLOCK_BATCH); - if blocks.is_empty() { - mempool - } else { - let mut payload = vec![GossipEnvelope::Blocks { blocks }]; - payload.extend(mempool); - payload - } + let blocks = + byte_bounded_block_page(node.blocks_from(peer_status.height + 1, MAX_BLOCK_BATCH)); + return (!blocks.is_empty()) + .then_some(GossipEnvelope::Blocks { blocks }) + .into_iter() + .collect(); } else if peer_status.height == local_status.height && peer_status.tip_hash != local_status.tip_hash { - let mut payload = vec![GossipEnvelope::ChainSnapshot(node.chain_snapshot())]; - payload.extend(mempool); - payload - } else { - mempool + return Vec::new(); } + node.mempool_gossip() } pub(super) async fn apply_peer_list( @@ -166,24 +168,22 @@ pub(super) async fn envelopes_for_peer( .collect(); } if peer_status.height < local_status.height { - let mut payload = vec![GossipEnvelope::Blocks { - blocks: node.blocks_from(peer_status.height + 1, MAX_BLOCK_BATCH), - }]; - payload.extend( - envelopes - .iter() - .filter(|envelope| !matches!(envelope, GossipEnvelope::Block(_))) - .cloned(), - ); - return payload; + let blocks = + byte_bounded_block_page(node.blocks_from(peer_status.height + 1, MAX_BLOCK_BATCH)); + return (!blocks.is_empty()) + .then_some(GossipEnvelope::Blocks { blocks }) + .into_iter() + .collect(); } if peer_status.height == local_status.height && peer_status.tip_hash != local_status.tip_hash { - return vec![GossipEnvelope::ChainSnapshot(node.chain_snapshot())]; + return Vec::new(); } - if peer_needs_snapshot(peer_status.height, envelopes) { - return vec![GossipEnvelope::ChainSnapshot(node.chain_snapshot())]; + if peer_has_block_gap(peer_status.height, envelopes) { + let blocks = + byte_bounded_block_page(node.blocks_from(peer_status.height + 1, MAX_BLOCK_BATCH)); + return vec![GossipEnvelope::Blocks { blocks }]; } envelopes diff --git a/src/adapters/p2p/tests.rs b/src/adapters/p2p/tests.rs @@ -70,7 +70,7 @@ async fn full_outbound_queue_is_metric_not_peer_error() { } #[tokio::test] -async fn single_block_fork_error_requests_chain_snapshot() { +async fn single_block_fork_error_requests_blocks_by_locator() { let alice = Wallet::from_seed("single-block-fork-alice"); let allocations = allocations(std::slice::from_ref(&alice), 1_000); let mut local_node = node( @@ -129,10 +129,18 @@ async fn single_block_fork_error_requests_chain_snapshot() { .unwrap() .unwrap() .unwrap(); - assert!(matches!( - super::parse_envelope(&line).unwrap(), - GossipEnvelope::ChainSnapshotRequest - )); + let GossipEnvelope::BlockLocatorRequest { locator, limit } = + super::parse_envelope(&line).unwrap() + else { + panic!("expected block locator request"); + }; + let fork_blocks = remote_node.blocks_after_locator(&locator, limit); + assert_eq!(fork_blocks.len(), 2); + let local_ledger = network.inner.node.lock().await.clone_ledger(); + let adopted = super::validate_blocks_extension(local_ledger, fork_blocks, crate::app::now_ms()) + .await + .unwrap(); + assert_eq!(adopted.tip_hash(), remote_node.ledger().tip_hash()); } #[tokio::test] @@ -294,7 +302,7 @@ async fn invalid_block_batch_is_rejected_atomically_without_partial_import() { } #[tokio::test] -async fn chain_snapshot_request_only_writes_snapshot_without_mutating_local_state() { +async fn chain_bootstrap_request_only_writes_bootstrap_without_mutating_local_state() { let alice = Wallet::from_seed("snapshot-request-spam-alice"); let allocations = allocations(std::slice::from_ref(&alice), 1_000); let mut local_node = node("snapshot-request-spam", alice.clone(), allocations); @@ -326,7 +334,7 @@ async fn chain_snapshot_request_only_writes_snapshot_without_mutating_local_stat &mut server_writer, remote_addr, &mut known_peer, - GossipEnvelope::ChainSnapshotRequest, + GossipEnvelope::ChainBootstrapRequest, ) .await .unwrap(); @@ -336,10 +344,12 @@ async fn chain_snapshot_request_only_writes_snapshot_without_mutating_local_stat .unwrap() .unwrap() .unwrap(); - assert_eq!( - super::parse_envelope(&line).unwrap(), - GossipEnvelope::ChainSnapshot(before.clone()) - ); + let GossipEnvelope::ChainBootstrap(bootstrap) = super::parse_envelope(&line).unwrap() else { + panic!("expected chain bootstrap"); + }; + assert_eq!(bootstrap.genesis_block, before.blocks[0]); + assert_eq!(bootstrap.height, before.blocks.last().unwrap().height); + assert_eq!(bootstrap.tip_hash, before.blocks.last().unwrap().hash); assert_eq!(network.inner.node.lock().await.chain_snapshot(), before); assert!(peers.lock().await.list().is_empty()); assert!(known_peer.is_none()); @@ -759,7 +769,7 @@ async fn inbound_verification_only_session_closes_after_response() { } #[tokio::test] -async fn setup_placeholder_accepts_remote_genesis_and_adopts_snapshot() { +async fn setup_placeholder_accepts_remote_genesis_and_adopts_bootstrap() { let local_wallet = Wallet::from_seed("setup-placeholder-local"); let local_ledger = Ledger::new(BTreeMap::new(), 1); let local_node = Arc::new(tokio::sync::Mutex::new(NodeCore::from_ledger( @@ -784,13 +794,13 @@ async fn setup_placeholder_accepts_remote_genesis_and_adopts_snapshot() { }; let remote_wallet = Wallet::from_seed("setup-placeholder-remote"); - let remote_snapshot = node( + let remote_node = node( "remote", remote_wallet.clone(), allocations(std::slice::from_ref(&remote_wallet), 1_000), - ) - .chain_snapshot(); - let remote_genesis = remote_snapshot.blocks[0].hash.clone(); + ); + let remote_genesis = remote_node.ledger().genesis_hash().to_string(); + let remote_bootstrap = remote_node.chain_bootstrap(); let hello = ProtocolHello { protocol_version: PROTOCOL_VERSION, network_id: NETWORK_ID.to_string(), @@ -811,8 +821,8 @@ async fn setup_placeholder_accepts_remote_genesis_and_adopts_snapshot() { .await .unwrap(); - assert!(peer_status.request_snapshot); - assert!(!peer_status.push_snapshot); + assert!(peer_status.request_bootstrap); + assert!(!peer_status.push_bootstrap); assert_eq!(known_peer.as_deref(), Some("iuna.jhx.app:9444")); let listed = network.inner.peers.lock().await.list(); assert_eq!(listed.len(), 1); @@ -823,10 +833,9 @@ async fn setup_placeholder_accepts_remote_genesis_and_adopts_snapshot() { assert_eq!(peer.misbehavior_score, 0); assert!(!peer.is_banned_at(crate::app::now_ms())); - let adopted = - super::validate_snapshot_extension(local_ledger, remote_snapshot, crate::app::now_ms()) - .await - .unwrap(); + let adopted = super::validate_chain_bootstrap(remote_bootstrap, crate::app::now_ms()) + .await + .unwrap(); assert_eq!(adopted.genesis_hash(), remote_genesis); assert!( network @@ -844,7 +853,7 @@ async fn setup_placeholder_accepts_remote_genesis_and_adopts_snapshot() { } #[tokio::test] -async fn real_node_accepts_setup_placeholder_peer_and_pushes_snapshot() { +async fn real_node_accepts_setup_placeholder_peer_and_pushes_bootstrap() { let wallet = Wallet::from_seed("setup-placeholder-peer-real-node"); let node = Arc::new(tokio::sync::Mutex::new(node( "real", @@ -885,12 +894,38 @@ async fn real_node_accepts_setup_placeholder_peer_and_pushes_snapshot() { .await .unwrap(); - assert!(!peer_status.request_snapshot); - assert!(peer_status.push_snapshot); + assert!(!peer_status.request_bootstrap); + assert!(peer_status.push_bootstrap); let payload = super::catchup_payload_for_peer(&node, &peer_status).await; assert!(matches!( payload.as_slice(), - [GossipEnvelope::ChainSnapshot(_)] + [GossipEnvelope::ChainBootstrap(_)] + )); +} + +#[tokio::test] +async fn lagging_peer_receives_block_pages_before_mempool() { + let wallet = Wallet::from_seed("lagging-peer-blocks-before-mempool"); + let mut local = node( + "lagging-peer-source", + wallet.clone(), + allocations(std::slice::from_ref(&wallet), 1_000), + ); + let genesis_hash = local.ledger().genesis_hash().to_string(); + queue_plaintext_burn(&mut local, &wallet, 1); + local.drain_outbox(); + local.mine_one_at(1).unwrap(); + local.drain_outbox(); + queue_plaintext_burn(&mut local, &wallet, 1); + assert!(!local.ledger().pending().is_empty()); + let local = Arc::new(tokio::sync::Mutex::new(local)); + let peer_status = super::PeerStatus::new(0, genesis_hash); + + let payload = super::catchup_payload_for_peer(&local, &peer_status).await; + + assert!(matches!( + payload.as_slice(), + [GossipEnvelope::Blocks { .. }] )); } diff --git a/src/adapters/p2p/writer.rs b/src/adapters/p2p/writer.rs @@ -31,3 +31,81 @@ pub(super) async fn write_envelope( writer.write_all(b"\n").await?; Ok(()) } + +pub(super) fn byte_bounded_block_page( + blocks: Vec<crate::domain::Block>, +) -> Vec<crate::domain::Block> { + let mut page = Vec::new(); + let mut encoded_len = serde_json::to_vec(&GossipEnvelope::Blocks { blocks: Vec::new() }) + .expect("block envelope serialization cannot fail") + .len(); + for block in blocks { + let block_len = serde_json::to_vec(&block) + .expect("block serialization cannot fail") + .len(); + let separator_len = usize::from(!page.is_empty()); + let Some(candidate_len) = encoded_len + .checked_add(separator_len) + .and_then(|len| len.checked_add(block_len)) + else { + break; + }; + if candidate_len > MAX_GOSSIP_LINE_BYTES { + break; + } + encoded_len = candidate_len; + page.push(block); + } + page +} + +#[cfg(test)] +mod tests { + use crate::{ + app::GossipEnvelope, + domain::{Block, BurnBundleSection, FinalizerMode}, + }; + + use super::{MAX_GOSSIP_LINE_BYTES, byte_bounded_block_page}; + + fn large_block(height: u64) -> Block { + Block { + height, + prev_hash: format!("{:064x}", height.saturating_sub(1)), + timestamp_ms: height, + miner: "0".repeat(64), + finalizer_mode: FinalizerMode::Ticket, + finalizer_rank: 0, + reward: 0, + vdf_rounds: 0, + vdf_output: "x".repeat(512 * 1024), + leader_proof: None, + burn_bundle_section: BurnBundleSection::default(), + transactions: Vec::new(), + hash: format!("{height:064x}"), + } + } + + #[test] + fn block_page_is_bounded_by_wire_bytes_not_only_item_count() { + let blocks = (1..=32).map(large_block).collect::<Vec<_>>(); + let page = byte_bounded_block_page(blocks.clone()); + let encoded = serde_json::to_vec(&GossipEnvelope::Blocks { + blocks: page.clone(), + }) + .unwrap(); + + assert!(!page.is_empty()); + assert!(page.len() < blocks.len()); + assert!(encoded.len() <= MAX_GOSSIP_LINE_BYTES); + + let mut one_more = page; + one_more.push(blocks[one_more.len()].clone()); + assert!( + serde_json::to_vec(&GossipEnvelope::Blocks { blocks: one_more }) + .unwrap() + .len() + > MAX_GOSSIP_LINE_BYTES + ); + } +} diff --git a/src/app.rs b/src/app.rs @@ -28,8 +28,9 @@ mod wallet; pub use in_memory_network::InMemoryNetwork; pub use peer_book::{PeerBook, PeerDirection, PeerInfo}; pub use types::{ - AutoMineOutcome, AutoMinePlan, BlockInventory, ExternalMineJob, FeeEstimate, GossipEnvelope, - LaunchProfileStatus, MiningStatus, NodeConfig, NodeStatus, ProtocolHello, StratumStatus, + AutoMineOutcome, AutoMinePlan, BlockInventory, ChainBootstrap, ExternalMineJob, FeeEstimate, + GossipEnvelope, LaunchProfileStatus, MiningStatus, NodeConfig, NodeStatus, ProtocolHello, + StratumStatus, }; use wallet::NodeWallet; diff --git a/src/app/automatic_mining.rs b/src/app/automatic_mining.rs @@ -869,7 +869,7 @@ mod tests { envelope, GossipEnvelope::Block(_) | GossipEnvelope::Blocks { .. } - | GossipEnvelope::ChainSnapshot(_) + | GossipEnvelope::ChainBootstrap(_) )) }) .unwrap(); diff --git a/src/app/gossip.rs b/src/app/gossip.rs @@ -1,8 +1,8 @@ use crate::domain::{Block, ChainSnapshot}; use super::{ - BLOCK_REQUEST_LIMIT, GossipEnvelope, NETWORK_ID, NodeCore, PROTOCOL_VERSION, ProtocolHello, - TRANSACTION_BATCH_LIMIT, now_ms, types::BlockInventory, + BLOCK_REQUEST_LIMIT, ChainBootstrap, GossipEnvelope, NETWORK_ID, NodeCore, PROTOCOL_VERSION, + ProtocolHello, TRANSACTION_BATCH_LIMIT, now_ms, types::BlockInventory, }; impl NodeCore { @@ -30,6 +30,26 @@ impl NodeCore { self.ledger.snapshot() } + pub fn chain_bootstrap(&self) -> ChainBootstrap { + let snapshot = self.ledger.genesis_snapshot(); + ChainBootstrap { + genesis_allocations: snapshot.genesis_allocations, + vdf_rounds: snapshot.vdf_rounds, + launch_profile: snapshot.launch_profile, + genesis_block: snapshot.blocks[0].clone(), + height: self.ledger.height(), + tip_hash: self.ledger.tip_hash().to_string(), + } + } + + pub fn block_locator(&self) -> Vec<String> { + self.ledger.block_locator() + } + + pub fn blocks_after_locator(&self, locator: &[String], limit: usize) -> Vec<Block> { + self.ledger.blocks_after_locator(locator, limit) + } + pub fn hello(&self, listen_addr: Option<String>, node_id: Option<String>) -> GossipEnvelope { GossipEnvelope::Hello(ProtocolHello { protocol_version: PROTOCOL_VERSION, diff --git a/src/app/receive.rs b/src/app/receive.rs @@ -64,7 +64,9 @@ impl NodeCore { match envelope { GossipEnvelope::Hello(_) | GossipEnvelope::PeerStatus { .. } - | GossipEnvelope::ChainSnapshotRequest + | GossipEnvelope::ChainBootstrapRequest + | GossipEnvelope::ChainBootstrap(_) + | GossipEnvelope::BlockLocatorRequest { .. } | GossipEnvelope::BlockRangeRequest { .. } | GossipEnvelope::BlockRequest { .. } | GossipEnvelope::Inventory { .. } => Ok(()), @@ -111,7 +113,6 @@ impl NodeCore { } Ok(()) } - GossipEnvelope::ChainSnapshot(snapshot) => self.import_chain_snapshot(snapshot), GossipEnvelope::PeerAnnouncement { .. } | GossipEnvelope::PeerVerificationChallenge { .. } | GossipEnvelope::PeerVerificationResponse { .. } diff --git a/src/app/types.rs b/src/app/types.rs @@ -3,7 +3,7 @@ use std::collections::BTreeMap; use serde::{Deserialize, Serialize}; use crate::domain::{ - Amount, Block, BurnBundle, ChainSnapshot, ChainStatus, PreparedBlock, StratumMineTemplate, + Amount, Block, BurnBundle, ChainStatus, LaunchProfile, PreparedBlock, StratumMineTemplate, Transaction, Wallet, }; @@ -39,7 +39,12 @@ pub enum GossipEnvelope { #[serde(default)] time_ms: u64, }, - ChainSnapshotRequest, + ChainBootstrapRequest, + ChainBootstrap(ChainBootstrap), + BlockLocatorRequest { + locator: Vec<String>, + limit: usize, + }, BlockRangeRequest { from_height: u64, limit: usize, @@ -67,7 +72,6 @@ pub enum GossipEnvelope { Blocks { blocks: Vec<Block>, }, - ChainSnapshot(ChainSnapshot), PeerAnnouncement { address: String, #[serde(default)] @@ -89,6 +93,16 @@ pub enum GossipEnvelope { } #[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)] +pub struct ChainBootstrap { + pub genesis_allocations: BTreeMap<String, Amount>, + pub vdf_rounds: u64, + pub launch_profile: LaunchProfile, + pub genesis_block: Block, + pub height: u64, + pub tip_hash: String, +} + +#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)] pub struct ProtocolHello { pub protocol_version: u32, pub network_id: String, diff --git a/src/domain/ledger_queries.rs b/src/domain/ledger_queries.rs @@ -77,6 +77,15 @@ impl Ledger { } } + pub(crate) fn genesis_snapshot(&self) -> ChainSnapshot { + ChainSnapshot { + genesis_allocations: self.genesis_allocations.clone(), + vdf_rounds: self.initial_vdf_rounds, + launch_profile: self.launch_profile.clone(), + blocks: vec![self.chain[0].clone()], + } + } + pub fn status(&self) -> ChainStatus { self.status_with_balances(true) } @@ -429,6 +438,35 @@ impl Ledger { .collect() } + pub(crate) fn block_locator(&self) -> Vec<String> { + let mut locator = Vec::new(); + let mut index = self.chain.len().saturating_sub(1); + let mut step = 1_usize; + loop { + locator.push(self.chain[index].hash.clone()); + if index == 0 { + break; + } + index = index.saturating_sub(step); + if locator.len() > 10 { + step = step.saturating_mul(2); + } + } + locator + } + + pub(crate) fn blocks_after_locator(&self, locator: &[String], limit: usize) -> Vec<Block> { + let common_height = locator.iter().find_map(|hash| { + self.chain + .iter() + .find(|block| block.hash == *hash) + .map(|block| block.height) + }); + common_height + .map(|height| self.blocks_from(height.saturating_add(1), limit)) + .unwrap_or_default() + } + pub fn block_by_hash(&self, hash: &str) -> Option<Block> { self.chain.iter().find(|block| block.hash == hash).cloned() }