diff --git a/Cargo.lock b/Cargo.lock index f5dfa89f..6b75f7f4 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1981,7 +1981,6 @@ dependencies = [ "hex", "libssz", "libssz-types", - "rayon", "serde", "spawned-concurrency", "thiserror 2.0.18", diff --git a/bin/ethlambda/src/benchmark/keys.rs b/bin/ethlambda/src/benchmark/keys.rs index c116549d..e49d2457 100644 --- a/bin/ethlambda/src/benchmark/keys.rs +++ b/bin/ethlambda/src/benchmark/keys.rs @@ -11,7 +11,7 @@ use std::collections::HashMap; use std::path::Path; use std::time::Instant; -use ethlambda_blockchain::key_manager::{KeyManager, ValidatorKeyPair}; +use ethlambda_blockchain::key_manager::{KeyManager, KeyRole, ValidatorKeyPair}; use ethlambda_crypto::signature::ValidatorSecretKey; use ethlambda_types::state::ValidatorPubkeyBytes; use eyre::WrapErr as _; @@ -19,22 +19,6 @@ use rayon::prelude::*; const PUBKEY_LEN: usize = size_of::(); -#[derive(Debug, Clone, Copy)] -#[repr(u64)] -enum Role { - Attestation = 0, - Proposal = 1, -} - -impl Role { - fn tag(self) -> &'static str { - match self { - Role::Attestation => "attestation", - Role::Proposal => "proposal", - } - } -} - struct Key { pubkey: ValidatorPubkeyBytes, secret: ValidatorSecretKey, @@ -75,8 +59,8 @@ impl KeySet { // Every key is independent and deterministic in (seed, index, role), so // they are produced in parallel; `collect` keeps the job order. let start = Instant::now(); - let jobs: Vec<(u64, Role)> = (0..num_validators) - .flat_map(|index| [(index, Role::Attestation), (index, Role::Proposal)]) + let jobs: Vec<(u64, KeyRole)> = (0..num_validators) + .flat_map(|index| [(index, KeyRole::Attestation), (index, KeyRole::Proposal)]) .collect(); let keys: Vec = jobs .into_par_iter() @@ -118,7 +102,7 @@ impl KeySet { fn load_or_generate( seed: u64, index: u64, - role: Role, + role: KeyRole, num_active_epochs: usize, cache: Option<&Path>, ) -> eyre::Result { @@ -126,7 +110,7 @@ fn load_or_generate( dir.join(format!( "xmss-{}-seed{seed}-v{index}-{}-w{num_active_epochs}.bin", env!("ETHLAMBDA_LEANVM_REV"), - role.tag() + role.name() )) }); if let Some(file) = &file @@ -135,7 +119,13 @@ fn load_or_generate( return load_cached(file); } - let key_seed = seed ^ (index << 1 | role as u64).rotate_left(32); + // Spelled out rather than `role as u64`, so reordering `KeyRole` cannot + // change the keys a seed derives. + let role_bit = match role { + KeyRole::Attestation => 0, + KeyRole::Proposal => 1, + }; + let key_seed = seed ^ (index << 1 | role_bit).rotate_left(32); // leanVM seeds a key with 32 bytes; the run's `u64` fills the low end and // the rest stays zero, which keeps the derivation reproducible without // pretending to more entropy than the seed carries. @@ -198,13 +188,13 @@ mod tests { #[test] fn keys_are_deterministic_per_seed_and_distinct_per_role() { - let a = load_or_generate(7, 3, Role::Attestation, 2, None).unwrap(); - let b = load_or_generate(7, 3, Role::Attestation, 2, None).unwrap(); + let a = load_or_generate(7, 3, KeyRole::Attestation, 2, None).unwrap(); + let b = load_or_generate(7, 3, KeyRole::Attestation, 2, None).unwrap(); assert_eq!(a.pubkey, b.pubkey); assert_eq!(a.secret.to_bytes().unwrap(), b.secret.to_bytes().unwrap()); - let proposal = load_or_generate(7, 3, Role::Proposal, 2, None).unwrap(); + let proposal = load_or_generate(7, 3, KeyRole::Proposal, 2, None).unwrap(); assert_ne!(a.pubkey, proposal.pubkey); - let other_seed = load_or_generate(8, 3, Role::Attestation, 2, None).unwrap(); + let other_seed = load_or_generate(8, 3, KeyRole::Attestation, 2, None).unwrap(); assert_ne!(a.pubkey, other_seed.pubkey); } diff --git a/bin/ethlambda/src/main.rs b/bin/ethlambda/src/main.rs index 496d8dc9..e50defb1 100644 --- a/bin/ethlambda/src/main.rs +++ b/bin/ethlambda/src/main.rs @@ -38,7 +38,7 @@ use cli::NodeOptions; use command::Command; use ethlambda_blockchain::block_builder::ProposerConfig; -use ethlambda_blockchain::key_manager::ValidatorKeyPair; +use ethlambda_blockchain::key_manager::{KeyRole, ValidatorKeyPair}; use ethlambda_crypto::signature::ValidatorSecretKey; use ethlambda_network_api::{InitBlockChain, InitP2P, ToBlockChainToP2PRef, ToP2PToBlockChainRef}; use ethlambda_p2p::{ @@ -619,18 +619,12 @@ where Ok(pubkey) } -#[derive(Debug)] -enum ValidatorKeyRole { - Attestation, - Proposal, -} - /// Classify a privkey file as attestation or proposal based on the filename. /// /// Matches zeam's (`pkgs/cli/src/node.zig:540`) and lantern's /// (`client_keys.c:606`) routing, which lets all three clients share the /// `lean-quickstart` generator output unchanged. -fn classify_role(file: &Path) -> Result { +fn classify_role(file: &Path) -> Result { let name = file .file_name() .and_then(|n| n.to_str()) @@ -638,8 +632,8 @@ fn classify_role(file: &Path) -> Result { let is_attester = name.contains("attester"); let is_proposer = name.contains("proposer"); match (is_attester, is_proposer) { - (true, false) => Ok(ValidatorKeyRole::Attestation), - (false, true) => Ok(ValidatorKeyRole::Proposal), + (true, false) => Ok(KeyRole::Attestation), + (false, true) => Ok(KeyRole::Proposal), (false, false) => Err(format!( "filename '{name}' must contain 'attester' or 'proposer'" )), @@ -695,8 +689,8 @@ fn read_validator_keys( let path = resolve_path(&entry.privkey_file); let slots = grouped.entry(entry.index).or_default(); let target = match role { - ValidatorKeyRole::Attestation => &mut slots.attestation, - ValidatorKeyRole::Proposal => &mut slots.proposal, + KeyRole::Attestation => &mut slots.attestation, + KeyRole::Proposal => &mut slots.proposal, }; if target.is_some() { eyre::bail!("validator {}: duplicate {role:?} entry", entry.index); diff --git a/crates/blockchain/Cargo.toml b/crates/blockchain/Cargo.toml index 342ae0b0..b577ba33 100644 --- a/crates/blockchain/Cargo.toml +++ b/crates/blockchain/Cargo.toml @@ -27,7 +27,6 @@ spawned-concurrency.workspace = true tokio.workspace = true tokio-util = { version = "0.7", default-features = false } -rayon.workspace = true serde.workspace = true thiserror.workspace = true tracing.workspace = true diff --git a/crates/blockchain/src/key_manager.rs b/crates/blockchain/src/key_manager.rs index 26196ec4..ae6cf60f 100644 --- a/crates/blockchain/src/key_manager.rs +++ b/crates/blockchain/src/key_manager.rs @@ -1,4 +1,6 @@ use std::collections::HashMap; +use std::num::NonZeroUsize; +use std::thread; use std::time::Instant; use ethlambda_crypto::signature::{ValidatorSecretKey, ValidatorSignature}; @@ -38,6 +40,30 @@ pub struct KeyManager { keys: HashMap, } +/// Which of a validator's two XMSS keys: each signs one kind of message. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum KeyRole { + Attestation, + Proposal, +} + +impl KeyRole { + /// Lowercase name, as it appears in key file names. + pub fn name(self) -> &'static str { + match self { + Self::Attestation => "attestation", + Self::Proposal => "proposal", + } + } +} + +/// Fewest keys worth a helper thread of [`KeyManager::prepare_keys_for`]. +/// +/// leanVM rebuilds a subtree sequentially, so the only parallelism is across +/// keys. A batch this size still has a few rebuilds to spread its spawn cost +/// over, and a key set no larger than one batch warms without spawning. +const MIN_KEYS_PER_WARM_THREAD: usize = 4; + impl KeyManager { pub fn new(keys: HashMap) -> Self { Self { keys } @@ -48,21 +74,48 @@ impl KeyManager { self.keys.keys().copied().collect() } - /// Warms every validator's signing cache for `slot`. + /// Warms the signing caches the duties at `slot` sign with: every + /// attestation key, plus `proposer`'s proposal key when one of our + /// validators proposes at `slot`. /// /// Pure latency shifting: the key rebuilds the same bottom Merkle subtree /// inside `sign` on a miss, so this only moves that cost off the duty's - /// critical path. Called one slot ahead, so a miss here is not yet a - /// failure to sign. - pub fn prepare_keys_for(&self, slot: u32) { - for (validator_id, key_pair) in &self.keys { - let _ = prepare_key(&key_pair.attestation_key, slot).inspect_err( - |err| warn!(validator_id, slot, %err, "Failed to warm attestation key signing cache"), - ); - let _ = prepare_key(&key_pair.proposal_key, slot).inspect_err( - |err| warn!(validator_id, slot, %err, "Failed to warm proposal key signing cache"), - ); - } + /// critical path. A key caches a single subtree, though, so warming it + /// evicts the subtree an earlier slot signs with whenever the two slots + /// straddle a subtree boundary. Call this only once nothing is left to sign + /// before `slot`, or that signature rebuilds the evicted subtree on its own + /// critical path. + /// + /// The other proposal keys are left cold: nothing signs with them before + /// their own turn to propose, and warming all of them rebuilds one subtree + /// per validator at every boundary. + /// + /// Blocks until every key is warm, with the rebuilds spread over at most + /// one thread per core (see [`for_each_batched`]). Each key knows which + /// subtree it holds and rebuilds only on a miss, so a repeat call for a + /// slot already warmed costs a lock per key and the thread spawns. + pub fn prepare_keys_for(&self, slot: u32, proposer: Option) { + let keys = self.keys_to_warm(proposer); + let start = Instant::now(); + let max_threads = thread::available_parallelism().map_or(1, NonZeroUsize::get); + for_each_batched(&keys, max_threads, |&(validator_id, role, key)| { + prepare_key(validator_id, role, key, slot) + }); + trace!(slot, keys = keys.len(), elapsed = ?start.elapsed(), "Warmed XMSS signing caches"); + } + + /// The keys [`Self::prepare_keys_for`] warms: every attestation key, then + /// `proposer`'s proposal key if that validator is ours. + fn keys_to_warm(&self, proposer: Option) -> Vec<(u64, KeyRole, &ValidatorSecretKey)> { + let attestation_keys = self + .keys + .iter() + .map(|(&id, pair)| (id, KeyRole::Attestation, &pair.attestation_key)); + let proposal_key = proposer.and_then(|id| { + let pair = self.keys.get(&id)?; + Some((id, KeyRole::Proposal, &pair.proposal_key)) + }); + attestation_keys.chain(proposal_key).collect() } /// Signs an attestation using the validator's attestation key. @@ -162,16 +215,56 @@ fn signable_at( } /// Warm one key's signing cache, timing the miss that rebuilds a subtree. -fn prepare_key(key: &ValidatorSecretKey, slot: u32) -> Result<(), KeyManagerError> { +/// +/// A failure only warns: `sign` rebuilds the subtree itself on a miss. +fn prepare_key(validator_id: u64, role: KeyRole, key: &ValidatorSecretKey, slot: u32) { let start = Instant::now(); - key.prepare(slot) - .map_err(|err| KeyManagerError::SigningError(err.to_string()))?; - trace!(slot, elapsed = ?start.elapsed(), "Warmed XMSS signing cache"); - Ok(()) + let result = key.prepare(slot); + let elapsed = start.elapsed(); + let _ = result + .inspect(|()| trace!(slot, validator_id, ?role, ?elapsed, "Prepared XMSS key")) + .inspect_err( + |err| warn!(slot, validator_id, ?role, %err, "Failed to warm XMSS signing cache"), + ); +} + +/// Runs `f` on every item, in batches of at least [`MIN_KEYS_PER_WARM_THREAD`] +/// items spread over at most `max_threads` threads, and returns once all are +/// done. +/// +/// The calling thread takes the first batch, so a single batch spawns nothing. +/// A batch whose thread cannot be spawned also runs on the calling thread, +/// which costs latency instead of panicking the caller. +fn for_each_batched(items: &[T], max_threads: usize, f: impl Fn(&T) + Sync) { + let batch_len = items + .len() + .div_ceil(max_threads.max(1)) + .max(MIN_KEYS_PER_WARM_THREAD); + let mut batches = items.chunks(batch_len); + let Some(first) = batches.next() else { + return; + }; + let f = &f; + thread::scope(|scope| { + let mut inline = vec![first]; + for batch in batches { + let spawned = thread::Builder::new() + .name("xmss-warm".into()) + .spawn_scoped(scope, move || batch.iter().for_each(f)); + if let Err(err) = spawned { + warn!(%err, "Failed to spawn an XMSS warm thread, warming its keys inline"); + inline.push(batch); + } + } + inline.into_iter().flatten().for_each(f); + }); } #[cfg(test)] mod tests { + use std::collections::HashSet; + use std::sync::Mutex; + use super::*; #[test] @@ -206,4 +299,112 @@ mod tests { Err(KeyManagerError::ValidatorKeyNotFound(123)) )); } + + /// A key manager for validators `0..count`, each key covering slots 0 and + /// 1 only, which keeps generation cheap. + fn tiny_key_manager(count: u64) -> KeyManager { + // Seeded from the whole `(id, role)` pair, so no two keys share a seed. + let key = |id: u64, role: KeyRole| { + let mut seed = [0; 32]; + seed[..8].copy_from_slice(&id.to_le_bytes()); + seed[8] = role as u8; + ValidatorSecretKey::generate_from_seed(seed, 0..=1).unwrap() + }; + let keys = (0..count) + .map(|id| { + let pair = ValidatorKeyPair { + attestation_key: key(id, KeyRole::Attestation), + proposal_key: key(id, KeyRole::Proposal), + }; + (id, pair) + }) + .collect(); + KeyManager::new(keys) + } + + fn warmed(key_manager: &KeyManager, proposer: Option) -> Vec<(u64, KeyRole)> { + let mut entries: Vec<_> = key_manager + .keys_to_warm(proposer) + .into_iter() + .map(|(id, role, _)| (id, role)) + .collect(); + entries.sort_by_key(|&(id, role)| (role == KeyRole::Proposal, id)); + entries + } + + #[test] + fn keys_to_warm_takes_every_attestation_key_and_only_the_proposers_proposal_key() { + let key_manager = tiny_key_manager(3); + let attestation_keys: Vec<_> = (0..3).map(|id| (id, KeyRole::Attestation)).collect(); + + assert_eq!(warmed(&key_manager, None), attestation_keys); + + let mut with_proposer = attestation_keys.clone(); + with_proposer.push((1, KeyRole::Proposal)); + assert_eq!(warmed(&key_manager, Some(1)), with_proposer); + + // A proposer that is not one of ours adds nothing. + assert_eq!(warmed(&key_manager, Some(99)), attestation_keys); + } + + /// Runs [`for_each_batched`] over `0..len` and returns, per item, the + /// threads that visited it. + fn visits(len: usize, max_threads: usize) -> Vec> { + let items: Vec = (0..len).collect(); + let visits = Mutex::new(vec![Vec::new(); len]); + for_each_batched(&items, max_threads, |&item| { + visits.lock().unwrap()[item].push(thread::current().id()); + }); + visits.into_inner().unwrap() + } + + #[test] + fn for_each_batched_visits_every_item_once() { + let batch = MIN_KEYS_PER_WARM_THREAD; + for len in [0, 1, batch, batch + 1, 3 * batch + 1, 40] { + for max_threads in [0, 1, 3, 64] { + let visits = visits(len, max_threads); + assert!( + visits.iter().all(|threads| threads.len() == 1), + "len {len}, max_threads {max_threads}: {visits:?}" + ); + } + } + } + + #[test] + fn for_each_batched_caps_threads_and_keeps_one_batch_inline() { + let batch = MIN_KEYS_PER_WARM_THREAD; + let threads_used = |len, max_threads| -> HashSet<_> { + visits(len, max_threads).into_iter().flatten().collect() + }; + + // One batch's worth runs on the calling thread alone. + let caller = HashSet::from([thread::current().id()]); + assert_eq!(threads_used(batch, 64), caller); + // More batches than threads: the batches grow instead. + assert_eq!(threads_used(40, 3).len(), 3); + // Fewer batches than threads: one thread per batch. + assert_eq!(threads_used(3 * batch, 64).len(), 3); + } + + #[test] + fn prepare_keys_for_leaves_every_key_signable() { + // More keys than one batch holds, so the warm spans several threads. + let count = 2 * MIN_KEYS_PER_WARM_THREAD as u64 + 1; + let mut key_manager = tiny_key_manager(count); + + key_manager.prepare_keys_for(1, Some(0)); + key_manager.prepare_keys_for(1, Some(0)); + // A slot outside every key's range only warns. + key_manager.prepare_keys_for(7, None); + + let message = H256::default(); + for id in 0..count { + key_manager + .sign_with_attestation_key(id, 1, &message) + .unwrap(); + } + key_manager.sign_block_root(0, 1, &message).unwrap(); + } } diff --git a/crates/blockchain/src/lib.rs b/crates/blockchain/src/lib.rs index 039f3021..9292d5c4 100644 --- a/crates/blockchain/src/lib.rs +++ b/crates/blockchain/src/lib.rs @@ -185,15 +185,7 @@ impl BlockChain { let genesis_time = time_config.genesis_time; let key_manager = key_manager::KeyManager::new(validator_keys); - // Warm the XMSS signing caches for the current slot before the first tick. - // store.time() doesn't work here: after an offline gap it lags wall-clock by - // exactly the gap the first duty will be at - let now_ms = unix_now_ms(); - let current_slot = (now_ms.saturating_sub(time_config.genesis_time_ms()) - / time_config.milliseconds_per_slot) as u32; - key_manager.prepare_keys_for(current_slot); - - let handle = BlockChainServer { + let server = BlockChainServer { store, p2p: None, key_manager, @@ -211,8 +203,40 @@ impl BlockChain { sync_status: SyncStatusTracker::new(gate_duties), sync_status_controller, events, + }; + + // Warm the XMSS signing caches for the next duties before the first + // tick, which fires right away and runs the current interval's duty. + // store.time() doesn't work here: after an offline gap it lags + // wall-clock by exactly the gap the first duty will be at. + let ms_since_genesis = unix_now_ms().saturating_sub(time_config.genesis_time_ms()); + let current_slot = ms_since_genesis / time_config.milliseconds_per_slot; + match SlotInterval::from_ms_since_genesis(ms_since_genesis, &time_config) { + // The first tick still attests at the current slot. No proposal + // key: the current slot's block was due at the previous slot's + // interval 4, before we started, and the interval-1 tick warms the + // next slot's. + SlotInterval::BlockPublication | SlotInterval::AttestationProduction => { + server + .key_manager + .prepare_keys_for(current_slot as u32, None); + } + // This slot's attestations are behind us, so the next signatures + // are the next slot's block, built at this slot's interval 4, and + // that slot's attestations. + SlotInterval::Aggregation + | SlotInterval::SafeTargetUpdate + | SlotInterval::EndOfSlot => { + let num_validators = server.store.head_state().validators.len() as u64; + let next_slot = current_slot + 1; + let proposer = server.get_our_proposer(next_slot, num_validators); + server + .key_manager + .prepare_keys_for(next_slot as u32, proposer); + } } - .start(); + + let handle = server.start(); let time_until_genesis = (SystemTime::UNIX_EPOCH + Duration::from_secs(genesis_time)) .duration_since(SystemTime::now()) .unwrap_or_default(); @@ -343,8 +367,10 @@ impl BlockChainServer { } // Fail fast: a state with zero validators is invalid and would cause - // panics in proposer selection and attestation processing. - if self.store.head_state().validators.is_empty() { + // panics in proposer selection and attestation processing. Read once + // per tick, since `head_state` clones the whole state. + let num_validators = self.store.head_state().validators.len() as u64; + if num_validators == 0 { error!("Head state has no validators, skipping tick"); return; } @@ -382,7 +408,7 @@ impl BlockChainServer { // Whether one of our validators proposes this slot. Drives the store's // interval-0 attestation acceptance. let is_proposer = (interval == SlotInterval::BlockPublication && slot > 0) - .then(|| self.get_our_proposer(slot)) + .then(|| self.get_our_proposer(slot, num_validators)) .flatten() .is_some(); @@ -445,6 +471,19 @@ impl BlockChainServer { EarlyAggregationCheck, ); } + + // Warm the XMSS signing caches for the next slot so the signing + // paths don't have to, now that this slot's attestations are + // signed: a key caches one bottom subtree, so warming any earlier + // evicts the subtree this slot's attestation signs with whenever + // the two slots straddle a subtree boundary. This lands before + // interval 4 signs the next slot's block. A skipped interval-1 + // tick costs only latency, since `sign` rebuilds the subtree + // itself on a miss. + let next_slot = slot + 1; + let proposer = self.get_our_proposer(next_slot, num_validators); + self.key_manager + .prepare_keys_for(next_slot as u32, proposer); } // ==== interval 2 ==== @@ -482,7 +521,7 @@ impl BlockChainServer { SlotInterval::EndOfSlot => { let next_slot = slot + 1; let next_proposer = self - .get_our_proposer(next_slot) + .get_our_proposer(next_slot, num_validators) .filter(|_| self.sync_status.duties_allowed()); if let Some(validator_id) = next_proposer { @@ -495,9 +534,6 @@ impl BlockChainServer { metrics::update_safe_target_slot(self.store.safe_target_slot()); // Update head slot metric (head may change when attestations are promoted at intervals 0/4) metrics::update_head_slot(self.store.head_slot()); - - // Warm the XMSS signing caches for the next slot so the signing paths don't have to - self.key_manager.prepare_keys_for((slot + 1) as u32); } /// Kick off a committee-signature aggregation session: @@ -534,8 +570,9 @@ impl BlockChainServer { // Limit ourselves to a single round of aggregation if we propose next round. // This buys us time to build the block before the next slot's interval-0 tick. + let num_validators = self.store.head_state().validators.len() as u64; let next_proposer = self - .get_our_proposer(slot + 1) + .get_our_proposer(slot + 1, num_validators) .filter(|_| self.sync_status.duties_allowed()); let max_jobs = if next_proposer.is_some() { 1 @@ -685,10 +722,7 @@ impl BlockChainServer { } /// Returns the validator ID if any of our validators is the proposer for this slot. - fn get_our_proposer(&self, slot: u64) -> Option { - let head_state = self.store.head_state(); - let num_validators = head_state.validators.len() as u64; - + fn get_our_proposer(&self, slot: u64, num_validators: u64) -> Option { self.key_manager .validator_ids() .into_iter() diff --git a/crates/common/crypto/src/signature.rs b/crates/common/crypto/src/signature.rs index fd8a8ae3..4b7b8b68 100644 --- a/crates/common/crypto/src/signature.rs +++ b/crates/common/crypto/src/signature.rs @@ -212,8 +212,9 @@ impl ValidatorSecretKey { /// /// The key holds one cached subtree, so this is worth calling only for the /// slot about to be signed, and it is pure latency shifting: a miss inside - /// `sign` rebuilds the same subtree. Errors only when `slot` is outside - /// [`Self::signable_slots`]. + /// `sign` rebuilds the same subtree. A slot the cached subtree already + /// covers returns without rebuilding, so repeating a call is cheap. Errors + /// only when `slot` is outside [`Self::signable_slots`]. pub fn prepare(&self, slot: u32) -> Result<(), XmssSignError> { self.inner.prepare(slot) }