diff --git a/crates/config/src/keys.rs b/crates/config/src/keys.rs index ebb5bc7e2..c58e27718 100644 --- a/crates/config/src/keys.rs +++ b/crates/config/src/keys.rs @@ -23,7 +23,7 @@ use sha2::Sha256; use std::sync::Arc; use tn_types::{ construct_proof_of_possession_message, Address, BlsKeypair, BlsPublicKey, BlsSignature, - BlsSigner, DefaultHashFunction, NetworkKeypair, NetworkPublicKey, Signer, + BlsSigner, DefaultHashFunction, NetworkKeypair, NetworkPublicKey, Signer, WorkerId, }; use zeroize::Zeroizing; @@ -208,15 +208,27 @@ fn warn_if_key_permissions_are_loose( ) { } -#[derive(Debug)] +/// Private key material and derivation inputs shared by a key manager. struct KeyConfigInner { - // DO NOT expose the private key to other code. Tests that need this will provide a primary - // key. Use the BlsSigner trait for signing for the primary. + /// DO NOT expose the private key to other code. Tests provide their own primary key. + /// Use the BlsSigner trait for signing for the primary. primary_keypair: BlsKeypair, - // Derived from the primary_keypair. + /// Derived from the primary keypair. primary_network_keypair: NetworkKeypair, - // Derived from the primary_keypair. - worker_network_keypair: NetworkKeypair, + /// Seed string for worker network keypairs. Per-worker keypairs are derived on demand from + /// the primary keypair and this seed; see [`KeyConfig::worker_network_keypair`]. + worker_network_seed: String, +} + +impl std::fmt::Debug for KeyConfigInner { + fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + formatter + .debug_struct("KeyConfigInner") + .field("primary_keypair", &self.primary_keypair) + .field("primary_network_keypair", &self.primary_network_keypair) + .field("worker_network_seed", &"[REDACTED]") + .finish() + } } /// Basic implementation of a key manager. This version will read a BLS key @@ -225,7 +237,7 @@ struct KeyConfigInner { /// It should NOT expose the BLS private key, even though it is currently read /// from a file this will not always be the case and all code needing signatures /// MUST go through KeyConfig. -/// NOTE: The two network keys (primary and worker) are derived from the BLS key +/// NOTE: The network keys (primary and per-worker) are derived from the BLS key /// and are exposed to other code. This is required to work with libp2p which /// wants the actual private key. This method of deriving the key is an attempt /// to provide some protection to the key- even though it will exist in memory it @@ -363,12 +375,11 @@ impl KeyConfig { }; let primary_network_keypair = Self::generate_network_keypair(&primary_keypair, &primary_seed); - let worker_network_keypair = Self::generate_network_keypair(&primary_keypair, &worker_seed); Ok(Self { inner: Arc::new(KeyConfigInner { primary_keypair, primary_network_keypair, - worker_network_keypair, + worker_network_seed: worker_seed, }), }) } @@ -429,7 +440,6 @@ impl KeyConfig { let worker_seed = "worker network keypair"; let primary_network_keypair = Self::generate_network_keypair(&primary_keypair, primary_seed); - let worker_network_keypair = Self::generate_network_keypair(&primary_keypair, worker_seed); // Make sure we have the validator dir, owner-only. // Don't error out if path exists. create_keys_dir(&tn_datadir.node_keys_path())?; @@ -458,7 +468,7 @@ impl KeyConfig { inner: Arc::new(KeyConfigInner { primary_keypair, primary_network_keypair, - worker_network_keypair, + worker_network_seed: worker_seed.to_string(), }), }) } @@ -467,13 +477,11 @@ impl KeyConfig { pub fn new_with_testing_key(primary_keypair: BlsKeypair) -> Self { let primary_network_keypair = Self::generate_network_keypair(&primary_keypair, "primary network keypair"); - let worker_network_keypair = - Self::generate_network_keypair(&primary_keypair, "worker network keypair"); Self { inner: Arc::new(KeyConfigInner { primary_keypair, primary_network_keypair, - worker_network_keypair, + worker_network_seed: "worker network keypair".to_string(), }), } } @@ -494,15 +502,30 @@ impl KeyConfig { self.primary_network_keypair().public().clone().into() } - /// Provide the keypair (with private key) for the worker network. + /// Provide the keypair (with private key) for the network of `worker_id`. /// Allows building the libp2p worker network. - pub fn worker_network_keypair(&self) -> &NetworkKeypair { - &self.inner.worker_network_keypair + /// + /// Worker 0 derives from the stored seed exactly as before per-worker swarms existed. This + /// keeps worker 0's PeerId stable for deployed nodes: that network identity is advertised + /// on-chain and cached in peers' kad stores, so it must not change. Worker ids above 0 + /// append the id to the seed to get a distinct keypair per swarm. + pub fn worker_network_keypair(&self, worker_id: WorkerId) -> NetworkKeypair { + if worker_id == 0 { + Self::generate_network_keypair( + &self.inner.primary_keypair, + &self.inner.worker_network_seed, + ) + } else { + Self::generate_network_keypair( + &self.inner.primary_keypair, + &format!("{} {worker_id}", self.inner.worker_network_seed), + ) + } } - /// The [NetworkPublicKey] for the worker network. - pub fn worker_network_public_key(&self) -> NetworkPublicKey { - self.worker_network_keypair().public().into() + /// The [NetworkPublicKey] for the network of `worker_id`. + pub fn worker_network_public_key(&self, worker_id: WorkerId) -> NetworkPublicKey { + self.worker_network_keypair(worker_id).public().into() } /// Creates a proof that the authority account address is owned by the @@ -989,7 +1012,22 @@ mod tests { let config = KeyConfig::new_with_testing_key(keypair); let rendered = format!("{config:?}"); - assert!(rendered.contains("[REDACTED]"), "BLS private half must be redacted: {rendered}"); + assert!( + rendered.contains("private: \"[REDACTED]\""), + "BLS private half must be redacted: {rendered}" + ); + assert!( + rendered.contains(&format!("public: {:?}", config.primary_public_key())), + "BLS public key should still be shown: {rendered}" + ); + assert!( + !rendered.contains(&config.inner.worker_network_seed), + "worker network seed must not be shown" + ); + assert!( + rendered.contains("worker_network_seed: \"[REDACTED]\""), + "worker network seed field must remain present and redacted: {rendered}" + ); let assert_secret_absent = |bytes: &[u8], what: &str| { assert!(!rendered.contains(&hex::encode(bytes)), "{what} leaked as hex"); @@ -1017,14 +1055,18 @@ mod tests { ed25519_secret(config.primary_network_keypair()).as_ref(), "primary network secret", ); + // Per-worker network keypairs are derived on demand from the primary key and the stored + // seed (#555), so `KeyConfigInner` stores no worker keypair. Worker 0 is the legacy + // derivation; check it in case a future field caches derived keypairs. assert_secret_absent( - ed25519_secret(config.worker_network_keypair()).as_ref(), + ed25519_secret(&config.worker_network_keypair(0)).as_ref(), "worker network secret", ); - // Positive anchors: the network fields must actually render their public halves, + // Positive anchor: the primary network field must actually render its public half, // otherwise the negative checks above pass vacuously once `KeyConfigInner`'s Debug - // stops printing the network keypairs at all. + // stops printing the network keypair at all. The worker seed's redacted field is + // anchored above; no worker keypair is stored. let ed25519_public_rendered = |net: &NetworkKeypair| { let ed25519: libp2p::identity::ed25519::Keypair = net.clone().try_into().expect("network keypairs are ed25519"); @@ -1034,9 +1076,21 @@ mod tests { rendered.contains(&ed25519_public_rendered(config.primary_network_keypair())), "primary network public key should still be shown: {rendered}" ); - assert!( - rendered.contains(&ed25519_public_rendered(config.worker_network_keypair())), - "worker network public key should still be shown: {rendered}" - ); + } + + /// Worker 0 must keep the legacy bare-seed derivation (its PeerId is advertised on-chain), + /// worker 1 must get a distinct keypair, and derivation must be deterministic per id. + #[test] + fn test_worker_network_keypair_per_id_derivation() { + let kc = KeyConfig::new_with_testing_key(random_keypair()); + let legacy: NetworkPublicKey = KeyConfig::generate_network_keypair( + &kc.inner.primary_keypair, + "worker network keypair", + ) + .public() + .into(); + assert_eq!(kc.worker_network_public_key(0), legacy); + assert_ne!(kc.worker_network_public_key(1), kc.worker_network_public_key(0)); + assert_eq!(kc.worker_network_public_key(1), kc.worker_network_public_key(1)); } } diff --git a/crates/network-libp2p/src/consensus.rs b/crates/network-libp2p/src/consensus.rs index 338ae455e..9cc6710d6 100644 --- a/crates/network-libp2p/src/consensus.rs +++ b/crates/network-libp2p/src/consensus.rs @@ -365,7 +365,7 @@ where external_addr: Multiaddr, rpc: Option, ) -> NetworkResult { - let network_key = key_config.worker_network_keypair().clone(); + let network_key = key_config.worker_network_keypair(worker_id); Self::new( network_config, event_stream, diff --git a/crates/network-libp2p/src/kad.rs b/crates/network-libp2p/src/kad.rs index 64e003ee9..70920aef8 100644 --- a/crates/network-libp2p/src/kad.rs +++ b/crates/network-libp2p/src/kad.rs @@ -19,7 +19,7 @@ use tn_config::KeyConfig; use tn_storage::tables::{ KadProviderRecords, KadRecords, KadWorkerProviderRecords, KadWorkerRecords, }; -use tn_types::{decode, encode, try_decode, BlockHash, Database, DefaultHashFunction}; +use tn_types::{encode, try_decode, BlockHash, Database, DefaultHashFunction}; use tracing::{error, warn}; /// A record stored in the DHT. @@ -162,9 +162,9 @@ fn instant_to_system(expires: &Option) -> Option { /// Decode a stored provider-record blob, tolerating bytes that no longer decode. /// -/// The provider tables persist `Vec` values. A schema/version skew +/// Provider envelopes contain separately encoded `Vec` values. A schema skew /// across a restart, or on-disk corruption, can leave a row whose bytes the current -/// software can no longer decode. The panicking [`decode`] would turn that single bad +/// software can no longer decode. The panicking [`tn_types::decode`] would turn that single bad /// row into a crash of the whole `ConsensusNetwork` task on the first provider read /// after restart, and because the row is never purged the crash recurs on every /// restart. This returns `None` (logging a warning) instead, so the caller can skip @@ -179,6 +179,35 @@ fn decode_providers(key: &BlockHash, raw: &[u8]) -> Option, +} + +impl KadProviderRow { + /// Encode a provider set while preserving independently decodable row ownership. + fn encode(key: RecordKey, records: &[KadProviderRecord]) -> Vec { + encode(&Self { key, records: encode(&records) }) + } + + /// Decode a provider set only when every member belongs to this row's discovery key. + fn decode_records(&self, hash: &BlockHash) -> Option> { + decode_providers(hash, &self.records) + .filter(|records| records.iter().all(|record| record.key == self.key)) + } +} + /// Minimum spacing between saturated-table provider eviction scans. /// /// Provider records carry a 48h TTL, so once `num_providers` reaches @@ -192,7 +221,9 @@ const PROVIDER_EVICT_INTERVAL: Duration = Duration::from_secs(60); /// Wraps around the consensus DB. #[derive(Clone, Debug)] pub struct KadStore { + /// Shared database containing all swarm namespaces. db: DB, + /// Discovery key under which this node publishes its provider record. node_key: RecordKey, /// This node's libp2p peer id. /// @@ -206,9 +237,9 @@ pub struct KadStore { /// basically just here to prevent or mitigate attacks on the Kad store. /// Use the same settings as a Kad Memery store. config: MemoryStoreConfig, - /// Tracks to number of records in DB. + /// Number of persisted discovery records owned by this swarm, including expired rows. num_records: usize, - /// Tracks to number of provider records in DB. + /// Number of provider rows with a readable ownership envelope belonging to this swarm. num_providers: usize, /// Last time the saturated-table provider eviction scan ran, so a full table /// cannot be turned into a full-table decode scan per inbound `AddProvider`. @@ -234,25 +265,18 @@ impl KadStore { let node_key = RecordKey::new(&encode(&key_config.primary_public_key())); // Defaults for sanity. let config = MemoryStoreConfig::default(); - let (num_records, num_providers) = match kad_type { - NetworkType::Primary => { - (db.iter::().count(), db.iter::().count()) - } - NetworkType::Worker(_) => ( - db.iter::().count(), - db.iter::().count(), - ), - }; - let store = Self { + let mut store = Self { db, node_key, local_peer_id, config, - num_records, - num_providers, + num_records: 0, + num_providers: 0, last_provider_evict: None, kad_type, }; + store.num_records = store.owned_records().count(); + store.num_providers = store.owned_provider_rows().count(); metrics::describe_counter!( "tn_network.kad_provider_write_failures_total", metrics::Unit::Count, @@ -280,44 +304,71 @@ impl KadStore { .increment(1); } + /// Namespace shared discovery tables by the same role and worker id as `RecordDomain`. fn key_to_hash(&self, key: &RecordKey) -> BlockHash { + let (role, worker_id): (u8, tn_types::WorkerId) = match self.kad_type { + NetworkType::Primary => (0, 0), + NetworkType::Worker(id) => (1, id), + }; let mut h = DefaultHashFunction::new(); + h.update(&[role]); + h.update(&worker_id.to_le_bytes()); h.update(encode(key).as_ref()); BlockHash::from_slice(h.finalize().as_bytes()) } + /// Whether a persisted discovery key rederives this store's row hash. + fn owns(&self, key: &RecordKey, hash: &BlockHash) -> bool { + self.key_to_hash(key) == *hash + } + + /// Decode a discovery row only when its persisted key matches this store's namespace. + fn decode_record(&self, hash: &BlockHash, raw: &[u8]) -> Option { + try_decode::(raw).ok().filter(|record| self.owns(&record.key, hash)) + } + + /// Enumerate this swarm's rows, including expired records needed for accounting. + fn owned_records(&self) -> impl Iterator + '_ { + let rows = match self.kad_type { + NetworkType::Primary => self.db.iter::(), + NetworkType::Worker(_) => self.db.iter::(), + }; + rows.filter_map(move |(hash, raw)| { + self.decode_record(&hash, &raw).map(|record| (hash, record)) + }) + } + + /// Read row ownership before inspecting provider payloads from a shared table. + fn decode_provider_row(&self, hash: &BlockHash, raw: &[u8]) -> Option { + try_decode::(raw).ok().filter(|row| self.owns(&row.key, hash)) + } + + /// Enumerate only this swarm's provider rows, even when their payloads are malformed. + fn owned_provider_rows(&self) -> impl Iterator + '_ { + let rows = match self.kad_type { + NetworkType::Primary => self.db.iter::(), + NetworkType::Worker(_) => self.db.iter::(), + }; + rows.filter_map(move |(hash, raw)| { + self.decode_provider_row(&hash, &raw).map(|row| (hash, row)) + }) + } + /// Scan the records table and remove any rows whose expiry has passed. /// Returns the number of rows actually removed and updates `num_records`. fn evict_expired_records(&mut self) -> usize { let now = SystemTime::now(); - let expired_keys: Vec = match self.kad_type { - NetworkType::Primary => self - .db - .iter::() - .filter_map(|(k, v)| { - let r: KadRecord = decode(v.as_ref()); - r.is_expired(now).then_some(k) - }) - .collect(), - NetworkType::Worker(_) => self - .db - .iter::() - .filter_map(|(k, v)| { - let r: KadRecord = decode(v.as_ref()); - r.is_expired(now).then_some(k) - }) - .collect(), - }; - let mut evicted = 0; - for k in &expired_keys { - let ok = match self.kad_type { + let expired_keys: Vec = self + .owned_records() + .filter_map(|(hash, record)| record.is_expired(now).then_some(hash)) + .collect(); + let evicted = expired_keys + .iter() + .filter(|k| match self.kad_type { NetworkType::Primary => self.db.remove::(k).is_ok(), NetworkType::Worker(_) => self.db.remove::(k).is_ok(), - }; - if ok { - evicted += 1; - } - } + }) + .count(); self.num_records = self.num_records.saturating_sub(evicted); self.update_records_gauge(); evicted @@ -327,42 +378,22 @@ impl KadStore { /// Returns the number of keys removed and updates `num_providers`. fn evict_expired_providers(&mut self) -> usize { let now = SystemTime::now(); - let drop_keys: Vec = match self.kad_type { - NetworkType::Primary => self - .db - .iter::() - .filter_map(|(k, v)| { - // An undecodable row is unusable; treat it as droppable so eviction - // frees the slot and purges it instead of panicking (issue #999). - let should_drop = decode_providers(&k, v.as_ref()) - .map(|recs| !recs.is_empty() && recs.iter().all(|r| r.is_expired(now))) - .unwrap_or(true); - should_drop.then_some(k) - }) - .collect(), - NetworkType::Worker(_) => self - .db - .iter::() - .filter_map(|(k, v)| { - // An undecodable row is unusable; treat it as droppable so eviction - // frees the slot and purges it instead of panicking (issue #999). - let should_drop = decode_providers(&k, v.as_ref()) - .map(|recs| !recs.is_empty() && recs.iter().all(|r| r.is_expired(now))) - .unwrap_or(true); - should_drop.then_some(k) - }) - .collect(), - }; - let mut evicted = 0; - for k in &drop_keys { - let ok = match self.kad_type { + let drop_keys: Vec = self + .owned_provider_rows() + .filter_map(|(hash, row)| { + // Ownership is known even for empty or malformed provider payloads. + row.decode_records(&hash) + .is_none_or(|records| records.iter().all(|record| record.is_expired(now))) + .then_some(hash) + }) + .collect(); + let evicted = drop_keys + .iter() + .filter(|k| match self.kad_type { NetworkType::Primary => self.db.remove::(k).is_ok(), NetworkType::Worker(_) => self.db.remove::(k).is_ok(), - }; - if ok { - evicted += 1; - } - } + }) + .count(); self.num_providers = self.num_providers.saturating_sub(evicted); evicted } @@ -372,20 +403,13 @@ impl KadStore { /// `consensus.rs`. Without this, a schema/version skew or on-disk corruption leaves a /// row that panics the whole `ConsensusNetwork` task on the first provider read after /// restart, and because the row is never purged the panic recurs on every restart. - /// Returns the number of rows removed and keeps `num_providers` in step (issue #999). + /// Only rows with a readable ownership envelope are considered; unknown ownership is never + /// guessed. Returns the number of rows removed and keeps `num_providers` in step (issue #999). pub fn scrub_corrupt_providers(&mut self) -> usize { - let corrupt: Vec = match self.kad_type { - NetworkType::Primary => self - .db - .iter::() - .filter_map(|(k, v)| decode_providers(&k, v.as_ref()).is_none().then_some(k)) - .collect(), - NetworkType::Worker(_) => self - .db - .iter::() - .filter_map(|(k, v)| decode_providers(&k, v.as_ref()).is_none().then_some(k)) - .collect(), - }; + let corrupt: Vec = self + .owned_provider_rows() + .filter_map(|(hash, row)| row.decode_records(&hash).is_none().then_some(hash)) + .collect(); let evicted = corrupt .iter() .filter(|k| match self.kad_type { @@ -428,7 +452,8 @@ impl KadStore { /// Iterator of KAD records. pub struct RecordIter<'a> { - iter: Box)> + 'a>, + /// Decoded rows whose hashes match the originating store's namespace. + iter: Box + 'a>, } impl<'a> std::fmt::Debug for RecordIter<'a> { @@ -442,14 +467,9 @@ impl<'a> Iterator for RecordIter<'a> { fn next(&mut self) -> Option { let now = SystemTime::now(); - loop { - let (_, raw) = self.iter.next()?; - let r: KadRecord = decode(raw.as_ref()); - if r.is_expired(now) { - continue; - } - return Some(Cow::Owned(r.into())); - } + self.iter + .find(|(_, record)| !record.is_expired(now)) + .map(|(_, record)| Cow::Owned(record.into())) } } @@ -472,11 +492,10 @@ impl RecordStore for KadStore { }) .ok()?; let raw = record?; - let r: KadRecord = decode(raw.as_ref()); - if r.is_expired(SystemTime::now()) { - return None; - } - Some(Cow::Owned(r.into())) + try_decode::(&raw) + .ok() + .filter(|record| record.key == *k && !record.is_expired(SystemTime::now())) + .map(|record| Cow::Owned(record.into())) } fn put(&mut self, r: Record) -> libp2p::kad::store::Result<()> { @@ -486,28 +505,23 @@ impl RecordStore for KadStore { let key = self.key_to_hash(&r.key); let kr: KadRecord = r.into(); - // Are we adding a new record or replacing an existing? - let new_record = match self.kad_type { + let stored = match self.kad_type { NetworkType::Primary => self.db.get::(&key), NetworkType::Worker(_) => self.db.get::(&key), } .map_err(|error| { error!(target: "network-kad", ?error, kad_type = ?self.kad_type, "failed to read Kademlia record before insert"); Error::ValueTooLarge - })? - .is_none(); - // We have a new record so go ahead and inc num_records. - // Should be safe since a failure to insert indicates a fatal DB condition. - if new_record { + })?; + // Startup excludes unreadable records, so repairing one is an insertion for capacity + // accounting. Replacing a readable owned row keeps the existing count. + let new_record = stored.as_deref().and_then(|raw| self.decode_record(&key, raw)).is_none(); + if new_record && self.num_records >= self.config.max_records { + // Try to free a slot by evicting any records whose TTL has passed. + self.evict_expired_records(); if self.num_records >= self.config.max_records { - // Try to free a slot by evicting any records whose TTL has passed. - self.evict_expired_records(); - if self.num_records >= self.config.max_records { - return Err(Error::MaxRecords); - } + return Err(Error::MaxRecords); } - self.num_records += 1; - self.update_records_gauge(); } match self.kad_type { NetworkType::Primary => self.db.insert::(&key, &encode(&kr)), @@ -517,34 +531,40 @@ impl RecordStore for KadStore { error!(target: "network-kad", ?error, kad_type = ?self.kad_type, "failed to insert Kademlia record"); Error::ValueTooLarge })?; + if new_record { + self.num_records += 1; + self.update_records_gauge(); + } Ok(()) } fn remove(&mut self, k: &RecordKey) { let key = self.key_to_hash(k); + let row_counted = match self.kad_type { + NetworkType::Primary => self.db.get::(&key), + NetworkType::Worker(_) => self.db.get::(&key), + } + .ok() + .flatten() + .and_then(|raw| self.decode_record(&key, &raw)) + .is_some(); if match self.kad_type { NetworkType::Primary => self.db.remove::(&key), NetworkType::Worker(_) => self.db.remove::(&key), } .is_ok() + && row_counted { - // Record was removed so dec num_records. Saturate to match the eviction - // siblings (`evict_expired_records`): on MDBX `db.remove` returns `Ok` even - // when the key was absent, so a `remove` for an uncounted key (a double - // `remove`, or one for a row already dropped by eviction) would otherwise - // drive this `usize` below zero and wrap to `usize::MAX`, permanently - // wedging `put` behind the `num_records >= max_records` cap. + // Only readable owned rows contribute to startup accounting. An absent or + // malformed row cannot uncount another row even when MDBX removal returns Ok. + // Saturation also tolerates a preexisting stale count without wrapping capacity. self.num_records = self.num_records.saturating_sub(1); self.update_records_gauge(); } } fn records(&self) -> Self::RecordsIter<'_> { - let iter = match self.kad_type { - NetworkType::Primary => self.db.iter::(), - NetworkType::Worker(_) => self.db.iter::(), - }; - RecordIter { iter } + RecordIter { iter: Box::new(self.owned_records()) } } fn add_provider(&mut self, record: ProviderRecord) -> libp2p::kad::store::Result<()> { @@ -564,7 +584,8 @@ impl RecordStore for KadStore { Error::ValueTooLarge })?; - let key = self.key_to_hash(&record.key); + let record_key = record.key.clone(); + let key = self.key_to_hash(&record_key); let kr: KadProviderRecord = record.into(); let stored = match self.kad_type { NetworkType::Primary => self.db.get::(&key), @@ -575,11 +596,11 @@ impl RecordStore for KadStore { Error::ValueTooLarge })?; - // A present-but-undecodable row is treated as absent for the merge: the new - // provider set overwrites (purges) it instead of the read panicking. Unlike a - // genuinely new key it is already counted in `num_providers`, so only a missing - // row bumps the gauge (issue #999). - let row_exists = stored.is_some(); + // A readable ownership envelope is counted even when the provider payload is corrupt. + // An unreadable envelope was excluded by startup accounting, so its replacement must + // pass the capacity check and increment the count like a new key. + let stored_row = stored.as_deref().and_then(|raw| self.decode_provider_row(&key, raw)); + let row_exists = stored_row.is_some(); // The capacity check applies only to a brand-new key, mirroring `put`'s // `new_record` gate: an overwrite of an existing, already-counted row cannot @@ -606,18 +627,16 @@ impl RecordStore for KadStore { .ok_or(Error::MaxProvidedKeys)?; } - let merged = stored - .as_deref() - .and_then(|raw| decode_providers(&key, raw)) + let merged = stored_row + .and_then(|row| row.decode_records(&key)) .map(|existing| self.merge_provider(existing, kr.clone())) .transpose()?; let records: Vec = merged.unwrap_or_else(|| vec![kr]); + let encoded = KadProviderRow::encode(record_key, &records); match self.kad_type { - NetworkType::Primary => self.db.insert::(&key, &encode(&records)), - NetworkType::Worker(_) => { - self.db.insert::(&key, &encode(&records)) - } + NetworkType::Primary => self.db.insert::(&key, &encoded), + NetworkType::Worker(_) => self.db.insert::(&key, &encoded), } .map_err(|error| { error!(target: "network-kad", ?error, kad_type = ?self.kad_type, "failed to insert Kademlia provider records"); @@ -647,7 +666,8 @@ impl RecordStore for KadStore { }) .ok() .flatten() - .and_then(|recs| decode_providers(&hash, &recs)) + .and_then(|raw| self.decode_provider_row(&hash, &raw)) + .and_then(|row| row.decode_records(&hash)) .map(|records| { records.into_iter().filter(|r| !r.is_expired(now)).map(Into::into).collect() }) @@ -684,10 +704,13 @@ impl RecordStore for KadStore { .ok() .flatten(); - if let Some(raw) = stored { - // An undecodable row is purged wholesale (it is unusable and would otherwise - // panic every read); a decodable row keeps every provider except `p`. - let remaining: Vec = decode_providers(&hash, &raw) + if stored.is_some() { + let row = stored.as_deref().and_then(|raw| self.decode_provider_row(&hash, raw)); + let row_counted = row.is_some(); + // Direct lookup establishes the namespace even if the envelope itself is corrupt. + // Preserve the count when removing an unreadable envelope excluded at startup. + let remaining: Vec = row + .and_then(|row| row.decode_records(&hash)) .map(|records| records.into_iter().filter(|r| r.provider != *p).collect()) .unwrap_or_default(); if remaining.is_empty() { @@ -696,18 +719,17 @@ impl RecordStore for KadStore { NetworkType::Worker(_) => self.db.remove::(&hash), } .is_ok(); - if removed { + if removed && row_counted { // The key holds no providers now (all filtered out, or the row was // purged): drop the count once, saturating to avoid an underflow panic. self.num_providers = self.num_providers.saturating_sub(1); } } else { + let encoded = KadProviderRow::encode(key.clone(), &remaining); let _ = match self.kad_type { - NetworkType::Primary => { - self.db.insert::(&hash, &encode(&remaining)) - } + NetworkType::Primary => self.db.insert::(&hash, &encoded), NetworkType::Worker(_) => { - self.db.insert::(&hash, &encode(&remaining)) + self.db.insert::(&hash, &encoded) } } .inspect_err(|error| { @@ -870,6 +892,141 @@ mod test { ); } + /// Sibling workers retain separate records, provider sets, startup counts and removals. + #[test] + fn test_kad_worker_store_isolation() -> eyre::Result<()> { + let tmp_dir = TempDir::new()?; + let db = open_db(tmp_dir.path()); + let key_config = test_key_config(); + // Sharing even the local peer id must not allow `provided()` to cross namespaces. + let local_peer_id = PeerId::random(); + let mut worker_0 = + KadStore::new(db.clone(), local_peer_id, &key_config, NetworkType::Worker(0)); + let mut worker_1 = + KadStore::new(db.clone(), local_peer_id, &key_config, NetworkType::Worker(1)); + let record_0 = Record { + key: RecordKey::new(&b"shared-worker-record"), + value: vec![0], + publisher: None, + expires: None, + }; + let record_1 = Record { value: vec![1], ..record_0.clone() }; + worker_0.put(record_0.clone())?; + assert!(worker_1.get(&record_0.key).is_none()); + worker_1.put(record_1.clone())?; + assert_eq!(worker_0.get(&record_0.key).map(|record| record.value.clone()), Some(vec![0])); + assert_eq!(worker_1.get(&record_1.key).map(|record| record.value.clone()), Some(vec![1])); + assert_eq!(worker_0.records().count(), 1); + assert_eq!(worker_1.records().count(), 1); + + let provider_0 = ProviderRecord { + key: worker_0.node_key.clone(), + provider: local_peer_id, + expires: None, + addresses: vec!["/ip4/127.0.0.1/tcp/1000".parse()?], + }; + let provider_1 = ProviderRecord { + addresses: vec!["/ip4/127.0.0.1/tcp/1001".parse()?], + ..provider_0.clone() + }; + worker_0.add_provider(provider_0.clone())?; + assert!(worker_1.providers(&provider_0.key).is_empty()); + worker_1.add_provider(provider_1.clone())?; + assert_eq!(worker_0.providers(&provider_0.key), vec![provider_0.clone()]); + assert_eq!(worker_1.providers(&provider_1.key), vec![provider_1.clone()]); + assert_eq!( + worker_0.provided().map(Cow::into_owned).collect::>(), + vec![provider_0.clone()] + ); + assert_eq!( + worker_1.provided().map(Cow::into_owned).collect::>(), + vec![provider_1.clone()] + ); + + let restarted_0 = + KadStore::new(db.clone(), local_peer_id, &key_config, NetworkType::Worker(0)); + let restarted_1 = KadStore::new(db, local_peer_id, &key_config, NetworkType::Worker(1)); + assert_eq!((restarted_0.num_records, restarted_0.num_providers), (1, 1)); + assert_eq!((restarted_1.num_records, restarted_1.num_providers), (1, 1)); + + // Persist both namespaces so the deletion checks also exercise the disk fallback. + worker_0.db.sync_persist(); + worker_0.remove(&record_0.key); + worker_0.remove_provider(&provider_0.key, &local_peer_id); + // Layered storage requires a persistence barrier before reading a deleted key. + worker_0.db.sync_persist(); + assert!(worker_0.get(&record_0.key).is_none()); + assert!(worker_0.providers(&provider_0.key).is_empty()); + assert_eq!(worker_1.get(&record_1.key).map(|record| record.value.clone()), Some(vec![1])); + assert_eq!(worker_1.providers(&provider_1.key), vec![provider_1]); + Ok(()) + } + + /// Repairing or removing a row excluded at startup must preserve record capacity accounting. + #[test] + fn test_kad_worker_corrupt_record_accounting() -> eyre::Result<()> { + let tmp_dir = TempDir::new()?; + let db = open_db(tmp_dir.path()); + let key_config = test_key_config(); + let mut seed = + KadStore::new(db.clone(), PeerId::random(), &key_config, NetworkType::Worker(0)); + let good_record = test_record(false); + seed.put(good_record.clone())?; + let corrupt_record = test_record(false); + let corrupt_hash = seed.key_to_hash(&corrupt_record.key); + db.insert::(&corrupt_hash, &vec![0xff])?; + + let mut worker_0 = + KadStore::new(db.clone(), PeerId::random(), &key_config, NetworkType::Worker(0)); + worker_0.config.max_records = 1; + assert_eq!(worker_0.num_records, 1); + assert!(matches!(worker_0.put(corrupt_record.clone()), Err(Error::MaxRecords))); + assert_eq!(worker_0.records().count(), 1); + assert_eq!(worker_0.num_records, 1); + worker_0.remove(&corrupt_record.key); + assert_eq!( + worker_0.num_records, 1, + "removing unreadable data must not uncount a valid row" + ); + assert!(worker_0.get(&good_record.key).is_some()); + worker_0.remove(&good_record.key); + assert_eq!(worker_0.num_records, 0); + db.insert::(&corrupt_hash, &vec![0xff])?; + worker_0.put(corrupt_record.clone())?; + assert_eq!(worker_0.num_records, 1, "repairing unreadable data consumes capacity"); + assert!(worker_0.get(&corrupt_record.key).is_some()); + assert_eq!(db.iter::().count(), 1); + Ok(()) + } + + /// A sibling's eviction must leave the owner's expired rows and counters in step. + #[test] + fn test_kad_worker_sibling_eviction_preserves_capacity() -> eyre::Result<()> { + let tmp_dir = TempDir::new()?; + let db = open_db(tmp_dir.path()); + let key_config = test_key_config(); + let mut worker_0 = + KadStore::new(db.clone(), PeerId::random(), &key_config, NetworkType::Worker(0)); + let mut worker_1 = KadStore::new(db, PeerId::random(), &key_config, NetworkType::Worker(1)); + worker_0.config.max_records = 1; + worker_0.config.max_provided_keys = 1; + worker_0.put(test_record(true))?; + worker_0.add_provider(expired_provider_under(&fresh_record_key()))?; + + assert_eq!(worker_1.evict_expired_records(), 0); + assert_eq!(worker_1.evict_expired_providers(), 0); + assert_eq!(worker_0.db.iter::().count(), 1); + assert_eq!(worker_0.db.iter::().count(), 1); + let fresh = test_record(false); + worker_0.put(fresh.clone())?; + worker_0.add_provider(live_provider_under(&fresh.key))?; + assert_eq!((worker_0.num_records, worker_0.num_providers), (1, 1)); + assert!(worker_0.get(&fresh.key).is_some()); + assert_eq!(worker_0.providers(&fresh.key).len(), 1); + assert_eq!((worker_1.num_records, worker_1.num_providers), (0, 0)); + Ok(()) + } + #[test] fn test_kad_store() { let tmp_dir = TempDir::new().expect("temp dir"); @@ -1448,19 +1605,23 @@ mod test { RecordKey::new(&encode(&test_key_config().primary_public_key())) } + /// Corrupt a primary provider payload while keeping its ownership envelope readable. fn inject_corrupt_primary_provider(store: &KadStore, key: &RecordKey) { let hash = store.key_to_hash(key); + let row = KadProviderRow { key: key.clone(), records: CORRUPT_PROVIDER_BYTES.to_vec() }; store .db - .insert::(&hash, &CORRUPT_PROVIDER_BYTES.to_vec()) + .insert::(&hash, &encode(&row)) .expect("inject corrupt provider row"); } + /// Corrupt a worker provider payload while keeping its ownership envelope readable. fn inject_corrupt_worker_provider(store: &KadStore, key: &RecordKey) { let hash = store.key_to_hash(key); + let row = KadProviderRow { key: key.clone(), records: CORRUPT_PROVIDER_BYTES.to_vec() }; store .db - .insert::(&hash, &CORRUPT_PROVIDER_BYTES.to_vec()) + .insert::(&hash, &encode(&row)) .expect("inject corrupt worker provider row"); } @@ -1588,6 +1749,73 @@ mod test { assert_eq!(store.num_providers, 0, "worker count reflects the purge"); } + /// Provider scrubbing and expiry respect ownership even for malformed or empty payloads. + #[test] + fn test_kad_worker_provider_corruption_isolation() -> eyre::Result<()> { + let tmp_dir = TempDir::new()?; + let db = open_db(tmp_dir.path()); + let key_config = test_key_config(); + let mut seed_0 = + KadStore::new(db.clone(), PeerId::random(), &key_config, NetworkType::Worker(0)); + let mut seed_1 = + KadStore::new(db.clone(), PeerId::random(), &key_config, NetworkType::Worker(1)); + let good_key = fresh_record_key(); + seed_0.add_provider(live_provider_under(&good_key))?; + seed_1.add_provider(live_provider_under(&good_key))?; + let corrupt_key = fresh_record_key(); + inject_corrupt_worker_provider(&seed_1, &corrupt_key); + let empty_key = fresh_record_key(); + db.insert::( + &seed_1.key_to_hash(&empty_key), + &KadProviderRow::encode(empty_key.clone(), &[]), + )?; + let mixed_key = fresh_record_key(); + db.insert::( + &seed_1.key_to_hash(&mixed_key), + &KadProviderRow::encode(mixed_key.clone(), &[live_provider_under(&good_key).into()]), + )?; + let unknown_key = fresh_record_key(); + let unknown_hash = seed_1.key_to_hash(&unknown_key); + db.insert::(&unknown_hash, &CORRUPT_PROVIDER_BYTES.to_vec())?; + + let mut worker_0 = + KadStore::new(db.clone(), PeerId::random(), &key_config, NetworkType::Worker(0)); + let mut worker_1 = + KadStore::new(db.clone(), PeerId::random(), &key_config, NetworkType::Worker(1)); + assert_eq!(worker_0.num_providers, 1); + assert_eq!(worker_1.num_providers, 4, "unreadable envelopes have no assumed owner"); + assert!(worker_0.providers(&corrupt_key).is_empty()); + assert!(worker_1.providers(&mixed_key).is_empty(), "mismatched payload key is rejected"); + assert!(worker_1.providers(&unknown_key).is_empty()); + assert_eq!(worker_0.scrub_corrupt_providers(), 0); + assert_eq!(worker_0.evict_expired_providers(), 0); + assert_eq!(db.iter::().count(), 6); + + assert_eq!(worker_1.scrub_corrupt_providers(), 2); + assert_eq!(worker_1.evict_expired_providers(), 1, "empty owned payload frees its slot"); + assert_eq!(worker_1.num_providers, 1); + assert_eq!(worker_0.num_providers, 1); + assert_eq!(worker_0.providers(&good_key).len(), 1); + assert_eq!(worker_1.providers(&good_key).len(), 1); + assert!(db.get::(&unknown_hash)?.is_some()); + + worker_1.config.max_provided_keys = 1; + assert!( + matches!( + worker_1.add_provider(live_provider_under(&unknown_key)), + Err(Error::MaxProvidedKeys) + ), + "replacing an uncounted envelope must still enforce capacity" + ); + // Exercise deletion of an on-disk row, then wait for the queued removal to persist. + db.sync_persist(); + worker_1.remove_provider(&unknown_key, &PeerId::random()); + db.sync_persist(); + assert_eq!(worker_1.num_providers, 1, "removing an uncounted envelope preserves the count"); + assert!(db.get::(&unknown_hash)?.is_none()); + Ok(()) + } + // ---- issue #1185: a provider record's address list is capped ---- /// A live provider record under `key` that carries `count` distinct addresses. diff --git a/crates/network-libp2p/src/lib.rs b/crates/network-libp2p/src/lib.rs index c4a691d5e..8b4ab4b9e 100644 --- a/crates/network-libp2p/src/lib.rs +++ b/crates/network-libp2p/src/lib.rs @@ -46,3 +46,6 @@ pub use libp2p::{ #[cfg(test)] #[path = "./tests/common.rs"] pub(crate) mod common; +#[cfg(test)] +#[path = "tests/fixture_tests.rs"] +mod fixture_tests; diff --git a/crates/network-libp2p/src/metrics.rs b/crates/network-libp2p/src/metrics.rs index 4f8be7f14..4966e0620 100644 --- a/crates/network-libp2p/src/metrics.rs +++ b/crates/network-libp2p/src/metrics.rs @@ -1,7 +1,7 @@ //! Prometheus metrics for the libp2p consensus networks. //! //! Both the primary and worker networks instantiate the same types, so every series -//! carries a `network` label ({`primary`, `worker`}) set at construction. +//! carries a `network` label (`primary` or `worker-{id}`) set at construction. use crate::{peers::Penalty, types::NetworkType}; use reth_metrics::{ @@ -10,10 +10,10 @@ use reth_metrics::{ }; /// Map a [`NetworkType`] to its metric label value. -pub(crate) fn network_label(network_type: &NetworkType) -> &'static str { +pub(crate) fn network_label(network_type: &NetworkType) -> String { match network_type { - NetworkType::Primary => "primary", - NetworkType::Worker(_) => "worker", + NetworkType::Primary => "primary".to_owned(), + NetworkType::Worker(id) => format!("worker-{id}"), } } @@ -41,14 +41,17 @@ pub(crate) struct SwarmMetrics { /// The derive-backed handles. handles: SwarmMetricHandles, /// The network label value for per-event labeled counters. - network: &'static str, + network: String, } impl SwarmMetrics { /// Create the swarm metric handles for `network_type`. pub(crate) fn new_for(network_type: &NetworkType) -> Self { let network = network_label(network_type); - Self { handles: SwarmMetricHandles::new_with_labels(&[("network", network)]), network } + Self { + handles: SwarmMetricHandles::new_with_labels(&[("network", network.clone())]), + network, + } } /// Record a successfully published gossip message. @@ -81,7 +84,7 @@ impl SwarmMetrics { pub(crate) fn record_outbound_failure(&self, kind: &'static str) { metrics::counter!( "tn_network.outbound_request_failures_total", - "network" => self.network, + "network" => self.network.clone(), "kind" => kind, ) .increment(1); @@ -117,7 +120,7 @@ pub(crate) struct PeerManagerMetrics { /// The derive-backed handles. handles: PeerManagerMetricHandles, /// The network label value for per-event labeled counters. - network: &'static str, + network: String, } impl PeerManagerMetrics { @@ -125,7 +128,7 @@ impl PeerManagerMetrics { pub(crate) fn new_for(network_type: &NetworkType) -> Self { let network = network_label(network_type); Self { - handles: PeerManagerMetricHandles::new_with_labels(&[("network", network)]), + handles: PeerManagerMetricHandles::new_with_labels(&[("network", network.clone())]), network, } } @@ -148,7 +151,7 @@ impl PeerManagerMetrics { pub(crate) fn record_connection_established(&self, direction: &'static str) { metrics::counter!( "tn_network.connections_established_total", - "network" => self.network, + "network" => self.network.clone(), "direction" => direction, ) .increment(1); @@ -179,7 +182,7 @@ impl PeerManagerMetrics { }; metrics::counter!( "tn_network.peer_penalties_total", - "network" => self.network, + "network" => self.network.clone(), "severity" => severity, ) .increment(1); @@ -196,6 +199,7 @@ mod tests { use super::*; use metrics_util::debugging::{DebugValue, DebuggingRecorder}; + /// Primary and worker metrics register their expected labels and update every handle. #[test] fn test_metrics_register_and_update() { let recorder = DebuggingRecorder::new(); @@ -233,7 +237,7 @@ mod tests { let (key, _, _, value) = find("tn_network.connected_peers"); assert!(matches!(value, DebugValue::Gauge(g) if g.0 == 4.0)); - assert!(key.key().labels().any(|l| l.key() == "network" && l.value() == "worker")); + assert!(key.key().labels().any(|l| l.key() == "network" && l.value() == "worker-0")); let (key, _, _, _) = find("tn_network.outbound_request_failures_total"); assert!(key.key().labels().any(|l| l.key() == "kind" && l.value() == "timeout")); @@ -252,4 +256,75 @@ mod tests { find("tn_network.dial_failures_total"); find("tn_network.connections_closed_total"); } + + /// Worker swarms must retain independent gauges and counters in the shared recorder. + #[test] + fn test_worker_metrics_are_isolated() { + let recorder = DebuggingRecorder::new(); + let snapshotter = recorder.snapshotter(); + + metrics::with_local_recorder(&recorder, || { + let first = SwarmMetrics::new_for(&NetworkType::Worker(0)); + let second = SwarmMetrics::new_for(&NetworkType::Worker(1)); + first.set_pending(3, 3); + second.set_pending(7, 7); + [&first, &second].into_iter().for_each(|swarm| { + swarm.record_gossip_published(); + swarm.record_outbound_failure("timeout"); + }); + + let first = PeerManagerMetrics::new_for(&NetworkType::Worker(0)); + let second = PeerManagerMetrics::new_for(&NetworkType::Worker(1)); + first.set_peer_counts(3, 3, 3, 3); + second.set_peer_counts(7, 7, 7, 7); + [&first, &second].into_iter().for_each(|peers| { + peers.record_connection_established("in"); + peers.record_penalty(&Penalty::Severe); + }); + }); + + let snapshot = snapshotter.snapshot().into_vec(); + [("worker-0", 3.0), ("worker-1", 7.0)].into_iter().for_each(|(network, expected)| { + let value = |name| { + snapshot + .iter() + .find(|(key, ..)| { + key.key().name() == name + && key + .key() + .labels() + .any(|label| label.key() == "network" && label.value() == network) + }) + .map(|(_, _, _, value)| value) + }; + [ + "tn_network.px_disconnects_pending", + "tn_network.outbound_requests_pending", + "tn_network.connected_peers", + "tn_network.known_peers", + "tn_network.discovery_peers", + "tn_network.banned_peers", + ] + .into_iter() + .for_each(|name| { + assert!( + matches!(value(name), Some(DebugValue::Gauge(g)) if g.0 == expected), + "{name} must retain {network}'s gauge value" + ); + }); + [ + "tn_network.gossip_published_total", + "tn_network.outbound_request_failures_total", + "tn_network.connections_established_total", + "tn_network.peer_penalties_total", + ] + .into_iter() + .for_each(|name| { + assert!( + matches!(value(name), Some(DebugValue::Counter(1))), + "{name} must count {network}'s events separately" + ); + }); + }); + } } diff --git a/crates/network-libp2p/src/tests/fixture_tests.rs b/crates/network-libp2p/src/tests/fixture_tests.rs new file mode 100644 index 000000000..8e603d596 --- /dev/null +++ b/crates/network-libp2p/src/tests/fixture_tests.rs @@ -0,0 +1,31 @@ +//! Regression tests for network identities exposed by committee fixtures. + +use tn_storage::mem_db::MemDatabase; +use tn_test_utils::CommitteeFixture; +use tn_types::{NetworkPublicKey, DEFAULT_WORKER_ID}; + +/// Every authority fixture exposes worker zero with the key advertised for that worker. +#[test] +fn every_authority_fixture_matches_advertised_worker_zero() -> Result<(), &'static str> { + let fixture = CommitteeFixture::builder(MemDatabase::default).build(); + let committee = fixture.committee(); + let bootstrap_servers = committee.bootstrap_servers(); + assert_eq!(fixture.num_authorities(), 4, "exercise authorities beyond the first position"); + assert_eq!(committee.number_of_workers(), 1); + + fixture.authorities().try_for_each(|authority| { + let worker = authority.worker(); + let advertised = bootstrap_servers + .get(&authority.primary_public_key()) + .and_then(|server| server.worker(DEFAULT_WORKER_ID)) + .ok_or("every authority must advertise worker zero")?; + assert_eq!(worker.id, DEFAULT_WORKER_ID, "worker id must not depend on authority order"); + let fixture_key: NetworkPublicKey = worker.keypair().public().into(); + assert_eq!( + fixture_key, advertised.network_key, + "fixture must authenticate as its authority's advertised worker zero" + ); + Ok::<(), &'static str>(()) + })?; + Ok(()) +} diff --git a/crates/network-libp2p/src/tests/network_tests.rs b/crates/network-libp2p/src/tests/network_tests.rs index 9fb842681..8a6e60cd9 100644 --- a/crates/network-libp2p/src/tests/network_tests.rs +++ b/crates/network-libp2p/src/tests/network_tests.rs @@ -723,7 +723,7 @@ async fn test_primary_worker_protocol_isolation() -> eyre::Result<()> { config_2.network_config(), tx2, config_2.key_config().clone(), - config_2.key_config().worker_network_keypair().clone(), + config_2.key_config().worker_network_keypair(DEFAULT_WORKER_ID), MemDatabase::default(), task_manager.get_spawner(), NetworkType::Worker(0), @@ -753,7 +753,7 @@ async fn test_primary_worker_protocol_isolation() -> eyre::Result<()> { primary .add_explicit_peer( worker_bls, - config_2.key_config().worker_network_public_key(), + config_2.key_config().worker_network_public_key(DEFAULT_WORKER_ID), worker_addr, ) .await?; @@ -853,7 +853,7 @@ async fn test_unsupported_protocol_does_not_penalize() -> eyre::Result<()> { config_2.network_config(), tx2, config_2.key_config().clone(), - config_2.key_config().worker_network_keypair().clone(), + config_2.key_config().worker_network_keypair(DEFAULT_WORKER_ID), MemDatabase::default(), task_manager.get_spawner(), NetworkType::Worker(0), @@ -878,7 +878,7 @@ async fn test_unsupported_protocol_does_not_penalize() -> eyre::Result<()> { primary .add_explicit_peer( worker_bls, - config_2.key_config().worker_network_public_key(), + config_2.key_config().worker_network_public_key(DEFAULT_WORKER_ID), worker_addr, ) .await?; @@ -3362,6 +3362,100 @@ async fn test_startup_scrubs_legacy_and_corrupt_kad_records() -> eyre::Result<() Ok(()) } +/// Startup verifies only its worker namespace and preserves sibling records signed by the same key. +#[tokio::test] +async fn test_worker_startup_preserves_sibling_kad_records() -> eyre::Result<()> { + use libp2p::kad; + use tn_config::KeyConfig; + use tn_storage::tables::KadWorkerRecords; + + let key_config = KeyConfig::new_with_testing_key(BlsKeypair::generate(&mut rand::rng())); + let publisher_key_config = + KeyConfig::new_with_testing_key(BlsKeypair::generate(&mut rand::rng())); + let owner_bls = publisher_key_config.primary_public_key(); + let key = kad::RecordKey::new(&owner_bls); + let network_config = NetworkConfig::default(); + let chain_id = network_config.libp2p_config().chain_id; + let task_manager = TaskManager::default(); + let db = MemDatabase::default(); + let local_network_key = key_config.worker_network_keypair(0); + let worker_0_key = publisher_key_config.worker_network_keypair(0); + let worker_1_key = publisher_key_config.worker_network_keypair(1); + let worker_0_address = create_multiaddr(None); + let worker_0_record = NodeRecord::build( + RecordDomain::new(chain_id, NetworkType::Worker(0)), + worker_0_key.public().into(), + worker_0_address.clone(), + None, + |data| publisher_key_config.request_signature_direct(data), + ); + let worker_1_record = NodeRecord::build( + RecordDomain::new(chain_id, NetworkType::Worker(1)), + worker_1_key.public().into(), + create_multiaddr(None), + None, + |data| publisher_key_config.request_signature_direct(data), + ); + let mut store_0 = KadStore::new( + db.clone(), + local_network_key.public().into(), + &key_config, + NetworkType::Worker(0), + ); + let mut store_1 = KadStore::new( + db.clone(), + key_config.worker_network_keypair(1).public().into(), + &key_config, + NetworkType::Worker(1), + ); + store_0.put(kad::Record { + key: key.clone(), + value: encode(&worker_0_record), + publisher: None, + expires: None, + })?; + store_1.put(kad::Record { + key: key.clone(), + value: encode(&worker_1_record), + publisher: None, + expires: None, + })?; + let persisted = db.iter::().collect::>(); + assert_eq!(persisted.len(), 2); + + let (tx, _network_events) = mpsc::channel(10); + let network = ConsensusNetwork::< + TestWorkerRequest, + TestWorkerResponse, + MemDatabase, + mpsc::Sender>, + >::new( + &network_config, + tx, + key_config.clone(), + local_network_key, + db.clone(), + task_manager.get_spawner(), + NetworkType::Worker(0), + worker_0_address, + None, + )?; + assert_eq!(db.iter::().collect::>(), persisted); + assert_eq!( + store_0.get(&key).map(|record| record.value.clone()), + Some(encode(&worker_0_record)) + ); + assert_eq!( + store_1.get(&key).map(|record| record.value.clone()), + Some(encode(&worker_1_record)) + ); + assert_eq!( + network.swarm.behaviour().peer_manager.auth_to_peer(owner_bls).map(|(peer_id, _)| peer_id), + Some(worker_0_key.public().into()), + ); + Ok(()) +} + /// Records restored from the persisted kad store at startup are UNPINNED: the store legitimately /// holds arbitrary signature-valid third-party records (DHT storage duty), so a restart must not /// convert them into permanently pinned `known_peers` entries. Restored records resolve until the diff --git a/crates/network-libp2p/src/tests/types.rs b/crates/network-libp2p/src/tests/types.rs index 21ba6f2a6..8241fde76 100644 --- a/crates/network-libp2p/src/tests/types.rs +++ b/crates/network-libp2p/src/tests/types.rs @@ -178,6 +178,27 @@ fn test_cross_role_replay_rejected() { assert!(NodeRecord::decode_and_verify(&bytes, primary_domain, &pubkey).is_none()); } +/// Sibling workers must reject each other's records even with the same chain and BLS key. +#[test] +fn test_cross_worker_replay_rejected() { + let key_config = KeyConfig::new_with_testing_key(BlsKeypair::generate(&mut rand::rng())); + let pubkey = key_config.primary_public_key(); + let worker_0 = RecordDomain::new(2017, NetworkType::Worker(0)); + let worker_1 = RecordDomain::new(2017, NetworkType::Worker(1)); + let record = NodeRecord::build( + worker_0, + key_config.primary_network_public_key(), + create_multiaddr(None), + None, + |data| key_config.request_signature_direct(data), + ); + assert!(record.clone().verify(worker_0, &pubkey).is_some()); + assert!(record.clone().verify(worker_1, &pubkey).is_none()); + let bytes = tn_types::encode(&record); + assert!(NodeRecord::decode_and_verify(&bytes, worker_0, &pubkey).is_some()); + assert!(NodeRecord::decode_and_verify(&bytes, worker_1, &pubkey).is_none()); +} + /// GHSA-cc64-wfq5-56ph cross-CHAIN replay: a record signed for one chain /// verifies under that chain but is REJECTED under a different chain id (same /// role, same BLS key), on both the in-memory and bytes paths. diff --git a/crates/node/src/manager/node.rs b/crates/node/src/manager/node.rs index b99045312..04e32f9ae 100644 --- a/crates/node/src/manager/node.rs +++ b/crates/node/src/manager/node.rs @@ -32,8 +32,8 @@ use tn_types::{ gas_accumulator::{entry_fee_for_worker, GasAccumulator}, repack_monitor::RepackMonitor, BlsPublicKey, BootstrapServer, Committee, ConsensusHeader, ConsensusHeaderDigest, - ConsensusNumHash, ConsensusOutput, Database as TNDatabase, EngineUpdate, Epoch, SealedHeader, - ShutdownNotifier, TaskError, TaskManager, TaskSpawner, TimestampSec, WorkerId, + ConsensusNumHash, ConsensusOutput, Database as TNDatabase, EngineUpdate, Epoch, P2pNode, + SealedHeader, ShutdownNotifier, TaskError, TaskManager, TaskSpawner, TimestampSec, WorkerId, DEFAULT_WORKER_ID, }; // The canonical worker-attribution helper lives in `tn-types` (one implementation, no drift); @@ -73,6 +73,81 @@ const EXEX_EVENT_CAPACITY: usize = 16; /// manager's forwarder instead of consuming memory. const TO_ENGINE_CAPACITY: usize = 64; +/// Inputs for one worker swarm, validated before any process-lifetime network is spawned. +struct PreparedWorkerNetwork { + /// Worker identity shared by its key derivation, protocols, and network handle. + worker_id: WorkerId, + /// This worker's advertised address and optional RPC endpoint. + p2p: P2pNode, + /// The persistent event stream at the same index as the worker configuration. + event_stream: Events, +} + +/// Require local swarm configuration to match the authoritative count for the entering epoch. +fn check_configured_worker_count( + epoch: Epoch, + on_chain_workers: usize, + configured_workers: usize, +) -> eyre::Result<()> { + eyre::ensure!( + configured_workers != 0, + "node config `node_info.p2p_info.workers` must configure at least one worker" + ); + eyre::ensure!( + configured_workers == on_chain_workers, + "node config `node_info.p2p_info.workers` lists {configured_workers} workers but the \ + chain-derived count for epoch {epoch} is {on_chain_workers}: every validator must run \ + the worker count the committee carries" + ); + Ok(()) +} + +/// Validate every worker and pair its configuration with its persistent event stream. +/// +/// Runs before either primary or worker swarm construction. Checking the complete layout first +/// prevents a bad later RPC endpoint, count, or worker id from leaving partially spawned networks. +fn prepare_worker_networks( + workers: &[P2pNode], + event_streams: &[Events], + epoch: Epoch, + on_chain_workers: usize, +) -> eyre::Result>> { + let configured = workers.len(); + eyre::ensure!( + configured <= 1 || tn_types::forks::multi_workers_fork_active(epoch), + "node config `node_info.p2p_info.workers` lists {configured} workers but the \ + multi-workers fork is not active at epoch {epoch}: configure exactly one worker" + ); + check_configured_worker_count(epoch, on_chain_workers, configured)?; + eyre::ensure!( + configured <= usize::from(WorkerId::MAX) + 1, + "node config lists {configured} workers, exceeding the WorkerId range" + ); + eyre::ensure!( + configured == event_streams.len(), + "node config lists {configured} workers but has {} worker event streams", + event_streams.len() + ); + + workers + .iter() + .zip(event_streams) + .zip(0..=WorkerId::MAX) + .map(|((p2p, event_stream), worker_id)| { + p2p.rpc.as_ref().map(tn_types::RpcInfo::validate).transpose().wrap_err_with(|| { + format!( + "invalid `node_info.p2p_info.workers[{worker_id}].rpc` endpoint in node config" + ) + })?; + Ok(PreparedWorkerNetwork { + worker_id, + p2p: p2p.clone(), + event_stream: event_stream.clone(), + }) + }) + .collect() +} + /// The long-running owner that oversees epoch transitions. /// /// One instance exists for the lifetime of the process. It holds the resources that must survive @@ -90,8 +165,9 @@ pub(crate) struct EpochManager { tn_datadir: P, /// Primary network handle. primary_network_handle: Option, - /// Worker network handle. - worker_network_handle: Option, + /// Worker network handles, indexed by [`WorkerId`](tn_types::WorkerId). Empty until + /// [`spawn_node_networks`](Self::spawn_node_networks) runs. + worker_network_handles: Vec, /// Key config - loaded once for application lifetime. key_config: KeyConfig, /// The epoch manager's [ShutdownNotifier] to shutdown all node processes. @@ -124,9 +200,10 @@ pub(crate) struct EpochManager { /// Application-scoped consensus bus. Survives epoch boundaries and is reset between epochs via /// `reset_for_epoch`; carries `recent_blocks`, node mode, and other cross-component state. consensus_bus: ConsensusBusApp, - /// Persistent event stream for the long-running worker network. Outlives any single epoch so - /// the worker swarm does not have to be rebuilt on each transition. - worker_event_stream: QueChannel>, + /// Persistent event streams for the long-running worker networks, one per configured worker + /// and indexed by [`WorkerId`](tn_types::WorkerId). Outlive any single epoch so the worker + /// swarms do not have to be rebuilt on each transition. + worker_event_streams: Vec>>, /// Final consensus header of the epoch that just closed, carried into the next epoch so it can /// be used as the starting point for the new epoch's chain. @@ -610,7 +687,10 @@ where // Don't risk keeping the default CVV active mode... consensus_bus.node_mode().send_replace(NodeMode::Observer); } - let worker_event_stream = QueChannel::new(); + // one event stream per configured worker, indexed by worker id + let worker_event_streams = (0..builder.tn_config.node_info.p2p_info.num_workers()) + .map(|_| QueChannel::new()) + .collect(); let bootstrap_servers = if let Ok(committee_zero) = Config::load_from_path_or_default::( tn_datadir.committee_path(), @@ -639,7 +719,7 @@ where builder, tn_datadir, primary_network_handle: None, - worker_network_handle: None, + worker_network_handles: Vec::new(), key_config, node_shutdown, epoch_boundary: Default::default(), @@ -647,7 +727,7 @@ where reth_db, consensus_db, consensus_bus, - worker_event_stream, + worker_event_streams, last_consensus_header: None, last_forwarded_consensus_number: 0, consensus_chain, @@ -745,8 +825,13 @@ where let reth_env = engine.get_reth_env().await; reth_env.heal_finalized_to_persisted_tip()?; // retrieve epoch information from canonical tip on startup - let EpochState { epoch, .. } = engine.epoch_state_from_canonical_tip().await?; + let EpochState { epoch, epoch_info, .. } = engine.epoch_state_from_canonical_tip().await?; debug!(target: "epoch-manager", ?epoch, "retrieved epoch state from canonical tip"); + // Read the raw count independently: fresh genesis has no finalized header, so catchup + // leaves the accumulator at its initial size. The accumulator also clamps zero to one. + // Both startup validation and epoch entry must use the authoritative closing-block count. + let on_chain_workers = + read_num_workers_at_epoch_entry(&reth_env, epoch_info.blockHeight).await?; // The canonical epoch cross-checks the finalized header catchup pins its reads to. catchup_accumulator(reth_env, &gas_accumulator, &mut self.consensus_chain, epoch).await?; self.try_restore_state(&engine).await?; @@ -757,7 +842,8 @@ where // network builder, the gossip handles, and the gossip-validation handlers. let mut network_config = NetworkConfig::read_config(&self.tn_datadir)?; network_config.set_chain_id(self.builder.tn_config.genesis().config.chain_id); - self.spawn_node_networks(node_task_spawner, &network_config, epoch).await?; + self.spawn_node_networks(node_task_spawner, &network_config, epoch, on_chain_workers) + .await?; let primary_network_handle = self.primary_network_handle.as_ref().expect("primary network").clone(); // `epoch_vote_topic` and `consensus_output_topic` are committee-only publish topics, so @@ -994,14 +1080,25 @@ where /// Spawn the process-lifetime primary and worker [`ConsensusNetwork`] swarms. /// /// Each swarm runs as a critical task until node shutdown. The resulting network handles are - /// stored on the manager for use by every epoch; the worker handle is seeded with the starting - /// `epoch` and its task spawner is refreshed on each epoch transition. + /// stored on the manager for use by every epoch; the worker handles are seeded with the + /// starting `epoch` and their task spawners are refreshed on each epoch transition. + /// The configured worker count must match the raw chain count at the previous epoch's closing + /// block (genesis for epoch 0) before any swarm is created. This includes fresh genesis where + /// accumulator catchup is a no-op. Worker RPC descriptors and event streams validate together. async fn spawn_node_networks( &mut self, node_task_spawner: TaskSpawner, network_config: &NetworkConfig, epoch: Epoch, + on_chain_workers: usize, ) -> eyre::Result<()> { + let workers = prepare_worker_networks( + &self.builder.tn_config.node_info.p2p_info.workers, + &self.worker_event_streams, + epoch, + on_chain_workers, + )?; + // Reject an invalid peer-score config before it is installed into the process-global, // first-write-wins `GLOBAL_SCORE_CONFIG` by the `PeerManager` built below // (`init_peer_score_config`). This is the boot-path install funnel, so validating here @@ -1046,59 +1143,55 @@ where self.primary_network_handle = Some(PrimaryNetworkHandle::new(primary_network_handle, network_config.chain_id())); - // pass through the worker's RPC descriptor so peers can discover this - // validator's JSON-RPC endpoint via kademlia. validators that did not - // configure RPC leave the descriptor `None`. fail fast on a misconfigured - // endpoint rather than advertising something peers will reject. - let worker_p2p = self - .builder - .tn_config - .node_info - .p2p_info - .worker(DEFAULT_WORKER_ID) - .ok_or_else(|| eyre!("no worker {DEFAULT_WORKER_ID} in node info"))? - .clone(); - let worker_rpc = worker_p2p.rpc; - if let Some(rpc) = &worker_rpc { - rpc.validate() - .wrap_err("invalid `node_info.p2p_info.workers[0].rpc` endpoint in node config")?; - } - - // create long-running network task for worker - let worker_network = ConsensusNetwork::new_for_worker( - DEFAULT_WORKER_ID, - network_config, - self.worker_event_stream.clone(), - self.key_config.clone(), - self.consensus_db.clone(), - node_task_spawner.clone(), - worker_p2p.network_address, - worker_rpc, - )?; - let worker_network_handle = worker_network.network_handle(); - let node_shutdown = self.node_shutdown.subscribe(); + // + //=== WORKERS + // - // spawn long-running primary network task - node_task_spawner.spawn_critical_task("Worker Network", async move { - tokio::select!( - _ = &node_shutdown => { - Ok(()) - } - res = worker_network.run() => { - warn!(target: "epoch-manager", ?res, "worker network stopped"); - Ok(res?) - } - ) - }); + // create one long-running swarm per configured worker + // the per-epoch code still drives worker 0 only (#557 loops over worker components) + self.worker_network_handles = workers + .into_iter() + .map(|PreparedWorkerNetwork { worker_id, p2p, event_stream }| { + // create long-running network task for this worker + let worker_network = ConsensusNetwork::new_for_worker( + worker_id, + network_config, + event_stream, + self.key_config.clone(), + self.consensus_db.clone(), + node_task_spawner.clone(), + p2p.network_address, + p2p.rpc, + )?; + let worker_network_handle = worker_network.network_handle(); + let node_shutdown = self.node_shutdown.subscribe(); + + // spawn long-running worker network task + node_task_spawner.spawn_critical_task( + format!("Worker Network {worker_id}"), + async move { + tokio::select!( + _ = &node_shutdown => { + Ok(()) + } + res = worker_network.run() => { + warn!(target: "epoch-manager", ?res, "worker network stopped"); + Ok(res?) + } + ) + }, + ); - // set temporary task spawner - this is updated with each epoch - self.worker_network_handle = Some(WorkerNetworkHandle::new( - worker_network_handle, - node_task_spawner.clone(), - DEFAULT_WORKER_ID, - epoch, - network_config.chain_id(), - )); + // set temporary task spawner - this is updated with each epoch + Ok(WorkerNetworkHandle::new( + worker_network_handle, + node_task_spawner.clone(), + worker_id, + epoch, + network_config.chain_id(), + )) + }) + .collect::>>()?; Ok(()) } @@ -1341,7 +1434,172 @@ fn check_restore_consistency( #[cfg(test)] mod tests { use super::*; - use tn_types::{ExecHeader, B256}; + use rand::{rngs::StdRng, SeedableRng as _}; + use tn_types::{BlsKeypair, ExecHeader, RpcInfo, TnReceiver as _, TnSender as _, B256}; + + /// Reproducible keys for checking the identity assigned to each prepared swarm. + fn worker_key_config() -> KeyConfig { + KeyConfig::new_with_testing_key(BlsKeypair::generate(&mut StdRng::seed_from_u64(1315))) + } + + /// Give each worker a distinct advertised address, derived network key, and RPC endpoint. + fn worker_p2p(keys: &KeyConfig, worker_id: WorkerId) -> eyre::Result { + Ok(P2pNode { + network_address: format!("/ip4/127.0.0.1/udp/{}/quic-v1", 9000 + u32::from(worker_id)) + .parse()?, + network_key: keys.worker_network_public_key(worker_id), + rpc: Some(RpcInfo { + http: format!("https://worker-{worker_id}.example.com/").parse()?, + ws: None, + }), + }) + } + + /// Every prepared swarm keeps the key, endpoint, worker id, and original event stream aligned. + #[test] + fn prepare_worker_networks_preserves_worker_identity_and_event_streams() -> eyre::Result<()> { + let keys = worker_key_config(); + let workers = [worker_p2p(&keys, 0)?, worker_p2p(&keys, 1)?, worker_p2p(&keys, 2)?]; + let streams: Vec> = + workers.iter().map(|_| QueChannel::new()).collect(); + let receivers: Vec<_> = streams.iter().map(QueChannel::subscribe).collect(); + let prepared = prepare_worker_networks(&workers, &streams, Epoch::MAX, 3)?; + + assert_eq!(prepared.len(), 3); + assert_ne!(keys.worker_network_public_key(0), keys.worker_network_public_key(1)); + prepared.iter().zip(&workers).zip([0, 1, 2]).try_for_each( + |((worker, configured), expected_id)| -> eyre::Result<()> { + assert_eq!(worker.worker_id, expected_id); + assert_eq!(worker.p2p, *configured); + assert_eq!( + worker.p2p.network_key, + keys.worker_network_public_key(worker.worker_id) + ); + worker.event_stream.try_send(worker.worker_id)?; + Ok(()) + }, + )?; + receivers.into_iter().zip([0, 1, 2]).try_for_each( + |(mut receiver, expected_id)| -> eyre::Result<()> { + assert_eq!(receiver.try_recv()?, expected_id); + assert!(receiver.try_recv().is_err()); + Ok(()) + }, + )?; + Ok(()) + } + + /// Both missing and surplus event streams are rejected instead of silently truncating `zip`. + #[test] + fn prepare_worker_networks_rejects_event_stream_count_mismatch() -> eyre::Result<()> { + let workers = [worker_p2p(&worker_key_config(), 0)?]; + [Vec::new(), vec![(), ()]].into_iter().try_for_each(|streams| -> eyre::Result<()> { + let error = prepare_worker_networks(&workers, &streams, 0, 1) + .err() + .ok_or_else(|| eyre!("expected worker event stream count mismatch"))?; + assert!(error.to_string().contains("worker event streams")); + Ok(()) + }) + } + + /// A bad RPC endpoint on a later worker fails preparation of the entire swarm set. + #[test] + fn prepare_worker_networks_rejects_later_invalid_rpc() -> eyre::Result<()> { + let keys = worker_key_config(); + let worker = worker_p2p(&keys, 1)?; + let workers = [ + worker_p2p(&keys, 0)?, + P2pNode { + rpc: Some(RpcInfo { http: "ftp://worker-1.example.com/".parse()?, ws: None }), + ..worker + }, + ]; + let error = prepare_worker_networks(&workers, &[(), ()], Epoch::MAX, 2) + .err() + .ok_or_else(|| eyre!("expected invalid RPC endpoint"))?; + assert!(error.to_string().contains("node_info.p2p_info.workers[1].rpc")); + Ok(()) + } + + /// Startup requires the local count to match the raw chain-derived count, in both directions. + #[test] + fn prepare_worker_networks_rejects_chain_count_mismatch() -> eyre::Result<()> { + let keys = worker_key_config(); + let workers = [worker_p2p(&keys, 0)?, worker_p2p(&keys, 1)?]; + let single_worker = workers.get(..1).ok_or_else(|| eyre!("expected worker zero"))?; + [(workers.as_slice(), 1), (single_worker, 2)].into_iter().try_for_each( + |(configured, chain_count)| -> eyre::Result<()> { + let streams = vec![(); configured.len()]; + let error = prepare_worker_networks(configured, &streams, Epoch::MAX, chain_count) + .err() + .ok_or_else(|| eyre!("expected chain worker count mismatch"))?; + assert!(error.to_string().contains("chain-derived count")); + Ok(()) + }, + )?; + Ok(()) + } + + /// Raw zero counts must fail even though the gas accumulator clamps a zero resize to one. + #[test] + fn prepare_worker_networks_rejects_zero_worker_counts() -> eyre::Result<()> { + let workers = [worker_p2p(&worker_key_config(), 0)?]; + let chain_error = prepare_worker_networks(&workers, &[()], 0, 0) + .err() + .ok_or_else(|| eyre!("expected zero chain count rejection"))?; + assert!(chain_error.to_string().contains("chain-derived count for epoch 0 is 0")); + let config_error = prepare_worker_networks::<()>(&[], &[], 0, 0) + .err() + .ok_or_else(|| eyre!("expected empty worker config rejection"))?; + assert!(config_error.to_string().contains("at least one worker")); + Ok(()) + } + + /// Fresh multi-worker genesis uses the chain count even before accumulator catchup can run. + #[cfg(not(feature = "adiri"))] + #[test] + fn prepare_worker_networks_accepts_multi_worker_genesis() -> eyre::Result<()> { + let keys = worker_key_config(); + let workers = [worker_p2p(&keys, 0)?, worker_p2p(&keys, 1)?]; + let prepared = prepare_worker_networks(&workers, &[(), ()], 0, 2)?; + assert_eq!(prepared.len(), 2); + Ok(()) + } + + /// The full WorkerId range is accepted, and the next configured worker is rejected. + #[test] + fn prepare_worker_networks_enforces_worker_id_bound() -> eyre::Result<()> { + let max_workers = usize::from(WorkerId::MAX) + 1; + let worker = P2pNode { rpc: None, ..worker_p2p(&worker_key_config(), 0)? }; + let workers = vec![worker; max_workers + 1]; + let streams = vec![(); max_workers + 1]; + let prepared = prepare_worker_networks( + &workers[..max_workers], + &streams[..max_workers], + Epoch::MAX, + max_workers, + )?; + assert_eq!(prepared.len(), max_workers); + assert_eq!(prepared.last().map(|worker| worker.worker_id), Some(WorkerId::MAX)); + let error = prepare_worker_networks(&workers, &streams, Epoch::MAX, max_workers + 1) + .err() + .ok_or_else(|| eyre!("expected WorkerId overflow"))?; + assert!(error.to_string().contains("WorkerId range")); + Ok(()) + } + + /// A multi-worker local config cannot start any swarms before the fork activates. + #[cfg(feature = "adiri")] + #[test] + fn prepare_worker_networks_rejects_pre_fork_multiple_workers() -> eyre::Result<()> { + let keys = worker_key_config(); + let workers = [worker_p2p(&keys, 0)?, worker_p2p(&keys, 1)?]; + let error = prepare_worker_networks(&workers, &[(), ()], 0, 2) + .err() + .ok_or_else(|| eyre!("expected pre-fork multi-worker rejection"))?; + assert!(error.to_string().contains("multi-workers fork is not active")); + Ok(()) + } /// A tip sealed header at `number` whose nonce encodes `epoch` (upper 32 bits), matching the /// payload builder's `nonce = epoch << 32 | round` layout that `deconstruct_nonce` reads back. diff --git a/crates/node/src/manager/node/start_epoch.rs b/crates/node/src/manager/node/start_epoch.rs index deaf04a86..1802694e6 100644 --- a/crates/node/src/manager/node/start_epoch.rs +++ b/crates/node/src/manager/node/start_epoch.rs @@ -231,7 +231,8 @@ where authority_id: public_key.into(), execution_address: self.builder.tn_config.node_info.execution_address, primary_network_key: self.key_config.primary_network_public_key(), - worker_network_key: self.key_config.worker_network_public_key(), + // the node record only advertises worker 0 for now (#557) + worker_network_key: self.key_config.worker_network_public_key(DEFAULT_WORKER_ID), primary_external_address: self .builder .tn_config @@ -361,7 +362,11 @@ where .ok_or_else(|| eyre!("on-chain WorkerConfigs reports zero workers")) }) .wrap_err("failed to read the committee worker count from chain")?; - check_committee_worker_count(epoch, num_workers)?; + check_committee_worker_count( + epoch, + num_workers, + self.builder.tn_config.node_info.p2p_info.num_workers(), + )?; // the network must be live let committee = if epoch == 0 { @@ -434,10 +439,10 @@ where /// Construct the epoch's [`WorkerNode`] and bring up its [`WorkerNetwork`]. /// - /// Only worker id [`tn_types::DEFAULT_WORKER_ID`] is supported. The shared - /// [`WorkerNetworkHandle`] on the [`EpochManager`] is re-pointed at this epoch's task - /// spawner and epoch number before anything else, so batch reporting runs under the - /// epoch-scoped lifetime. + /// Only worker id [`tn_types::DEFAULT_WORKER_ID`] is driven for now (#557 adds the loop). + /// That worker's [`WorkerNetworkHandle`] on the [`EpochManager`] is re-pointed at this + /// epoch's task spawner and epoch number before anything else, so batch reporting runs + /// under the epoch-scoped lifetime. /// /// The engine's worker components are initialized on the initial epoch, and also whenever /// the engine reports no workers yet — the latter covers the case where the first epoch @@ -471,9 +476,9 @@ where // update the network handle's task spawner for reporting batches in the epoch { let network_handle = self - .worker_network_handle - .as_mut() - .ok_or_eyre("worker network handle missing from epoch manager")?; + .worker_network_handles + .get_mut(usize::from(worker_id)) + .ok_or_else(|| eyre!("no network handle for worker {worker_id}"))?; network_handle.update_task_spawner(epoch_task_spawner.clone()); network_handle.update_epoch(consensus_config.committee().epoch()); @@ -523,9 +528,9 @@ where }); let network_handle = self - .worker_network_handle - .as_ref() - .ok_or_eyre("worker network handle missing from epoch manager")? + .worker_network_handles + .get(usize::from(worker_id)) + .ok_or_else(|| eyre!("no network handle for worker {worker_id}"))? .clone(); let validator = engine @@ -804,7 +809,11 @@ where previous_committee_keys: HashSet, ) -> eyre::Result<()> { // get event streams for the worker network handler - let rx_event_stream = self.worker_event_stream.subscribe(); + let rx_event_stream = self + .worker_event_streams + .get(usize::from(*worker_id)) + .ok_or_else(|| eyre!("no event stream for worker {worker_id}"))? + .subscribe(); debug!(target: "epoch-manager", "spawning worker network for epoch"); let committee_keys: HashSet = consensus_config @@ -818,9 +827,10 @@ where .committee() .bootstrap_servers() .iter() - // worker 0 always exists (the non-empty list invariant is enforced at deserialize), so - // this `filter_map` cannot drop a peer - .filter_map(|(k, v)| v.worker(DEFAULT_WORKER_ID).cloned().map(|worker| (*k, worker))) + // worker 0 always exists (the non-empty list invariant is enforced at deserialize). + // for higher ids a missing entry drops the peer, which is correct: a peer that runs + // fewer workers has no swarm for this id + .filter_map(|(k, v)| v.worker(*worker_id).cloned().map(|worker| (*k, worker))) .collect(); let next_committee_keys: HashSet = consensus_config.next_committee_keys().iter().copied().collect(); @@ -836,17 +846,24 @@ where // start listening if the network needs to be initialized if initial_epoch { - let worker_address = Self::parse_listener_address_for_swarm( - "WORKER_LISTENER_MULTIADDR", - consensus_config.primary_networkkey(), - consensus_config - .worker_address(DEFAULT_WORKER_ID) - .ok_or_eyre("no worker network address in node info")?, - )?; + let configured_address = consensus_config + .worker_address(*worker_id) + .ok_or_else(|| eyre!("no network address for worker {worker_id} in node info"))?; + // the env override applies to worker 0 only: one env var cannot name N distinct + // listeners, so higher ids always bind their configured address + let worker_address = if *worker_id == DEFAULT_WORKER_ID { + Self::parse_listener_address_for_swarm( + "WORKER_LISTENER_MULTIADDR", + consensus_config.primary_networkkey(), + configured_address, + )? + } else { + configured_address + }; network_handle.inner_handle().start_listening(worker_address).await?; } - let worker_address = consensus_config.worker_address(DEFAULT_WORKER_ID); + let worker_address = consensus_config.worker_address(*worker_id); // always attempt to dial peers for the new epoch // the network's peer manager will intercept dial attempts for peers that are already @@ -875,10 +892,8 @@ where // later epoch unless the subscription is explicitly dropped. Skipping alone would also // skip the only refresh of this topic's authorized-publisher allowlist, freezing it on // the committee that was current when the node last subscribed. - let batch_topic = tn_config::LibP2pConfig::worker_batch_topic( - consensus_config.chain_id(), - DEFAULT_WORKER_ID, - ); + let batch_topic = + tn_config::LibP2pConfig::worker_batch_topic(consensus_config.chain_id(), *worker_id); let mode = self.consensus_bus.current_node_mode(); if should_subscribe_batch_topic(mode) { debug!(target: "epoch-manager", ?mode, "subscribing to worker batch topic"); @@ -1107,19 +1122,25 @@ fn should_subscribe_batch_topic(mode: NodeMode) -> bool { /// Whether `epoch` may be entered with an on-chain worker count of `num_workers`. /// -/// A single worker is always fine. Above one, the answer depends on the multi-workers fork -/// ([`multi_workers_fork_active`]), evaluated at the epoch being entered - the same epoch carried -/// inside the [`Committee`] this count is about to be stamped onto, so the gate here and the gate -/// the encoder consults cannot disagree: +/// The configured swarm count must match chain state, including after governance changes at an +/// epoch boundary. A matching single worker is always fine. Above one, the answer depends on the +/// multi-workers fork ([`multi_workers_fork_active`]), evaluated at the epoch being entered - the +/// same epoch carried inside the [`Committee`] this count is about to be stamped onto, so the gate +/// here and the gate the encoder consults cannot disagree: /// /// - pre-fork the legacy committee layout has no field to carry a worker count, so the encoder /// refuses the value. Halting here turns that into a diagnosable epoch-entry failure instead of a /// panic from the first pack write, which is the only thing the node could do about it anyway: /// the count is chain state and cannot be talked down locally. -/// - post-fork the count is representable and epoch entry proceeds. It still only warns, because -/// this node version spawns worker [`DEFAULT_WORKER_ID`] alone: header payloads keyed to higher -/// worker ids validate, but nothing local produces them. -fn check_committee_worker_count(epoch: Epoch, num_workers: NonZeroUsize) -> eyre::Result<()> { +/// - post-fork the count is representable and epoch entry proceeds. It still warns because this +/// node version starts epoch components only for worker [`DEFAULT_WORKER_ID`]: header payloads +/// keyed to higher worker ids validate, but nothing local produces them yet. +fn check_committee_worker_count( + epoch: Epoch, + num_workers: NonZeroUsize, + configured_workers: usize, +) -> eyre::Result<()> { + super::check_configured_worker_count(epoch, num_workers.get(), configured_workers)?; if num_workers.get() == 1 { return Ok(()); } @@ -1137,7 +1158,8 @@ fn check_committee_worker_count(epoch: Epoch, num_workers: NonZeroUsize) -> eyre epoch, num_workers, spawned_worker = DEFAULT_WORKER_ID, - "committee runs multiple workers but this node version spawns only worker {DEFAULT_WORKER_ID}: \ + "committee runs multiple workers but this node version starts epoch components only for \ + worker {DEFAULT_WORKER_ID}: \ ids >= 1 are accepted by header validation but not produced locally" ); Ok(()) @@ -1179,11 +1201,10 @@ mod tests { /// One worker is representable in both committee layouts, so entry never blocks on it. #[test] - fn single_worker_epoch_entry_is_always_allowed() { - for epoch in [0, 1, 407, u32::MAX] { - check_committee_worker_count(epoch, NonZeroUsize::MIN) - .expect("one worker is representable at every epoch"); - } + fn single_worker_epoch_entry_is_always_allowed() -> eyre::Result<()> { + [0, 1, 407, u32::MAX] + .into_iter() + .try_for_each(|epoch| check_committee_worker_count(epoch, NonZeroUsize::MIN, 1)) } /// Pre-fork the legacy committee layout cannot carry a worker count, so entry halts rather than @@ -1195,18 +1216,42 @@ mod tests { /// `OnceLock` is process-wide and the whole test binary shares one process. #[cfg(feature = "adiri")] #[test] - fn pre_fork_epoch_entry_rejects_multiple_workers() { - let err = check_committee_worker_count(0, NonZeroUsize::new(2).expect("2 is not 0")) - .expect_err("a pre-fork multi-worker committee cannot be encoded"); + fn pre_fork_epoch_entry_rejects_multiple_workers() -> eyre::Result<()> { + let count = NonZeroUsize::new(2).ok_or_else(|| eyre::eyre!("nonzero worker count"))?; + let err = check_committee_worker_count(0, count, 2) + .err() + .ok_or_else(|| eyre::eyre!("expected pre-fork multi-worker rejection"))?; assert!(err.to_string().contains("multi-workers fork is not active"), "{err}"); + Ok(()) } /// Default builds have the multi-worker layout active from genesis, so a count above one is - /// representable and entry proceeds (with a warning that this node still spawns one worker). + /// representable and entry proceeds (with a warning about worker-0-only epoch components). #[cfg(not(feature = "adiri"))] #[test] - fn post_fork_epoch_entry_allows_multiple_workers() { - check_committee_worker_count(0, NonZeroUsize::new(2).expect("2 is not 0")) - .expect("the post-fork layout holds a worker count"); + fn post_fork_epoch_entry_allows_multiple_workers() -> eyre::Result<()> { + let count = NonZeroUsize::new(2).ok_or_else(|| eyre::eyre!("nonzero worker count"))?; + check_committee_worker_count(0, count, 2) + } + + /// A single-worker chain must reject extra local swarms before its early return. + #[test] + fn single_worker_epoch_entry_rejects_extra_configured_workers() -> eyre::Result<()> { + let result = check_committee_worker_count(0, NonZeroUsize::MIN, 3); + assert!(result.is_err()); + let error = result.err().ok_or_else(|| eyre::eyre!("expected worker count mismatch"))?; + assert!(error.to_string().contains("chain-derived count for epoch 0 is 1")); + Ok(()) + } + + /// Governance cannot increase the committee count while the process retains fewer swarms. + #[test] + fn epoch_entry_rejects_changed_worker_count() -> eyre::Result<()> { + let count = NonZeroUsize::new(2).ok_or_else(|| eyre::eyre!("nonzero worker count"))?; + let result = check_committee_worker_count(7, count, 1); + assert!(result.is_err()); + let error = result.err().ok_or_else(|| eyre::eyre!("expected worker count mismatch"))?; + assert!(error.to_string().contains("chain-derived count for epoch 7 is 2")); + Ok(()) } } diff --git a/crates/storage/src/lib.rs b/crates/storage/src/lib.rs index cb2580a56..257e61555 100644 --- a/crates/storage/src/lib.rs +++ b/crates/storage/src/lib.rs @@ -65,10 +65,14 @@ const NODE_BATCHES_CACHE_CF: &str = "node_batches_cache"; const OUR_NODE_BATCHES_CACHE_CF: &str = "our_node_batches_cache"; const CONSENSUS_OUTPUT_CACHE_CF: &str = "consensus_output_cache"; -const KAD_RECORD_CF: &str = "kad_record"; -const KAD_PROVIDER_RECORD_CF: &str = "kad_provider_record"; -const KAD_WORKER_RECORD_CF: &str = "kad_worker_record"; -const KAD_WORKER_PROVIDER_RECORD_CF: &str = "kad_worker_provider_record"; +/// Discovery records with role and worker id in their row hashes. +const KAD_RECORD_CF: &str = "kad_record_v2"; +/// Provider rows with a separately decodable ownership key. +const KAD_PROVIDER_RECORD_CF: &str = "kad_provider_record_v2"; +/// Worker discovery records isolated by worker id. +const KAD_WORKER_RECORD_CF: &str = "kad_worker_record_v2"; +/// Worker provider rows isolated by worker id and ownership key. +const KAD_WORKER_PROVIDER_RECORD_CF: &str = "kad_worker_provider_record_v2"; macro_rules! tables { ( $($table:ident;$name:expr;$hint:expr;<$K:ty, $V:ty>),*) => { diff --git a/crates/telcoin-network-cli/src/keytool/generate.rs b/crates/telcoin-network-cli/src/keytool/generate.rs index 8609dc0ec..24d5f13e9 100644 --- a/crates/telcoin-network-cli/src/keytool/generate.rs +++ b/crates/telcoin-network-cli/src/keytool/generate.rs @@ -192,8 +192,8 @@ impl KeygenArgs { info!(target: "tn::generate_keys", primary=?node_info.p2p_info.primary.network_address, "updating primary external network address"); - // network keypair for workers (the key config holds worker 0's key) - let network_publickey = key_config.worker_network_public_key(); + // network keypair for workers (keytool still generates a single worker entry) + let network_publickey = key_config.worker_network_public_key(DEFAULT_WORKER_ID); let worker = node_info .p2p_info .worker_mut(DEFAULT_WORKER_ID) diff --git a/crates/test-utils-committee/src/authority.rs b/crates/test-utils-committee/src/authority.rs index edfab5f22..3d006bead 100644 --- a/crates/test-utils-committee/src/authority.rs +++ b/crates/test-utils-committee/src/authority.rs @@ -166,8 +166,10 @@ impl AuthorityFixture { // These key updates don't return errors... let _ = config.update_protocol_key(key_config.primary_public_key()); let _ = config.update_primary_network_key(key_config.primary_network_public_key()); - let _ = config - .update_worker_network_key(DEFAULT_WORKER_ID, key_config.worker_network_public_key()); + let _ = config.update_worker_network_key( + DEFAULT_WORKER_ID, + key_config.worker_network_public_key(DEFAULT_WORKER_ID), + ); let consensus_config = ConsensusConfig::new_with_committee_and_prior_epoch_record_for_test( config, diff --git a/crates/test-utils-committee/src/builder.rs b/crates/test-utils-committee/src/builder.rs index 460c8ea07..81f9935c5 100644 --- a/crates/test-utils-committee/src/builder.rs +++ b/crates/test-utils-committee/src/builder.rs @@ -9,7 +9,7 @@ use tn_config::{KeyConfig, NetworkConfig, Parameters}; use tn_types::{ get_available_udp_port, test_genesis, Address, Authority, AuthorityIdentifier, BlsKeypair, BootstrapServer, Committee, Database, Epoch, EpochDigest, Multiaddr, NetworkKeypair, P2pNode, - TimestampSec, DEFAULT_WORKER_PORT, + TimestampSec, DEFAULT_WORKER_ID, DEFAULT_WORKER_PORT, }; /// The committee builder for tests. @@ -61,7 +61,7 @@ where /// Set the number of workers every authority runs (defaults to one). /// /// Worker 0 uses the authority's [KeyConfig] worker network key; every further worker - /// gets a fresh network keypair, since [KeyConfig] holds a single worker key. + /// gets a fresh network keypair. pub fn number_of_workers(mut self, number_of_workers: NonZeroUsize) -> Self { self.number_of_workers = number_of_workers; self @@ -133,12 +133,12 @@ where DB: Database, F: Fn() -> DB, { + /// Build the committee and each authority's worker-zero fixture. pub fn build(mut self) -> CommitteeFixture { let committee_size = self.committee_size.get(); let network_config = self.network_config.unwrap_or_default(); let mut rng = StdRng::from_rng(&mut self.rng); - let mut committee_info = Vec::with_capacity(committee_size); #[allow(clippy::mutable_key_type)] let mut authorities = BTreeMap::new(); let mut bootstrap_servers = BTreeMap::new(); @@ -167,7 +167,7 @@ where let worker_nodes: Vec = (0..self.number_of_workers.get()) .map(|worker_id| { let key = if worker_id == 0 { - key_config.worker_network_public_key() + key_config.worker_network_public_key(DEFAULT_WORKER_ID) } else { NetworkKeypair::generate_ed25519().public().into() }; @@ -190,18 +190,19 @@ where (primary_keypair, key_config, authority.clone()), ); } - // Reset the authority ids so they are in sort order. Some tests require this. - for (i, (_, (primary_keypair, key_config, authority))) in authorities.iter_mut().enumerate() - { - let worker = WorkerFixture::generate(key_config.clone(), i as u16); - committee_info.push(( - primary_keypair.copy(), - key_config.clone(), - authority.clone(), - worker, - network_config.clone(), - )); - } + // Every authority fixture represents its own worker 0, independent of authority order. + let committee_info: Vec<_> = authorities + .values() + .map(|(primary_keypair, key_config, authority)| { + ( + primary_keypair.copy(), + key_config.clone(), + authority.clone(), + WorkerFixture::generate(key_config.clone(), DEFAULT_WORKER_ID), + network_config.clone(), + ) + }) + .collect(); // Make the committee so we can give it the AuthorityFixtures below. let committee = Committee::new_for_test( authorities.into_iter().map(|(k, (_, _, a))| (k, a)).collect(), diff --git a/crates/test-utils-committee/src/worker.rs b/crates/test-utils-committee/src/worker.rs index 269075730..1d7391970 100644 --- a/crates/test-utils-committee/src/worker.rs +++ b/crates/test-utils-committee/src/worker.rs @@ -9,15 +9,19 @@ use tn_types::{NetworkKeypair, WorkerId}; /// [WorkerFixture] holds keypairs and should not be used in production. #[derive(Debug)] pub struct WorkerFixture { + /// Key manager deriving this worker's network identity. key_config: KeyConfig, + /// Worker id within its authority, independent of the authority's committee position. pub id: WorkerId, } impl WorkerFixture { - pub fn keypair(&self) -> &NetworkKeypair { - self.key_config.worker_network_keypair() + /// The derived network keypair for this fixture's worker id. + pub fn keypair(&self) -> NetworkKeypair { + self.key_config.worker_network_keypair(self.id) } + /// Create a worker fixture with an id scoped to its authority. pub fn generate(key_config: KeyConfig, id: WorkerId) -> Self { Self { key_config, id } }