diff --git a/Cargo.lock b/Cargo.lock index 6cdc450e9..84ba40655 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -13494,6 +13494,7 @@ dependencies = [ "p256", "rand 0.8.7", "reqwest 0.13.4", + "reth-chain-state", "reth-chainspec", "reth-consensus", "reth-eth-wire-types", @@ -13555,6 +13556,7 @@ dependencies = [ "commonware-runtime", "commonware-utils", "const-hex", + "derive_more", "eyre", "serde", "thiserror 2.0.18", diff --git a/Cargo.toml b/Cargo.toml index 6854d960f..d7ecad336 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -129,6 +129,7 @@ zone-sequencer = { path = "crates/sequencer" } # reth reth-basic-payload-builder = { git = "https://github.com/paradigmxyz/reth", rev = "1bf2384" } +reth-chain-state = { git = "https://github.com/paradigmxyz/reth", rev = "1bf2384" } reth-chainspec = { git = "https://github.com/paradigmxyz/reth", rev = "1bf2384" } reth-cli = { git = "https://github.com/paradigmxyz/reth", rev = "1bf2384" } reth-cli-util = { git = "https://github.com/paradigmxyz/reth", rev = "1bf2384" } diff --git a/crates/node/Cargo.toml b/crates/node/Cargo.toml index 1ccd4481f..605e97f93 100644 --- a/crates/node/Cargo.toml +++ b/crates/node/Cargo.toml @@ -31,6 +31,7 @@ zone-rpc.workspace = true zone-sequencer.workspace = true # reth +reth-chain-state.workspace = true reth-chainspec.workspace = true reth-consensus = { workspace = true, optional = true } reth-eth-wire-types.workspace = true @@ -59,6 +60,7 @@ alloy-genesis.workspace = true alloy-network.workspace = true alloy-primitives.workspace = true alloy-provider = { workspace = true, features = ["ws"] } +alloy-rlp.workspace = true alloy-rpc-client.workspace = true alloy-rpc-types-engine.workspace = true alloy-rpc-types-eth.workspace = true @@ -88,12 +90,11 @@ alloy = { workspace = true, features = [ alloy-contract.workspace = true alloy-eips.workspace = true alloy-evm.workspace = true -alloy-rlp.workspace = true alloy-signer.workspace = true base64.workspace = true +const-hex.workspace = true commonware-codec.workspace = true commonware-cryptography.workspace = true -const-hex.workspace = true p256.workspace = true rand.workspace = true reth-ethereum = { workspace = true, features = ["node", "test-utils"] } diff --git a/crates/node/src/cli.rs b/crates/node/src/cli.rs index b4904d896..eb1c805ae 100644 --- a/crates/node/src/cli.rs +++ b/crates/node/src/cli.rs @@ -129,6 +129,12 @@ fn run_node(mut cli: Cli) -> eyre::Result<()> { } let manifest_mode = p2p_config.is_some(); + if manifest_mode { + // Replicate only durable blocks. Persist every block immediately so followers can + // acknowledge each block without waiting for Reth's in-memory buffer to fill. + builder.config_mut().engine.persistence_threshold = 0; + builder.config_mut().engine.memory_block_buffer_target = 0; + } let should_sequence_blocks = sequencer_enabled(args.enable_sequencer, manifest_role); let sequencer_signer = (should_sequence_blocks || manifest_mode) .then(|| { diff --git a/crates/node/src/lib.rs b/crates/node/src/lib.rs index adeea050d..711ae41bb 100644 --- a/crates/node/src/lib.rs +++ b/crates/node/src/lib.rs @@ -12,6 +12,7 @@ pub mod dev; pub mod engine; pub mod genesis; pub mod node; +mod replication; pub mod rpc; pub use engine::ZoneEngine; diff --git a/crates/node/src/node.rs b/crates/node/src/node.rs index d8cf98910..0f1038e9e 100644 --- a/crates/node/src/node.rs +++ b/crates/node/src/node.rs @@ -5,6 +5,7 @@ use crate::{ ZoneEngine, + replication::{broadcast_persisted_blocks, import_leader_blocks}, rpc::{ZoneRpc, ZoneRpcApi, rpc_connection_config, start_private_rpc}, }; use alloy_primitives::Address; @@ -67,7 +68,7 @@ use zone_l1::{ spawn_policy_resolution_task, spawn_pool_prefetch_task, }, }; -use zone_p2p::{P2pConfig, P2pNetworkId, spawn_p2p}; +use zone_p2p::{P2pConfig, P2pNetworkId, Role, spawn_p2p}; use zone_payload::{ DEFAULT_WITHDRAWAL_BATCH_INTERVAL_BLOCKS, WithdrawalRevealEncryptor, ZonePayloadAttributes, ZonePayloadFactory, ZonePayloadTypes, @@ -466,14 +467,30 @@ where .erased(); self.resolve_and_seed_tokens(&l1_provider).await?; - self.spawn_l1_subscriber(&ctx); + let p2p_role = self.p2p_config.as_ref().map(P2pConfig::role); + if p2p_role == Some(Role::Follower) { + // TODO(multi-sequencer): Split L1 observation/cache updates from deposit + // enqueueing. Followers import complete blocks from the leader and do not consume + // DepositQueue; starting the unified subscriber here would grow that queue forever. + // On promotion/restart the subscriber resumes from the tempoBlockNumber persisted in + // the follower's imported zone state. + info!(target: "reth::cli", "Skipping L1 deposit subscriber on follower"); + } else { + self.spawn_l1_subscriber(&ctx); + } self.spawn_policy_tasks(&l1_provider, &ctx); let task_executor = ctx.node.task_executor().clone(); if let Some(config) = self.p2p_config.take() { let network_id = P2pNetworkId::new(l1_provider.get_chain_id().await?, self.portal_address); - Self::launch_p2p(config, network_id, &task_executor)?; + Self::launch_p2p( + config, + network_id, + &task_executor, + ctx.node.provider().clone(), + ctx.beacon_engine_handle.clone(), + )?; } if let Some(ref config) = self.sequencer_config { @@ -527,16 +544,39 @@ where config: P2pConfig, network_id: P2pNetworkId, task_executor: &reth_tasks::TaskExecutor, + provider: N::Provider, + engine: reth_node_builder::ConsensusEngineHandle, ) -> eyre::Result<()> { + let role = config.role(); let handle = spawn_p2p(config, network_id)?; let zone_p2p::P2pHandleParts { shutdown: shutdown_token, mut stopped, thread, - commands: _commands, - events: _events, + commands, + events, } = handle.into_parts(); + match role { + Role::Leader => { + task_executor.spawn_critical_task( + "zone-p2p-block-broadcast", + broadcast_persisted_blocks(provider, commands), + ); + // Leaders do not receive block messages. Dropping this receiver is harmless: the + // runtime only emits BlockReceived on followers. + drop(events); + } + Role::Follower => { + // Keep the command sender alive so the runtime's command loop remains available + // for later ACK/backfill commands even though followers send nothing in this PR. + task_executor.spawn_critical_task( + "zone-p2p-block-import", + import_leader_blocks(provider, engine, events, commands), + ); + } + } + task_executor.spawn_critical_with_graceful_shutdown_signal( "zone-p2p", |shutdown| async move { diff --git a/crates/node/src/replication.rs b/crates/node/src/replication.rs new file mode 100644 index 000000000..da3ecd453 --- /dev/null +++ b/crates/node/src/replication.rs @@ -0,0 +1,318 @@ +//! Node-side leader block replication and follower import. + +use alloy_consensus::BlockHeader as _; +use alloy_primitives::B256; +use alloy_rlp::Decodable as _; +use alloy_rpc_types_engine::ForkchoiceState; +use futures::{StreamExt as _, stream::BoxStream}; +use reth_chain_state::PersistedBlockSubscriptions; +use reth_node_api::PayloadTypes as _; +use reth_node_builder::ConsensusEngineHandle; +use reth_primitives_traits::SealedBlock; +use reth_storage_api::{BlockNumReader, BlockReader, HeaderProvider}; +use tempo_primitives::{Block, TempoHeader}; +use tokio::sync::mpsc; +use tracing::{debug, info}; +use zone_p2p::{P2pCommand, P2pEvent}; +use zone_payload::ZonePayloadTypes; + +#[derive(Debug, Clone, Copy)] +pub(crate) struct PersistedTip { + number: u64, + hash: B256, +} + +pub(crate) struct EncodedPersistedBlock { + number: u64, + hash: B256, + encoded: Vec, +} + +/// Interface used by the replication task to keep track of blocks that are persisted vs broadcast +pub(crate) trait PersistedBlockSource: Clone + Send + Sync + 'static { + fn last_block_number(&self) -> eyre::Result; + fn persisted_block_stream(&self) -> BoxStream<'static, PersistedTip>; + fn encoded_block_by_number(&self, number: u64) -> eyre::Result; +} + +impl

PersistedBlockSource for P +where + P: PersistedBlockSubscriptions + BlockReader + Clone + Send + Sync + 'static, +{ + fn last_block_number(&self) -> eyre::Result { + Ok(BlockNumReader::last_block_number(self)?) + } + + fn persisted_block_stream(&self) -> BoxStream<'static, PersistedTip> { + PersistedBlockSubscriptions::persisted_block_stream(self) + .map(|tip| PersistedTip { + number: tip.number, + hash: tip.hash, + }) + .boxed() + } + + fn encoded_block_by_number(&self, number: u64) -> eyre::Result { + let block = self + .block_by_number(number)? + .ok_or_else(|| eyre::eyre!("persisted zone block {number} is missing"))?; + let sealed = SealedBlock::seal_slow(block); + Ok(EncodedPersistedBlock { + number: sealed.number(), + hash: sealed.hash(), + encoded: alloy_rlp::encode(sealed.into_block()), + }) + } +} + +/// Broadcast every newly persisted leader block in canonical order. +pub(crate) async fn broadcast_persisted_blocks

(provider: P, commands: mpsc::Sender) +where + P: PersistedBlockSource, +{ + // Handle race conditions carefully at startup. Read before subscribing, then reconcile after subscribing. + // This closes both startup windows: a block persisted before the subscription is found by the + // second read, while a block persisted after the subscription is retained by the stream. + let mut last_broadcast = match provider.last_block_number() { + Ok(number) => number, + Err(err) => { + tracing::error!(target: "zone::p2p", %err, "Failed reading persisted zone head"); + return; + } + }; + let mut persisted = provider.persisted_block_stream(); + let startup_tip = match provider.last_block_number() { + Ok(number) => number, + Err(err) => { + tracing::error!(target: "zone::p2p", %err, "Failed reconciling persisted zone head"); + return; + } + }; + + if let Err(err) = + broadcast_persisted_range(&provider, &commands, &mut last_broadcast, startup_tip, None) + .await + { + tracing::error!(target: "zone::p2p", %err, "Failed broadcasting persisted zone blocks"); + return; + } + + while let Some(persisted_tip) = persisted.next().await { + if persisted_tip.number < last_broadcast { + tracing::error!( + target: "zone::p2p", + persisted = persisted_tip.number, + last_broadcast, + "Persisted zone head moved backwards" + ); + return; + } + + if let Err(err) = broadcast_persisted_range( + &provider, + &commands, + &mut last_broadcast, + persisted_tip.number, + Some(persisted_tip.hash), + ) + .await + { + tracing::error!(target: "zone::p2p", %err, "Failed broadcasting persisted zone blocks"); + return; + } + } + debug!(target: "zone::p2p", "Persisted block stream closed"); +} + +async fn broadcast_persisted_range

( + provider: &P, + commands: &mpsc::Sender, + last_broadcast: &mut u64, + tip_number: u64, + expected_tip_hash: Option, +) -> eyre::Result<()> +where + P: PersistedBlockSource, +{ + for number in last_broadcast.saturating_add(1)..=tip_number { + let block = provider.encoded_block_by_number(number)?; + let number = block.number; + let hash = block.hash; + if number == tip_number + && let Some(expected) = expected_tip_hash + && hash != expected + { + eyre::bail!( + "persisted zone block hash does not match notification at height {number}: expected={expected}, actual={hash}" + ); + } + commands + .send(P2pCommand::BroadcastBlock(block.encoded)) + .await + .map_err(|_| eyre::eyre!("P2P command channel closed"))?; + debug!(target: "zone::p2p", number, ?hash, "Queued persisted block for followers"); + *last_broadcast = number; + } + Ok(()) +} + +/// Decode, fully execute, and canonicalize blocks received by a follower. +pub(crate) async fn import_leader_blocks

( + provider: P, + engine: ConsensusEngineHandle, + mut events: mpsc::Receiver, + _commands: mpsc::Sender, +) where + P: BlockNumReader + HeaderProvider

+ Clone + Send + Sync + 'static, +{ + while let Some(event) = events.recv().await { + let P2pEvent::BlockReceived { block, .. } = event else { + continue; + }; + if let Err(err) = import_leader_block(&provider, &engine, &block).await { + tracing::error!(target: "zone::p2p", %err, "Rejected leader block"); + } + } + debug!(target: "zone::p2p", "P2P event channel closed"); +} + +async fn import_leader_block

( + provider: &P, + engine: &ConsensusEngineHandle, + encoded: &[u8], +) -> eyre::Result<()> +where + P: BlockNumReader + HeaderProvider

+ Clone + Send + Sync + 'static, +{ + // Check the received block + let mut input = encoded; + let block = Block::decode(&mut input) + .map_err(|err| eyre::eyre!("invalid RLP-encoded zone block: {err}"))?; + if !input.is_empty() { + eyre::bail!("encoded zone block has {} trailing bytes", input.len()); + } + + let block = SealedBlock::seal_slow(block); + let block_number = block.number(); + let hash = block.hash(); + let best_block = provider.best_block_number()?; + + // 1. Block number is correct + if block_number <= best_block { + let existing = provider.sealed_header(block_number)?.ok_or_else(|| { + eyre::eyre!("missing local canonical header at height {block_number}") + })?; + if existing.hash() == hash { + debug!(target: "zone::p2p", block_number, ?hash, "Ignoring duplicate leader block"); + return Ok(()); + } + eyre::bail!( + "leader block conflicts with canonical block at height {block_number}: local={}, received={hash}", + existing.hash() + ); + } + + let expected_number = best_block.saturating_add(1); + if block_number != expected_number { + eyre::bail!( + "leader block gap: local head is {best_block}, received height {block_number}, expected {expected_number}; backfill is not implemented yet" + ); + } + + // 2. Block's parent hash is correct + let parent = provider + .sealed_header(best_block)? + .ok_or_else(|| eyre::eyre!("missing local canonical head at height {best_block}"))?; + if block.parent_hash() != parent.hash() { + eyre::bail!( + "leader block parent mismatch at height {block_number}: local={}, received={}", + parent.hash(), + block.parent_hash() + ); + } + + // 3. All txns in the block execute properly + let payload = ZonePayloadTypes::block_to_payload(block, None); + let status = engine.new_payload(payload).await?; + if !status.is_valid() { + eyre::bail!("execution engine rejected leader block {block_number} ({hash}): {status:?}"); + } + + // 4. Forkchoice + let forkchoice = ForkchoiceState::same_hash(hash); + let result = engine.fork_choice_updated(forkchoice, None).await?; + if !result.is_valid() { + eyre::bail!( + "execution engine rejected forkchoice for block {block_number} ({hash}): {result:?}" + ); + } + + info!(target: "zone::p2p", block_number, ?hash, "Imported canonical leader block"); + Ok(()) +} + +#[cfg(test)] +mod tests { + use std::sync::{ + Arc, + atomic::{AtomicUsize, Ordering}, + }; + + use futures::{StreamExt as _, stream}; + + use super::{ + EncodedPersistedBlock, PersistedBlockSource, PersistedTip, broadcast_persisted_blocks, + }; + use alloy_primitives::B256; + use zone_p2p::P2pCommand; + + #[derive(Clone)] + struct StartupRaceSource { + reads: Arc, + tip: PersistedTip, + } + + impl PersistedBlockSource for StartupRaceSource { + fn last_block_number(&self) -> eyre::Result { + let read = self.reads.fetch_add(1, Ordering::SeqCst); + Ok(if read == 0 { + self.tip.number - 1 + } else { + self.tip.number + }) + } + + fn persisted_block_stream(&self) -> futures::stream::BoxStream<'static, PersistedTip> { + stream::iter([self.tip]).boxed() + } + + fn encoded_block_by_number(&self, number: u64) -> eyre::Result { + assert_eq!(number, self.tip.number); + Ok(EncodedPersistedBlock { + number, + hash: self.tip.hash, + encoded: vec![number as u8], + }) + } + } + + #[tokio::test] + async fn broadcasts_block_persisted_during_startup_reconciliation_once() { + let source = StartupRaceSource { + reads: Arc::new(AtomicUsize::new(0)), + tip: PersistedTip { + number: 1, + hash: B256::repeat_byte(0x11), + }, + }; + let (commands, mut command_rx) = tokio::sync::mpsc::channel(4); + + broadcast_persisted_blocks(source, commands).await; + + assert_eq!( + command_rx.recv().await, + Some(P2pCommand::BroadcastBlock(vec![1])) + ); + assert_eq!(command_rx.recv().await, None); + } +} diff --git a/crates/node/tests/it/e2e.rs b/crates/node/tests/it/e2e.rs index f8e5db088..35910531a 100644 --- a/crates/node/tests/it/e2e.rs +++ b/crates/node/tests/it/e2e.rs @@ -5,7 +5,7 @@ //! subscriber retries a dummy URL in the background, but L2 execution is fully //! exercised via queue injection (with the L1 state cache seeded for precompile reads). -use std::net::TcpListener; +use std::{net::TcpListener, time::Duration}; use alloy::primitives::{Address, B256, Bytes, TxKind, U256, address}; use alloy_consensus::Transaction; @@ -24,11 +24,54 @@ use zone_l1::ChainTempoStateExt; use crate::utils::{ DEFAULT_POLL, DEFAULT_TIMEOUT, L1Fixture, WITHDRAWAL_TX_GAS, ZoneTestNode, approve_outbox, leader_p2p_config, local_dev_zone_account, poll_until, seed_fixture_for_zone, - start_chain_id_rpc, start_local_zone_with_fixture, + start_chain_id_rpc, start_local_p2p_pair, start_local_zone_with_fixture, }; const CONTRACT_CREATION_TX_GAS: u64 = 1_000_000; +/// A follower imports the leader's executed block and exposes the resulting state over RPC. +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn test_p2p_follower_tracks_leader_balance() -> eyre::Result<()> { + reth_tracing::init_test_tracing(); + + let (leader, follower, mut fixture) = start_local_p2p_pair(10).await?; + + // Commonware deliberately drops messages for offline peers. Wait for + // peer dial/handshake (loopback dials every 500ms) before producing the + // first block. + tokio::time::sleep(Duration::from_secs(1)).await; + + fixture.inject_empty_block(leader.deposit_queue()); + leader.wait_for_block_number(1, DEFAULT_TIMEOUT).await?; + follower.wait_for_block_number(1, DEFAULT_TIMEOUT).await?; + + let depositor = address!("0x0000000000000000000000000000000000001234"); + let recipient = address!("0x0000000000000000000000000000000000005678"); + let amount = 1_000_000_u128; + let deposit = fixture.make_deposit(PATH_USD_ADDRESS, depositor, recipient, amount); + fixture.inject_deposits(leader.deposit_queue(), vec![deposit]); + + leader + .wait_for_balance( + PATH_USD_ADDRESS, + recipient, + U256::from(amount), + DEFAULT_TIMEOUT, + ) + .await?; + let follower_balance = follower + .wait_for_balance( + PATH_USD_ADDRESS, + recipient, + U256::from(amount), + DEFAULT_TIMEOUT, + ) + .await?; + follower.wait_for_block_number(2, DEFAULT_TIMEOUT).await?; + assert_eq!(follower_balance, U256::from(amount)); + Ok(()) +} + /// A P2P bind failure is fatal rather than leaving the node running without P2P. #[tokio::test(flavor = "multi_thread", worker_threads = 4)] async fn test_p2p_listener_failure_stops_node() -> eyre::Result<()> { diff --git a/crates/node/tests/it/utils.rs b/crates/node/tests/it/utils.rs index d4cf3903d..ccdc2e70b 100644 --- a/crates/node/tests/it/utils.rs +++ b/crates/node/tests/it/utils.rs @@ -41,6 +41,7 @@ use tempo_contracts::precompiles::{ use tempo_precompiles::{PATH_USD_ADDRESS, tip403_registry::ALLOW_ALL_POLICY_ID}; use tempo_primitives::{TempoHeader, transaction::tt_signature::TempoSignature}; use tempo_zone_contracts::{ZONE_FACTORY_ADDRESS, ZONE_OUTBOX_ADDRESS}; +use tokio::io::{AsyncReadExt as _, AsyncWriteExt as _}; use zone_chainspec::ZoneChainSpec; use zone_l1::{ Deposit, DepositQueue, EnabledToken, EncryptedDeposit, L1Deposit, L1PortalEvents, L1StateCache, @@ -492,6 +493,31 @@ impl ZoneTestNode { .await } + /// Wait for the zone L2 RPC head to reach at least `target`. + /// + /// This polls `eth_blockNumber`, which is useful when a test needs to assert + /// that a follower imported leader-produced zone blocks over P2P. + pub(crate) async fn wait_for_block_number( + &self, + target: u64, + timeout: Duration, + ) -> eyre::Result { + let provider = self.provider(); + poll_until( + timeout, + DEFAULT_POLL, + &format!("eth_blockNumber >= {target}"), + || { + let provider = &provider; + async move { + let n = provider.get_block_number().await?; + if n >= target { Ok(Some(n)) } else { Ok(None) } + } + }, + ) + .await + } + /// Read a TIP-20 token balance on this zone (single-shot, no polling). pub(crate) async fn balance_of(&self, token: Address, account: Address) -> eyre::Result { use tempo_contracts::precompiles::ITIP20; @@ -606,6 +632,7 @@ impl ZoneTestNode { 8, initial_tokens, None, + true, ) .await } @@ -630,6 +657,7 @@ impl ZoneTestNode { withdrawal_batch_interval_blocks, Some(vec![]), None, + true, ) .await } @@ -705,6 +733,7 @@ impl ZoneTestNode { 8, Some(vec![]), Some(p2p_config), + true, ) .await } @@ -747,6 +776,7 @@ impl ZoneTestNode { 8, Some(vec![]), None, + true, ) .await } @@ -762,6 +792,7 @@ impl ZoneTestNode { withdrawal_batch_interval_blocks: u64, initial_tokens: Option>, p2p_config: Option, + spawn_engine: bool, ) -> eyre::Result { let tasks = Runtime::test(); let is_local_dummy_l1 = l1_ws_url == DUMMY_L1_URL; @@ -788,6 +819,7 @@ impl ZoneTestNode { if let Some(initial_tokens) = initial_tokens { zone_node = zone_node.with_initial_tokens(initial_tokens); } + let p2p_enabled = p2p_config.is_some(); if let Some(p2p_config) = p2p_config { zone_node = zone_node.with_p2p(p2p_config); } @@ -805,6 +837,10 @@ impl ZoneTestNode { ) .apply(|mut c| { c.network.discovery.disable_discovery = true; + if p2p_enabled { + c.engine.persistence_threshold = 0; + c.engine.memory_block_buffer_target = 0; + } c }); @@ -821,34 +857,36 @@ impl ZoneTestNode { .launch_with_debug_capabilities() .await?; - let l1_provider = ProviderBuilder::new_with_network::() - .connect(&l1_provider_url) - .await? - .erased(); - let policy_provider = zone_l1::PolicyProvider::new( - policy_cache.clone(), - l1_provider, - tokio::runtime::Handle::current(), - ); - let provider = node_handle.node.provider(); - let last_header = provider - .sealed_header(provider.best_block_number()?)? - .ok_or_else(|| eyre::eyre!("no latest block header"))?; - let engine = zone_node::ZoneEngine::new( - provider.chain_spec(), - node_handle.node.add_ons_handle.beacon_engine_handle.clone(), - node_handle.node.payload_builder_handle.clone(), - deposit_queue.clone(), - last_header, - sequencer_signer.address(), - SecretKey::from(sequencer_signer.credential()), - portal_address, - policy_provider, - ); - node_handle - .node - .task_executor - .spawn_critical_task("zone-engine", engine.run()); + if spawn_engine { + let l1_provider = ProviderBuilder::new_with_network::() + .connect(&l1_provider_url) + .await? + .erased(); + let policy_provider = zone_l1::PolicyProvider::new( + policy_cache.clone(), + l1_provider, + tokio::runtime::Handle::current(), + ); + let provider = node_handle.node.provider(); + let last_header = provider + .sealed_header(provider.best_block_number()?)? + .ok_or_else(|| eyre::eyre!("no latest block header"))?; + let engine = zone_node::ZoneEngine::new( + provider.chain_spec(), + node_handle.node.add_ons_handle.beacon_engine_handle.clone(), + node_handle.node.payload_builder_handle.clone(), + deposit_queue.clone(), + last_header, + sequencer_signer.address(), + SecretKey::from(sequencer_signer.credential()), + portal_address, + policy_provider, + ); + node_handle + .node + .task_executor + .spawn_critical_task("zone-engine", engine.run()); + } let http_url: url::Url = node_handle .node @@ -2547,6 +2585,147 @@ pub(crate) async fn start_local_zone_with_fixture( Ok((zone, fixture)) } +/// Start a leader and follower with identical genesis state and authenticated P2P identities. +pub(crate) async fn start_local_p2p_pair( + seed_blocks: u64, +) -> eyre::Result<(ZoneTestNode, ZoneTestNode, L1Fixture)> { + fn available_address() -> eyre::Result { + let listener = TcpListener::bind("127.0.0.1:0")?; + Ok(listener.local_addr()?) + } + + async fn spawn_test_l1_rpc() -> eyre::Result { + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await?; + let address = listener.local_addr()?; + tokio::spawn(async move { + loop { + let Ok((mut stream, _)) = listener.accept().await else { + return; + }; + tokio::spawn(async move { + let mut request = vec![0_u8; 16 * 1024]; + let Ok(read) = stream.read(&mut request).await else { + return; + }; + let request = String::from_utf8_lossy(&request[..read]); + let body = request.split("\r\n\r\n").nth(1).unwrap_or_default(); + let value: serde_json::Value = serde_json::from_str(body).unwrap_or_default(); + let id = value.get("id").cloned().unwrap_or(serde_json::Value::Null); + let result = match value.get("method").and_then(|method| method.as_str()) { + Some("eth_chainId") => serde_json::json!("0x539"), + Some("eth_blockNumber") => serde_json::json!("0x0"), + _ => serde_json::Value::Null, + }; + let response_body = serde_json::json!({ + "jsonrpc": "2.0", + "id": id, + "result": result, + }) + .to_string(); + let response = format!( + "HTTP/1.1 200 OK\r\ncontent-type: application/json\r\ncontent-length: {}\r\nconnection: close\r\n\r\n{}", + response_body.len(), + response_body + ); + let _ = stream.write_all(response.as_bytes()).await; + }); + } + }); + Ok(format!("http://{address}")) + } + + let addresses = [ + available_address()?, + available_address()?, + available_address()?, + ]; + let identities = [ + Ed25519PrivateKey::from_seed(101), + Ed25519PrivateKey::from_seed(102), + Ed25519PrivateKey::from_seed(103), + ]; + let public_keys = identities.each_ref().map(|key| key.public_key()); + + let unique = NEXT_CHAIN_ID.fetch_add(1, Ordering::Relaxed); + let config_dir = std::env::temp_dir().join(format!( + "tempo-zone-p2p-test-{}-{unique}", + std::process::id() + )); + std::fs::create_dir_all(&config_dir)?; + let manifest_path = config_dir.join("manifest.toml"); + let mut manifest = format!( + "zone_id = 0\nleader_ed25519_public_key = \"{}\"\n", + const_hex::encode_prefixed(public_keys[0].as_ref()) + ); + for (index, (public_key, address)) in public_keys.iter().zip(addresses).enumerate() { + manifest.push_str(&format!( + "\n[[nodes]]\nname = \"node-{index}\"\ned25519_public_key = \"{}\"\naddress = \"{address}\"\n", + const_hex::encode_prefixed(public_key.as_ref()) + )); + } + std::fs::write(&manifest_path, manifest)?; + let mut configs = Vec::with_capacity(2); + for (index, role) in [(0, Role::Leader), (1, Role::Follower)] { + let key_path = config_dir.join(format!("node-{index}.key")); + std::fs::write( + &key_path, + const_hex::encode_prefixed(identities[index].encode().as_ref()), + )?; + configs.push(P2pConfig::load( + &manifest_path, + &key_path, + addresses[index], + false, + 0, + Some(role), + )?); + } + let _ = std::fs::remove_dir_all(&config_dir); + + let chain_id = next_unique_chain_id(); + let l1_rpc_url = spawn_test_l1_rpc().await?; + let genesis: Genesis = serde_json::from_str(zone_node::genesis::GENESIS_TEMPLATE_JSON)?; + let signer = l1_dev_signer(); + let leader = ZoneTestNode::launch_with_genesis_and_withdrawal_batch_interval( + l1_rpc_url.clone(), + Address::ZERO, + None, + chain_id, + Some(genesis.clone()), + signer.clone(), + 8, + Some(vec![]), + Some(configs.remove(0)), + true, + ) + .await?; + let follower = ZoneTestNode::launch_with_genesis_and_withdrawal_batch_interval( + l1_rpc_url, + Address::ZERO, + None, + chain_id, + Some(genesis), + signer, + 8, + Some(vec![]), + Some(configs.remove(0)), + false, + ) + .await?; + + let fixture = L1Fixture::new(); + for zone in [&leader, &follower] { + seed_local_policy_cache(zone.policy_cache()); + fixture.seed_l1_cache( + zone.l1_state_cache(), + Address::ZERO, + Address::ZERO, + seed_blocks, + ); + } + Ok((leader, follower, fixture)) +} + pub(crate) fn leader_p2p_config(listen: SocketAddr) -> eyre::Result { fn available_address() -> eyre::Result { Ok(TcpListener::bind("127.0.0.1:0")?.local_addr()?) diff --git a/crates/p2p/Cargo.toml b/crates/p2p/Cargo.toml index 3b34eae3a..058bf4c6c 100644 --- a/crates/p2p/Cargo.toml +++ b/crates/p2p/Cargo.toml @@ -21,6 +21,7 @@ commonware-runtime = { workspace = true, features = ["external"] } commonware-utils.workspace = true const-hex.workspace = true +derive_more.workspace = true eyre.workspace = true serde = { workspace = true, features = ["derive"] } thiserror.workspace = true diff --git a/crates/p2p/src/lib.rs b/crates/p2p/src/lib.rs index 6be0a1bab..a26cd1b0c 100644 --- a/crates/p2p/src/lib.rs +++ b/crates/p2p/src/lib.rs @@ -3,7 +3,6 @@ mod identity; mod manifest; -mod messages; mod network; mod runtime; diff --git a/crates/p2p/src/manifest.rs b/crates/p2p/src/manifest.rs index a1f049930..1229aa71a 100644 --- a/crates/p2p/src/manifest.rs +++ b/crates/p2p/src/manifest.rs @@ -9,10 +9,13 @@ use commonware_codec::DecodeExt as _; use commonware_cryptography::ed25519::PublicKey; use commonware_p2p::{Address, Ingress}; use commonware_utils::Hostname; +use derive_more::{Display, FromStr}; use serde::Deserialize; /// The role assigned to a node by the manifest. -#[derive(Debug, Clone, Copy, PartialEq, Eq)] +#[derive(Debug, Clone, Copy, PartialEq, Eq, Display, FromStr)] +#[display(rename_all = "lowercase")] +#[from_str(rename_all = "lowercase")] pub enum Role { /// Builds blocks and runs the existing sequencer settlement tasks. Leader, @@ -20,27 +23,6 @@ pub enum Role { Follower, } -impl fmt::Display for Role { - fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { - match self { - Self::Leader => f.write_str("leader"), - Self::Follower => f.write_str("follower"), - } - } -} - -impl std::str::FromStr for Role { - type Err = String; - - fn from_str(value: &str) -> Result { - match value { - "leader" => Ok(Self::Leader), - "follower" => Ok(Self::Follower), - _ => Err(format!("expected `leader` or `follower`, got `{value}`")), - } - } -} - /// A validated P2P address from the manifest. #[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord)] pub enum ManifestAddress { diff --git a/crates/p2p/src/messages.rs b/crates/p2p/src/messages.rs deleted file mode 100644 index cbc2c622e..000000000 --- a/crates/p2p/src/messages.rs +++ /dev/null @@ -1,67 +0,0 @@ -use commonware_codec::{Error, FixedSize, Read, ReadExt as _, Write}; -use commonware_runtime::{Buf, BufMut}; - -const HEARTBEAT_TAG: u8 = 0; -const HEARTBEAT_ACK_TAG: u8 = 1; - -/// PoC messages exchanged on the P2P control channel. -/// -/// TODO: This heartbeat exchange exists only to exercise authenticated -/// message transport. Replace it with the v0 block, ACK/signature, transaction-forwarding, -/// and backfill protocols. It is not a leader-election heartbeat or a source of leadership -/// or finality, and will be removed with the next PR. -#[derive(Debug, Clone, Copy, PartialEq, Eq)] -pub(crate) enum ControlMessage { - Heartbeat { nonce: u64 }, - HeartbeatAck { nonce: u64 }, -} - -impl Write for ControlMessage { - fn write(&self, buffer: &mut impl BufMut) { - match self { - Self::Heartbeat { nonce } => { - HEARTBEAT_TAG.write(buffer); - nonce.write(buffer); - } - Self::HeartbeatAck { nonce } => { - HEARTBEAT_ACK_TAG.write(buffer); - nonce.write(buffer); - } - } - } -} - -impl FixedSize for ControlMessage { - const SIZE: usize = u8::SIZE + u64::SIZE; -} - -impl Read for ControlMessage { - type Cfg = (); - - fn read_cfg(buffer: &mut impl Buf, _cfg: &Self::Cfg) -> Result { - let tag = u8::read(buffer)?; - let nonce = u64::read(buffer)?; - match tag { - HEARTBEAT_TAG => Ok(Self::Heartbeat { nonce }), - HEARTBEAT_ACK_TAG => Ok(Self::HeartbeatAck { nonce }), - other => Err(Error::InvalidEnum(other)), - } - } -} - -#[cfg(test)] -mod tests { - use commonware_codec::{DecodeExt as _, Encode as _}; - - use super::ControlMessage; - - #[test] - fn control_messages_round_trip() { - for message in [ - ControlMessage::Heartbeat { nonce: 42 }, - ControlMessage::HeartbeatAck { nonce: u64::MAX }, - ] { - assert_eq!(ControlMessage::decode(message.encode()).unwrap(), message); - } - } -} diff --git a/crates/p2p/src/network.rs b/crates/p2p/src/network.rs index 8b423abe4..3b366adf7 100644 --- a/crates/p2p/src/network.rs +++ b/crates/p2p/src/network.rs @@ -9,10 +9,12 @@ use eyre::WrapErr as _; use crate::ZoneManifest; -/// The final block/ack/tx/backfill protocol reserves channel IDs 0 through 3. -pub(crate) const CONTROL_CHANNEL: u64 = 4; -pub(crate) const CONTROL_BACKLOG: usize = 128; -pub(crate) const MAX_MESSAGE_SIZE: u32 = 10 * 1024 * 1024; +/// Leader-to-follower sealed block replication channel. +pub(crate) const BLOCK_CHANNEL: u64 = 0; +pub(crate) const BLOCK_BACKLOG: usize = 128; + +// At 30M gas, calldata is bounded below 7.5 MiB; leave headroom for block overhead. +pub(crate) const MAX_MESSAGE_SIZE: u32 = 20 * 1024 * 1024; /// Version of the Tempo Zone P2P wire protocol. pub(crate) const WIRE_PROTOCOL_VERSION: u8 = 0; @@ -90,8 +92,8 @@ fn namespace(zone_id: u32, network_id: P2pNetworkId) -> Vec { namespace } -pub(crate) fn control_quota() -> Quota { - Quota::per_second(NZU32!(4)) +pub(crate) fn block_quota() -> Quota { + Quota::per_second(NZU32!(128)) } #[cfg(test)] diff --git a/crates/p2p/src/runtime.rs b/crates/p2p/src/runtime.rs index b3f8ce1c5..8abb150eb 100644 --- a/crates/p2p/src/runtime.rs +++ b/crates/p2p/src/runtime.rs @@ -1,6 +1,5 @@ -use std::{collections::BTreeSet, net::SocketAddr, path::Path, sync::Arc, time::Duration}; +use std::{net::SocketAddr, path::Path, sync::Arc, time::Duration}; -use commonware_codec::{DecodeExt as _, Encode as _}; use commonware_cryptography::ed25519::PublicKey; use commonware_p2p::{ AddressableManager as _, Receiver as _, Recipients, Sender as _, authenticated::lookup, @@ -8,17 +7,17 @@ use commonware_p2p::{ use commonware_runtime::{Runner as _, Spawner as _}; use tokio::sync::{mpsc, oneshot}; use tokio_util::sync::CancellationToken; -use tracing::{debug, info, warn}; +use tracing::{debug, error, info, warn}; use crate::{ P2pNetworkId, Role, ZoneManifest, identity::Ed25519Identity, - messages::ControlMessage, - network::{self, CONTROL_BACKLOG, CONTROL_CHANNEL}, + network::{self, BLOCK_BACKLOG, BLOCK_CHANNEL, MAX_MESSAGE_SIZE}, }; -const HEARTBEAT_INTERVAL: Duration = Duration::from_secs(5); const SHUTDOWN_TIMEOUT: Duration = Duration::from_secs(5); +const BROADCAST_RETRY_INTERVAL: Duration = Duration::from_secs(1); +const BROADCAST_RETRY_TIMEOUT: Duration = Duration::from_secs(30); const COMMAND_BACKLOG: usize = 128; const EVENT_BACKLOG: usize = 128; @@ -100,59 +99,25 @@ fn validate_ip_check_configuration( } /// Outbound protocol commands accepted by the dedicated P2P runtime. -/// -/// The Commonware sender remains owned by the dedicated runtime. Callers communicate with it -/// through the bounded Tokio channel exposed by [`P2pHandle`]. These PoC variants will be -/// replaced by the block, ACK/signature, transaction-forwarding, and backfill commands. #[derive(Debug, Clone, PartialEq, Eq)] pub enum P2pCommand { - /// Send a PoC heartbeat to one authenticated peer. - Heartbeat { - /// Intended peer's Ed25519 Commonware identity. - recipient: PublicKey, - /// Heartbeat nonce. - nonce: u64, - }, - /// Acknowledge a PoC heartbeat from one authenticated peer. - HeartbeatAck { - /// Intended peer's Ed25519 Commonware identity. - recipient: PublicKey, - /// Acknowledged nonce. - nonce: u64, - }, + /// Broadcast one RLP-encoded sealed zone block to all configured followers. + BroadcastBlock(Vec), } -impl P2pCommand { - fn into_wire_parts(self) -> (PublicKey, ControlMessage) { - match self { - Self::Heartbeat { recipient, nonce } => { - (recipient, ControlMessage::Heartbeat { nonce }) - } - Self::HeartbeatAck { recipient, nonce } => { - (recipient, ControlMessage::HeartbeatAck { nonce }) - } - } - } -} - -/// Observable lifecycle and heartbeat events emitted by the P2P runtime. +/// Observable lifecycle and block events emitted by the P2P runtime. #[derive(Debug, Clone, PartialEq, Eq)] pub enum P2pEvent { - /// The network and control channel were started. + /// The network and block channel were started. Started { role: Role, ed25519_public_key: PublicKey, listen: SocketAddr, }, - /// A leader received a heartbeat from a follower. - HeartbeatReceived { - follower_ed25519_public_key: PublicKey, - nonce: u64, - }, - /// A follower received the leader's acknowledgement. - HeartbeatAcknowledged { + /// A follower received an encoded sealed block from its configured leader. + BlockReceived { leader_ed25519_public_key: PublicKey, - nonce: u64, + block: Vec, }, } @@ -228,27 +193,19 @@ async fn join_runtime_thread(thread: std::thread::JoinHandle<()>) -> eyre::Resul .map_err(|_| eyre::eyre!("P2P runtime thread panicked")) } -/// Starts Commonware and the role-specific PoC heartbeat actor on a dedicated OS thread. +/// Starts Commonware block transport on a dedicated OS thread. pub fn spawn_p2p(config: P2pConfig, network_id: P2pNetworkId) -> eyre::Result { let shutdown = CancellationToken::new(); let thread_shutdown = shutdown.clone(); let (stopped_tx, stopped) = oneshot::channel(); let (commands, command_rx) = mpsc::channel(COMMAND_BACKLOG); - let runtime_commands = commands.clone(); let (events_tx, events) = mpsc::channel(EVENT_BACKLOG); let thread = std::thread::Builder::new() .name(format!("zone-p2p-{}", config.role())) .spawn(move || { - let result = run( - config, - network_id, - thread_shutdown, - runtime_commands, - command_rx, - events_tx, - ) - .map_err(|err| format!("{err:?}")); + let result = run(config, network_id, thread_shutdown, command_rx, events_tx) + .map_err(|err| format!("{err:?}")); let _ = stopped_tx.send(result); }) .map_err(|err| eyre::eyre!("failed spawning P2P runtime thread: {err}"))?; @@ -268,7 +225,6 @@ fn run( config: P2pConfig, network_id: P2pNetworkId, shutdown: CancellationToken, - commands: mpsc::Sender, command_rx: mpsc::Receiver, events: mpsc::Sender, ) -> eyre::Result<()> { @@ -287,11 +243,8 @@ fn run( network_id, )?; oracle.track(0, peers).await; - let (sender, receiver) = commonware.register( - CONTROL_CHANNEL, - network::control_quota(), - CONTROL_BACKLOG, - ); + let (sender, receiver) = + commonware.register(BLOCK_CHANNEL, network::block_quota(), BLOCK_BACKLOG); let mut network_task = commonware.start(); if config.bypass_ip_check { @@ -317,26 +270,33 @@ fn run( }) .await; - let command_loop = run_commands(sender, command_rx); + let followers = config + .manifest + .nodes() + .iter() + .filter(|node| config.manifest.role_of(node.ed25519_public_key()) == Some(Role::Follower)) + .map(|node| node.ed25519_public_key().clone()) + .collect(); + let command_loop = run_commands(config.role, followers, sender, command_rx); tokio::pin!(command_loop); - let heartbeat = run_heartbeat( + let receive_loop = run_block_receiver( config.role, config.manifest, - commands, receiver, events, ); - tokio::pin!(heartbeat); + tokio::pin!(receive_loop); let result = tokio::select! { + biased; () = shutdown.cancelled() => Ok(()), network_result = &mut network_task => match network_result { Ok(()) => Err(eyre::eyre!("Commonware network stopped unexpectedly")), Err(err) => Err(eyre::eyre!("Commonware network failed: {err}")), }, result = &mut command_loop => result, - result = &mut heartbeat => result, + result = &mut receive_loop => result, }; context @@ -348,149 +308,89 @@ fn run( } async fn run_commands( + role: Role, + followers: Vec, mut sender: lookup::Sender, mut commands: mpsc::Receiver, ) -> eyre::Result<()> { while let Some(command) = commands.recv().await { - let (recipient, message) = command.into_wire_parts(); - let sent = sender - .send(Recipients::One(recipient.clone()), message.encode(), true) - .await - .map_err(|err| eyre::eyre!("failed sending P2P control message: {err}"))?; - if sent.is_empty() { - debug!(target: "zone::p2p", %recipient, ?message, "Peer is not connected; control message was not sent"); + if role != Role::Leader { + warn!(target: "zone::p2p", ?command, "Ignoring outbound block command on follower"); + continue; + } + let P2pCommand::BroadcastBlock(block) = command; + if block.len() > MAX_MESSAGE_SIZE as usize { + error!( + target: "zone::p2p", + block_size_bytes = block.len(), + max_message_size_bytes = MAX_MESSAGE_SIZE, + "Canonical block exceeds the P2P message size limit; block was not broadcast" + ); + continue; + } + let sent = tokio::time::timeout(BROADCAST_RETRY_TIMEOUT, async { + loop { + let sent = sender + .send(Recipients::Some(followers.clone()), block.clone(), true) + .await + .map_err(|err| eyre::eyre!("failed broadcasting zone block: {err}"))?; + if !sent.is_empty() || followers.is_empty() { + return Ok::<_, eyre::Report>(sent); + } + debug!( + target: "zone::p2p", + "No followers are connected; retrying canonical block broadcast" + ); + tokio::time::sleep(BROADCAST_RETRY_INTERVAL).await; + } + }) + .await; + let sent = match sent { + Ok(sent) => sent?, + Err(_) => { + warn!( + target: "zone::p2p", + timeout_secs = BROADCAST_RETRY_TIMEOUT.as_secs(), + "No followers connected before block broadcast timed out" + ); + continue; + } + }; + if sent.len() != followers.len() { + debug!(target: "zone::p2p", connected = sent.len(), configured = followers.len(), "Some followers are not connected; block was not sent to them"); } } Err(eyre::eyre!("P2P command channel closed unexpectedly")) } -async fn run_heartbeat( +async fn run_block_receiver( role: Role, manifest: Arc, - commands: mpsc::Sender, - receiver: lookup::Receiver, - events: mpsc::Sender, -) -> eyre::Result<()> { - match role { - Role::Leader => run_leader_heartbeat(manifest, commands, receiver, events).await, - Role::Follower => run_follower_heartbeat(manifest, commands, receiver, events).await, - } -} - -async fn run_leader_heartbeat( - manifest: Arc, - commands: mpsc::Sender, mut receiver: lookup::Receiver, events: mpsc::Sender, ) -> eyre::Result<()> { - let mut seen = BTreeSet::new(); + let leader = manifest.leader_ed25519_public_key().clone(); loop { let (peer, bytes) = receiver .recv() .await - .map_err(|err| eyre::eyre!("control channel receive failed: {err}"))?; - let message = match ControlMessage::decode(bytes) { - Ok(message) => message, - Err(err) => { - warn!(target: "zone::p2p", %peer, %err, "Ignoring invalid control message"); - continue; - } - }; - if manifest.role_of(&peer) != Some(Role::Follower) { - warn!(target: "zone::p2p", %peer, ?message, "Ignoring control message from non-follower"); + .map_err(|err| eyre::eyre!("block channel receive failed: {err}"))?; + if role == Role::Leader { + warn!(target: "zone::p2p", %peer, "Leader received an unexpected block message"); continue; } - - match message { - ControlMessage::Heartbeat { nonce } => { - if seen.insert(peer.clone()) { - info!(target: "zone::p2p", follower = %peer, "Received first follower heartbeat"); - } else { - debug!(target: "zone::p2p", follower = %peer, nonce, "Received follower heartbeat"); - } - let _ = events - .send(P2pEvent::HeartbeatReceived { - follower_ed25519_public_key: peer.clone(), - nonce, - }) - .await - .ok(); - commands - .send(P2pCommand::HeartbeatAck { - recipient: peer, - nonce, - }) - .await - .map_err(|_| eyre::eyre!("P2P command channel closed"))?; - } - ControlMessage::HeartbeatAck { nonce } => { - warn!(target: "zone::p2p", follower = %peer, nonce, "Leader received unexpected heartbeat acknowledgement"); - } - } - } -} - -async fn run_follower_heartbeat( - manifest: Arc, - commands: mpsc::Sender, - mut receiver: lookup::Receiver, - events: mpsc::Sender, -) -> eyre::Result<()> { - let leader = manifest.leader_ed25519_public_key().clone(); - let mut interval = tokio::time::interval(HEARTBEAT_INTERVAL); - interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip); - let mut next_nonce = 0_u64; - let mut acknowledged = false; - - loop { - tokio::select! { - _ = interval.tick() => { - let nonce = next_nonce; - next_nonce = next_nonce.wrapping_add(1); - commands - .send(P2pCommand::Heartbeat { - recipient: leader.clone(), - nonce, - }) - .await - .map_err(|_| eyre::eyre!("P2P command channel closed"))?; - } - received = receiver.recv() => { - let (peer, bytes) = received - .map_err(|err| eyre::eyre!("control channel receive failed: {err}"))?; - let message = match ControlMessage::decode(bytes) { - Ok(message) => message, - Err(err) => { - warn!(target: "zone::p2p", %peer, %err, "Ignoring invalid control message"); - continue; - } - }; - if peer != leader { - warn!(target: "zone::p2p", %peer, ?message, "Ignoring control message from non-leader"); - continue; - } - match message { - ControlMessage::HeartbeatAck { nonce } => { - if !acknowledged { - acknowledged = true; - info!(target: "zone::p2p", %leader, "Established heartbeat exchange with leader"); - } else { - debug!(target: "zone::p2p", %leader, nonce, "Leader acknowledged heartbeat"); - } - let _ = events - .send(P2pEvent::HeartbeatAcknowledged { - leader_ed25519_public_key: leader.clone(), - nonce, - }) - .await; - } - ControlMessage::Heartbeat { nonce } => { - warn!(target: "zone::p2p", %leader, nonce, "Follower received unexpected heartbeat request"); - } - } - } + if peer != leader { + warn!(target: "zone::p2p", %peer, "Ignoring block from non-leader"); + continue; } + events + .send(P2pEvent::BlockReceived { + leader_ed25519_public_key: peer, + block: bytes.into(), + }) + .await + .map_err(|_| eyre::eyre!("P2P event channel closed"))?; } } @@ -506,8 +406,8 @@ mod tests { use commonware_codec::Encode as _; use commonware_cryptography::{Signer as _, ed25519::PrivateKey}; - use super::{P2pConfig, P2pEvent, spawn_p2p, validate_ip_check_configuration}; - use crate::{P2pNetworkId, ZoneManifest, identity::Ed25519Identity}; + use super::{P2pCommand, P2pConfig, P2pEvent, spawn_p2p, validate_ip_check_configuration}; + use crate::{P2pNetworkId, ZoneManifest, identity::Ed25519Identity, network::MAX_MESSAGE_SIZE}; fn available_address() -> SocketAddr { let listener = TcpListener::bind("127.0.0.1:0").unwrap(); @@ -544,7 +444,7 @@ mod tests { } #[tokio::test(flavor = "multi_thread", worker_threads = 2)] - async fn three_nodes_exchange_heartbeats() { + async fn leader_broadcasts_blocks_to_followers() { let addresses = [ available_address(), available_address(), @@ -587,20 +487,40 @@ mod tests { }) .collect::>(); + let block = vec![0xf8, 0x01, 0x80]; + let commands = handles[0].parts.as_ref().unwrap().commands.clone(); + let oversized_block = vec![0; MAX_MESSAGE_SIZE as usize + 1]; + commands + .send(P2pCommand::BroadcastBlock(oversized_block)) + .await + .expect("P2P command channel should remain open"); + let broadcast_block = block.clone(); + let broadcaster = tokio::spawn(async move { + loop { + commands + .send(P2pCommand::BroadcastBlock(broadcast_block.clone())) + .await + .unwrap(); + tokio::time::sleep(Duration::from_millis(100)).await; + } + }); + for handle in handles.iter_mut().skip(1) { tokio::time::timeout(Duration::from_secs(15), async { loop { - if matches!( - handle.events_mut().recv().await, - Some(P2pEvent::HeartbeatAcknowledged { .. }) - ) { - break; + if let Some(P2pEvent::BlockReceived { + block: received, .. + }) = handle.events_mut().recv().await + { + assert_eq!(received, block); + return; } } }) .await - .expect("follower did not receive heartbeat acknowledgement"); + .expect("follower did not receive block"); } + broadcaster.abort(); for handle in handles { tokio::time::timeout(Duration::from_secs(10), handle.shutdown())