From c41deeb0a13f34e1fe7e033290fab612ef96ff34 Mon Sep 17 00:00:00 2001 From: Pablo Deymonnaz Date: Tue, 29 Sep 2026 15:47:23 -0300 Subject: [PATCH] Serve POST /eth/v1/validator/liveness/{epoch} from ethlambda beacon, the endpoint a validator client's doppelganger protection calls before signing. A validator is live in an epoch if the head state credits it for that epoch (a non-zero participation byte, covering what blocks have already included) or if this node observed it act in that epoch, which covers what no block has included yet. The observations live in a new ObservedLiveness on the Store, one bitset per epoch for the newest three, because three places write them and all three already hold a Store clone: P2P when it accepts a gossip aggregate (its aggregator and every attester its signature verified) or subnet attestation, the chain actor when it imports a beacon block (its proposer, so range-synced and self-published blocks count), and the RPC for attestations and aggregates submitted through pool/attestations and aggregate_and_proofs, since gossip never delivers a node its own messages. The endpoint answers the store clock's previous, current and next epoch and is a 400 for any other epoch or an unknown index, and a 503 while syncing. --- crates/blockchain/src/lib.rs | 9 + crates/net/p2p/src/beacon/verdict.rs | 84 ++++++++++ crates/net/rpc/src/beacon/pool.rs | 72 +++++++- crates/net/rpc/src/beacon/validator.rs | 218 ++++++++++++++++++++++++- crates/storage/src/lib.rs | 2 + crates/storage/src/liveness.rs | 142 ++++++++++++++++ crates/storage/src/store.rs | 22 +++ docs/rpc.md | 12 ++ 8 files changed, 558 insertions(+), 3 deletions(-) create mode 100644 crates/storage/src/liveness.rs diff --git a/crates/blockchain/src/lib.rs b/crates/blockchain/src/lib.rs index 0044889d..994d0f67 100644 --- a/crates/blockchain/src/lib.rs +++ b/crates/blockchain/src/lib.rs @@ -2772,6 +2772,15 @@ impl BlockChainServer { "Block imported successfully" ); + // The proposer acted in this block's epoch, for the Beacon + // API's liveness endpoint. Recorded here rather than where + // gossip accepts a block, so a block that arrived by range + // sync or through this node's own API counts too. + if self.store.chain() == Chain::Beacon { + let epoch = ethlambda_types::beacon::signing::compute_epoch_at_slot(slot); + self.store.observed_liveness().record(epoch, proposer); + } + // Recover per-attestation single-message aggregates from the // block's merged multi-message aggregate and fold them into // the local pool. `Some` only for a lean block imported while diff --git a/crates/net/p2p/src/beacon/verdict.rs b/crates/net/p2p/src/beacon/verdict.rs index 6d7c3136..2d335ca8 100644 --- a/crates/net/p2p/src/beacon/verdict.rs +++ b/crates/net/p2p/src/beacon/verdict.rs @@ -147,7 +147,14 @@ impl Validated { /// [`pool_aggregator_attestation`]. An accepted aggregate goes into the /// pool too, whatever else happens to it, so block production can pack /// other nodes' votes; see [`pool_gossip_aggregate`]. + /// + /// Every accepted aggregate or subnet attestation also marks the + /// validators it names as live for its target epoch; see + /// [`record_liveness`]. fn forward(self, server: &P2PServer, received_at: Instant, outcome: Outcome) { + if outcome == Outcome::Accept { + record_liveness(server, &self); + } if let Self::Aggregate { aggregate, .. } = &self && outcome == Outcome::Accept { @@ -201,6 +208,33 @@ impl Validated { } } +/// Mark the validators an accepted aggregate or subnet attestation names as +/// live for its target epoch, for `/eth/v1/validator/liveness`. +/// +/// An aggregate names its aggregator and every attester its aggregate +/// signature verified, which is what makes an aggregate worth more here than +/// the attestation subnets alone: this node joins only a few of those. A +/// block's proposer is recorded where the chain actor imports it instead, so +/// a block that reached it by range sync or through this node's own API +/// counts too. +fn record_liveness(server: &P2PServer, object: &Validated) { + let observed = server.store.observed_liveness(); + match object { + Validated::Aggregate { + aggregate, + attesting_indices, + } => { + let (epoch, _root) = aggregate.target(); + let aggregator = std::iter::once(aggregate.aggregator_index()); + observed.record_all(epoch, aggregator.chain(attesting_indices.iter().copied())); + } + Validated::Attestation { attestation, .. } => { + observed.record(attestation.data.target.epoch, attestation.attester_index); + } + Validated::Block { .. } | Validated::Column(_) => {} + } +} + /// Pool an accepted gossip aggregate for block production to pack. /// /// Pooled here, on `Accept`, because this is where all three of its @@ -737,6 +771,56 @@ mod tests { ); } + /// An accepted aggregate names its aggregator and every attester its + /// signature verified; all of them count as live for its target epoch. + #[tokio::test] + async fn an_accepted_aggregate_marks_its_aggregator_and_attesters_live() { + let server = unconnected_beacon_server(Config::mainnet(), 0).await; + let slot = 40; + Validated::Aggregate { + aggregate: Box::new(electra_aggregate(slot, 7)), + attesting_indices: vec![11, 12], + } + .forward(&server, Instant::now(), Outcome::Accept); + + let observed = server.store.observed_liveness(); + let epoch = slot / 32; + for validator in [7, 11, 12] { + assert!(observed.is_live(epoch, validator), "{validator}"); + } + assert!(!observed.is_live(epoch, 13)); + } + + #[tokio::test] + async fn an_accepted_subnet_attestation_marks_its_attester_live() { + let server = unconnected_beacon_server(Config::mainnet(), 0).await; + Validated::Attestation { + attestation: Box::new(electra_attestation(40, 5)), + subnet_id: 0, + } + .forward(&server, Instant::now(), Outcome::Accept); + assert!(server.store.observed_liveness().is_live(40 / 32, 5)); + } + + /// Liveness rests on the same verdict the pool does: an object that was + /// not accepted proves nothing about the validators it names. + #[tokio::test] + async fn an_object_that_was_not_accepted_marks_nobody_live() { + let server = unconnected_beacon_server(Config::mainnet(), 0).await; + Validated::Aggregate { + aggregate: Box::new(electra_aggregate(40, 7)), + attesting_indices: vec![11], + } + .forward( + &server, + Instant::now(), + Outcome::Ignore(IgnoreReason::Overloaded), + ); + let observed = server.store.observed_liveness(); + assert!(!observed.is_live(1, 7)); + assert!(!observed.is_live(1, 11)); + } + /// Only `Accept` means the signatures were verified; anything else must /// stay out of the pool, since one unverified attestation fails the /// whole block it is packed into. diff --git a/crates/net/rpc/src/beacon/pool.rs b/crates/net/rpc/src/beacon/pool.rs index 86809218..aced62b8 100644 --- a/crates/net/rpc/src/beacon/pool.rs +++ b/crates/net/rpc/src/beacon/pool.rs @@ -114,6 +114,11 @@ async fn post_pool_attestations( // messages, so without this an aggregator served by this node would // be missing its own validator client's votes. let published = checked.and_then(|checked| { + // Live for liveness too: gossip never delivers this node its own + // validator clients' messages, so this is where it sees them. + store + .observed_liveness() + .record(attestation.data.target.epoch, validator); pool.lock().expect("attestation pool lock poisoned").insert( &attestation, checked.committee_position, @@ -301,13 +306,19 @@ async fn post_aggregate_and_proofs( let aggregator = aggregate.aggregator_index(); let seen = aggregate::SeenAggregates::new(capacity, capacity); let checked = aggregate::cheap_checks(&seen, &store, &aggregate, now_ms) - .and_then(|()| aggregate::stateful_checks(&store, &aggregate).map(|_| ())); + .and_then(|()| aggregate::stateful_checks(&store, &aggregate)); let published = checked .map_err(|outcome: Outcome| { warn!(%slot, aggregator, ?outcome, "Refused a submitted aggregate"); "aggregate failed validation" }) - .and_then(|()| { + .and_then(|attesting_indices| { + // The aggregator and every attester its signature verified are + // live, as when P2P accepts a gossip aggregate. + let (epoch, _root) = aggregate.target(); + store + .observed_liveness() + .record_all(epoch, std::iter::once(aggregator).chain(attesting_indices)); // Recorded for block production, which packs the aggregates // this node has validated. pool.lock() @@ -728,6 +739,63 @@ mod tests { assert!(fixture.network.aggregates.lock().unwrap().is_empty()); } + /// Gossip never delivers a node its own validator clients' messages, so a + /// submission through this API is where the node sees them act. + #[tokio::test] + async fn a_published_attestation_marks_its_attester_live() { + let fixture = fixture(); + let attestation = attestation(&fixture, 0, 0); + let epoch = attestation.data.target.epoch; + submit(&fixture, std::slice::from_ref(&attestation)).await; + let observed = fixture.store.observed_liveness(); + assert!(observed.is_live(epoch, attestation.attester_index)); + } + + #[tokio::test] + async fn a_refused_attestation_marks_nobody_live() { + let fixture = fixture(); + let mut forged = attestation(&fixture, 0, 0); + forged.signature = attestation(&fixture, 0, 1).signature; + let epoch = forged.data.target.epoch; + submit(&fixture, std::slice::from_ref(&forged)).await; + assert!( + !fixture + .store + .observed_liveness() + .is_live(epoch, forged.attester_index) + ); + } + + /// The aggregate is built on one node and submitted to a fresh one, so its + /// attesters can only have been marked live by the aggregate itself, not + /// by their own votes. + #[tokio::test] + async fn a_published_aggregate_marks_its_aggregator_and_attesters_live() { + let builder = fixture(); + let slot = builder.state.slot(); + let committee = get_beacon_committee(&builder.state, slot, 0).unwrap(); + let votes: Vec = (0..committee.len()) + .map(|position| attestation(&builder, 0, position)) + .collect(); + submit(&builder, &votes).await; + let aggregate = builder + .pool + .lock() + .unwrap() + .aggregate(votes[0].data.hash_tree_root(), slot, 0) + .unwrap(); + + let fresh = fixture(); + let signed = signed_aggregate(&fresh, committee[0], aggregate); + let (status, json) = submit_aggregates(&fresh, std::slice::from_ref(&signed)).await; + assert_eq!(status, StatusCode::OK, "{json}"); + let epoch = compute_epoch_at_slot(slot); + let observed = fresh.store.observed_liveness(); + for validator in &committee { + assert!(observed.is_live(epoch, *validator), "{validator}"); + } + } + #[tokio::test] async fn a_pre_electra_fork_header_is_refused() { let fixture = fixture(); diff --git a/crates/net/rpc/src/beacon/validator.rs b/crates/net/rpc/src/beacon/validator.rs index 4b740688..7b218c74 100644 --- a/crates/net/rpc/src/beacon/validator.rs +++ b/crates/net/rpc/src/beacon/validator.rs @@ -32,7 +32,7 @@ use serde::{Deserialize, Serialize}; use ethlambda_network_api::RpcToP2PRef; use ethlambda_state_transition::beacon::{ - fork_choice::checkpoint_state, + fork_choice::{checkpoint_state, get_current_store_epoch}, gossip::attestation::compute_subnet_for_attestation, helpers::accessors::{CommitteeCacheExt as _, get_block_root_at_slot}, helpers::altair::compute_sync_committee_period, @@ -54,6 +54,7 @@ pub(crate) fn routes() -> Router { "/eth/v1/validator/duties/sync/{epoch}", post(post_sync_duties), ) + .route("/eth/v1/validator/liveness/{epoch}", post(post_liveness)) .route( "/eth/v1/validator/attestation_data", get(get_attestation_data), @@ -163,6 +164,94 @@ fn sync_duties( })) } +#[derive(Debug, Serialize)] +struct Liveness { + #[serde(with = "ethlambda_types::beacon::serde_helpers::quoted_or_bare")] + index: ValidatorIndex, + is_live: bool, +} + +/// `POST /eth/v1/validator/liveness/{epoch}`: whether this node saw each +/// validator act in `epoch`, which is what a validator client's doppelganger +/// protection asks before it signs anything. +/// +/// The Beacon API leaves the source to the node's own view. A validator is +/// live if either: +/// - the head state credits it for `epoch` (a non-zero participation byte), +/// which covers everything already included on chain; or +/// - the node observed it act in `epoch`: an accepted gossip aggregate or +/// subnet attestation, an imported block's proposer, or a submission +/// through this API. That covers what no block has included yet, most of +/// the current epoch. See [`ethlambda_storage::ObservedLiveness`]. +/// +/// Answered for the store clock's previous, current and next epoch; the next +/// one is always `false`, and is accepted because a doppelganger check made +/// at an epoch boundary can land on it. Anything else, and an index outside +/// the head state's registry, is a `400`; the node syncing is a `503`. +async fn post_liveness( + Path(epoch): Path, + State(store): State, + Extension(sync_status): Extension, + Json(indices): Json>, +) -> Response { + if sync_status.get() == SyncStatus::Syncing { + return ApiError::ServiceUnavailable("the node is syncing").into_response(); + } + match liveness(&store, &epoch, &indices) { + Ok(body) => crate::json_response(body), + Err(err) => err.into_response(), + } +} + +fn liveness(store: &Store, epoch: &str, indices: &[String]) -> Result { + let epoch = parse_epoch(epoch)?; + let indices = indices + .iter() + .map(|index| index.parse::()) + .collect::, _>>() + .map_err(|_| ApiError::BadRequest("invalid validator index"))?; + + let current = get_current_store_epoch(store, &store.config()); + if epoch + 1 < current || epoch > current + 1 { + return Err(ApiError::BadRequest( + "epoch is not the previous, current or next epoch", + )); + } + + let (_head_root, state) = head(store)?; + let state_epoch = compute_epoch_at_slot(state.slot()); + // The head state's flags for `epoch`, if it keeps them: its own epoch's + // and the one before. `None` before altair, which keeps no flags. + let participation = state + .altair_validator_lists() + .ok() + .and_then(|(previous, current, _)| { + if epoch == state_epoch { + Some(current) + } else if epoch + 1 == state_epoch { + Some(previous) + } else { + None + } + }); + + let observed = store.observed_liveness(); + let mut data = Vec::with_capacity(indices.len()); + for index in indices { + state + .validator(index) + .map_err(|_| ApiError::BadRequest("unknown validator index"))?; + let credited = participation + .and_then(|flags| flags.get(index as usize)) + .is_some_and(|flags| *flags != 0); + data.push(Liveness { + index, + is_live: credited || observed.is_live(epoch, index), + }); + } + Ok(serde_json::json!({ "data": data })) +} + /// One entry of `beacon_committee_subscriptions`. Parsed so a malformed body /// is refused, though `validator_index` is never read. #[derive(Debug, Deserialize)] @@ -1052,4 +1141,131 @@ mod tests { assert_eq!(status, StatusCode::SERVICE_UNAVAILABLE); } } + + // --- liveness -------------------------------------------------------- + + mod liveness { + use super::*; + + /// The epoch every test's state and store clock sit in: far enough + /// from genesis that the epoch two before it exists. + const EPOCH: u64 = 5; + + /// A fulu state at [`EPOCH`]'s first slot in which validator 2 is + /// credited for this epoch and validator 3 for the previous one. + fn credited_state() -> BeaconState { + let mut state = fulu_state(); + let BeaconState::Fulu(fulu) = &mut state else { + unreachable!("built as fulu") + }; + fulu.slot = compute_start_slot_at_epoch(EPOCH); + fulu.current_epoch_participation[2] = 0b001; + fulu.previous_epoch_participation[3] = 0b111; + state + } + + /// `state`'s store, its clock moved to `state`'s slot, as the chain + /// actor's tick keeps it. + fn store_for(state: BeaconState) -> Store { + let slot = state.slot(); + let (mut store, _root) = beacon_store_at(state); + let config = store.config(); + let now = config.genesis_time_ms() + slot * config.slot_duration_ms; + store.set_time_ms(now).unwrap(); + store + } + + async fn post_liveness( + store: Store, + epoch: u64, + indices: &[&str], + sync_status: SyncStatusController, + ) -> (StatusCode, serde_json::Value) { + let request = Request::post(format!("/eth/v1/validator/liveness/{epoch}")) + .header("content-type", "application/json") + .body(Body::from(serde_json::json!(indices).to_string())) + .unwrap(); + let app = routes().with_state(store).layer(Extension(sync_status)); + let response = app.oneshot(request).await.unwrap(); + let status = response.status(); + let body = response.into_body().collect().await.unwrap().to_bytes(); + (status, serde_json::from_slice(&body).unwrap_or_default()) + } + + /// `(index, is_live)` pairs, in the order answered. + fn answers(json: &serde_json::Value) -> Vec<(String, bool)> { + json["data"] + .as_array() + .unwrap() + .iter() + .map(|entry| { + ( + entry["index"].as_str().unwrap().to_owned(), + entry["is_live"].as_bool().unwrap(), + ) + }) + .collect() + } + + #[tokio::test] + async fn a_participation_flag_makes_a_validator_live() { + let store = store_for(credited_state()); + let (status, json) = + post_liveness(store.clone(), EPOCH, &["2", "3"], Default::default()).await; + assert_eq!(status, StatusCode::OK); + assert_eq!( + answers(&json), + [("2".into(), true), ("3".into(), false)], + "2 is credited for this epoch, 3 only for the previous one" + ); + + let (_, json) = post_liveness(store, EPOCH - 1, &["2", "3"], Default::default()).await; + assert_eq!(answers(&json), [("2".into(), false), ("3".into(), true)]); + } + + /// What no block has included yet: a validator the node saw act is + /// live without any flag. + #[tokio::test] + async fn an_observed_validator_is_live_without_a_flag() { + let store = store_for(credited_state()); + store.observed_liveness().record(EPOCH, 9); + let (status, json) = + post_liveness(store, EPOCH, &["9", "10"], Default::default()).await; + assert_eq!(status, StatusCode::OK); + assert_eq!(answers(&json), [("9".into(), true), ("10".into(), false)]); + } + + #[tokio::test] + async fn the_next_epoch_is_answered_and_nobody_is_live_in_it() { + let store = store_for(credited_state()); + let (status, json) = post_liveness(store, EPOCH + 1, &["2"], Default::default()).await; + assert_eq!(status, StatusCode::OK); + assert_eq!(answers(&json), [("2".into(), false)]); + } + + #[tokio::test] + async fn epochs_outside_the_window_are_a_400() { + for epoch in [EPOCH - 2, EPOCH + 2] { + let store = store_for(credited_state()); + let (status, _) = post_liveness(store, epoch, &["2"], Default::default()).await; + assert_eq!(status, StatusCode::BAD_REQUEST, "epoch {epoch}"); + } + } + + #[tokio::test] + async fn an_unknown_validator_is_a_400() { + let store = store_for(credited_state()); + let unknown = (COUNT as u64).to_string(); + let (status, _) = post_liveness(store, EPOCH, &[&unknown], Default::default()).await; + assert_eq!(status, StatusCode::BAD_REQUEST); + } + + #[tokio::test] + async fn a_syncing_node_answers_503() { + let store = store_for(credited_state()); + let syncing = SyncStatusController::new(SyncStatus::Syncing); + let (status, _) = post_liveness(store, EPOCH, &["2"], syncing).await; + assert_eq!(status, StatusCode::SERVICE_UNAVAILABLE); + } + } } diff --git a/crates/storage/src/lib.rs b/crates/storage/src/lib.rs index 8f1fc61b..759f3ca4 100644 --- a/crates/storage/src/lib.rs +++ b/crates/storage/src/lib.rs @@ -3,6 +3,7 @@ pub mod backend; mod beacon_state_delta; mod committee_cache; mod error; +mod liveness; mod metrics; mod state_codec; mod state_diff; @@ -14,6 +15,7 @@ pub use committee_cache::{CommitteeCache, Lookup, ShufflingKey}; /// Error type returned by the fallible [`Store`] operations, exported so /// callers can match on it (e.g. to distinguish [`Error::DbVersionMismatch`]). pub use error::Error; +pub use liveness::ObservedLiveness; // `CacheKey` lives in `state_writer` (beside the `StateCache` it keys), not // `store`; re-exported here so the public path (`ethlambda_storage::CacheKey`) // is unaffected by which module owns it. diff --git a/crates/storage/src/liveness.rs b/crates/storage/src/liveness.rs new file mode 100644 index 00000000..b9aff663 --- /dev/null +++ b/crates/storage/src/liveness.rs @@ -0,0 +1,142 @@ +//! Validators this node has seen act, by epoch, for the Beacon API's +//! `POST /eth/v1/validator/liveness/{epoch}`. +//! +//! The endpoint's answer is this node's own view, which the Beacon API allows +//! to come from the network, the chain or the API. The head state's +//! participation flags already cover the chain, but only once a block has +//! included a vote; this set covers what the node sees before that: accepted +//! gossip, imported blocks' proposers, and what validator clients submit +//! through this node's own API (gossip never delivers a node its own +//! messages). The endpoint ORs the two. +//! +//! Held by the `Store` because P2P, the chain actor and the RPC all write or +//! read it, and all three already hold a clone. + +use std::collections::BTreeMap; +use std::sync::Mutex; + +use ethlambda_types::beacon::primitives::{Epoch, ValidatorIndex}; + +/// How many epochs are kept, counting back from the newest one recorded. +/// +/// The endpoint answers the previous, current and next epoch; the next one has +/// nothing to observe yet, so the newest epoch and the two before it cover +/// every epoch it can be asked about, with one to spare at a boundary. +const RETAINED_EPOCHS: u64 = 3; + +/// The largest validator index recorded, exclusive. +/// +/// Every writer records an index that passed validation against a state, so +/// this is a guard rather than a limit anything should reach: it bounds one +/// epoch's bitset at 2 MiB whatever an index turns out to be. Mainnet has +/// about 2.4 million validators, well under it. +const MAX_TRACKED_INDEX: ValidatorIndex = 1 << 24; + +/// One bitset per retained epoch, indexed by validator index. +/// +/// A bitset rather than a set of indices: on mainnet an epoch sees most of +/// the registry act, which is about 300 KB as bits and tens of megabytes as a +/// hash set. +#[derive(Default)] +pub struct ObservedLiveness(Mutex>>); + +impl ObservedLiveness { + /// Record that `validator` did something in `epoch`. + pub fn record(&self, epoch: Epoch, validator: ValidatorIndex) { + self.record_all(epoch, [validator]); + } + + /// Record every one of `validators` for `epoch`, under one lock. + /// + /// An epoch older than the retained window is ignored rather than + /// recorded and immediately pruned: a block imported during range sync + /// names an epoch nobody will ask about. + pub fn record_all(&self, epoch: Epoch, validators: impl IntoIterator) { + let mut epochs = self.0.lock().expect("liveness lock poisoned"); + let newest = epochs + .keys() + .next_back() + .copied() + .unwrap_or(epoch) + .max(epoch); + let floor = newest.saturating_sub(RETAINED_EPOCHS - 1); + if epoch < floor { + return; + } + let bits = epochs.entry(epoch).or_default(); + for validator in validators { + if validator >= MAX_TRACKED_INDEX { + continue; + } + let word = (validator / 64) as usize; + if bits.len() <= word { + bits.resize(word + 1, 0); + } + bits[word] |= 1 << (validator % 64); + } + epochs.retain(|&kept, _| kept >= floor); + } + + /// Whether `validator` was recorded for `epoch`. + pub fn is_live(&self, epoch: Epoch, validator: ValidatorIndex) -> bool { + let epochs = self.0.lock().expect("liveness lock poisoned"); + epochs + .get(&epoch) + .and_then(|bits| bits.get((validator / 64) as usize)) + .is_some_and(|word| word & (1 << (validator % 64)) != 0) + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn a_recorded_validator_is_live_in_that_epoch_only() { + let observed = ObservedLiveness::default(); + observed.record(5, 70); + assert!(observed.is_live(5, 70)); + assert!(!observed.is_live(5, 71), "a neighbouring bit"); + assert!(!observed.is_live(4, 70), "another epoch"); + assert!(!observed.is_live(5, 6_000), "past the recorded words"); + } + + #[test] + fn record_all_sets_every_index() { + let observed = ObservedLiveness::default(); + observed.record_all(5, [0, 63, 64, 1_000_000]); + for index in [0, 63, 64, 1_000_000] { + assert!(observed.is_live(5, index), "{index}"); + } + assert!(!observed.is_live(5, 1)); + } + + #[test] + fn epochs_older_than_the_window_are_pruned() { + let observed = ObservedLiveness::default(); + observed.record(10, 1); + observed.record(11, 1); + observed.record(12, 1); + assert!(observed.is_live(10, 1), "10, 11 and 12 fit the window"); + + observed.record(13, 1); + assert!(!observed.is_live(10, 1), "13 pushes 10 out"); + assert!(observed.is_live(11, 1)); + } + + #[test] + fn an_epoch_below_the_window_is_not_recorded() { + let observed = ObservedLiveness::default(); + observed.record(20, 1); + observed.record(5, 2); + assert!(!observed.is_live(5, 2)); + assert!(observed.is_live(20, 1), "and nothing newer is disturbed"); + } + + #[test] + fn an_index_past_the_guard_is_ignored() { + let observed = ObservedLiveness::default(); + observed.record(5, MAX_TRACKED_INDEX); + assert!(!observed.is_live(5, MAX_TRACKED_INDEX)); + } +} diff --git a/crates/storage/src/store.rs b/crates/storage/src/store.rs index 96addfbd..1fe155e1 100644 --- a/crates/storage/src/store.rs +++ b/crates/storage/src/store.rs @@ -7,6 +7,7 @@ use lru::LruCache; use crate::api::{StorageBackend, StorageReadView, StorageWriteBatch, Table}; use crate::committee_cache::CommitteeCache; use crate::error::Error; +use crate::liveness::ObservedLiveness; use ethlambda_crypto::signature::ValidatorSignature; use ethlambda_types::{ @@ -882,6 +883,12 @@ pub struct Store { /// /// Always empty on lean, which has no beacon committees. committee_cache: Arc, + /// Validators this node has seen act, by epoch, for the Beacon API's + /// liveness endpoint. Written by P2P (accepted gossip), the chain actor + /// (imported blocks' proposers) and the RPC (submissions through this + /// node's own API), which is why it lives on the `Store` all three share. + /// See [`ObservedLiveness`]. Always empty on lean. + observed_liveness: 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 +1519,7 @@ impl Store { state_cache, pending_states, committee_cache: Arc::new(CommitteeCache::default()), + observed_liveness: Arc::new(ObservedLiveness::default()), beacon: Default::default(), state_writer, } @@ -2658,6 +2666,12 @@ impl Store { Arc::clone(&self.committee_cache) } + /// The validators this node has seen act, shared by every clone of this + /// `Store`. See [`ObservedLiveness`]. + pub fn observed_liveness(&self) -> &ObservedLiveness { + &self.observed_liveness + } + /// Returns whether a state is available for the given block root. /// /// True if `pending_states` holds the state, a snapshot exists, or the @@ -5483,6 +5497,14 @@ mod tests { assert!(store.cached_state(key).is_some()); } + #[test] + fn observed_liveness_is_shared_across_store_clones() { + let store = beacon_test_store(Arc::new(InMemoryBackend::new())); + let clone = store.clone(); + clone.observed_liveness().record(3, 7); + assert!(store.observed_liveness().is_live(3, 7)); + } + #[test] fn the_committee_cache_is_shared_across_store_clones() { let store = beacon_test_store(Arc::new(InMemoryBackend::new())); diff --git a/docs/rpc.md b/docs/rpc.md index 6d007e06..6da301fa 100644 --- a/docs/rpc.md +++ b/docs/rpc.md @@ -242,6 +242,7 @@ surface rather than sitting beside it; a `/lean/v0` path on a beacon node is a | `GET` | `/eth/v1/validator/duties/proposer/{epoch}` | JSON | Proposers for the head's epoch or the next | | `POST` | `/eth/v1/validator/duties/attester/{epoch}` | JSON | Committee assignments for the given indices | | `POST` | `/eth/v1/validator/duties/sync/{epoch}` | JSON | Sync committee seats for the given indices, in the head's current or next period | +| `POST` | `/eth/v1/validator/liveness/{epoch}` | JSON | Whether this node saw each validator act in the epoch (doppelganger protection) | | `GET` | `/eth/v1/validator/attestation_data` | JSON | What to attest to at `slot` | | `POST` | `/eth/v2/beacon/pool/attestations` | *(status only)* | Validate and gossip `SingleAttestation`s | | `POST` | `/eth/v1/validator/beacon_committee_subscriptions` | *(status only)* | Aggregators' entries join their committee's subnet | @@ -272,6 +273,17 @@ the chain actor writes, so no request waits on the actor. `503` while the node is syncing. This node serves no sync committee message or contribution endpoint yet, so a validator client that gets duties here cannot publish what they ask for. +- **Liveness** is this node's own view, which the Beacon API allows. A + validator is live in an epoch if the head state credits it for that epoch (a + non-zero participation byte, so anything a block already included), **or** + the node observed it act: an accepted gossip aggregate (its aggregator and + every attester its signature verified) or subnet attestation, an imported + block's proposer, or a submission through `pool/attestations` or + `aggregate_and_proofs`. The observed half covers what no block has included + yet, most of the current epoch; it keeps the newest three epochs recorded, + as one bitset each. The window is the store clock's previous, current and + next epoch (the next is always `false`); anything else, and an unknown + index, is a `400`, and the endpoint is a `503` while the node is syncing. - **`attestation_data`** follows phase0's `validator.md`: the head block, the epoch's boundary block as target, and as source the current justified checkpoint of the head state advanced to the slot's epoch (through fork