diff --git a/lez/sequencer/core/src/config.rs b/lez/sequencer/core/src/config.rs index a3ae4d0d0..ba8ec7ad2 100644 --- a/lez/sequencer/core/src/config.rs +++ b/lez/sequencer/core/src/config.rs @@ -127,41 +127,3 @@ const fn default_metrics_address() -> Option { pub const fn default_priority_fee() -> u64 { logos_blockchain_zone_sdk::sequencer::FundingConfig::DEFAULT_PRIORITY_FEE } - -#[cfg(test)] -mod tests { - use super::*; - - #[test] - fn gossip_config_defaults() { - let config: GossipConfig = serde_json::from_str("{}").unwrap(); - assert_eq!(config.listen_addr, "/ip4/0.0.0.0/udp/0/quic-v1"); - assert!(config.bootstrap_peers.is_empty()); - } - - #[test] - fn debug_config_parses_with_gossip() { - let config = SequencerConfig::from_path(Path::new(concat!( - env!("CARGO_MANIFEST_DIR"), - "/../service/configs/debug/sequencer_config.json" - ))) - .unwrap(); - assert!(config.gossip.is_some()); - } - - #[test] - fn sequencer_config_gossip_defaults_to_none() { - let json = r#"{ - "home": ".", "max_num_tx_in_block": 1, "mempool_max_size": 1, - "block_create_timeout": "1s", "retry_pending_blocks_timeout": "1s", - "signing_key": [0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0], - "bedrock_config": { - "channel_id": "0101010101010101010101010101010101010101010101010101010101010101", - "node_url": "http://localhost:18080", - "funding_key": "2e03b2eff5a45478e7e79668d2a146cf2c5c7925bce927f2b1c67f2ab4fc0d26" - } - }"#; - let config: SequencerConfig = serde_json::from_str(json).unwrap(); - assert!(config.gossip.is_none()); - } -} diff --git a/lez/sequencer/core/src/gossip/accreditation/mod.rs b/lez/sequencer/core/src/gossip/accreditation/mod.rs index fbdd4f58d..82bbf484a 100644 --- a/lez/sequencer/core/src/gossip/accreditation/mod.rs +++ b/lez/sequencer/core/src/gossip/accreditation/mod.rs @@ -18,8 +18,9 @@ use logos_blockchain_zone_sdk::{ use crate::config::BedrockConfig; pub trait AccreditedKeysProvider: Send + 'static { - /// The channel's current accredited Ed25519 keys. An empty set is valid - /// (channel does not exist yet); errors keep the caller's last set. + /// The channel's current accredited Ed25519 keys. + /// + /// An empty set is valid, meaning that the channel does not exist yet. fn accredited_keys(&self) -> impl Future>> + Send; } @@ -51,6 +52,7 @@ impl AccreditedKeysProvider for NodeKeysProvider { .channel_state(self.channel_id) .await .context("Failed to read channel state for accredited keys")?; + Ok(state .map(|state| { state diff --git a/lez/sequencer/core/src/gossip/network.rs b/lez/sequencer/core/src/gossip/network.rs index 1a3ca4a1b..466d2df9b 100644 --- a/lez/sequencer/core/src/gossip/network.rs +++ b/lez/sequencer/core/src/gossip/network.rs @@ -1,9 +1,3 @@ -//! The libp2p swarm and its drive task. -//! -//! The drive task owns the [`Swarm`]; [`Libp2pNetwork`] is the handle. A -//! drive-task failure after startup is log-and-continue: the node keeps -//! running L1-only (see the gossip spec). - use std::{ collections::{HashMap, HashSet}, sync::{ @@ -25,7 +19,6 @@ use libp2p::{ multiaddr::Protocol, swarm::{NetworkBehaviour, Swarm, SwarmEvent}, }; -use log::{debug, error, info, warn}; use mempool::MemPoolHandle; use tokio::sync::{mpsc, watch}; use tokio_util::sync::CancellationToken; @@ -61,8 +54,7 @@ pub trait PeerNetworkTrait { fn connected_peers(&self) -> Vec<[u8; 32]>; /// Cancelled when the drive task terminates. Unlike the publisher's - /// token, observers must NOT halt the node on it — gossip is an - /// optimization. + /// token, observers must NOT halt the node on it. fn driver_cancellation(&self) -> CancellationToken; } @@ -88,7 +80,7 @@ pub struct TxPublisher(mpsc::Sender); impl TxPublisher { pub fn publish(&self, tx: LeeTransaction) { if let Err(err) = self.0.try_send(tx) { - debug!("Dropping local tx publish: outbound gossip channel full or closed: {err}"); + log::debug!("Dropping local tx publish: outbound gossip channel full or closed: {err}"); } } } @@ -123,6 +115,7 @@ impl Libp2pNetwork { }) .collect::>()?; + // TODO: use a helper for this let topic = gossipsub::IdentTopic::new(format!("/lez/{}/v1/txs", hex::encode(channel_id))); let message_id_fn = |msg: &gossipsub::Message| { @@ -151,6 +144,7 @@ impl Libp2pNetwork { .max_transmit_size(max_transmit_size) .build() .map_err(|err| anyhow!("Failed to build gossipsub config: {err}"))?; + let gossipsub_behaviour = gossipsub::Behaviour::new( gossipsub::MessageAuthenticity::Signed(keypair.clone()), gossipsub_config, @@ -182,6 +176,7 @@ impl Libp2pNetwork { .with_swarm_config(|cfg| cfg.with_idle_connection_timeout(Duration::from_secs(60))) .build(); + // swarm .behaviour_mut() .gossipsub @@ -194,7 +189,7 @@ impl Libp2pNetwork { // Fail fast on bind errors: wait for the first listen address. let listen_addrs = wait_for_listen_addr(&mut swarm).await?; - info!("Gossip listening on {listen_addrs:?} as {local_peer_id}"); + log::info!("Gossip listening on {listen_addrs:?} as {local_peer_id}"); // Seed Kademlia with bootstrap peers that carry an embedded peer id; // dial the rest directly, since Kademlia can't route to an address @@ -212,11 +207,11 @@ impl Libp2pNetwork { continue; } if let Err(err) = swarm.dial(addr.clone()) { - warn!("Failed to dial gossip bootstrap peer {addr}: {err}"); + log::warn!("Failed to dial gossip bootstrap peer {addr}: {err}"); } } if let Err(err) = swarm.behaviour_mut().kademlia.bootstrap() { - debug!("Kademlia bootstrap skipped (no known peers yet): {err}"); + log::debug!("Kademlia bootstrap skipped (no known peers yet): {err}"); } let (connected_tx, connected_rx) = watch::channel(Vec::new()); @@ -358,7 +353,7 @@ impl DriveTask { GossipBehaviourEvent::Mdns(mdns::Event::Discovered(peers)) => { for (peer_id, addr) in peers { if let Err(err) = self.swarm.dial(addr) { - debug!("Failed to dial mdns-discovered peer {peer_id}: {err}"); + log::debug!("Failed to dial mdns-discovered peer {peer_id}: {err}"); } } } @@ -395,7 +390,7 @@ impl DriveTask { let acceptance = match evaluate_transaction(data, self.max_block_size) { TxEvaluation::Reject(reason) => { - debug!("Rejecting gossiped tx from {source}: {reason}"); + log::debug!("Rejecting gossiped tx from {source}: {reason}"); gossipsub::MessageAcceptance::Reject } TxEvaluation::Accept(tx) => { @@ -403,22 +398,22 @@ impl DriveTask { if self.seen.contains(&hash) { gossipsub::MessageAcceptance::Ignore } else { - // Mark seen only on a successful push. A full local mempool - // still forwards to peers (who may have room) without - // admitting or marking seen, so a later re-gossip can be - // admitted once we drain. match self.mempool.try_push((TransactionOrigin::Gossip, tx)) { Ok(()) => { + // mark seen only on succesful pushes, so that if mempool is full we can later + // receive the same tx from gossip and try pushing it again self.seen.insert(hash); } Err(_) => { - debug!("Gossip mempool full; forwarding tx {hash:?} without admitting"); + log::debug!("Mempool full; forwarding tx {hash:?} without admitting"); } } + gossipsub::MessageAcceptance::Accept } } }; + _ = self .swarm .behaviour_mut() @@ -437,7 +432,7 @@ impl DriveTask { .gossipsub .publish(self.topic.clone(), bytes) { - debug!("Skipping local tx publish {hash:?}: {err}"); + log::debug!("Skipping local tx publish {hash:?}: {err}"); } } } @@ -457,8 +452,7 @@ fn is_unspecified(addr: &Multiaddr) -> bool { }) } -/// Derives the libp2p `PeerId` an Ed25519 public key produces. `None` only -/// for byte strings that are not a valid curve point. +/// Derives the libp2p `PeerId` an Ed25519 public key produces. #[cfg_attr( not(test), expect( @@ -466,9 +460,10 @@ fn is_unspecified(addr: &Multiaddr) -> bool { reason = "unused by the mesh until a later gossip task; exercised by the identity test below" ) )] -pub(crate) fn peer_id_from_ed25519(pubkey: &[u8; 32]) -> Option { +pub(crate) fn peer_id_from_ed25519( + pubkey: &[u8; 32], +) -> Result { libp2p::identity::ed25519::PublicKey::try_from_bytes(pubkey) - .ok() .map(|key| libp2p::identity::PublicKey::from(key).to_peer_id()) } @@ -488,15 +483,15 @@ async fn wait_for_listen_addr(swarm: &mut Swarm) -> Result match event { SwarmEvent::NewListenAddr { address, .. } => return Ok(vec![address]), SwarmEvent::ListenerError { error, .. } => { - return Err(anyhow!("Gossip listener error during startup: {error}")); + anyhow::bail!("Gossip listener error during startup: {error}"); } SwarmEvent::ListenerClosed { reason, .. } => { - return Err(anyhow!("Gossip listener closed during startup: {reason:?}")); + anyhow::bail!("Gossip listener closed during startup: {reason:?}"); } _ => {} }, () = &mut deadline => { - return Err(anyhow!("Timed out waiting for gossip listen address")); + anyhow::bail!("Timed out waiting for gossip listen address"); } } } @@ -528,7 +523,7 @@ fn spawn_death_reminder(cancellation: CancellationToken, graceful_shutdown: Arc< return; } loop { - error!( + log::error!( "Sequencer gossip network is down; continuing L1-only. \ Restart the node to restore p2p." ); diff --git a/lez/sequencer/core/src/gossip/seen_cache.rs b/lez/sequencer/core/src/gossip/seen_cache.rs index 69fc318c6..ad467391e 100644 --- a/lez/sequencer/core/src/gossip/seen_cache.rs +++ b/lez/sequencer/core/src/gossip/seen_cache.rs @@ -1,14 +1,13 @@ -//! Bounded, FIFO-eviction membership cache over transaction hashes. -//! -//! The mempool is a plain channel with no dedup, so the gossip layer tracks -//! recently seen transactions here to avoid re-admitting duplicates that -//! arrive from multiple peers or echo back after a local publish. Counters -//! are exposed for a future metrics surface. - +use common::HashType; use std::collections::{HashSet, VecDeque}; -use common::HashType; - +/// Bounded, FIFO-eviction membership cache over transaction hashes. +/// +/// The mempool is a plain channel with no dedup, so the gossip layer tracks +/// recently seen transactions here to avoid re-admitting duplicates that +/// arrive from multiple peers or echo back after a local publish. +/// +/// TODO: Counters are exposed for a future metrics surface. pub struct SeenCache { capacity: usize, order: VecDeque, diff --git a/lez/sequencer/core/src/gossip/tests.rs b/lez/sequencer/core/src/gossip/tests.rs index e2b11d81f..bf621bbdb 100644 --- a/lez/sequencer/core/src/gossip/tests.rs +++ b/lez/sequencer/core/src/gossip/tests.rs @@ -1,5 +1,3 @@ -//! Multi-node integration tests over real QUIC on 127.0.0.1. - use std::time::{Duration, Instant}; use common::transaction::LeeTransaction;