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] 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.