From c150c75ff4e9e21c3a36210b93d90f1b30f1f7df Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Tom=C3=A1s=20Gr=C3=BCner?= <47506558+MegaRedHand@users.noreply.github.com> Date: Fri, 2 Oct 2026 11:42:18 -0300 Subject: [PATCH 1/2] perf(beacon): weigh fork-choice votes from a justified-balances snapshot Since the validator registry moved onto a persistent tree, every `state.validator(i)` in the fork-choice vote loop is a tree descent, and `get_proposer_score` rescans the registry through `get_total_active_balance` (an active-index Vec, then one descent per active index). Both run on every `get_head`, so the head recompute went from flat-array reads to millions of descents per call at mainnet scale. Flatten the justified checkpoint state once into `JustifiedBalances`: a per-validator weight (zero when inactive or slashed) and the spec-exact total active balance (slashed-but-active validators count, floored at one increment), computed in the same in-order pass. It is keyed by the justified checkpoint itself, so every writer of that checkpoint is covered without a hook, and its only source is `checkpoint_state(justified)`, the state `get_weight` reads. `get_weight` stays spec-literal and is the oracle: the fork-choice fixture runner now compares `compute_weights` with it for every block of the filtered tree at each `checks` step. `calculate_committee_fraction` is split so `is_head_weak`, `is_parent_strong` and the boost share the snapshot's total. Adds lean_beacon_justified_balances_lookups_total and lean_beacon_justified_balances_build_seconds. --- .../src/beacon/fork_choice.rs | 297 ++++++++++++++++-- .../state_transition/src/metrics.rs | 34 ++ .../tests/beacon_spec/fork_choice.rs | 35 +++ crates/common/types/src/beacon/fork_choice.rs | 69 +++- crates/storage/src/store.rs | 32 +- docs/beacon_stf.md | 15 + docs/metrics.md | 19 ++ 7 files changed, 468 insertions(+), 33 deletions(-) diff --git a/crates/blockchain/state_transition/src/beacon/fork_choice.rs b/crates/blockchain/state_transition/src/beacon/fork_choice.rs index 5560782b..50d136d2 100644 --- a/crates/blockchain/state_transition/src/beacon/fork_choice.rs +++ b/crates/blockchain/state_transition/src/beacon/fork_choice.rs @@ -185,6 +185,7 @@ use crate::beacon::primitives::{ ValidatorIndex, }; use crate::beacon::stf; +use crate::metrics; // --------------------------------------------------------------------------- // LatestMessage, PowBlock, PayloadStatusV1 @@ -195,7 +196,7 @@ use crate::beacon::stf; // Re-exported at the paths they had when they were defined here, so [`Store`] // and its callers are unchanged. pub use ethlambda_types::beacon::fork_choice::{ - LatestMessage, PayloadStatusEnum, PayloadStatusV1, PowBlock, + JustifiedBalances, LatestMessage, PayloadStatusEnum, PayloadStatusV1, PowBlock, }; // --------------------------------------------------------------------------- @@ -1076,8 +1077,18 @@ pub fn get_ancestor(index: &HashMap, root: Root, slot: Slot) /// values it is called with are themselves already expressed on a 0-100 /// scale. pub fn calculate_committee_fraction(state: &BeaconState, committee_percent: u64) -> Result { - let committee_weight = get_total_active_balance(state)? / preset::SLOTS_PER_EPOCH; - Ok(committee_weight.saturating_mul(committee_percent) / 100) + Ok(committee_fraction( + get_total_active_balance(state)?, + committee_percent, + )) +} + +/// The arithmetic of [`calculate_committee_fraction`], for a caller that +/// already holds the total active balance (the [`JustifiedBalances`] snapshot's +/// own). +pub fn committee_fraction(total_active_balance: Gwei, committee_percent: u64) -> Gwei { + let committee_weight = total_active_balance / preset::SLOTS_PER_EPOCH; + committee_weight.saturating_mul(committee_percent) / 100 } /// The checkpoint block for `epoch`, on `root`'s chain: the ancestor of `root` @@ -1098,10 +1109,67 @@ pub fn get_checkpoint_block( /// See [`calculate_committee_fraction`] for why this divides by a bare /// `100` rather than [`constants::BASIS_POINTS`]. pub fn get_proposer_score(store: &Store, config: &Config) -> Result { - let justified_checkpoint = store.beacon_justified_checkpoint(); - let justified_state = checkpoint_state(store, &justified_checkpoint, config)?; - let committee_weight = get_total_active_balance(&justified_state)? / preset::SLOTS_PER_EPOCH; - Ok(committee_weight.saturating_mul(config.proposer_score_boost) / 100) + let balances = justified_balances(store, config)?; + Ok(committee_fraction( + balances.total_active_balance(), + config.proposer_score_boost, + )) +} + +/// The flat balances of the store's justified checkpoint state, built on first +/// use after the checkpoint moves and shared until it moves again. +/// +/// Keyed by the checkpoint itself, which every writer of the justified +/// checkpoint (`update_checkpoints` from three handlers, store construction, +/// restart) already changes, so none of them needs a hook. The only source is +/// [`checkpoint_state`], the state `get_weight` reads, so the snapshot is +/// exactly what the specification's per-root definition sees; a later +/// post-state of the same epoch would not be, since `slashed` can differ. +/// A failed `checkpoint_state` fails the call, as it does for every other +/// fork-choice reader, and the head stays where it is. +pub fn justified_balances(store: &Store, config: &Config) -> Result> { + let checkpoint = store.beacon_justified_checkpoint(); + if let Some(balances) = store.justified_balances(&checkpoint) { + metrics::inc_justified_balances_lookups("hit"); + return Ok(balances); + } + metrics::inc_justified_balances_lookups("miss"); + let state = checkpoint_state(store, &checkpoint, config)?; + let balances = Arc::new(build_justified_balances(checkpoint, &state)); + store.set_justified_balances(Arc::clone(&balances)); + Ok(balances) +} + +/// One in-order pass over `state`'s registry: each validator's vote weight +/// (zero if inactive or slashed) and, in the same pass, the total active +/// balance with exactly the semantics of `get_total_active_balance` +/// (slashed-but-active validators count, saturating, floored at one increment). +fn build_justified_balances(checkpoint: Checkpoint, state: &BeaconState) -> JustifiedBalances { + let _timing = metrics::time_justified_balances_build(); + // Activity at the state's own epoch, as `get_weight` reads it; not asserted + // equal to `checkpoint.epoch`, which the unit-test stores do not keep. + let epoch = get_current_epoch(state); + let mut total: Gwei = 0; + let balances = state + .validators() + .iter() + .map(|validator| { + if !is_active_validator(validator, epoch) { + return 0; + } + total = total.saturating_add(validator.effective_balance); + if validator.slashed { + 0 + } else { + validator.effective_balance + } + }) + .collect(); + JustifiedBalances::new( + checkpoint, + balances, + total.max(preset::EFFECTIVE_BALANCE_INCREMENT), + ) } /// The LMD GHOST weight of `root`: the effective balance of every @@ -1191,8 +1259,7 @@ pub fn compute_weights( config: &Config, ) -> Result> { let justified_checkpoint = store.beacon_justified_checkpoint(); - let state = checkpoint_state(store, &justified_checkpoint, config)?; - let current_epoch = get_current_epoch(&state); + let balances = justified_balances(store, config)?; // Keyed on the voted block itself; the fold below turns these into subtree // totals in place. @@ -1201,19 +1268,15 @@ pub fn compute_weights( // `for_each_non_equivocating_latest_message` for why asking it per voter // from in here would deadlock. store.for_each_non_equivocating_latest_message(|validator_index, message| { - // Not `get_active_validator_indices`: that allocates the whole active - // set (~2 million entries on mainnet) to answer a membership question, - // and an index past this state's registry is a validator that did not - // exist yet at the justified checkpoint, which is a skip rather than an - // error. - let Ok(validator) = state.validator(validator_index) else { - return; - }; - if validator.slashed || !is_active_validator(validator, current_epoch) { + // An index past the snapshot is a validator that did not exist at the + // justified checkpoint, and an inactive or slashed one is zero in it: + // both are a skip rather than an error. + let balance = balances.get(validator_index); + if balance == 0 { return; } let entry = weights.entry(message.root).or_default(); - *entry = entry.saturating_add(validator.effective_balance); + *entry = entry.saturating_add(balance); }); // Highest slot first: see above for why that is a topological order. @@ -1245,7 +1308,8 @@ pub fn compute_weights( let justified_slot = index .get(&justified_checkpoint.root) .map_or(0, |(slot, _)| *slot); - let proposer_score = get_proposer_score(store, config)?; + let proposer_score = + committee_fraction(balances.total_active_balance(), config.proposer_score_boost); let mut cursor = boost_root; while let Some((slot, parent_root)) = index.get(&cursor).copied() { let entry = weights.entry(cursor).or_default(); @@ -1605,20 +1669,22 @@ pub fn is_proposing_on_time(store: &Store, config: &Config) -> bool { /// proposer's own boost, i.e. reorging it out would not be fighting an /// already-decisive lead. pub fn is_head_weak(store: &Store, head_root: Root, config: &Config) -> Result { - let justified_checkpoint = store.beacon_justified_checkpoint(); - let justified_state = checkpoint_state(store, &justified_checkpoint, config)?; - let reorg_threshold = - calculate_committee_fraction(&justified_state, config.reorg_head_weight_threshold)?; + // The snapshot is built from `checkpoint_state(justified)`, the state this + // read from before, so its total is the same number. + let reorg_threshold = committee_fraction( + justified_balances(store, config)?.total_active_balance(), + config.reorg_head_weight_threshold, + ); Ok(get_weight(store, &store.block_index(), head_root, config)? < reorg_threshold) } /// Whether `parent_root` already has enough votes of its own that the missing /// votes are assigned to it rather than being hoarded elsewhere. pub fn is_parent_strong(store: &Store, parent_root: Root, config: &Config) -> Result { - let justified_checkpoint = store.beacon_justified_checkpoint(); - let justified_state = checkpoint_state(store, &justified_checkpoint, config)?; - let parent_threshold = - calculate_committee_fraction(&justified_state, config.reorg_parent_weight_threshold)?; + let parent_threshold = committee_fraction( + justified_balances(store, config)?.total_active_balance(), + config.reorg_parent_weight_threshold, + ); Ok(get_weight(store, &store.block_index(), parent_root, config)? > parent_threshold) } @@ -2746,7 +2812,12 @@ mod tests { /// [`anchor_pair`] over a registry of `count` validators, for the weight /// tests, which name a voter per validator index. fn anchor_pair_with(count: usize) -> (BeaconState, SignedBeaconBlock) { - let mut state = test_state::with_validators(count); + anchor_pair_from(test_state::with_validators(count)) + } + + /// [`anchor_pair`] over a caller-built `state`, for tests that need a + /// registry with slashed or inactive validators in it. + fn anchor_pair_from(mut state: BeaconState) -> (BeaconState, SignedBeaconBlock) { let parent_root = state.latest_block_header().parent_root; let mut signed = block(state.slot(), parent_root); @@ -2795,7 +2866,12 @@ mod tests { /// The anchor is what `checkpoint_state` resolves the justified checkpoint /// to, which is the one thing both weight functions need from a real store. fn anchored_store(count: usize) -> (Store, Root, Slot) { - let (anchor_state, anchor_block) = anchor_pair_with(count); + anchored_store_from(test_state::with_validators(count)) + } + + /// [`anchored_store`] over a caller-built state. + fn anchored_store_from(state: BeaconState) -> (Store, Root, Slot) { + let (anchor_state, anchor_block) = anchor_pair_from(state); let anchor_slot = anchor_state.slot(); let anchor_root = anchor_block.message_hash_tree_root(); let store = get_forkchoice_store( @@ -2808,13 +2884,30 @@ mod tests { (store, anchor_root, anchor_slot) } + /// A `count`-validator state (`count >= 9`) whose registry has every kind + /// of validator the weight paths must tell apart, at the state's own epoch: + /// validator 6 is slashed but active, 7 activates one epoch later, 8 exited + /// at this epoch, and the rest are plain active ones. + fn state_with_slashed_and_inactive_validators(count: usize) -> BeaconState { + let mut state = test_state::with_validators(count); + let epoch = get_current_epoch(&state); + state.validator_mut(6).expect("registered").slashed = true; + state.validator_mut(7).expect("registered").activation_epoch = epoch + 1; + state.validator_mut(8).expect("registered").exit_epoch = epoch; + state.apply_pending_mutations(); + state + } + /// `compute_weights` is the specification's `get_weight` for every root at /// once, so the two have to agree root by root: over a fork, over voters /// spread across both branches, and with the proposer boost applied. #[test] fn the_single_pass_weights_match_the_specifications_per_root_weight() { let config = Config::active(); - let (mut store, anchor_root, anchor_slot) = anchored_store(8); + // Validators 6 (slashed), 7 (not yet active) and 8 (exited) vote below + // and must weigh nothing; validator 40 is past the registry. + let (mut store, anchor_root, anchor_slot) = + anchored_store_from(state_with_slashed_and_inactive_validators(12)); // anchor -> a -> {b, c}: a fork whose two leaves split the vote, so a // wrong fold shows up as a leaf carrying its sibling's balance. @@ -2849,6 +2942,9 @@ mod tests { }, ); store.insert_equivocating_index(5); + for (validator_index, root) in [(6, c_root), (7, c_root), (8, c_root), (40, c_root)] { + store.set_latest_message(validator_index, LatestMessage { epoch: 0, root }); + } store.set_proposer_boost_root(b_root); let weights = compute_weights(&store, &index, &config).expect("the anchor state is there"); @@ -2864,6 +2960,145 @@ mod tests { weights[&b_root] > weights[&c_root], "three voters and the boost must outweigh one voter" ); + assert_eq!( + weights[&c_root], + preset::MAX_EFFECTIVE_BALANCE, + "only validator 3 counts on c: the slashed, inactive, exited and unknown voters weigh nothing" + ); + assert_eq!( + weights[&b_root], + 3 * preset::MAX_EFFECTIVE_BALANCE + + get_proposer_score(&store, &config).expect("the anchor state is there"), + ); + } + + /// The boost's committee weight divides the total active balance, which + /// counts a slashed validator that is still active and leaves out one that + /// is not active yet or already exited: the snapshot's per-vote zeros for + /// the first kind must not leak into the total. + #[test] + fn the_proposer_score_counts_a_slashed_but_active_validator() { + let config = Config::active(); + let (store, _anchor_root, _anchor_slot) = + anchored_store_from(state_with_slashed_and_inactive_validators(12)); + + // 12 validators, 2 of them (7 and 8) inactive; the slashed one stays. + let expected = committee_fraction( + 10 * preset::MAX_EFFECTIVE_BALANCE, + config.proposer_score_boost, + ); + assert_eq!( + get_proposer_score(&store, &config).expect("the anchor state is there"), + expected + ); + assert_ne!( + expected, + committee_fraction( + 9 * preset::MAX_EFFECTIVE_BALANCE, + config.proposer_score_boost + ), + "a total that dropped the slashed validator would differ" + ); + // And it is the state's own definition, not a second one. + let state = checkpoint_state(&store, &store.beacon_justified_checkpoint(), &config) + .expect("the anchor state is there"); + assert_eq!( + expected, + calculate_committee_fraction(&state, config.proposer_score_boost).expect("total"), + ); + } + + #[test] + fn building_the_snapshot_handles_activation_exit_slashing_and_an_empty_set() { + let mut state = test_state::with_validators(6); + let epoch = get_current_epoch(&state); + state.validator_mut(1).expect("registered").activation_epoch = epoch; + state.validator_mut(2).expect("registered").activation_epoch = epoch + 1; + state.validator_mut(3).expect("registered").exit_epoch = epoch; + state.validator_mut(4).expect("registered").slashed = true; + state.apply_pending_mutations(); + let checkpoint = Checkpoint { + epoch, + root: Root::repeat_byte(1), + }; + + let snapshot = build_justified_balances(checkpoint, &state); + let full = preset::MAX_EFFECTIVE_BALANCE; + // Activation at the epoch counts; exit at the epoch does not; a slashed + // but active validator weighs zero as a voter and counts in the total. + assert_eq!( + [0, 1, 2, 3, 4, 5].map(|index| snapshot.get(index)), + [full, full, 0, 0, 0, full] + ); + assert_eq!(snapshot.total_active_balance(), 4 * full); + assert_eq!( + snapshot.total_active_balance(), + get_total_active_balance(&state).expect("total") + ); + assert_eq!(snapshot.get(6), 0, "past the registry reads zero"); + assert_eq!(snapshot.checkpoint(), checkpoint); + + // No active validator at all: the total is floored like the spec's. + for index in 0..6 { + state.validator_mut(index).expect("registered").exit_epoch = epoch; + } + state.apply_pending_mutations(); + let empty = build_justified_balances(checkpoint, &state); + assert_eq!( + empty.total_active_balance(), + preset::EFFECTIVE_BALANCE_INCREMENT + ); + assert_eq!( + empty.total_active_balance(), + get_total_active_balance(&state).expect("total") + ); + assert_eq!(empty.get(0), 0); + } + + #[test] + fn a_new_justified_checkpoint_rebuilds_the_snapshot() { + let config = Config::active(); + let (mut store, anchor_root, _anchor_slot) = anchored_store(8); + + let first = justified_balances(&store, &config).expect("the anchor state is there"); + let again = justified_balances(&store, &config).expect("cached"); + assert!( + Arc::ptr_eq(&first, &again), + "same checkpoint, same snapshot" + ); + assert_eq!(first.get(0), preset::MAX_EFFECTIVE_BALANCE); + + // A later justified checkpoint over a state with different balances, + // planted where `checkpoint_state` finds it. + let next = Checkpoint { + epoch: first.checkpoint().epoch + 1, + root: Root::repeat_byte(0x77), + }; + let mut state = test_state::with_validators(8); + *state.slot_mut() = compute_start_slot_at_epoch(next.epoch); + state + .validator_mut(0) + .expect("registered") + .effective_balance = 5; + state.validator_mut(1).expect("registered").slashed = true; + state.apply_pending_mutations(); + store.cache_state( + CacheKey::CheckpointState { + epoch: next.epoch, + root: next.root, + }, + Arc::new(state), + ); + let finalized = store.beacon_finalized_checkpoint(); + update_checkpoints(&mut store, next, finalized); + assert_ne!(store.beacon_justified_checkpoint().root, anchor_root); + + let rebuilt = justified_balances(&store, &config).expect("planted state"); + assert!(!Arc::ptr_eq(&first, &rebuilt)); + assert_eq!(rebuilt.checkpoint(), next); + assert_eq!(rebuilt.get(0), 5); + assert_eq!(rebuilt.get(1), 0, "slashed in the new state"); + assert_eq!(rebuilt.get(2), preset::MAX_EFFECTIVE_BALANCE); } /// The failure a live mainnet follower hit: `promote_beacon_anchor` prunes diff --git a/crates/blockchain/state_transition/src/metrics.rs b/crates/blockchain/state_transition/src/metrics.rs index e7da0ce8..aacf9acf 100644 --- a/crates/blockchain/state_transition/src/metrics.rs +++ b/crates/blockchain/state_transition/src/metrics.rs @@ -62,6 +62,40 @@ pub fn inc_committee_cache_lookups(result: &str) { .inc(); } +static LEAN_BEACON_JUSTIFIED_BALANCES_LOOKUPS_TOTAL: LazyLock = LazyLock::new( + || { + register_int_counter_vec!( + "lean_beacon_justified_balances_lookups_total", + "Beacon fork-choice justified-balances snapshot lookups, by whether the cached snapshot served them", + &["result"] + ) + .unwrap() + }, +); + +/// Count one justified-balances lookup: `hit` (the cached snapshot matched the +/// justified checkpoint) or `miss` (it was rebuilt from the checkpoint state). +pub fn inc_justified_balances_lookups(result: &str) { + LEAN_BEACON_JUSTIFIED_BALANCES_LOOKUPS_TOTAL + .with_label_values(&[result]) + .inc(); +} + +static LEAN_BEACON_JUSTIFIED_BALANCES_BUILD_SECONDS: LazyLock = LazyLock::new(|| { + register_histogram!( + "lean_beacon_justified_balances_build_seconds", + "Duration of one justified-balances snapshot build (one pass over the checkpoint state's registry)", + vec![0.001, 0.005, 0.01, 0.025, 0.05, 0.1, 0.25, 0.5, 1.0] + ) + .unwrap() +}); + +/// Time one justified-balances snapshot build, excluding the checkpoint state +/// lookup it starts from. +pub fn time_justified_balances_build() -> TimingGuard { + TimingGuard::new(&LEAN_BEACON_JUSTIFIED_BALANCES_BUILD_SECONDS) +} + static LEAN_STATE_TRANSITION_TIME_SECONDS: LazyLock = LazyLock::new(|| { register_histogram!( "lean_state_transition_time_seconds", diff --git a/crates/blockchain/state_transition/tests/beacon_spec/fork_choice.rs b/crates/blockchain/state_transition/tests/beacon_spec/fork_choice.rs index 87eb3ff9..f9fc680b 100644 --- a/crates/blockchain/state_transition/tests/beacon_spec/fork_choice.rs +++ b/crates/blockchain/state_transition/tests/beacon_spec/fork_choice.rs @@ -698,6 +698,41 @@ pub(super) fn apply_checks( if let Some(expected) = &checks.should_override_forkchoice_update { check_should_override_forkchoice_update(expected, store, config)?; } + check_weights_against_the_spec(store, config) +} + +/// Oracle for the fork-choice weights, run at every `checks` step. +/// +/// No fixture checks weights directly, but `get_head` descends on +/// `compute_weights` (one pass over the votes, balances from the justified- +/// balances snapshot), while `get_weight` is the specification's per-root +/// definition and reads the justified checkpoint state itself. Comparing the +/// two for every block in the filtered tree puts the snapshot against the spec +/// path on every fixture, whatever the fixture asserts. +fn check_weights_against_the_spec(store: &Store, config: &Config) -> Result<(), String> { + let index = store.block_index(); + let tree = fork_choice::get_filtered_block_tree(store, &index, config) + .map_err(|err| format!("get_filtered_block_tree: {err:?}"))?; + let weights = fork_choice::compute_weights(store, &index, config) + .map_err(|err| format!("compute_weights: {err:?}"))?; + for root in tree.keys() { + let spec = match fork_choice::get_weight(store, &index, *root, config) { + Ok(spec) => spec, + // A vote for a block that invalidation removed from the index + // (`sync/optimistic`): the specification's `get_weight` raises on + // it, and `compute_weights` drops such a vote on purpose, so the + // two are not comparable for that case. + Err(err) if format!("{err:?}").contains("root in store.blocks") => continue, + Err(err) => return Err(format!("get_weight(0x{}): {err:?}", hex::encode(root.0))), + }; + let single_pass = weights.get(root).copied().unwrap_or_default(); + if spec != single_pass { + return Err(format!( + "weight of 0x{}: spec get_weight {spec}, compute_weights {single_pass}", + hex::encode(root.0) + )); + } + } Ok(()) } diff --git a/crates/common/types/src/beacon/fork_choice.rs b/crates/common/types/src/beacon/fork_choice.rs index 7002466d..8919a15e 100644 --- a/crates/common/types/src/beacon/fork_choice.rs +++ b/crates/common/types/src/beacon/fork_choice.rs @@ -12,7 +12,8 @@ use libssz_derive::{HashTreeRoot, SszDecode, SszEncode}; -use crate::beacon::primitives::{Epoch, ExecutionBlockHash, Root, Uint256}; +use crate::beacon::containers::Checkpoint; +use crate::beacon::primitives::{Epoch, ExecutionBlockHash, Gwei, Root, Uint256, ValidatorIndex}; /// One validator's most recent attestation: the epoch it targeted, and the /// block it attested to (the LMD GHOST vote). @@ -100,12 +101,78 @@ pub struct PayloadStatusV1 { pub validation_error: Option, } +/// The justified checkpoint state's balances, flattened for the fork-choice +/// vote loop: `store.checkpoint_states[checkpoint]` as `get_weight` and +/// `get_proposer_score` read it. +/// +/// A tree descent per vote (`state.validator(i)`) is an order of magnitude +/// dearer than an array read, and the loop runs once per `get_head` over every +/// voter. This is built once per justified checkpoint and keyed by it, so a +/// reader compares [`Self::checkpoint`] with the store's current one and +/// rebuilds on a mismatch, with no hook on the code that moves the checkpoint. +/// +/// Held in the store, hence here rather than in `ethlambda-state-transition`; +/// the builder lives there, next to the state accessors. +#[derive(Debug, PartialEq, Eq)] +pub struct JustifiedBalances { + checkpoint: Checkpoint, + /// Zero unless the validator is active at `checkpoint.epoch` and unslashed, + /// the two conditions under which a vote weighs anything. + balances: Box<[Gwei]>, + /// `get_total_active_balance` of the checkpoint state: unlike `balances` + /// it counts slashed validators, and it is floored at one increment. + total_active_balance: Gwei, +} + +impl JustifiedBalances { + /// Wraps already-computed values; see the field docs for what they mean. + pub fn new(checkpoint: Checkpoint, balances: Box<[Gwei]>, total_active_balance: Gwei) -> Self { + Self { + checkpoint, + balances, + total_active_balance, + } + } + + /// The checkpoint these balances were derived from. + pub fn checkpoint(&self) -> Checkpoint { + self.checkpoint + } + + /// The weight of `index`'s vote: zero when inactive or slashed, and also + /// past the end, since that validator did not exist at the checkpoint. + pub fn get(&self, index: ValidatorIndex) -> Gwei { + usize::try_from(index) + .ok() + .and_then(|index| self.balances.get(index)) + .copied() + .unwrap_or_default() + } + + /// The checkpoint state's total active balance, as the specification's + /// `get_total_active_balance` defines it. + pub fn total_active_balance(&self) -> Gwei { + self.total_active_balance + } +} + #[cfg(test)] mod tests { use libssz::{SszDecode as _, SszEncode as _}; use super::*; + #[test] + fn justified_balances_read_zero_past_the_registry() { + let balances = JustifiedBalances::new(Checkpoint::default(), vec![7, 0, 9].into(), 16); + assert_eq!(balances.get(0), 7); + assert_eq!(balances.get(1), 0); + assert_eq!(balances.get(2), 9); + assert_eq!(balances.get(3), 0); + assert_eq!(balances.get(u64::MAX), 0); + assert_eq!(balances.total_active_balance(), 16); + } + #[test] fn a_pow_block_round_trips_through_ssz() { // The store persists these, so the derive has to survive the move out diff --git a/crates/storage/src/store.rs b/crates/storage/src/store.rs index 96addfbd..497acaf4 100644 --- a/crates/storage/src/store.rs +++ b/crates/storage/src/store.rs @@ -19,7 +19,7 @@ use ethlambda_types::{ config::Config, containers::{BeaconState, Checkpoint as BeaconCheckpoint, SignedBeaconBlock}, fork::ForkName, - fork_choice::{LatestMessage, PayloadStatusV1, PowBlock}, + fork_choice::{JustifiedBalances, LatestMessage, PayloadStatusV1, PowBlock}, preset::{Preset, SLOTS_PER_EPOCH}, primitives::ExecutionBlockHash, }, @@ -667,6 +667,10 @@ pub(crate) struct BeaconScratch { pub(crate) block_timeliness: HashMap, pub(crate) equivocating_indices: HashSet, pub(crate) latest_messages: HashMap, + /// The justified checkpoint state's balances, flattened for the vote loop + /// and keyed by their own checkpoint. A derived cache: a miss is rebuilt + /// from `checkpoint_states`, so nothing needs to persist it. + pub(crate) justified_balances: Option>, pub(crate) pow_blocks: HashMap, pub(crate) unrealized_justifications: HashMap, /// Beacon roots imported on an execution client's `NOT_VALIDATED` answer, @@ -3210,6 +3214,32 @@ impl Store { } } + /// The cached justified-balances snapshot, if it was built for exactly + /// `checkpoint`. + /// + /// The checkpoint is the whole key, so a caller never has to know which + /// code moved the justified checkpoint: a stale snapshot just misses. + pub fn justified_balances( + &self, + checkpoint: &BeaconCheckpoint, + ) -> Option> { + self.beacon + .lock() + .unwrap() + .justified_balances + .as_ref() + .filter(|balances| balances.checkpoint() == *checkpoint) + .cloned() + } + + /// Replaces the cached justified-balances snapshot. + /// + /// Takes `&self`, like [`Self::cache_state`]: the read-only fork-choice + /// helpers fill it on a miss. + pub fn set_justified_balances(&self, balances: Arc) { + self.beacon.lock().unwrap().justified_balances = Some(balances); + } + /// Looks up a PoW block by its own hash, standing in for the /// specification's `get_pow_block(hash)`. pub fn beacon_pow_block(&self, hash: H256) -> Option { diff --git a/docs/beacon_stf.md b/docs/beacon_stf.md index 17ebe686..4eb511b5 100644 --- a/docs/beacon_stf.md +++ b/docs/beacon_stf.md @@ -293,6 +293,21 @@ over the registry should walk `validators().iter()`, zipped with candidate for computing once per epoch rather than per call, once it is shown that no block operation changes it mid-epoch. +Fork choice has the same problem on a longer loop: `get_head` weighs every +validator's latest vote by the voter's balance at the justified checkpoint, and +the boost needs that state's total active balance. Both read one +`JustifiedBalances` snapshot (`ethlambda-types`, held in the store's +`BeaconScratch`) instead of a `validator(i)` descent per vote and a registry +scan per boost. It is built from `checkpoint_state(justified)` only, keyed by +the checkpoint it came from, so a moved justified checkpoint is a miss with no +hook on the writers; it holds a zero for validators that are inactive or +slashed, and a separate total that still counts slashed-but-active validators +(the specification's `get_total_active_balance`, floored at one increment). +Equivocations stay out of it: they are store-level and can change at any time, +so they are filtered per vote. `get_weight` stays as the specification's +per-root definition, and the fork-choice fixture runner checks +`compute_weights` against it at every `checks` step. + ## Macros and traits Two `macro_rules!` in the whole crate, both local, both replacing boilerplate that diff --git a/docs/metrics.md b/docs/metrics.md index bf215a3c..960346c8 100644 --- a/docs/metrics.md +++ b/docs/metrics.md @@ -422,6 +422,25 @@ is a registry scan plus a whole-epoch shuffle on the import thread. `unkeyable` is a lookup the cache could not key at all, mostly the genesis state asking about its own first epochs, and should be zero on a checkpoint-synced follower. +### Beacon Justified Balances + +Beacon fork choice weighs every vote by the voter's effective balance at the +justified checkpoint. `justified_balances` (in +`crates/blockchain/state_transition/src/beacon/fork_choice.rs`) flattens that +state into one array, keyed by the justified checkpoint and rebuilt on the first +`get_head` after the checkpoint moves. This is ethlambda-specific, not part of +the leanMetrics spec. + +| Name | Type | Usage | Sample collection event | Labels | +|------|------|-------|-------------------------|--------| +| `lean_beacon_justified_balances_lookups_total` | Counter | Snapshot lookups, by whether the cached snapshot served them | On every `justified_balances` call: `get_head`'s weights, the proposer boost, and the reorg helpers | result=hit,miss | +| `lean_beacon_justified_balances_build_seconds` | Histogram | Time to build one snapshot from the checkpoint state | On every miss, around the registry pass (the checkpoint state lookup is not included) | | + +**Read the miss rate against the justified-checkpoint rate**: about one miss +per justified checkpoint change, so a handful per hour on a healthy chain. +Misses that track `get_head` calls mean the justified checkpoint is flapping +between branches, or the snapshot is being replaced between two readers. + ### Beacon Pubkey Cache Every BLS signature check `ethlambda beacon` runs (block import, fork choice, From 6a6160ef25ed4438156122eb37f13c17a52e0a58 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Tom=C3=A1s=20Gr=C3=BCner?= <47506558+MegaRedHand@users.noreply.github.com> Date: Fri, 2 Oct 2026 11:52:30 -0300 Subject: [PATCH 2/2] perf(beacon): keep latest messages dense, indexed by validator The vote loop walked a HashMap in hash order and probed the equivocator set once per vote, then read the snapshot balance at a random offset. Store the votes as a Vec> indexed by validator instead, so votes and snapshot balances are read in index order and the equivocator probe is skipped while the set is empty. Saturating sums of non-negative balances give the same weights in any order, so the result is unchanged. The table grows to the highest voting index, so set_latest_message asserts (debug builds) that the index is plausible; indices only come from validated attestations, and a wild one would otherwise allocate a table to match. --- crates/storage/src/store.rs | 75 +++++++++++++++++++++++++++++++------ docs/beacon_stf.md | 6 +++ 2 files changed, 70 insertions(+), 11 deletions(-) diff --git a/crates/storage/src/store.rs b/crates/storage/src/store.rs index 497acaf4..f2dbecc2 100644 --- a/crates/storage/src/store.rs +++ b/crates/storage/src/store.rs @@ -636,6 +636,12 @@ impl GossipSignatureBuffer { } } +/// Bound on a validator index [`Store::set_latest_message`] accepts in debug +/// builds. Far above any registry this client will meet, far below what a +/// corrupt index would need to exhaust memory: the dense vote table has one slot +/// per index up to the highest seen. +const MAX_PLAUSIBLE_VALIDATOR_INDEX: u64 = 1 << 26; + /// Beacon fork-choice state that is per-slot or per-epoch scratch rather than /// chain history: nothing here survives a restart, and nothing here is worth /// the write amplification of persisting. @@ -666,7 +672,10 @@ pub(crate) struct BeaconScratch { pub(crate) proposer_boost_root: H256, pub(crate) block_timeliness: HashMap, pub(crate) equivocating_indices: HashSet, - pub(crate) latest_messages: HashMap, + /// Dense, indexed by validator index: the vote loop reads votes and + /// snapshot balances in index order, not hash order. `None` is a validator + /// that has not voted (or an index below the highest one seen). + pub(crate) latest_messages: Vec>, /// The justified checkpoint state's balances, flattened for the vote loop /// and keyed by their own checkpoint. A derived cache: a miss is rebuilt /// from `checkpoint_states`, so nothing needs to persist it. @@ -3184,21 +3193,33 @@ impl Store { .lock() .unwrap() .latest_messages - .get(&index) + .get(usize::try_from(index).ok()?) .copied() + .flatten() } /// Records the latest attestation for validator `index`. + /// + /// The table grows to hold `index`, so an index must come from a validated + /// attestation (a member of a committee of the registry), never from raw + /// input: a wild index would allocate a table to match. The assertion is + /// the tripwire for that in debug builds and fixtures. pub fn set_latest_message(&mut self, index: u64, message: LatestMessage) { - self.beacon - .lock() - .unwrap() - .latest_messages - .insert(index, message); + debug_assert!( + index < MAX_PLAUSIBLE_VALIDATOR_INDEX, + "validator index {index} would grow the dense vote table far past any registry" + ); + let slot = usize::try_from(index).expect("validator index fits in usize"); + let mut beacon = self.beacon.lock().unwrap(); + if beacon.latest_messages.len() <= slot { + beacon.latest_messages.resize(slot + 1, None); + } + beacon.latest_messages[slot] = Some(message); } /// Calls `f` with `(validator_index, latest_message)` for every latest - /// message whose validator has not been observed equivocating. + /// message whose validator has not been observed equivocating, in + /// ascending validator index. /// /// Takes a closure rather than returning an iterator or a cloned map: /// the data lives behind a mutex, so a borrow of it cannot escape the @@ -3207,10 +3228,17 @@ impl Store { /// let it count for either side of the fork it created. pub fn for_each_non_equivocating_latest_message(&self, mut f: impl FnMut(u64, LatestMessage)) { let beacon = self.beacon.lock().unwrap(); - for (&index, &message) in &beacon.latest_messages { - if !beacon.equivocating_indices.contains(&index) { - f(index, message); + // Most of the time nobody has equivocated, so skip the set probe per vote. + let any_equivocators = !beacon.equivocating_indices.is_empty(); + for (index, message) in beacon.latest_messages.iter().enumerate() { + let Some(message) = message else { + continue; + }; + let index = index as u64; + if any_equivocators && beacon.equivocating_indices.contains(&index) { + continue; } + f(index, *message); } } @@ -6675,6 +6703,31 @@ mod tests { assert_eq!(seen, vec![1]); } + #[test] + fn latest_messages_are_visited_in_validator_index_order_and_gaps_are_skipped() { + let mut store = Store::test_store(); + let message = |epoch| LatestMessage { + epoch, + root: H256::from([epoch as u8; 32]), + }; + + // Written out of order, with a gap, and one overwritten. + store.set_latest_message(9, message(1)); + store.set_latest_message(2, message(2)); + store.set_latest_message(5, message(3)); + store.set_latest_message(5, message(4)); + + let mut seen = Vec::new(); + store + .for_each_non_equivocating_latest_message(|index, vote| seen.push((index, vote.epoch))); + assert_eq!(seen, vec![(2, 2), (5, 4), (9, 1)]); + + assert_eq!(store.latest_message(5).map(|vote| vote.epoch), Some(4)); + assert_eq!(store.latest_message(3), None, "a gap is not a vote"); + assert_eq!(store.latest_message(1_000), None, "past the table"); + assert_eq!(store.latest_message(u64::MAX), None); + } + #[test] fn a_pow_block_is_looked_up_by_its_own_hash() { let mut store = Store::test_store(); diff --git a/docs/beacon_stf.md b/docs/beacon_stf.md index 4eb511b5..85b8e18d 100644 --- a/docs/beacon_stf.md +++ b/docs/beacon_stf.md @@ -308,6 +308,12 @@ so they are filtered per vote. `get_weight` stays as the specification's per-root definition, and the fork-choice fixture runner checks `compute_weights` against it at every `checks` step. +The latest votes are stored the same way, as a dense table indexed by validator +index rather than a hash map, so the vote loop reads votes and snapshot +balances in index order. The table grows to the highest voting index, which is +why `set_latest_message` only takes indices from validated attestations (a +debug assertion bounds it). + ## Macros and traits Two `macro_rules!` in the whole crate, both local, both replacing boilerplate that