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/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/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/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/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/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..a496582a 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,11 @@ 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/liveness/{epoch}", post(post_liveness)) .route( "/eth/v1/validator/attestation_data", get(get_attestation_data), @@ -62,6 +69,192 @@ 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, + })) +} + +#[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"))?; + + // 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", + )); + } + + 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)] @@ -810,4 +1003,303 @@ 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); + } + } + + // --- liveness -------------------------------------------------------- + + mod liveness { + use super::*; + + /// 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 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.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(), wall_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, wall_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(wall_epoch(), 9); + let (status, json) = + 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)]); + } + + #[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, 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 [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, wall_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, wall_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 cd0f2259..6da301fa 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,8 @@ 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 | +| `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 | @@ -260,6 +264,26 @@ 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. +- **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 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.