From b252fbca519ab3ba83536c59cd3cb4e38459935c Mon Sep 17 00:00:00 2001 From: Pablo Deymonnaz Date: Tue, 29 Sep 2026 15:37:16 -0300 Subject: [PATCH 1/2] 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 6fa0cbb5973aae2a681300dc9fbd91bc84049e8d 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:17:44 -0300 Subject: [PATCH 2/2] fix(rpc): bound sync duties by the wall clock POST /eth/v1/validator/duties/sync/{epoch} accepted only the head state's sync committee period and the next one. The Beacon API bounds the period by the current epoch's: at the first slot of period P with the head still in P-1 (its boundary block late or missing), a validator client asks for P+1 and got a 400. The upper bound is now max(wall-clock period, head period) + 1. A period past the head's next one is read from a copy of the head advanced, through fork choice's checkpoint-state cache on a blocking thread, to the first epoch of the period before it, whose next_sync_committee is the requested one. An earlier period than the head's, an unknown index and a syncing node keep their 400 and 503. Adds a controlled-clock store helper to the rpc test utilities, and adapts the period-after-next test to pin the clock to the head's period, since the bound is no longer head-relative. --- crates/net/rpc/src/beacon/validator.rs | 187 ++++++++++++++++++++++--- crates/net/rpc/src/lib.rs | 18 ++- docs/rpc.md | 11 +- docs/spec_deviations.md | 9 +- 4 files changed, 201 insertions(+), 24 deletions(-) diff --git a/crates/net/rpc/src/beacon/validator.rs b/crates/net/rpc/src/beacon/validator.rs index 4b740688..2cf49267 100644 --- a/crates/net/rpc/src/beacon/validator.rs +++ b/crates/net/rpc/src/beacon/validator.rs @@ -83,9 +83,13 @@ struct SyncDuty { /// /// 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. +/// period after. The upper bound is the wall clock's, as the Beacon API +/// defines it: the period after the current one. When the head lags the clock +/// across a period boundary, the later period is read from a copy of the head +/// advanced (through fork choice's checkpoint-state cache, on a blocking +/// thread) to the first epoch of the period before it. An earlier period than +/// the head's 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 @@ -99,12 +103,32 @@ async fn post_sync_duties( 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(), + let computed = tokio::task::spawn_blocking(move || sync_duties(&store, &epoch, &indices)).await; + match computed { + Ok(Ok(body)) => crate::json_response(body), + Ok(Err(err)) => err.into_response(), + Err(_) => ApiError::Internal("computing the duties failed").into_response(), } } +/// The head's post-state advanced through empty slots to the first slot of +/// `epoch`, from fork choice's checkpoint-state cache. +/// +/// Runs `process_slots` on a miss, which is seconds on a mainnet registry, so +/// callers run on a blocking thread. +fn state_for_epoch( + store: &Store, + head_root: H256, + epoch: Epoch, +) -> Result, ApiError> { + let target = Checkpoint { + epoch, + root: head_root, + }; + checkpoint_state(store, &target, &store.config()) + .map_err(|_| ApiError::Internal("advancing the head state failed")) +} + fn sync_duties( store: &Store, epoch: &str, @@ -116,21 +140,40 @@ fn sync_duties( .map(|index| index.parse::()) .collect::, _>>() .map_err(|_| ApiError::BadRequest("invalid validator index"))?; - let (head_root, state) = head(store)?; + let (head_root, head_state) = head(store)?; + + // Bounded by the wall clock, as the Beacon API defines it: up to the + // period after the current one. The head lags the clock at a period + // boundary whose block is late or missing, so the head's own period only + // sets the bound when it is ahead of the clock. + let head_period = compute_sync_committee_period(compute_epoch_at_slot(head_state.slot())); + let clock_period = + compute_sync_committee_period(compute_epoch_at_slot(crate::beacon::node::wall_slot(store))); + let requested_period = compute_sync_committee_period(epoch); + if requested_period < head_period || requested_period > head_period.max(clock_period) + 1 { + return Err(ApiError::BadRequest( + "epoch is outside the sync committee periods the node serves duties for", + )); + } + // The head state answers its own period and the next. A later one (the + // head is behind the clock) is read from a copy of the head advanced to + // the first epoch of the period before it, whose `next_sync_committee` is + // the requested one. + let state = if requested_period <= head_period + 1 { + head_state + } else { + let first_epoch = (requested_period - 1) * preset::EPOCHS_PER_SYNC_COMMITTEE_PERIOD; + state_for_epoch(store, head_root, first_epoch)? + }; + let state_period = compute_sync_committee_period(compute_epoch_at_slot(state.slot())); 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 { + let committee = if requested_period == state_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", - )); + next }; // One pass over the committee rather than one per requested validator: @@ -1012,15 +1055,127 @@ mod tests { ); } + /// `post_sync` on a store whose wall clock is at `clock_slot`. + async fn post_sync_at_clock( + state: BeaconState, + clock_slot: u64, + epoch: u64, + indices: &[&str], + ) -> (StatusCode, serde_json::Value) { + let (store, _root) = crate::test_utils::beacon_store_at_clock(state, clock_slot); + 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(SyncStatusController::default())); + 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()) + } + + /// Adapted from a head-relative bound: the wall clock now decides how + /// far a request may reach, so the clock is pinned to the head's + /// period, where the period after next is past it. #[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; + let clock_slot = state.slot(); + let (status, _) = post_sync_at_clock(state, clock_slot, epoch, &["3"]).await; + assert_eq!(status, StatusCode::BAD_REQUEST); + } + + /// The head in the last epoch of period 0, the wall clock at the first + /// slot of period 1: a validator client asks for period 2, which only + /// a state advanced across the boundary knows. + fn lagging_head() -> (BeaconState, u64) { + 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 - 1) * preset::SLOTS_PER_EPOCH; + let boundary_slot = + compute_start_slot_at_epoch(preset::EPOCHS_PER_SYNC_COMMITTEE_PERIOD); + (state, boundary_slot) + } + + #[tokio::test] + async fn a_head_lagging_a_period_boundary_serves_the_period_after_from_an_advanced_state() { + let (state, boundary_slot) = lagging_head(); + let mut advanced = state.clone(); + ethlambda_state_transition::beacon::stf::process_slots( + &mut advanced, + boundary_slot, + ðlambda_types::beacon::config::Config::mainnet(), + ) + .unwrap(); + let (_, expected) = advanced.sync_committees().unwrap(); + let (_, head_next) = state.sync_committees().unwrap(); + assert_ne!( + expected.pubkeys, head_next.pubkeys, + "the boundary must rotate the committee, or the test proves nothing" + ); + + let period_after = preset::EPOCHS_PER_SYNC_COMMITTEE_PERIOD * 2; + let indices: Vec = (0..COUNT as u64).map(|i| i.to_string()).collect(); + let refs: Vec<&str> = indices.iter().map(String::as_str).collect(); + let (status, json) = + post_sync_at_clock(state.clone(), boundary_slot, period_after, &refs).await; + assert_eq!(status, StatusCode::OK); + + let mut want = Vec::new(); + for index in 0..COUNT as u64 { + let pubkey = advanced.validator(index).unwrap().pubkey; + let held: Vec = expected + .pubkeys + .iter() + .enumerate() + .filter(|(_, key)| **key == pubkey) + .map(|(seat, _)| seat.to_string()) + .collect(); + if !held.is_empty() { + want.push((index.to_string(), held)); + } + } + assert!(!want.is_empty()); + let got: Vec<(String, Vec)> = json["data"] + .as_array() + .unwrap() + .iter() + .map(|duty| { + let seats = duty["validator_sync_committee_indices"] + .as_array() + .unwrap() + .iter() + .map(|seat| seat.as_str().unwrap().to_string()) + .collect(); + (duty["validator_index"].as_str().unwrap().to_string(), seats) + }) + .collect(); + assert_eq!(got, want); + } + + #[tokio::test] + async fn a_head_lagging_a_period_boundary_still_refuses_two_periods_past_the_clock() { + let (state, boundary_slot) = lagging_head(); + let epoch = preset::EPOCHS_PER_SYNC_COMMITTEE_PERIOD * 3; + let (status, _) = post_sync_at_clock(state, boundary_slot, epoch, &["3"]).await; assert_eq!(status, StatusCode::BAD_REQUEST); } + #[tokio::test] + async fn a_head_lagging_a_period_boundary_serves_the_clocks_own_period_from_the_head() { + let (state, boundary_slot) = lagging_head(); + let epoch = preset::EPOCHS_PER_SYNC_COMMITTEE_PERIOD; + let (status, json) = post_sync_at_clock(state, boundary_slot, epoch, &["9"]).await; + assert_eq!(status, StatusCode::OK); + assert_eq!(json["data"][0]["validator_index"], "9"); + } + /// An earlier period would need a historical state, which a validator /// client never asks for; refused rather than answered wrongly. #[tokio::test] diff --git a/crates/net/rpc/src/lib.rs b/crates/net/rpc/src/lib.rs index a2144f7c..e9a0217c 100644 --- a/crates/net/rpc/src/lib.rs +++ b/crates/net/rpc/src/lib.rs @@ -463,12 +463,28 @@ pub(crate) mod test_utils { /// phase0 one whatever `state`'s fork: these endpoints read the state and /// the block's root and slot, never the block's body. pub(crate) fn beacon_store_at(state: BeaconState) -> (Store, H256) { + beacon_store_with_genesis(state, 1_606_824_023) + } + + /// [`beacon_store_at`] on a wall clock whose slot `clock_slot` began a + /// second ago, for endpoints that bound a request by the current epoch. + pub(crate) fn beacon_store_at_clock(state: BeaconState, clock_slot: u64) -> (Store, H256) { + let slot_secs = Config::mainnet().slot_duration_ms / 1000; + let now_secs = std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .expect("the clock is after the epoch") + .as_secs(); + let genesis = now_secs - clock_slot * slot_secs - 1; + beacon_store_with_genesis(state, genesis) + } + + fn beacon_store_with_genesis(state: BeaconState, genesis_time: u64) -> (Store, H256) { let slot = state.slot(); let block = phase0_beacon_block(slot, H256::ZERO); let root = block.message_hash_tree_root(); let mut store = Store::init_beacon( Arc::new(InMemoryBackend::default()), - 1_606_824_023, + genesis_time, Config::mainnet(), root, Checkpoint { root, slot }, diff --git a/docs/rpc.md b/docs/rpc.md index 6d007e06..9cb95265 100644 --- a/docs/rpc.md +++ b/docs/rpc.md @@ -241,7 +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 | +| `POST` | `/eth/v1/validator/duties/sync/{epoch}` | JSON | Sync committee seats for the given indices, up to the wall clock's current period plus one | | `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 | @@ -265,8 +265,13 @@ the chain actor writes, so no request waits on the actor. 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 + after. The upper bound is the wall clock's, as the Beacon API defines it: the + period after the clock's current one. When the head lags a period boundary + (its block is late or missing), the later period is read from a copy of the + head advanced to the first epoch of the period before it, through fork + choice's checkpoint-state cache on a blocking thread. A period past that + bound or before the head's 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 diff --git a/docs/spec_deviations.md b/docs/spec_deviations.md index a6cf1ec2..3017fae4 100644 --- a/docs/spec_deviations.md +++ b/docs/spec_deviations.md @@ -99,11 +99,12 @@ 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 +## Sync duties are not served for a period before the head's -`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`. +`POST /eth/v1/validator/duties/sync/{epoch}` refuses a period before the head +state's own with a `400`. Later periods follow the wall clock, as the Beacon +API defines: up to the clock's current period plus one, read from the head or, +when the head lags a period boundary, from a copy advanced across it. - **Beacon API:** allows any period up to the current one plus one, so an earlier period is valid to ask about.