diff --git a/crates/blockchain/state_transition/src/beacon/block_production.rs b/crates/blockchain/state_transition/src/beacon/block_production.rs index 0c412cb4..d8dd1e28 100644 --- a/crates/blockchain/state_transition/src/beacon/block_production.rs +++ b/crates/blockchain/state_transition/src/beacon/block_production.rs @@ -38,8 +38,8 @@ use super::bls; use super::config::Config; use super::error::{Error, Result, verify}; use super::helpers::accessors::{ - CommitteeCache, get_beacon_proposer_index, get_block_root, get_current_epoch, - get_previous_epoch, get_randao_mix, + ActiveBalanceCache, CommitteeCache, get_beacon_proposer_index, get_block_root, + get_current_epoch, get_previous_epoch, get_randao_mix, }; use super::helpers::electra::{ get_attesting_indices, get_indexed_attestation, is_valid_indexed_attestation, @@ -343,6 +343,7 @@ pub fn assemble_block( config, &ExecutionEngine::valid(), &CommitteeCache::default(), + &ActiveBalanceCache::default(), )?; block.state_root = post.hash_tree_root(); Ok(block) @@ -435,6 +436,7 @@ mod tests { &Config::mainnet(), &ExecutionEngine::valid(), &CommitteeCache::default(), + &ActiveBalanceCache::default(), ) .unwrap(); assert_eq!(block.state_root, post.hash_tree_root()); diff --git a/crates/blockchain/state_transition/src/beacon/fork_choice.rs b/crates/blockchain/state_transition/src/beacon/fork_choice.rs index 5560782b..2915f32c 100644 --- a/crates/blockchain/state_transition/src/beacon/fork_choice.rs +++ b/crates/blockchain/state_transition/src/beacon/fork_choice.rs @@ -2370,8 +2370,18 @@ pub fn on_block( stf::ExecutionEngine::valid() } }; - let transition = - stf::state_transition(&mut state, &signed_block, true, config, &engine, committees); + // The store's own cache, so every block of this chain shares one entry per + // epoch; see `ActiveBalanceCacheExt`. + let active_balances = store.active_balance_cache(); + let transition = stf::state_transition( + &mut state, + &signed_block, + true, + config, + &engine, + committees, + &active_balances, + ); // `optimistic-sync.md`: a block deemed `INVALIDATED` MUST NOT be included // in the canonical chain. That is stated here, on the verdict, rather than diff --git a/crates/blockchain/state_transition/src/beacon/helpers/accessors.rs b/crates/blockchain/state_transition/src/beacon/helpers/accessors.rs index 677e895e..8985d957 100644 --- a/crates/blockchain/state_transition/src/beacon/helpers/accessors.rs +++ b/crates/blockchain/state_transition/src/beacon/helpers/accessors.rs @@ -11,7 +11,9 @@ use std::sync::Arc; // and so on), so every existing import of the type this module used to define // is unaffected by its move to `ethlambda-storage`; see [`CommitteeCacheExt`] // below for what this crate still contributes. -pub use ethlambda_storage::{CommitteeCache, Lookup, ShufflingKey}; +pub use ethlambda_storage::{ + ActiveBalanceCache, ActiveBalanceKey, CommitteeCache, Lookup, ShufflingKey, +}; use crate::beacon::config::Config; use crate::beacon::constants; @@ -301,6 +303,96 @@ impl CommitteeCacheExt for CommitteeCache { } } +/// The key `state`'s total active balance is cached under, or `None` if +/// `state` cannot name the block that fixes it: the state has not advanced +/// past that block's slot (the genesis state, asked about its own epoch), or +/// the slot has fallen out of its `SLOTS_PER_HISTORICAL_ROOT` window. `None` +/// means the lookup is computed for its caller alone and not cached. +/// +/// The total for a state in epoch `E` sums the effective balances of the +/// validators active at `E`, and both inputs are fixed once epoch `E - 1` +/// ends. Effective balances are written only by `process_effective_balance_updates`, +/// at the end of `E - 1`; an epoch's update is the last writer before `E`'s +/// states exist. The active set moves only through `activation_epoch` and +/// `exit_epoch`, and every assignment goes through `compute_activation_exit_epoch`, +/// at least `MAX_SEED_LOOKAHEAD` epochs ahead, so none can land in `E` for `E` +/// itself. Slashing sets `slashed` and queues an exit but leaves effective +/// balances alone, and the spec's total keeps slashed validators. Deposits +/// append validators that are not yet active, and the Electra upgrade zeroes +/// only never-activated validators. So the block root at the last slot of +/// `E - 1` identifies the history that determines the total: two states +/// agreeing on it agree on the total, however much they disagree after it. +/// +/// That is one epoch later than [`shuffling_key`]'s root (the last slot of +/// `E - 2`), because the shuffle's inputs are fixed `MIN_SEED_LOOKAHEAD` epochs +/// ahead of use and the total's are not: an effective balance written at the +/// end of `E - 1` is part of the total for `E`, so a key rooted at `E - 2` +/// would let two branches that diverge during `E - 1` share one entry. +/// +/// # Epoch 0 +/// +/// There is no epoch before it, so it takes the genesis block, at slot 0, as +/// its deciding block, as [`shuffling_key`] does for its first epochs. Nothing +/// after genesis can reach either input within epoch 0 (the effective-balance +/// update at its end writes epoch 1's), and a state from a different genesis +/// has a different genesis block root. +/// +/// # Callers must be block processing +/// +/// Between `process_effective_balance_updates` and the slot increment that +/// follows it, a state is still in epoch `E - 1` but already carries epoch +/// `E`'s effective balances, so its total would not match a key rooted at +/// `E - 2`'s end. The only callers routed through the cache are block +/// processing's (attestations, sync aggregate), which never run in that +/// window: a block is processed on a state that `process_slots` has already +/// advanced to the block's slot. Epoch processing keeps calling the +/// uncached [`get_total_active_balance`]. +fn active_balance_key(state: &BeaconState) -> Option { + let epoch = get_current_epoch(state); + let decision_slot = compute_start_slot_at_epoch(epoch).saturating_sub(1); + let decision_root = get_block_root_at_slot(state, decision_slot).ok()?; + Some(ActiveBalanceKey { + epoch, + decision_root, + }) +} + +/// Extends `ethlambda-storage`'s [`ActiveBalanceCache`] with the consensus +/// logic that keys and computes its entries, for the same reason +/// [`CommitteeCacheExt`] lives here and not in `ethlambda-storage`. +pub trait ActiveBalanceCacheExt { + /// `get_total_active_balance(state)`, computed once per key and then + /// served to every later caller naming the same one. Equal to the + /// uncached function for every state it is given; see [`active_balance_key`] + /// for why, and for who may call it. + fn total_active_balance(&self, state: &BeaconState) -> Result; +} + +impl ActiveBalanceCacheExt for ActiveBalanceCache { + fn total_active_balance(&self, state: &BeaconState) -> Result { + let Some(key) = active_balance_key(state) else { + crate::metrics::inc_total_active_balance_lookups("unkeyable"); + return get_total_active_balance(state); + }; + + let (total, lookup) = self.get_or_compute(key, || compute_total_active_balance(state)); + crate::metrics::inc_total_active_balance_lookups(match lookup { + Lookup::Hit => "hit", + Lookup::Miss => "miss", + }); + // The cross-check the key's soundness argument rests on: a stale hit + // would mean some writer moved an input mid-epoch. + if lookup == Lookup::Hit { + debug_assert_eq!( + total, + compute_total_active_balance(state), + "stale total active balance cache entry" + ); + } + Ok(total) + } +} + /// The committee at `slot` with index `index`. /// /// One epoch's active set is shuffled once and then split across every slot and @@ -398,9 +490,28 @@ pub fn get_total_balance(state: &BeaconState, indices: &[ValidatorIndex]) -> Res } /// The combined effective balance of the currently active validators. +/// +/// One in-order pass over the registry, summing as it goes: the same value as +/// [`get_total_balance`] over [`get_active_validator_indices`] (same +/// saturating sum, same one-increment floor), without the intermediate index +/// list or a tree descent per active validator. `state.validators().iter()` +/// walks the leaves, whereas `state.validator(i)` descends from the root each +/// time. pub fn get_total_active_balance(state: &BeaconState) -> Result { - let indices = get_active_validator_indices(state, get_current_epoch(state)); - get_total_balance(state, &indices) + Ok(compute_total_active_balance(state)) +} + +/// [`get_total_active_balance`]'s body, without the `Result` it never needed. +fn compute_total_active_balance(state: &BeaconState) -> Gwei { + let epoch = get_current_epoch(state); + let total = state + .validators() + .iter() + .filter(|validator| is_active_validator(validator, epoch)) + .fold(0, |sum: Gwei, validator| { + sum.saturating_add(validator.effective_balance) + }); + total.max(preset::EFFECTIVE_BALANCE_INCREMENT) } /// The signing domain for `domain_type` at `epoch`, or at the current epoch when @@ -649,6 +760,159 @@ mod tests { ); } + /// SplitMix64: a tiny deterministic generator, so the randomized tests + /// below need no dependency. + struct SplitMix64(u64); + + impl SplitMix64 { + fn next(&mut self) -> u64 { + self.0 = self.0.wrapping_add(0x9E37_79B9_7F4A_7C15); + let mut z = self.0; + z = (z ^ (z >> 30)).wrapping_mul(0xBF58_476D_1CE4_E5B9); + z = (z ^ (z >> 27)).wrapping_mul(0x94D0_49BB_1331_11EB); + z ^ (z >> 31) + } + } + + /// A state at `epoch` whose validators get random activation and exit + /// epochs around it and random effective balances. + fn random_registry_state(rng: &mut SplitMix64, count: usize, epoch: Epoch) -> BeaconState { + let mut state = with_validators(count); + *state.slot_mut() = compute_start_slot_at_epoch(epoch); + for index in 0..count { + let validator = state.validator_mut(index as ValidatorIndex).unwrap(); + validator.activation_epoch = rng.next() % (epoch + 3); + validator.exit_epoch = match rng.next() % 3 { + 0 => constants::FAR_FUTURE_EPOCH, + _ => rng.next() % (epoch + 3), + }; + validator.effective_balance = match rng.next() % 4 { + 0 => 0, + 1 => u64::MAX - rng.next() % 4, + _ => (rng.next() % 64) * preset::EFFECTIVE_BALANCE_INCREMENT, + }; + } + state.apply_pending_mutations(); + state + } + + /// The one-pass total is the spec's `get_total_balance` over the active + /// indices, including its saturation and its floor. + #[test] + fn the_one_pass_total_matches_the_spec_formulation() { + let mut rng = SplitMix64(0x5EED); + for round in 0..64 { + let count = (rng.next() % 40) as usize; + let epoch = rng.next() % 6; + let state = random_registry_state(&mut rng, count, epoch); + let indices = get_active_validator_indices(&state, get_current_epoch(&state)); + assert_eq!( + get_total_active_balance(&state).unwrap(), + get_total_balance(&state, &indices).unwrap(), + "round {round}" + ); + } + } + + /// An all-zero registry hits the floor, not zero. + #[test] + fn the_one_pass_total_is_floored_for_a_zero_registry() { + let mut state = with_validators(8); + for index in 0..8 { + state.validator_mut(index).unwrap().effective_balance = 0; + } + state.apply_pending_mutations(); + assert_eq!( + get_total_active_balance(&state).unwrap(), + preset::EFFECTIVE_BALANCE_INCREMENT + ); + } + + /// A state at the first slot of `epoch` whose every block root is `root`. + fn state_at_epoch_with_root(epoch: Epoch, root: u8) -> BeaconState { + let mut state = with_validators(16); + *state.slot_mut() = compute_start_slot_at_epoch(epoch); + for slot in 0..preset::SLOTS_PER_HISTORICAL_ROOT { + state.block_roots_mut()[slot] = Root::repeat_byte(root); + } + state + } + + /// A lookup returns the spec value and the second one is a hit. + #[test] + fn the_active_balance_cache_serves_the_spec_value() { + let state = state_at_epoch_with_root(3, 1); + let cache = ActiveBalanceCache::default(); + let expected = get_total_active_balance(&state).unwrap(); + + assert_eq!(cache.total_active_balance(&state).unwrap(), expected); + let key = active_balance_key(&state).unwrap(); + assert_eq!(cache.get(key), Some(expected)); + assert_eq!(cache.total_active_balance(&state).unwrap(), expected); + } + + /// Two sibling states in one epoch that disagree on the deciding block + /// keep separate entries, each with its own registry's total. + #[test] + fn sibling_states_with_different_decision_roots_get_different_entries() { + let cache = ActiveBalanceCache::default(); + let a = state_at_epoch_with_root(3, 1); + let mut b = state_at_epoch_with_root(3, 2); + // Sibling `b` has a validator the fork `a` never saw. + b.validator_mut(0).unwrap().effective_balance = 0; + b.apply_pending_mutations(); + + let total_a = cache.total_active_balance(&a).unwrap(); + let total_b = cache.total_active_balance(&b).unwrap(); + + assert_ne!(total_a, total_b); + assert_eq!(total_a, get_total_active_balance(&a).unwrap()); + assert_eq!(total_b, get_total_active_balance(&b).unwrap()); + assert_eq!(cache.get(active_balance_key(&a).unwrap()), Some(total_a)); + assert_eq!(cache.get(active_balance_key(&b).unwrap()), Some(total_b)); + } + + /// Across an epoch boundary the key changes, so the lookup misses and + /// fills a new entry instead of serving the old epoch's total. + #[test] + fn crossing_an_epoch_boundary_misses_and_refills() { + let cache = ActiveBalanceCache::default(); + let mut state = state_at_epoch_with_root(3, 1); + let before = cache.total_active_balance(&state).unwrap(); + let old_key = active_balance_key(&state).unwrap(); + + // The effective-balance update at the boundary, then the new epoch + // with its own deciding root. + state.validator_mut(0).unwrap().effective_balance = 0; + state.apply_pending_mutations(); + *state.slot_mut() = compute_start_slot_at_epoch(4); + let last_slot = + (compute_start_slot_at_epoch(4) - 1) as usize % preset::SLOTS_PER_HISTORICAL_ROOT; + state.block_roots_mut()[last_slot] = Root::repeat_byte(9); + + let new_key = active_balance_key(&state).unwrap(); + assert_ne!(old_key, new_key); + let after = cache.total_active_balance(&state).unwrap(); + assert_eq!(after, get_total_active_balance(&state).unwrap()); + assert_ne!(before, after); + assert_eq!(cache.get(old_key), Some(before)); + assert_eq!(cache.get(new_key), Some(after)); + } + + /// The genesis state cannot name the deciding block, so it is computed + /// for its caller alone and not cached. + #[test] + fn an_unkeyable_state_is_computed_but_not_cached() { + let mut state = with_validators(8); + *state.slot_mut() = 0; + assert!(active_balance_key(&state).is_none()); + let cache = ActiveBalanceCache::default(); + assert_eq!( + cache.total_active_balance(&state).unwrap(), + get_total_active_balance(&state).unwrap() + ); + } + #[test] fn proposer_is_drawn_from_the_active_set() { let state = crate::beacon::helpers::test_state::with_validators(32); diff --git a/crates/blockchain/state_transition/src/beacon/helpers/altair.rs b/crates/blockchain/state_transition/src/beacon/helpers/altair.rs index eaa008c2..3c45a7f1 100644 --- a/crates/blockchain/state_transition/src/beacon/helpers/altair.rs +++ b/crates/blockchain/state_transition/src/beacon/helpers/altair.rs @@ -182,11 +182,14 @@ pub fn get_next_sync_committee(state: &BeaconState) -> Result Result { - let total_active_balance = get_total_active_balance(state)?; - Ok( - preset::EFFECTIVE_BALANCE_INCREMENT * preset::BASE_REWARD_FACTOR - / integer_squareroot(total_active_balance), - ) + Ok(base_reward_per_increment(get_total_active_balance(state)?)) +} + +/// [`get_base_reward_per_increment`] for a total active balance the caller +/// already holds, such as one served by an `ActiveBalanceCache`. +pub fn base_reward_per_increment(total_active_balance: Gwei) -> Gwei { + preset::EFFECTIVE_BALANCE_INCREMENT * preset::BASE_REWARD_FACTOR + / integer_squareroot(total_active_balance) } /// The base reward for the validator at `index`. diff --git a/crates/blockchain/state_transition/src/beacon/stf/altair.rs b/crates/blockchain/state_transition/src/beacon/stf/altair.rs index 22b320f5..ef986e63 100644 --- a/crates/blockchain/state_transition/src/beacon/stf/altair.rs +++ b/crates/blockchain/state_transition/src/beacon/stf/altair.rs @@ -31,11 +31,12 @@ use crate::beacon::constants; use crate::beacon::containers::{BeaconState, altair, phase0}; use crate::beacon::error::{Error, Result, verify}; use crate::beacon::helpers::accessors::{ - CommitteeCache, CommitteeCacheExt, get_beacon_proposer_index, get_block_root_at_slot, - get_current_epoch, get_domain, get_previous_epoch, get_total_active_balance, + ActiveBalanceCache, ActiveBalanceCacheExt, CommitteeCache, CommitteeCacheExt, + get_beacon_proposer_index, get_block_root_at_slot, get_current_epoch, get_domain, + get_previous_epoch, }; use crate::beacon::helpers::altair::{ - add_flag, get_attestation_participation_flag_indices, get_base_reward_per_increment, has_flag, + add_flag, base_reward_per_increment, get_attestation_participation_flag_indices, has_flag, }; use crate::beacon::helpers::attestation::{get_indexed_attestation, is_valid_indexed_attestation}; use crate::beacon::helpers::misc::{compute_epoch_at_slot, compute_signing_root}; @@ -82,6 +83,7 @@ pub fn process_attestation( state: &mut BeaconState, attestation: &phase0::Attestation, committees: &CommitteeCache, + active_balances: &ActiveBalanceCache, ) -> Result<()> { let data = attestation.data; let current_epoch = get_current_epoch(state); @@ -160,7 +162,8 @@ pub fn process_attestation( }; // Hoisted: see the comment on the same line in `electra::process_attestation`. - let base_reward_per_increment = get_base_reward_per_increment(state)?; + let base_reward_per_increment = + base_reward_per_increment(active_balances.total_active_balance(state)?); let mut proposer_reward_numerator: Gwei = 0; let mut updates: Vec<(ValidatorIndex, ParticipationFlags)> = Vec::new(); @@ -175,6 +178,10 @@ pub fn process_attestation( })?; let mut new_flags: ParticipationFlags = 0; + // The attester's validator is read at most once, on the first flag + // it newly earns: each read is a tree descent, and the effective + // balance cannot change inside this read-only phase. + let mut attester_increments: Option = None; for &flag_index in &participation_flag_indices { if has_flag(current_flags, flag_index) { continue; @@ -186,8 +193,15 @@ pub fn process_attestation( // the result is bit-identical. See `electra::process_attestation` // for why the hoist is not a tidy-up but the difference between // importing a block at mainnet scale and not. - let increments = - state.validator(index)?.effective_balance / preset::EFFECTIVE_BALANCE_INCREMENT; + let increments = match attester_increments { + Some(increments) => increments, + None => { + let increments = state.validator(index)?.effective_balance + / preset::EFFECTIVE_BALANCE_INCREMENT; + attester_increments = Some(increments); + increments + } + }; let reward = (increments * base_reward_per_increment) .checked_mul(weight) .ok_or(Error::ArithmeticOverflow( @@ -289,6 +303,7 @@ pub fn process_attestation( pub fn process_sync_aggregate( state: &mut BeaconState, sync_aggregate: &altair::SyncAggregate, + active_balances: &ActiveBalanceCache, ) -> Result<()> { let committee_bits = &sync_aggregate.sync_committee_bits; @@ -358,9 +373,9 @@ pub fn process_sync_aggregate( }) .collect::>>()?; - let total_active_increments = - get_total_active_balance(state)? / preset::EFFECTIVE_BALANCE_INCREMENT; - let total_base_rewards = get_base_reward_per_increment(state)? + let total_active_balance = active_balances.total_active_balance(state)?; + let total_active_increments = total_active_balance / preset::EFFECTIVE_BALANCE_INCREMENT; + let total_base_rewards = base_reward_per_increment(total_active_balance) .checked_mul(total_active_increments) .ok_or(Error::ArithmeticOverflow( "get_base_reward_per_increment(state) * total_active_increments", @@ -407,7 +422,8 @@ mod tests { use super::*; use crate::beacon::containers::altair::SyncCommittee; use crate::beacon::containers::shared::{AttestationData, Checkpoint, Validator}; - use crate::beacon::helpers::accessors::get_beacon_committee; + use crate::beacon::helpers::accessors::{get_beacon_committee, get_total_active_balance}; + use crate::beacon::helpers::altair::get_base_reward_per_increment; use crate::beacon::primitives::{ BLS_SIGNATURE_SIZE, BlsPubkey, BlsSignature, Bytes32, HashTreeRoot as _, Root, }; @@ -577,7 +593,13 @@ mod tests { let proposer_index = get_beacon_proposer_index(&state).unwrap(); let balance_before = state.balance(proposer_index).unwrap(); - process_attestation(&mut state, &attestation, &CommitteeCache::default()).unwrap(); + process_attestation( + &mut state, + &attestation, + &CommitteeCache::default(), + &ActiveBalanceCache::default(), + ) + .unwrap(); let (previous_epoch_participation, _, _) = state.altair_validator_lists().unwrap(); for &index in &committee { @@ -607,7 +629,13 @@ mod tests { // so reprocessing the identical attestation must grant nothing new, and // the proposer's balance must not move. let balance_before_replay = state.balance(proposer_index).unwrap(); - process_attestation(&mut state, &attestation, &CommitteeCache::default()).unwrap(); + process_attestation( + &mut state, + &attestation, + &CommitteeCache::default(), + &ActiveBalanceCache::default(), + ) + .unwrap(); let balance_after_replay = state.balance(proposer_index).unwrap(); assert_eq!( balance_before_replay, balance_after_replay, @@ -683,7 +711,8 @@ mod tests { .map(|index| state.balance(index).unwrap()) .collect(); - process_sync_aggregate(&mut state, &sync_aggregate).unwrap(); + process_sync_aggregate(&mut state, &sync_aggregate, &ActiveBalanceCache::default()) + .unwrap(); // Independently derived expected rewards, sharing only the // already-tested building blocks (`get_total_active_balance`, @@ -748,7 +777,8 @@ mod tests { sync_committee_signature: BlsSignature::default(), }; assert!( - process_sync_aggregate(&mut state.clone(), &wrong).is_err(), + process_sync_aggregate(&mut state.clone(), &wrong, &ActiveBalanceCache::default()) + .is_err(), "the all-zero signature is not the point at infinity, so this must be rejected" ); @@ -757,7 +787,7 @@ mod tests { sync_committee_signature: g2_point_at_infinity(), }; let balance_before = state.balance(0).unwrap(); - process_sync_aggregate(&mut state, &correct).unwrap(); + process_sync_aggregate(&mut state, &correct, &ActiveBalanceCache::default()).unwrap(); // The sole validator holds every seat and none of them participated, // so it is penalized once per seat and the proposer (itself, the only diff --git a/crates/blockchain/state_transition/src/beacon/stf/bellatrix.rs b/crates/blockchain/state_transition/src/beacon/stf/bellatrix.rs index ddb32090..465cfa48 100644 --- a/crates/blockchain/state_transition/src/beacon/stf/bellatrix.rs +++ b/crates/blockchain/state_transition/src/beacon/stf/bellatrix.rs @@ -33,7 +33,9 @@ use crate::beacon::config::Config; use crate::beacon::constants; use crate::beacon::containers::{BeaconState, bellatrix}; use crate::beacon::error::{Error, Result, verify}; -use crate::beacon::helpers::accessors::{CommitteeCache, get_current_epoch, get_randao_mix}; +use crate::beacon::helpers::accessors::{ + ActiveBalanceCache, CommitteeCache, get_current_epoch, get_randao_mix, +}; use crate::beacon::preset; use crate::beacon::primitives::{ Bytes32, ExecutionAddress, ExecutionBlockHash, HashTreeRoot as _, Root, Slot, Uint256, @@ -67,6 +69,7 @@ pub fn process_block( config: &Config, engine: &ExecutionEngine, committees: &CommitteeCache, + active_balances: &ActiveBalanceCache, ) -> Result<()> { super::block::process_block_header( state, @@ -89,8 +92,9 @@ pub fn process_block( &block.body.voluntary_exits, config, committees, + active_balances, )?; - super::altair::process_sync_aggregate(state, &block.body.sync_aggregate)?; + super::altair::process_sync_aggregate(state, &block.body.sync_aggregate, active_balances)?; Ok(()) } diff --git a/crates/blockchain/state_transition/src/beacon/stf/block.rs b/crates/blockchain/state_transition/src/beacon/stf/block.rs index 234d1cef..c3813676 100644 --- a/crates/blockchain/state_transition/src/beacon/stf/block.rs +++ b/crates/blockchain/state_transition/src/beacon/stf/block.rs @@ -24,7 +24,8 @@ use crate::beacon::containers::{self, BeaconBlockHeader, BeaconState, Eth1Data, use crate::beacon::error::{Result, verify}; use crate::beacon::hash::hash; use crate::beacon::helpers::accessors::{ - CommitteeCache, get_beacon_proposer_index, get_current_epoch, get_domain, get_randao_mix, + ActiveBalanceCache, CommitteeCache, get_beacon_proposer_index, get_current_epoch, get_domain, + get_randao_mix, }; use crate::beacon::helpers::math::xor; use crate::beacon::helpers::misc::compute_signing_root; @@ -46,29 +47,55 @@ pub fn process_block( config: &Config, engine: &ExecutionEngine, committees: &CommitteeCache, + active_balances: &ActiveBalanceCache, ) -> Result<()> { match signed_block { containers::SignedBeaconBlock::Phase0(signed) => { - process_block_phase0(state, &signed.message, config, committees) + process_block_phase0(state, &signed.message, config, committees, active_balances) } containers::SignedBeaconBlock::Altair(signed) => { - process_block_altair(state, &signed.message, config, committees) - } - containers::SignedBeaconBlock::Bellatrix(signed) => { - bellatrix::process_block(state, &signed.message, config, engine, committees) - } - containers::SignedBeaconBlock::Capella(signed) => { - capella::process_block(state, &signed.message, config, engine, committees) - } - containers::SignedBeaconBlock::Deneb(signed) => { - deneb::process_block(state, &signed.message, config, engine, committees) - } - containers::SignedBeaconBlock::Electra(signed) => { - electra::process_block(state, &signed.message, config, engine, committees) - } - containers::SignedBeaconBlock::Fulu(signed) => { - fulu::process_block(state, &signed.message, config, engine, committees) + process_block_altair(state, &signed.message, config, committees, active_balances) } + containers::SignedBeaconBlock::Bellatrix(signed) => bellatrix::process_block( + state, + &signed.message, + config, + engine, + committees, + active_balances, + ), + containers::SignedBeaconBlock::Capella(signed) => capella::process_block( + state, + &signed.message, + config, + engine, + committees, + active_balances, + ), + containers::SignedBeaconBlock::Deneb(signed) => deneb::process_block( + state, + &signed.message, + config, + engine, + committees, + active_balances, + ), + containers::SignedBeaconBlock::Electra(signed) => electra::process_block( + state, + &signed.message, + config, + engine, + committees, + active_balances, + ), + containers::SignedBeaconBlock::Fulu(signed) => fulu::process_block( + state, + &signed.message, + config, + engine, + committees, + active_balances, + ), containers::SignedBeaconBlock::Lean(_) => { crate::beacon::lean_block_unreachable("process_block") } @@ -82,6 +109,7 @@ pub fn process_block_phase0( block: &phase0::BeaconBlock, config: &Config, committees: &CommitteeCache, + active_balances: &ActiveBalanceCache, ) -> Result<()> { process_block_header( state, @@ -101,6 +129,7 @@ pub fn process_block_phase0( &block.body.voluntary_exits, config, committees, + active_balances, )?; Ok(()) } @@ -116,6 +145,7 @@ pub fn process_block_altair( block: &altair::BeaconBlock, config: &Config, committees: &CommitteeCache, + active_balances: &ActiveBalanceCache, ) -> Result<()> { process_block_header( state, @@ -135,8 +165,9 @@ pub fn process_block_altair( &block.body.voluntary_exits, config, committees, + active_balances, )?; - super::altair::process_sync_aggregate(state, &block.body.sync_aggregate)?; + super::altair::process_sync_aggregate(state, &block.body.sync_aggregate, active_balances)?; Ok(()) } diff --git a/crates/blockchain/state_transition/src/beacon/stf/capella.rs b/crates/blockchain/state_transition/src/beacon/stf/capella.rs index b56639a9..2d2d82e4 100644 --- a/crates/blockchain/state_transition/src/beacon/stf/capella.rs +++ b/crates/blockchain/state_transition/src/beacon/stf/capella.rs @@ -51,7 +51,9 @@ use crate::beacon::containers::shared::{ use crate::beacon::containers::{BeaconState, capella, phase0}; use crate::beacon::error::{Error, Result, verify}; use crate::beacon::hash::hash; -use crate::beacon::helpers::accessors::{CommitteeCache, get_current_epoch, get_randao_mix}; +use crate::beacon::helpers::accessors::{ + ActiveBalanceCache, CommitteeCache, get_current_epoch, get_randao_mix, +}; use crate::beacon::helpers::capella::{ is_fully_withdrawable_validator, is_partially_withdrawable_validator, }; @@ -79,6 +81,7 @@ pub fn process_block( config: &Config, engine: &ExecutionEngine, committees: &CommitteeCache, + active_balances: &ActiveBalanceCache, ) -> Result<()> { super::block::process_block_header( state, @@ -101,8 +104,9 @@ pub fn process_block( &block.body.bls_to_execution_changes, config, committees, + active_balances, )?; - super::altair::process_sync_aggregate(state, &block.body.sync_aggregate)?; + super::altair::process_sync_aggregate(state, &block.body.sync_aggregate, active_balances)?; Ok(()) } @@ -374,6 +378,7 @@ pub fn process_operations( bls_to_execution_changes: &[capella::SignedBLSToExecutionChange], config: &Config, committees: &CommitteeCache, + active_balances: &ActiveBalanceCache, ) -> Result<()> { super::operations::process_operations( state, @@ -384,6 +389,7 @@ pub fn process_operations( voluntary_exits, config, committees, + active_balances, )?; for signed_change in bls_to_execution_changes { process_bls_to_execution_change(state, signed_change, config)?; diff --git a/crates/blockchain/state_transition/src/beacon/stf/deneb.rs b/crates/blockchain/state_transition/src/beacon/stf/deneb.rs index 528892b2..05045398 100644 --- a/crates/blockchain/state_transition/src/beacon/stf/deneb.rs +++ b/crates/blockchain/state_transition/src/beacon/stf/deneb.rs @@ -36,10 +36,11 @@ use crate::beacon::containers::{BeaconState, deneb, phase0}; use crate::beacon::error::{Error, Result, verify}; use crate::beacon::hash::hash; use crate::beacon::helpers::accessors::{ - CommitteeCache, CommitteeCacheExt, get_beacon_proposer_index, get_block_root, - get_block_root_at_slot, get_current_epoch, get_previous_epoch, get_randao_mix, + ActiveBalanceCache, ActiveBalanceCacheExt, CommitteeCache, CommitteeCacheExt, + get_beacon_proposer_index, get_block_root, get_block_root_at_slot, get_current_epoch, + get_previous_epoch, get_randao_mix, }; -use crate::beacon::helpers::altair::{add_flag, get_base_reward_per_increment, has_flag}; +use crate::beacon::helpers::altair::{add_flag, base_reward_per_increment, has_flag}; use crate::beacon::helpers::attestation::{get_indexed_attestation, is_valid_indexed_attestation}; use crate::beacon::helpers::math::integer_squareroot; use crate::beacon::helpers::misc::{compute_domain, compute_epoch_at_slot, compute_signing_root}; @@ -77,6 +78,7 @@ pub fn process_block( config: &Config, engine: &ExecutionEngine, committees: &CommitteeCache, + active_balances: &ActiveBalanceCache, ) -> Result<()> { super::block::process_block_header( state, @@ -105,8 +107,9 @@ pub fn process_block( &block.body.bls_to_execution_changes, config, committees, + active_balances, )?; - super::altair::process_sync_aggregate(state, &block.body.sync_aggregate)?; + super::altair::process_sync_aggregate(state, &block.body.sync_aggregate, active_balances)?; Ok(()) } @@ -147,6 +150,7 @@ fn process_operations( bls_to_execution_changes: &[SignedBLSToExecutionChange], config: &Config, committees: &CommitteeCache, + active_balances: &ActiveBalanceCache, ) -> Result<()> { let outstanding = state .eth1_data() @@ -167,7 +171,7 @@ fn process_operations( super::operations::process_attester_slashing(state, attester_slashing, config)?; } for attestation in attestations { - process_attestation(state, attestation, committees)?; + process_attestation(state, attestation, committees, active_balances)?; } for deposit in deposits { super::operations::process_deposit(state, deposit, config)?; @@ -281,6 +285,7 @@ pub fn process_attestation( state: &mut BeaconState, attestation: &phase0::Attestation, committees: &CommitteeCache, + active_balances: &ActiveBalanceCache, ) -> Result<()> { let data = attestation.data; let current_epoch = get_current_epoch(state); @@ -356,7 +361,8 @@ pub fn process_attestation( }; // Hoisted: see the comment on the same line in `electra::process_attestation`. - let base_reward_per_increment = get_base_reward_per_increment(state)?; + let base_reward_per_increment = + base_reward_per_increment(active_balances.total_active_balance(state)?); let mut proposer_reward_numerator: Gwei = 0; let mut updates: Vec<(ValidatorIndex, ParticipationFlags)> = Vec::new(); @@ -371,6 +377,10 @@ pub fn process_attestation( })?; let mut new_flags: ParticipationFlags = 0; + // The attester's validator is read at most once, on the first flag + // it newly earns: each read is a tree descent, and the effective + // balance cannot change inside this read-only phase. + let mut attester_increments: Option = None; for &flag_index in &participation_flag_indices { if has_flag(current_flags, flag_index) { continue; @@ -382,8 +392,15 @@ pub fn process_attestation( // the result is bit-identical. See `electra::process_attestation` // for why the hoist is not a tidy-up but the difference between // importing a block at mainnet scale and not. - let increments = - state.validator(index)?.effective_balance / preset::EFFECTIVE_BALANCE_INCREMENT; + let increments = match attester_increments { + Some(increments) => increments, + None => { + let increments = state.validator(index)?.effective_balance + / preset::EFFECTIVE_BALANCE_INCREMENT; + attester_increments = Some(increments); + increments + } + }; let reward = (increments * base_reward_per_increment) .checked_mul(weight) .ok_or(Error::ArithmeticOverflow( diff --git a/crates/blockchain/state_transition/src/beacon/stf/electra.rs b/crates/blockchain/state_transition/src/beacon/stf/electra.rs index b7152b36..188d1a38 100644 --- a/crates/blockchain/state_transition/src/beacon/stf/electra.rs +++ b/crates/blockchain/state_transition/src/beacon/stf/electra.rs @@ -56,10 +56,11 @@ use crate::beacon::containers::shared::{ use crate::beacon::containers::{BeaconState, capella, deneb, electra, fulu}; use crate::beacon::error::{Error, Result, verify}; use crate::beacon::helpers::accessors::{ - CommitteeCache, CommitteeCacheExt, get_beacon_proposer_index, get_block_root, - get_block_root_at_slot, get_current_epoch, get_previous_epoch, get_randao_mix, + ActiveBalanceCache, ActiveBalanceCacheExt, CommitteeCache, CommitteeCacheExt, + get_beacon_proposer_index, get_block_root, get_block_root_at_slot, get_current_epoch, + get_previous_epoch, get_randao_mix, }; -use crate::beacon::helpers::altair::{add_flag, get_base_reward_per_increment, has_flag}; +use crate::beacon::helpers::altair::{add_flag, base_reward_per_increment, has_flag}; use crate::beacon::helpers::electra::{ compute_exit_epoch_and_update_churn, electra_state, get_committee_indices, get_consolidation_churn_limit, get_indexed_attestation, get_max_effective_balance, @@ -657,6 +658,7 @@ pub fn process_attestation( state: &mut BeaconState, attestation: &electra::Attestation, committees: &CommitteeCache, + active_balances: &ActiveBalanceCache, ) -> Result<()> { let data = attestation.data; let current_epoch = get_current_epoch(state); @@ -761,7 +763,8 @@ pub fn process_attestation( // finished one. That was measured on a live mainnet run: the chain actor // sat at 98% CPU inside `get_active_validator_indices`, reached through // exactly this line, and the node's clock stopped advancing. - let base_reward_per_increment = get_base_reward_per_increment(state)?; + let base_reward_per_increment = + base_reward_per_increment(active_balances.total_active_balance(state)?); let mut proposer_reward_numerator: Gwei = 0; let mut updates: Vec<(ValidatorIndex, ParticipationFlags)> = Vec::new(); @@ -776,6 +779,10 @@ pub fn process_attestation( })?; let mut new_flags: ParticipationFlags = 0; + // The attester's validator is read at most once, on the first flag + // it newly earns: each read is a tree descent, and the effective + // balance cannot change inside this read-only phase. + let mut attester_increments: Option = None; for &flag_index in &participation_flag_indices { if has_flag(current_flags, flag_index) { continue; @@ -785,8 +792,15 @@ pub fn process_attestation( // `get_base_reward(state, index)` inlined against the hoisted // per-increment value above, keeping the helper's own order of // operations so the result is bit-identical. - let increments = - state.validator(index)?.effective_balance / preset::EFFECTIVE_BALANCE_INCREMENT; + let increments = match attester_increments { + Some(increments) => increments, + None => { + let increments = state.validator(index)?.effective_balance + / preset::EFFECTIVE_BALANCE_INCREMENT; + attester_increments = Some(increments); + increments + } + }; let reward = (increments * base_reward_per_increment) .checked_mul(weight) .ok_or(Error::ArithmeticOverflow( @@ -1699,6 +1713,7 @@ pub fn process_operations( body: &electra::BeaconBlockBody, config: &Config, committees: &CommitteeCache, + active_balances: &ActiveBalanceCache, ) -> Result<()> { let deposit_requests_start_index = block_ref(state, "process_operations")?.deposit_requests_start_index(); @@ -1724,7 +1739,7 @@ pub fn process_operations( } // [Modified in Electra:EIP7549] for attestation in body.attestations.iter() { - process_attestation(state, attestation, committees)?; + process_attestation(state, attestation, committees, active_balances)?; } for deposit in body.deposits.iter() { process_deposit(state, deposit, config)?; @@ -1765,6 +1780,7 @@ pub fn process_block( config: &Config, engine: &ExecutionEngine, committees: &CommitteeCache, + active_balances: &ActiveBalanceCache, ) -> Result<()> { super::block::process_block_header( state, @@ -1777,8 +1793,8 @@ pub fn process_block( process_execution_payload(state, &block.body, config, engine)?; super::block::process_randao(state, &block.body.randao_reveal)?; super::block::process_eth1_data(state, &block.body.eth1_data)?; - process_operations(state, &block.body, config, committees)?; - super::altair::process_sync_aggregate(state, &block.body.sync_aggregate)?; + process_operations(state, &block.body, config, committees, active_balances)?; + super::altair::process_sync_aggregate(state, &block.body.sync_aggregate, active_balances)?; Ok(()) } diff --git a/crates/blockchain/state_transition/src/beacon/stf/fulu.rs b/crates/blockchain/state_transition/src/beacon/stf/fulu.rs index 9cd8d326..0913a8db 100644 --- a/crates/blockchain/state_transition/src/beacon/stf/fulu.rs +++ b/crates/blockchain/state_transition/src/beacon/stf/fulu.rs @@ -41,7 +41,9 @@ use crate::beacon::config::Config; use crate::beacon::containers::{BeaconState, deneb, electra}; use crate::beacon::error::{Result, verify}; -use crate::beacon::helpers::accessors::{CommitteeCache, get_current_epoch, get_randao_mix}; +use crate::beacon::helpers::accessors::{ + ActiveBalanceCache, CommitteeCache, get_current_epoch, get_randao_mix, +}; use crate::beacon::helpers::fulu::{fulu_state, fulu_state_ref}; use crate::beacon::primitives::{Bytes32, HashTreeRoot as _}; @@ -61,6 +63,7 @@ pub fn process_block( config: &Config, engine: &ExecutionEngine, committees: &CommitteeCache, + active_balances: &ActiveBalanceCache, ) -> Result<()> { super::block::process_block_header( state, @@ -73,8 +76,8 @@ pub fn process_block( process_execution_payload(state, &block.body, config, engine)?; super::block::process_randao(state, &block.body.randao_reveal)?; super::block::process_eth1_data(state, &block.body.eth1_data)?; - super::electra::process_operations(state, &block.body, config, committees)?; - super::altair::process_sync_aggregate(state, &block.body.sync_aggregate)?; + super::electra::process_operations(state, &block.body, config, committees, active_balances)?; + super::altair::process_sync_aggregate(state, &block.body.sync_aggregate, active_balances)?; Ok(()) } diff --git a/crates/blockchain/state_transition/src/beacon/stf/mod.rs b/crates/blockchain/state_transition/src/beacon/stf/mod.rs index 74367286..31a1486a 100644 --- a/crates/blockchain/state_transition/src/beacon/stf/mod.rs +++ b/crates/blockchain/state_transition/src/beacon/stf/mod.rs @@ -66,7 +66,7 @@ pub mod operations; use crate::beacon::containers; use crate::beacon::containers::{BeaconState, phase0}; use crate::beacon::error::{Error, Result, verify}; -use crate::beacon::helpers::accessors::{CommitteeCache, get_domain}; +use crate::beacon::helpers::accessors::{ActiveBalanceCache, CommitteeCache, get_domain}; use crate::beacon::helpers::misc::compute_signing_root; use crate::beacon::preset; use crate::beacon::primitives::{HashTreeRoot as _, Slot}; @@ -132,6 +132,7 @@ pub fn state_transition( config: &Config, engine: &ExecutionEngine, committees: &CommitteeCache, + active_balances: &ActiveBalanceCache, ) -> Result<()> { process_slots(state, signed_block.slot(), config)?; @@ -153,7 +154,14 @@ pub fn state_transition( )?; } - block::process_block(state, signed_block, config, engine, committees)?; + block::process_block( + state, + signed_block, + config, + engine, + committees, + active_balances, + )?; // As in `process_slot`: the post-state's root, checked below and needed // again in the next slot, is then computed on the state's own nodes. state.apply_pending_mutations(); diff --git a/crates/blockchain/state_transition/src/beacon/stf/operations.rs b/crates/blockchain/state_transition/src/beacon/stf/operations.rs index c6fd6229..3cbd990e 100644 --- a/crates/blockchain/state_transition/src/beacon/stf/operations.rs +++ b/crates/blockchain/state_transition/src/beacon/stf/operations.rs @@ -23,8 +23,8 @@ use crate::beacon::containers::shared::{ use crate::beacon::error::{Error, Result, verify}; use crate::beacon::fork::ForkName; use crate::beacon::helpers::accessors::{ - CommitteeCache, CommitteeCacheExt, get_beacon_proposer_index, get_current_epoch, get_domain, - get_previous_epoch, + ActiveBalanceCache, CommitteeCache, CommitteeCacheExt, get_beacon_proposer_index, + get_current_epoch, get_domain, get_previous_epoch, }; use crate::beacon::helpers::attestation::{get_indexed_attestation, is_valid_indexed_attestation}; use crate::beacon::helpers::misc::{ @@ -72,6 +72,7 @@ pub fn process_operations( voluntary_exits: &[SignedVoluntaryExit], config: &Config, committees: &CommitteeCache, + active_balances: &ActiveBalanceCache, ) -> Result<()> { // `eth1_deposit_index` only ever advances by one per processed deposit, // and `deposit_count` only ever grows, so in a correctly-derived state the @@ -109,7 +110,12 @@ pub fn process_operations( if state.fork_name() == ForkName::Phase0 { process_attestation(state, attestation, config, committees)?; } else { - crate::beacon::stf::altair::process_attestation(state, attestation, committees)?; + crate::beacon::stf::altair::process_attestation( + state, + attestation, + committees, + active_balances, + )?; } } for deposit in deposits { diff --git a/crates/blockchain/state_transition/src/metrics.rs b/crates/blockchain/state_transition/src/metrics.rs index e7da0ce8..54862e98 100644 --- a/crates/blockchain/state_transition/src/metrics.rs +++ b/crates/blockchain/state_transition/src/metrics.rs @@ -62,6 +62,24 @@ pub fn inc_committee_cache_lookups(result: &str) { .inc(); } +static LEAN_BEACON_TOTAL_ACTIVE_BALANCE_LOOKUPS_TOTAL: LazyLock = + LazyLock::new(|| { + register_int_counter_vec!( + "lean_beacon_total_active_balance_lookups_total", + "Total active balance lookups, by whether the active balance cache served them", + &["result"] + ) + .unwrap() + }); + +/// Count one `ActiveBalanceCache` lookup: `hit`, `miss` (a registry pass was +/// run and cached), or `unkeyable` (computed for one caller, not cached). +pub fn inc_total_active_balance_lookups(result: &str) { + LEAN_BEACON_TOTAL_ACTIVE_BALANCE_LOOKUPS_TOTAL + .with_label_values(&[result]) + .inc(); +} + 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/operations.rs b/crates/blockchain/state_transition/tests/beacon_spec/operations.rs index e79d29f2..883d54f8 100644 --- a/crates/blockchain/state_transition/tests/beacon_spec/operations.rs +++ b/crates/blockchain/state_transition/tests/beacon_spec/operations.rs @@ -91,7 +91,7 @@ use ethlambda_state_transition::beacon::config::Config; use ethlambda_state_transition::beacon::containers::{ BeaconState, altair, bellatrix, capella, deneb, electra, phase0, shared, }; -use ethlambda_state_transition::beacon::helpers::accessors::CommitteeCache; +use ethlambda_state_transition::beacon::helpers::accessors::{ActiveBalanceCache, CommitteeCache}; use ethlambda_state_transition::beacon::primitives::HashTreeRoot as _; use ethlambda_state_transition::beacon::stf::altair as altair_stf; use ethlambda_state_transition::beacon::stf::bellatrix as bellatrix_stf; @@ -127,6 +127,7 @@ fn apply( // calls fits on one line; a case processes a single attestation, so there // is nothing for the cache to share between calls. let committees = CommitteeCache::default(); + let active_balances = ActiveBalanceCache::default(); let outcome = match handler { // Every fork's own `BeaconBlock` carries a different concrete body // (altair's adds a sync aggregate, bellatrix's an execution payload, @@ -216,20 +217,20 @@ fn apply( // so altair's own function serves both. ForkName::Altair | ForkName::Bellatrix | ForkName::Capella => { let attestation: phase0::Attestation = case.ssz("attestation"); - altair_stf::process_attestation(state, &attestation, &committees) + altair_stf::process_attestation(state, &attestation, &committees, &active_balances) } // Deneb widens the inclusion window and the timely-target // condition (EIP-7045); the container is still phase0's. ForkName::Deneb => { let attestation: phase0::Attestation = case.ssz("attestation"); - deneb_stf::process_attestation(state, &attestation, &committees) + deneb_stf::process_attestation(state, &attestation, &committees, &active_balances) } // Electra reshapes the container itself (EIP-7549's // `committee_bits`), and fulu's specification makes no further // change to either the container or the function. ForkName::Electra | ForkName::Fulu => { let attestation: electra::Attestation = case.ssz("attestation"); - electra_stf::process_attestation(state, &attestation, &committees) + electra_stf::process_attestation(state, &attestation, &committees, &active_balances) } ForkName::Lean => lean_is_not_a_fixture_fork("attestation"), }, @@ -307,7 +308,7 @@ fn apply( // further fork match. "sync_aggregate" => { let sync_aggregate: altair::SyncAggregate = case.ssz("sync_aggregate"); - altair_stf::process_sync_aggregate(state, &sync_aggregate) + altair_stf::process_sync_aggregate(state, &sync_aggregate, &active_balances) } // The fixture's operation file is a whole `BeaconBlockBody`, not a // bare payload: checking one needs fields that live on the body diff --git a/crates/blockchain/state_transition/tests/beacon_spec/sanity.rs b/crates/blockchain/state_transition/tests/beacon_spec/sanity.rs index eb18e67a..f097a598 100644 --- a/crates/blockchain/state_transition/tests/beacon_spec/sanity.rs +++ b/crates/blockchain/state_transition/tests/beacon_spec/sanity.rs @@ -17,7 +17,7 @@ use std::sync::Arc; use ethlambda_state_transition::beacon::config::Config; use ethlambda_state_transition::beacon::containers::{BeaconState, SignedBeaconBlock}; -use ethlambda_state_transition::beacon::helpers::accessors::CommitteeCache; +use ethlambda_state_transition::beacon::helpers::accessors::{ActiveBalanceCache, CommitteeCache}; use ethlambda_state_transition::beacon::stf::{self, ExecutionEngine}; use libtest_mimic::Trial; @@ -50,6 +50,7 @@ fn apply_blocks(case: &Case, state: &mut BeaconState, config: &Config) -> Result // One cache across the case's blocks, as the node holds one across its // imports, so consecutive blocks of an epoch share its shuffling. let committees = CommitteeCache::default(); + let active_balances = ActiveBalanceCache::default(); for index in 0..meta.blocks_count { let bytes = case.ssz_bytes_indexed("blocks", index); @@ -60,8 +61,16 @@ fn apply_blocks(case: &Case, state: &mut BeaconState, config: &Config) -> Result state.latest_block_header_mut().state_root = root; } - stf::state_transition(state, &block, true, config, &engine, &committees) - .map_err(|err| format!("block {index} rejected: {err:?}"))?; + stf::state_transition( + state, + &block, + true, + config, + &engine, + &committees, + &active_balances, + ) + .map_err(|err| format!("block {index} rejected: {err:?}"))?; previous_state_root = Some(block.state_root()); } diff --git a/crates/blockchain/state_transition/tests/beacon_spec/transition.rs b/crates/blockchain/state_transition/tests/beacon_spec/transition.rs index 7e4dd562..dc6d5a2f 100644 --- a/crates/blockchain/state_transition/tests/beacon_spec/transition.rs +++ b/crates/blockchain/state_transition/tests/beacon_spec/transition.rs @@ -69,7 +69,7 @@ use std::sync::Arc; use ethlambda_state_transition::beacon::ForkName; use ethlambda_state_transition::beacon::config::Config; use ethlambda_state_transition::beacon::containers::{BeaconState, SignedBeaconBlock}; -use ethlambda_state_transition::beacon::helpers::accessors::CommitteeCache; +use ethlambda_state_transition::beacon::helpers::accessors::{ActiveBalanceCache, CommitteeCache}; use ethlambda_state_transition::beacon::helpers::misc::compute_start_slot_at_epoch; use ethlambda_state_transition::beacon::primitives::Epoch; use ethlambda_state_transition::beacon::stf::{self, ExecutionEngine}; @@ -138,6 +138,7 @@ fn apply_blocks( // One cache across the case's blocks, as the node holds one across its // imports, so consecutive blocks of an epoch share its shuffling. let committees = CommitteeCache::default(); + let active_balances = ActiveBalanceCache::default(); for index in 0..meta.blocks_count { let fork = if index < first_post_fork_index { @@ -154,8 +155,16 @@ fn apply_blocks( state.latest_block_header_mut().state_root = root; } - stf::state_transition(state, &block, true, &config, &engine, &committees) - .map_err(|err| format!("block {index} rejected: {err:?}"))?; + stf::state_transition( + state, + &block, + true, + &config, + &engine, + &committees, + &active_balances, + ) + .map_err(|err| format!("block {index} rejected: {err:?}"))?; previous_state_root = Some(block.state_root()); } diff --git a/crates/storage/src/active_balance_cache.rs b/crates/storage/src/active_balance_cache.rs new file mode 100644 index 00000000..fecea5ee --- /dev/null +++ b/crates/storage/src/active_balance_cache.rs @@ -0,0 +1,171 @@ +//! The total-active-balance cache: one [`Gwei`] per epoch, shared across every +//! caller asking the same epoch's total of a state that agrees on the block +//! that fixes it. +//! +//! A sibling of [`crate::committee_cache::CommitteeCache`], held by the +//! `Store` for the same reasons (see that module's documentation): one table +//! shared by the chain actor's state transition and fork choice, and by any +//! other task that holds the store. A caller with no `Store` of its own (a +//! spec runner, a one-off lookup, a unit test, block production's scratch +//! state) holds a fresh [`ActiveBalanceCache::default`] for as long as that +//! work lasts instead. +//! +//! This module knows nothing about a `BeaconState`: it is handed an +//! [`ActiveBalanceKey`] and, on a miss, runs a closure that returns the total. +//! The consensus logic that derives the key from a state, and the one-pass +//! sum that computes the value, lives in `ethlambda-state-transition`'s +//! `beacon::helpers::accessors`, in the `ActiveBalanceCacheExt` trait this +//! type implements there; that crate depends on this one, not the other way +//! around. +//! +//! Unlike the committee cache there is no single-flight machinery: the value +//! is one registry pass, cheap next to a whole-epoch shuffle, so two callers +//! racing on the same missing key may both compute it. They compute the same +//! number (the key fixes it), so the second insert is a no-op. + +use std::sync::Mutex; + +use ethlambda_types::beacon::primitives::{Epoch, Gwei, Root}; + +pub use crate::committee_cache::Lookup; + +/// What pins the total active balance a state names for an epoch: the epoch +/// itself, and the last block root that could still have changed it. +/// +/// Opaque to this module: `epoch` and `decision_root` are read only for +/// equality and, for eviction, ordering by `epoch`. +/// `ethlambda-state-transition`'s `beacon::helpers::accessors::active_balance_key` +/// derives one from a state and documents why that root is sound; note that +/// it is the last slot of `epoch - 1`, one epoch later than a +/// [`crate::ShufflingKey`]'s, so the two key types are kept apart rather than +/// shared. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub struct ActiveBalanceKey { + pub epoch: Epoch, + pub decision_root: Root, +} + +/// How many distinct totals stay resident in an [`ActiveBalanceCache`]. +/// +/// A block asks for one epoch's total; fork choice and sibling branches ask +/// for a few more. Entries are 8 bytes of value plus a 40 byte key, so room +/// for many concurrent forks and several epochs each costs next to nothing, +/// and the bound is only there so a long-lived store cannot grow the table +/// without limit. +const ACTIVE_BALANCE_CACHE_CAPACITY: usize = 32; + +/// Totals shared across the calls asking for the same epoch's total active +/// balance, so a block derives it once instead of once per attestation. +/// +/// Internally synchronized, so every method takes `&self`. The lock is held +/// only for the table lookup or insert, never while the value is computed. +/// +/// # Eviction +/// +/// A miss on a full table drops the entry for the lowest epoch (the earliest +/// inserted among equals): an older epoch's total is less likely to be asked +/// for again than a newer one's, whichever branch either belongs to. Every +/// lookup is counted in `lean_beacon_total_active_balance_lookups_total`. +#[derive(Debug, Default)] +pub struct ActiveBalanceCache { + entries: Mutex>, +} + +impl ActiveBalanceCache { + /// `key`'s total, running `compute` on a miss and remembering the result. + /// + /// `compute` runs outside the lock. Concurrent misses on one key may each + /// run it; see the module documentation for why that is acceptable. + pub fn get_or_compute( + &self, + key: ActiveBalanceKey, + compute: impl FnOnce() -> Gwei, + ) -> (Gwei, Lookup) { + if let Some(total) = self.get(key) { + return (total, Lookup::Hit); + } + let total = compute(); + let mut entries = self.entries.lock().unwrap(); + if !entries.iter().any(|(entry_key, _)| *entry_key == key) { + if entries.len() >= ACTIVE_BALANCE_CACHE_CAPACITY { + let victim = entries + .iter() + .enumerate() + .min_by_key(|(_, (entry_key, _))| entry_key.epoch) + .map(|(position, _)| position); + if let Some(position) = victim { + entries.remove(position); + } + } + entries.push((key, total)); + } + (total, Lookup::Miss) + } + + /// The resident total for `key`, if any. + pub fn get(&self, key: ActiveBalanceKey) -> Option { + self.entries + .lock() + .unwrap() + .iter() + .find(|(entry_key, _)| *entry_key == key) + .map(|(_, total)| *total) + } +} + +#[cfg(test)] +mod tests { + use std::cell::Cell; + + use super::*; + + fn key(epoch: Epoch, decision_root: u8) -> ActiveBalanceKey { + ActiveBalanceKey { + epoch, + decision_root: Root::repeat_byte(decision_root), + } + } + + #[test] + fn a_repeat_lookup_is_served_from_the_cache() { + let cache = ActiveBalanceCache::default(); + let (first, first_lookup) = cache.get_or_compute(key(1, 0), || 7); + let (second, second_lookup) = cache.get_or_compute(key(1, 0), || panic!("rebuilt")); + assert_eq!((first, first_lookup), (7, Lookup::Miss)); + assert_eq!((second, second_lookup), (7, Lookup::Hit)); + } + + #[test] + fn keys_with_different_roots_or_epochs_are_distinct_entries() { + let cache = ActiveBalanceCache::default(); + cache.get_or_compute(key(1, 1), || 10); + let computed = Cell::new(0); + for (k, value) in [(key(1, 2), 20), (key(2, 1), 30)] { + let (total, lookup) = cache.get_or_compute(k, || { + computed.set(computed.get() + 1); + value + }); + assert_eq!((total, lookup), (value, Lookup::Miss)); + } + assert_eq!(computed.get(), 2); + assert_eq!(cache.get(key(1, 1)), Some(10)); + } + + #[test] + fn a_miss_on_a_full_cache_evicts_the_lowest_epoch() { + let cache = ActiveBalanceCache::default(); + // Newest first, so eviction by epoch and by insertion order disagree. + for epoch in (0..ACTIVE_BALANCE_CACHE_CAPACITY as Epoch).rev() { + cache.get_or_compute(key(epoch, 0), || epoch); + } + cache.get_or_compute(key(1000, 0), || 1000); + + assert_eq!(cache.get(key(0, 0)), None); + assert_eq!(cache.get(key(1, 0)), Some(1)); + assert_eq!(cache.get(key(1000, 0)), Some(1000)); + assert_eq!( + cache.entries.lock().unwrap().len(), + ACTIVE_BALANCE_CACHE_CAPACITY + ); + } +} diff --git a/crates/storage/src/lib.rs b/crates/storage/src/lib.rs index 8f1fc61b..6aca1cc0 100644 --- a/crates/storage/src/lib.rs +++ b/crates/storage/src/lib.rs @@ -1,3 +1,4 @@ +mod active_balance_cache; mod api; pub mod backend; mod beacon_state_delta; @@ -9,6 +10,7 @@ mod state_diff; mod state_writer; mod store; +pub use active_balance_cache::{ActiveBalanceCache, ActiveBalanceKey}; pub use api::{ALL_TABLES, StorageBackend, StorageReadView, StorageWriteBatch, Table}; pub use committee_cache::{CommitteeCache, Lookup, ShufflingKey}; /// Error type returned by the fallible [`Store`] operations, exported so diff --git a/crates/storage/src/store.rs b/crates/storage/src/store.rs index 96addfbd..2dc2155a 100644 --- a/crates/storage/src/store.rs +++ b/crates/storage/src/store.rs @@ -4,6 +4,7 @@ use std::sync::{Arc, LazyLock, Mutex}; use lru::LruCache; +use crate::active_balance_cache::ActiveBalanceCache; use crate::api::{StorageBackend, StorageReadView, StorageWriteBatch, Table}; use crate::committee_cache::CommitteeCache; use crate::error::Error; @@ -882,6 +883,13 @@ pub struct Store { /// /// Always empty on lean, which has no beacon committees. committee_cache: Arc, + /// The total active balance per epoch, shared the way + /// [`Self::committee_cache`] is: by the chain actor's state transition + /// and anything else holding a `Store` clone. See + /// [`ActiveBalanceCache`], and `ethlambda-state-transition`'s + /// `beacon::helpers::accessors::ActiveBalanceCacheExt` for the + /// state-aware half. Always empty on lean. + active_balance_cache: Arc, /// Beacon fork-choice scratch. Empty and untouched on a lean chain. pub(crate) beacon: Arc>, /// The background writer, joined when the last clone of this `Store` @@ -1512,6 +1520,7 @@ impl Store { state_cache, pending_states, committee_cache: Arc::new(CommitteeCache::default()), + active_balance_cache: Arc::new(ActiveBalanceCache::default()), beacon: Default::default(), state_writer, } @@ -2658,6 +2667,13 @@ impl Store { Arc::clone(&self.committee_cache) } + /// The total-active-balance cache shared by every clone of this store. + /// + /// An owned `Arc` for the same reason as [`Self::committee_cache`]. + pub fn active_balance_cache(&self) -> Arc { + Arc::clone(&self.active_balance_cache) + } + /// Returns whether a state is available for the given block root. /// /// True if `pending_states` holds the state, a snapshot exists, or the @@ -5504,6 +5520,24 @@ mod tests { assert!(Arc::ptr_eq(&first, &second)); } + #[test] + fn the_active_balance_cache_is_shared_across_store_clones() { + let store = beacon_test_store(Arc::new(InMemoryBackend::new())); + let clone = store.clone(); + let key = crate::ActiveBalanceKey { + epoch: 1, + decision_root: H256::from([1u8; 32]), + }; + + let (first, first_lookup) = clone.active_balance_cache().get_or_compute(key, || 42); + let (second, second_lookup) = store + .active_balance_cache() + .get_or_compute(key, || panic!("should not recompute")); + + assert_eq!((first, first_lookup), (42, Lookup::Miss)); + assert_eq!((second, second_lookup), (42, Lookup::Hit)); + } + #[test] fn the_state_cache_is_bounded() { let store = beacon_test_store(Arc::new(InMemoryBackend::new())); diff --git a/docs/beacon_stf.md b/docs/beacon_stf.md index 17ebe686..790569d6 100644 --- a/docs/beacon_stf.md +++ b/docs/beacon_stf.md @@ -283,15 +283,29 @@ lists lighthouse keeps its state in. The access pattern matters. `state.validator(i)` and `balances()[i]` are tree descents, cheap next to a hash but far from an array index, and they add up when a helper calls them once per validator: -`get_total_active_balance` builds the active-index `Vec` and then reads every -index back, and runs several times per block (once per attestation through -`get_base_reward_per_increment`, once per execution request through the churn -limits). In the 2026-09-28 import profile, those per-index reads and the -repeated whole-registry scans were the largest cost left after hashing. A loop -over the registry should walk `validators().iter()`, zipped with -`balances().iter()` where it needs both. The total active balance is the obvious -candidate for computing once per epoch rather than per call, once it is shown -that no block operation changes it mid-epoch. +`get_total_active_balance` used to build the active-index `Vec` and then read +every index back, and runs several times per block (once per attestation through +`get_base_reward_per_increment`, twice per sync aggregate, once per execution +request through the churn limits). In the 2026-09-28 import profile, those +per-index reads and the repeated whole-registry scans were the largest cost left +after hashing. A loop over the registry should walk `validators().iter()`, +zipped with `balances().iter()` where it needs both; `get_total_active_balance` +now does, in one pass. + +Block processing does not even pay that pass more than once per epoch. The +total is cached in `ActiveBalanceCache` (`crates/storage/src/active_balance_cache.rs`), +held by the `Store` beside the committee cache and consulted through +`ActiveBalanceCacheExt::total_active_balance`. The key is the epoch plus the +block root at the last slot of the epoch before it: effective balances are +written only by that epoch's `process_effective_balance_updates`, and every +activation or exit lands at least `MAX_SEED_LOOKAHEAD` epochs ahead, so no block +inside the epoch can move the total, and two states agreeing on that root agree +on it. The cache lives outside the state, so it adds no field to `BeaconState` +and no SSZ concern; callers with no store (spec runners, block production, +tests) hold a fresh `ActiveBalanceCache::default()`. Epoch processing, churn +helpers and fork choice keep calling the uncached one-pass +`get_total_active_balance`. In debug builds every hit is cross-checked against a +fresh computation. ## Macros and traits diff --git a/docs/metrics.md b/docs/metrics.md index bf215a3c..3a05e90a 100644 --- a/docs/metrics.md +++ b/docs/metrics.md @@ -422,6 +422,22 @@ 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 Total Active Balance Cache + +The total active balance (the divisor of every attestation's base reward and of +the sync aggregate's rewards) is one registry pass, and an ordinary block asks +for it about nine times. `ActiveBalanceCache` in the `Store` keeps it per epoch, +keyed by the epoch and the block root that fixes it. This is ethlambda-specific, +not part of the leanMetrics spec. + +| Name | Type | Usage | Sample collection event | Labels | +|------|------|-------|-------------------------|--------| +| `lean_beacon_total_active_balance_lookups_total` | Counter | Total active balance lookups, by whether the cache served them | On every `ActiveBalanceCacheExt::total_active_balance` call: block processing's attestations and sync aggregate | result=hit,miss,unkeyable | + +Expect one miss per epoch per chain, and hits otherwise. `unkeyable` is the +genesis state asking about its own epoch, and should be zero on a +checkpoint-synced follower. + ### Beacon Pubkey Cache Every BLS signature check `ethlambda beacon` runs (block import, fork choice,