From b252fbca519ab3ba83536c59cd3cb4e38459935c Mon Sep 17 00:00:00 2001 From: Pablo Deymonnaz Date: Tue, 29 Sep 2026 15:37:16 -0300 Subject: [PATCH 1/3] Serve GET /eth/v1/beacon/states/{state_id}/fork, GET /eth/v1/config/deposit_contract and POST /eth/v1/validator/duties/sync/{epoch} from ethlambda beacon, three of the four Beacon API endpoints validator clients call that were missing. fork returns the state's own Fork through the same state_id resolution as the other state endpoints, so a state root is the same 404. deposit_contract returns the Config's deposit chain id and contract address. duties/sync reads the head state's current_sync_committee for an epoch in the head's own sync committee period and next_sync_committee for the next one, matches each requested validator by pubkey and returns every seat it holds (the committee is drawn with replacement), leaves out validators with no seat, and answers 400 for an unknown index or any other period and 503 while syncing. An earlier period is refused rather than answered from a historical state, recorded in docs/spec_deviations.md. compute_sync_committee_period is added to the altair helpers as validator.md defines it. --- .../src/beacon/helpers/altair.rs | 19 ++ crates/net/rpc/src/beacon/config.rs | 42 ++- crates/net/rpc/src/beacon/states.rs | 53 ++++ crates/net/rpc/src/beacon/validator.rs | 242 ++++++++++++++++++ docs/rpc.md | 12 + docs/spec_deviations.md | 15 ++ 6 files changed, 375 insertions(+), 8 deletions(-) diff --git a/crates/blockchain/state_transition/src/beacon/helpers/altair.rs b/crates/blockchain/state_transition/src/beacon/helpers/altair.rs index eaa008c2..bb0ab6eb 100644 --- a/crates/blockchain/state_transition/src/beacon/helpers/altair.rs +++ b/crates/blockchain/state_transition/src/beacon/helpers/altair.rs @@ -170,6 +170,16 @@ pub fn get_next_sync_committee(state: &BeaconState) -> Result u64 { + epoch / preset::EPOCHS_PER_SYNC_COMMITTEE_PERIOD +} + /// The reward every one-increment slice of a validator's effective balance /// earns for a single timely, correct component of its attestation. /// @@ -592,4 +602,13 @@ mod tests { let elapsed = start.elapsed() / ITERATIONS; println!("get_flag_index_deltas, {VALIDATOR_COUNT} validators -> {elapsed:?}/call"); } + + #[test] + fn a_sync_committee_period_spans_its_epochs_and_no_more() { + let period = preset::EPOCHS_PER_SYNC_COMMITTEE_PERIOD; + assert_eq!(compute_sync_committee_period(0), 0); + assert_eq!(compute_sync_committee_period(period - 1), 0); + assert_eq!(compute_sync_committee_period(period), 1); + assert_eq!(compute_sync_committee_period(3 * period + 1), 3); + } } diff --git a/crates/net/rpc/src/beacon/config.rs b/crates/net/rpc/src/beacon/config.rs index 3ddb354b..11babc7e 100644 --- a/crates/net/rpc/src/beacon/config.rs +++ b/crates/net/rpc/src/beacon/config.rs @@ -1,4 +1,4 @@ -//! `/eth/v1/config/spec`. +//! `/eth/v1/config/spec` and `/eth/v1/config/deposit_contract`. //! //! The Beacon API asks for three things in one flat object: the network's //! configuration, the preset the node was built against, and the @@ -29,7 +29,22 @@ use ethlambda_types::beacon::{config::Config, constants, preset, serde_helpers:: use serde_json::{Map, Value}; pub(crate) fn routes() -> Router { - Router::new().route("/eth/v1/config/spec", get(get_spec)) + Router::new() + .route("/eth/v1/config/spec", get(get_spec)) + .route("/eth/v1/config/deposit_contract", get(get_deposit_contract)) +} + +/// `GET /eth/v1/config/deposit_contract`: the network's deposit contract, off +/// the same `Config` `/eth/v1/config/spec` reports it from (as +/// `DEPOSIT_CHAIN_ID` and `DEPOSIT_CONTRACT_ADDRESS`). +async fn get_deposit_contract(State(store): State) -> Response { + let config = store.config(); + crate::json_response(serde_json::json!({ + "data": { + "chain_id": config.deposit_chain_id.to_string(), + "address": HexPrefixed(&config.deposit_contract_address).to_string(), + } + })) } async fn get_spec(State(store): State) -> Response { @@ -202,15 +217,14 @@ mod tests { use tower::ServiceExt as _; async fn get_spec_json() -> serde_json::Value { + get_json("/eth/v1/config/spec").await + } + + async fn get_json(uri: &str) -> serde_json::Value { let fixture = beacon_fixture(64); let app = routes().with_state(fixture.store); let response = app - .oneshot( - Request::builder() - .uri("/eth/v1/config/spec") - .body(Body::empty()) - .unwrap(), - ) + .oneshot(Request::builder().uri(uri).body(Body::empty()).unwrap()) .await .unwrap(); @@ -219,6 +233,18 @@ mod tests { serde_json::from_slice(&body).unwrap() } + /// The fixture's store is bootstrapped with `Config::mainnet()`, whose + /// deposit contract is the one on Ethereum mainnet. + #[tokio::test] + async fn the_deposit_contract_is_mainnets() { + let json = get_json("/eth/v1/config/deposit_contract").await; + assert_eq!(json["data"]["chain_id"], "1"); + assert_eq!( + json["data"]["address"], + "0x00000000219ab540356cbb839cbe05303d7705fa" + ); + } + #[tokio::test] async fn the_spec_is_screaming_snake_case_with_quoted_values() { let json = get_spec_json().await; diff --git a/crates/net/rpc/src/beacon/states.rs b/crates/net/rpc/src/beacon/states.rs index 47b0e449..c16bd85b 100644 --- a/crates/net/rpc/src/beacon/states.rs +++ b/crates/net/rpc/src/beacon/states.rs @@ -37,6 +37,7 @@ use crate::{ pub(crate) fn routes() -> Router { Router::new() .route("/eth/v2/debug/beacon/states/{state_id}", get(get_state)) + .route("/eth/v1/beacon/states/{state_id}/fork", get(get_fork)) .route( "/eth/v1/beacon/states/{state_id}/finality_checkpoints", get(get_finality_checkpoints), @@ -100,6 +101,20 @@ async fn get_state( with_consensus_version(response, fork) } +/// `GET /eth/v1/beacon/states/{state_id}/fork`: the `Fork` the state carries, +/// which is what a validator client builds its signing domains from. +async fn get_fork(Path(state_id): Path, State(store): State) -> Response { + let (root, state) = match load(&store, &state_id) { + Ok(found) => found, + Err(err) => return err.into_response(), + }; + crate::json_response(serde_json::json!({ + "execution_optimistic": store.is_beacon_optimistic(root), + "finalized": is_finalized(&store, state.slot()), + "data": state.fork(), + })) +} + async fn get_finality_checkpoints( Path(state_id): Path, State(store): State, @@ -439,6 +454,44 @@ mod tests { } } + #[tokio::test] + async fn the_fork_is_the_one_the_state_carries() { + let fixture = beacon_fixture(ANCHOR_SLOT); + let head_state = fixture + .store + .get_state(&fixture.head_root) + .unwrap() + .unwrap(); + let expected = serde_json::to_value(head_state.fork()).unwrap(); + + let response = get("/eth/v1/beacon/states/head/fork", None).await; + assert_eq!(response.status(), StatusCode::OK); + let json = body_json(response).await; + assert_eq!(json["data"], expected); + assert!( + json["data"]["epoch"].is_string(), + "the epoch must be quoted" + ); + assert!(json["execution_optimistic"].is_boolean()); + assert!(json["finalized"].is_boolean()); + } + + #[tokio::test] + async fn the_finalized_states_fork_is_marked_finalized() { + let response = get("/eth/v1/beacon/states/finalized/fork", None).await; + assert_eq!(response.status(), StatusCode::OK); + assert_eq!(body_json(response).await["finalized"], true); + } + + /// The same refusal every other state endpoint gives: state roots are not + /// indexed, so a `0x` id is a 404 rather than a guess. + #[tokio::test] + async fn a_fork_by_state_root_is_a_404() { + let root = format!("0x{}", "ab".repeat(32)); + let response = get(&format!("/eth/v1/beacon/states/{root}/fork"), None).await; + assert_eq!(response.status(), StatusCode::NOT_FOUND); + } + mod validators { use super::*; use crate::test_utils::beacon_store_at; diff --git a/crates/net/rpc/src/beacon/validator.rs b/crates/net/rpc/src/beacon/validator.rs index abbbd95e..4b740688 100644 --- a/crates/net/rpc/src/beacon/validator.rs +++ b/crates/net/rpc/src/beacon/validator.rs @@ -14,6 +14,7 @@ use axum::{ response::{IntoResponse, Response}, routing::{get, post}, }; +use ethlambda_blockchain::{SyncStatusController, metrics::SyncStatus}; use ethlambda_storage::Store; use ethlambda_types::{ beacon::{ @@ -34,6 +35,7 @@ use ethlambda_state_transition::beacon::{ fork_choice::checkpoint_state, gossip::attestation::compute_subnet_for_attestation, helpers::accessors::{CommitteeCacheExt as _, get_block_root_at_slot}, + helpers::altair::compute_sync_committee_period, }; use crate::beacon::ApiError; @@ -48,6 +50,10 @@ pub(crate) fn routes() -> Router { "/eth/v1/validator/duties/attester/{epoch}", post(post_attester_duties), ) + .route( + "/eth/v1/validator/duties/sync/{epoch}", + post(post_sync_duties), + ) .route( "/eth/v1/validator/attestation_data", get(get_attestation_data), @@ -62,6 +68,101 @@ pub(crate) fn routes() -> Router { ) } +#[derive(Debug, Serialize)] +struct SyncDuty { + pubkey: BlsPubkey, + #[serde(with = "ethlambda_types::beacon::serde_helpers::quoted_or_bare")] + validator_index: ValidatorIndex, + /// Every position the validator holds in the committee, quoted. A + /// validator can hold more than one, since the committee is drawn with + /// replacement. + validator_sync_committee_indices: Vec, +} + +/// `POST /eth/v1/validator/duties/sync/{epoch}`. +/// +/// Answered from the head state: its `current_sync_committee` for an epoch in +/// the head's own sync committee period, its `next_sync_committee` for the +/// period after, which is as far ahead as the Beacon API allows. An earlier +/// period is refused rather than answered from a historical state, since a +/// validator client only ever asks about the current and next period. +/// +/// A requested validator that holds no seat is left out of `data`. The +/// answer is `503` while the node is syncing: the head state's committees are +/// not yet the chain's. +async fn post_sync_duties( + 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 sync_duties(&store, &epoch, &indices) { + Ok(body) => crate::json_response(body), + Err(err) => err.into_response(), + } +} + +fn sync_duties( + 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 (head_root, state) = head(store)?; + + let (current, next) = state + .sync_committees() + .map_err(|_| ApiError::BadRequest("sync committees start at altair"))?; + let head_period = compute_sync_committee_period(compute_epoch_at_slot(state.slot())); + let requested_period = compute_sync_committee_period(epoch); + let committee = if requested_period == head_period { + current + } else if requested_period == head_period + 1 { + next + } else { + return Err(ApiError::BadRequest( + "epoch is not in the head state's current or next sync committee period", + )); + }; + + // One pass over the committee rather than one per requested validator: + // the committee stores pubkeys, so that is what a validator is matched by. + let mut positions: HashMap> = HashMap::new(); + for (position, pubkey) in committee.pubkeys.iter().enumerate() { + positions + .entry(*pubkey) + .or_default() + .push(position.to_string()); + } + + let mut duties = Vec::new(); + for validator_index in indices { + let validator = state + .validator(validator_index) + .map_err(|_| ApiError::BadRequest("unknown validator index"))?; + if let Some(held) = positions.get(&validator.pubkey) { + duties.push(SyncDuty { + pubkey: validator.pubkey, + validator_index, + validator_sync_committee_indices: held.clone(), + }); + } + } + + Ok(serde_json::json!({ + "execution_optimistic": store.is_beacon_optimistic(head_root), + "data": duties, + })) +} + /// One entry of `beacon_committee_subscriptions`. Parsed so a malformed body /// is refused, though `validator_index` is never read. #[derive(Debug, Deserialize)] @@ -810,4 +911,145 @@ mod tests { .await; assert_eq!(status, StatusCode::BAD_REQUEST); } + + // --- duties/sync ----------------------------------------------------- + + mod sync_duties { + use super::*; + use ethlambda_types::beacon::containers::altair::SyncCommittee; + + /// A committee whose seat `i` belongs to validator `first + i % 8`, so + /// each of those eight holds `SYNC_COMMITTEE_SIZE / 8` seats and every + /// other validator holds none. + fn committee_of(state: &BeaconState, first: u64) -> SyncCommittee { + let pubkeys: Vec = (0..preset::SYNC_COMMITTEE_SIZE as u64) + .map(|seat| state.validator(first + seat % 8).unwrap().pubkey) + .collect(); + SyncCommittee { + aggregate_pubkey: pubkeys[0], + pubkeys: pubkeys.try_into().unwrap(), + } + } + + /// A fulu state whose current committee is validators 0-7 and whose + /// next committee is validators 8-15, so an answer drawn from the + /// wrong one shows up. + fn state_with_committees() -> BeaconState { + let mut state = fulu_state(); + let current = committee_of(&state, 0); + let next = committee_of(&state, 8); + let BeaconState::Fulu(fulu) = &mut state else { + unreachable!("built as fulu") + }; + fulu.current_sync_committee = current; + fulu.next_sync_committee = next; + state + } + + async fn post_sync( + state: BeaconState, + epoch: u64, + indices: &[&str], + sync_status: SyncStatusController, + ) -> (StatusCode, serde_json::Value) { + let (store, _root) = beacon_store_at(state); + let request = Request::post(format!("/eth/v1/validator/duties/sync/{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()) + } + + /// The seats `validator` holds in [`committee_of`]`(_, first)`. + fn seats(validator: u64, first: u64) -> Vec { + (0..preset::SYNC_COMMITTEE_SIZE as u64) + .filter(|seat| first + seat % 8 == validator) + .map(|seat| seat.to_string()) + .collect() + } + + #[tokio::test] + async fn the_current_period_reads_the_current_committee() { + let state = state_with_committees(); + let epoch = compute_epoch_at_slot(state.slot()); + let (status, json) = + post_sync(state.clone(), epoch, &["3", "9", "20"], Default::default()).await; + assert_eq!(status, StatusCode::OK); + + // Validator 3 sits in the current committee; 9 only in the next; + // 20 in neither, so only 3 is listed. + let duties = json["data"].as_array().unwrap(); + assert_eq!(duties.len(), 1); + assert_eq!(duties[0]["validator_index"], "3"); + let pubkey = state.validator(3).unwrap().pubkey; + assert_eq!(duties[0]["pubkey"], format!("0x{}", hex::encode(pubkey.0))); + assert_eq!( + duties[0]["validator_sync_committee_indices"], + serde_json::json!(seats(3, 0)) + ); + assert!(json["execution_optimistic"].is_boolean()); + } + + #[tokio::test] + async fn the_next_period_reads_the_next_committee() { + let state = state_with_committees(); + let next_period_epoch = preset::EPOCHS_PER_SYNC_COMMITTEE_PERIOD + * (compute_sync_committee_period(compute_epoch_at_slot(state.slot())) + 1); + let (status, json) = + post_sync(state, next_period_epoch, &["3", "9"], Default::default()).await; + assert_eq!(status, StatusCode::OK); + + let duties = json["data"].as_array().unwrap(); + assert_eq!(duties.len(), 1); + assert_eq!(duties[0]["validator_index"], "9"); + assert_eq!( + duties[0]["validator_sync_committee_indices"], + serde_json::json!(seats(9, 8)) + ); + } + + #[tokio::test] + async fn the_period_after_next_is_a_400() { + let state = state_with_committees(); + let period = compute_sync_committee_period(compute_epoch_at_slot(state.slot())); + let epoch = preset::EPOCHS_PER_SYNC_COMMITTEE_PERIOD * (period + 2); + let (status, _) = post_sync(state, epoch, &["3"], Default::default()).await; + assert_eq!(status, StatusCode::BAD_REQUEST); + } + + /// An earlier period would need a historical state, which a validator + /// client never asks for; refused rather than answered wrongly. + #[tokio::test] + async fn an_earlier_period_is_a_400() { + let mut state = state_with_committees(); + let BeaconState::Fulu(fulu) = &mut state else { + unreachable!("built as fulu") + }; + fulu.slot = preset::EPOCHS_PER_SYNC_COMMITTEE_PERIOD * preset::SLOTS_PER_EPOCH; + let (status, _) = post_sync(state, 0, &["3"], Default::default()).await; + assert_eq!(status, StatusCode::BAD_REQUEST); + } + + #[tokio::test] + async fn an_unknown_validator_is_a_400() { + let state = state_with_committees(); + let epoch = compute_epoch_at_slot(state.slot()); + let unknown = (COUNT as u64).to_string(); + let (status, _) = post_sync(state, epoch, &[&unknown], Default::default()).await; + assert_eq!(status, StatusCode::BAD_REQUEST); + } + + #[tokio::test] + async fn a_syncing_node_answers_503() { + let state = state_with_committees(); + let epoch = compute_epoch_at_slot(state.slot()); + let syncing = SyncStatusController::new(SyncStatus::Syncing); + let (status, _) = post_sync(state, epoch, &["3"], syncing).await; + assert_eq!(status, StatusCode::SERVICE_UNAVAILABLE); + } + } } diff --git a/docs/rpc.md b/docs/rpc.md index cd0f2259..6d007e06 100644 --- a/docs/rpc.md +++ b/docs/rpc.md @@ -230,8 +230,10 @@ surface rather than sitting beside it; a `/lean/v0` path on a beacon node is a | `GET` | `/eth/v1/beacon/headers/{block_id}` | JSON | `SignedBeaconBlockHeader`, plus `canonical` | | `GET` | `/eth/v2/debug/beacon/states/{state_id}` | JSON or SSZ | `BeaconState` at `state_id` | | `GET` | `/eth/v1/beacon/states/{state_id}/finality_checkpoints` | JSON | That state's three checkpoints | +| `GET` | `/eth/v1/beacon/states/{state_id}/fork` | JSON | That state's `Fork`: previous and current version, and the epoch it changed | | `GET` | `/eth/v1/beacon/genesis` | JSON | Genesis time, validators root, fork version | | `GET` | `/eth/v1/config/spec` | JSON | The store's `Config`, plus `PRESET_BASE`, `CONFIG_NAME`, the preset and the constants (see below) | +| `GET` | `/eth/v1/config/deposit_contract` | JSON | The `Config`'s deposit chain id and contract address | | `GET` | `/eth/v1/node/syncing` | JSON | Head slot, sync distance, optimistic flag | | `GET` | `/eth/v1/node/health` | *(status only)* | `200` caught up, `206` syncing | | `GET` | `/eth/v1/node/version` | JSON | Client version string | @@ -239,6 +241,7 @@ surface rather than sitting beside it; a `/lean/v0` path on a beacon node is a | `GET`, `POST` | `/eth/v1/beacon/states/{state_id}/validators` | JSON | Registry entries by index or pubkey, with status | | `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 | | `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 | @@ -260,6 +263,15 @@ the chain actor writes, so no request waits on the actor. which is as far as its shuffling is already fixed. Anything else is a `400`. `dependent_root` follows each endpoint's v1 definition. Attester duties walk every committee of the epoch, a full shuffle per request on mainnet. +- **Sync duties** read the head state's `current_sync_committee` for an epoch in + the head's own sync committee period and `next_sync_committee` for the one + after; any other period is a `400` (an earlier one would need a historical + state; see `docs/spec_deviations.md`). A validator is matched by pubkey and + gets every seat it holds, since the committee is drawn with replacement; one + with no seat is left out. An unknown index is a `400`, and the endpoint is a + `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. - **`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 diff --git a/docs/spec_deviations.md b/docs/spec_deviations.md index 7ff2d24e..a6cf1ec2 100644 --- a/docs/spec_deviations.md +++ b/docs/spec_deviations.md @@ -99,6 +99,21 @@ rather than populated, which is not spec-valid. work. Everything that reads `peer_id` is unaffected. Out of scope for the change that added the Beacon API surface; a follow-up exposes the record. +## Sync duties are served for the current and next period only + +`POST /eth/v1/validator/duties/sync/{epoch}` answers an `epoch` in the head +state's own sync committee period or the next one, and refuses an earlier +period with a `400`. + +- **Beacon API:** allows any period up to the current one plus one, so an + earlier period is valid to ask about. +- **ethlambda:** the head state carries only `current_sync_committee` and + `next_sync_committee`. Answering an earlier period means loading the state + at the start of that period, which is what Lighthouse does; Prysm also + serves it. A validator client only asks about the current and next period, + so this refuses rather than add a historical state lookup to a duty + endpoint (`sync_duties`, `crates/net/rpc/src/beacon/validator.rs`). + ## `block_id` cannot name `genesis`, and `state_id` cannot be a state root Two id forms the Beacon API defines return `404` here. From d973505de3362ff5085dbdd771db2b6423c65031 Mon Sep 17 00:00:00 2001 From: Pablo Deymonnaz Date: Tue, 29 Sep 2026 15:47:23 -0300 Subject: [PATCH 2/3] 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 From 95524f1d91a991ee82071f69ed8e4939c618e733 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Tom=C3=A1s=20Gr=C3=BCner?= <47506558+MegaRedHand@users.noreply.github.com> Date: Tue, 6 Oct 2026 00:03:46 -0300 Subject: [PATCH 3/3] fix(rpc): bound liveness by the wall clock The previous/current/next epoch window used the store's tick-driven clock, which still reads the old epoch until the slot tick runs. A validator client asking for the next epoch right after a boundary got a 400. Use the wall-clock epoch, as the duty endpoints do. --- crates/net/rpc/src/beacon/validator.rs | 64 ++++++++++++++++++++------ 1 file changed, 49 insertions(+), 15 deletions(-) diff --git a/crates/net/rpc/src/beacon/validator.rs b/crates/net/rpc/src/beacon/validator.rs index 7b218c74..a496582a 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, get_current_store_epoch}, + fork_choice::checkpoint_state, gossip::attestation::compute_subnet_for_attestation, helpers::accessors::{CommitteeCacheExt as _, get_block_root_at_slot}, helpers::altair::compute_sync_committee_period, @@ -211,7 +211,10 @@ fn liveness(store: &Store, epoch: &str, indices: &[String]) -> Result, _>>() .map_err(|_| ApiError::BadRequest("invalid validator index"))?; - let current = get_current_store_epoch(store, &store.config()); + // The wall clock, not the store's tick-driven one: the latter still reads + // the previous epoch until the slot tick runs, so a client asking for the + // next epoch right after a boundary would be refused. + let current = compute_epoch_at_slot(crate::beacon::node::wall_slot(store)); if epoch + 1 < current || epoch > current + 1 { return Err(ApiError::BadRequest( "epoch is not the previous, current or next epoch", @@ -1147,18 +1150,26 @@ mod tests { 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; + /// The wall clock's epoch. The store's mainnet genesis puts it far from + /// zero, so the epochs around it exist. + fn wall_epoch() -> u64 { + compute_epoch_at_slot(crate::beacon::node::wall_slot( + &beacon_store_at(fulu_state()).0, + )) + } - /// A fulu state at [`EPOCH`]'s first slot in which validator 2 is + /// A fulu state at the wall epoch's first slot in which validator 2 is /// credited for this epoch and validator 3 for the previous one. fn credited_state() -> BeaconState { + credited_state_at(wall_epoch()) + } + + fn credited_state_at(epoch: u64) -> 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.slot = compute_start_slot_at_epoch(epoch); fulu.current_epoch_participation[2] = 0b001; fulu.previous_epoch_participation[3] = 0b111; state @@ -1211,7 +1222,7 @@ mod tests { 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; + post_liveness(store.clone(), wall_epoch(), &["2", "3"], Default::default()).await; assert_eq!(status, StatusCode::OK); assert_eq!( answers(&json), @@ -1219,7 +1230,8 @@ mod tests { "2 is credited for this epoch, 3 only for the previous one" ); - let (_, json) = post_liveness(store, EPOCH - 1, &["2", "3"], Default::default()).await; + let (_, json) = + post_liveness(store, wall_epoch() - 1, &["2", "3"], Default::default()).await; assert_eq!(answers(&json), [("2".into(), false), ("3".into(), true)]); } @@ -1228,9 +1240,9 @@ mod tests { #[tokio::test] async fn an_observed_validator_is_live_without_a_flag() { let store = store_for(credited_state()); - store.observed_liveness().record(EPOCH, 9); + store.observed_liveness().record(wall_epoch(), 9); let (status, json) = - post_liveness(store, EPOCH, &["9", "10"], Default::default()).await; + post_liveness(store, wall_epoch(), &["9", "10"], Default::default()).await; assert_eq!(status, StatusCode::OK); assert_eq!(answers(&json), [("9".into(), true), ("10".into(), false)]); } @@ -1238,25 +1250,47 @@ mod tests { #[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; + let (status, json) = + post_liveness(store, wall_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] { + for epoch in [wall_epoch() - 2, wall_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}"); } } + /// The store's tick still reads epoch N while the wall clock has + /// entered N+1 (the slot tick has not run yet): the window follows the + /// wall clock, so a validator client asking for N+2 is answered. + #[tokio::test] + async fn the_window_follows_the_wall_clock_not_the_store_tick() { + let wall = wall_epoch(); + let store = store_for(credited_state_at(wall - 1)); + for (epoch, expected) in [ + (wall - 2, StatusCode::BAD_REQUEST), + (wall - 1, StatusCode::OK), + (wall, StatusCode::OK), + (wall + 1, StatusCode::OK), + (wall + 2, StatusCode::BAD_REQUEST), + ] { + let (status, _) = + post_liveness(store.clone(), epoch, &["2"], Default::default()).await; + assert_eq!(status, expected, "epoch {epoch}, wall epoch {wall}"); + } + } + #[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; + let (status, _) = + post_liveness(store, wall_epoch(), &[&unknown], Default::default()).await; assert_eq!(status, StatusCode::BAD_REQUEST); } @@ -1264,7 +1298,7 @@ mod tests { 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; + let (status, _) = post_liveness(store, wall_epoch(), &["2"], syncing).await; assert_eq!(status, StatusCode::SERVICE_UNAVAILABLE); } }