Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
297 changes: 266 additions & 31 deletions crates/blockchain/state_transition/src/beacon/fork_choice.rs

Large diffs are not rendered by default.

34 changes: 34 additions & 0 deletions crates/blockchain/state_transition/src/metrics.rs
Original file line number Diff line number Diff line change
Expand Up @@ -62,6 +62,40 @@ pub fn inc_committee_cache_lookups(result: &str) {
.inc();
}

static LEAN_BEACON_JUSTIFIED_BALANCES_LOOKUPS_TOTAL: LazyLock<IntCounterVec> = 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<Histogram> = 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<Histogram> = LazyLock::new(|| {
register_histogram!(
"lean_state_transition_time_seconds",
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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(())
}

Expand Down
69 changes: 68 additions & 1 deletion crates/common/types/src/beacon/fork_choice.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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).
Expand Down Expand Up @@ -100,12 +101,78 @@ pub struct PayloadStatusV1 {
pub validation_error: Option<String>,
}

/// 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
Expand Down
107 changes: 95 additions & 12 deletions crates/storage/src/store.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
},
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -666,7 +672,14 @@ pub(crate) struct BeaconScratch {
pub(crate) proposer_boost_root: H256,
pub(crate) block_timeliness: HashMap<H256, bool>,
pub(crate) equivocating_indices: HashSet<u64>,
pub(crate) latest_messages: HashMap<u64, LatestMessage>,
/// 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<Option<LatestMessage>>,
/// 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<Arc<JustifiedBalances>>,
pub(crate) pow_blocks: HashMap<H256, PowBlock>,
pub(crate) unrealized_justifications: HashMap<H256, BeaconCheckpoint>,
/// Beacon roots imported on an execution client's `NOT_VALIDATED` answer,
Expand Down Expand Up @@ -3180,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
Expand All @@ -3203,13 +3228,46 @@ 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);
}
}

/// 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<Arc<JustifiedBalances>> {
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<JustifiedBalances>) {
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<PowBlock> {
Expand Down Expand Up @@ -6645,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();
Expand Down
21 changes: 21 additions & 0 deletions docs/beacon_stf.md
Original file line number Diff line number Diff line change
Expand Up @@ -293,6 +293,27 @@ 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.

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
Expand Down
19 changes: 19 additions & 0 deletions docs/metrics.md
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
Loading