diff --git a/crates/net/rpc/src/beacon/ptc.rs b/crates/net/rpc/src/beacon/ptc.rs index e827d3e5..f8f1fd85 100644 --- a/crates/net/rpc/src/beacon/ptc.rs +++ b/crates/net/rpc/src/beacon/ptc.rs @@ -24,7 +24,7 @@ use ethlambda_blockchain::SyncStatusController; use ethlambda_blockchain::metrics::SyncStatus; use ethlambda_network_api::RpcToP2PRef; use ethlambda_state_transition::beacon::{ - fork_choice::{get_current_slot, get_payload_due_ms}, + fork_choice::get_payload_due_ms, gloas_block_production::aggregate_payload_attestations, gossip::{ Outcome, @@ -94,7 +94,8 @@ struct PtcDuty { /// seats still gets one duty per epoch, and one in no committee gets none, as /// does an unknown index. /// -/// `epoch` may be at most one past the current one. A gloas epoch is answered +/// `epoch` may be at most one past the wall clock's epoch (or the head's, if +/// that is later). A gloas epoch is answered /// from the head state's `ptc_window`, which `get_ptc` can read for the /// state's epoch, the one before and the next one. Anything else (the first /// gloas epoch while the head is still fulu, or an epoch the head has not @@ -152,11 +153,11 @@ fn ptc_duties( let config = store.config(); let (head_root, head_state) = head(store)?; let state_epoch = compute_epoch_at_slot(head_state.slot()); - // The wall clock can be ahead of the head (a node that missed slots) and, - // in a follower's store, behind it; either way an epoch more than one past - // the later of the two is not one a validator is asked about yet. - let clock_epoch = compute_epoch_at_slot(get_current_slot(store, &config)); - if epoch > state_epoch.max(clock_epoch) + 1 { + // Bounded by the wall clock, as the other duties are (see + // `validator::epoch_upper_bound`): the store's tick-driven clock still reads + // the previous epoch until the boundary slot's tick runs, which is when a + // validator client asks for the next epoch. + if epoch > crate::beacon::validator::epoch_upper_bound(store, state_epoch) { return Err(ApiError::BadRequest( "epoch is more than one past the current", )); @@ -572,6 +573,19 @@ mod tests { /// `commitments` blob commitments in its bid, on a clock that has that /// slot running now (the gossip checks only take the current slot). fn store_with_head(state: BeaconState, config: Config, commitments: usize) -> (Store, H256) { + let slot = state.slot(); + store_with_head_at_clock(state, config, commitments, slot) + } + + /// [`store_with_head`] whose wall clock is at `clock_slot` while the store's + /// own tick-driven time stays at the head's slot, as it does until the + /// boundary slot's tick runs. + fn store_with_head_at_clock( + state: BeaconState, + config: Config, + commitments: usize, + clock_slot: u64, + ) -> (Store, H256) { let slot = state.slot(); let mut block = gloas_beacon_block(slot, H256::ZERO, H256::ZERO, H256::repeat_byte(1)); let SignedBeaconBlock::Gloas(inner) = &mut block else { @@ -587,7 +601,7 @@ mod tests { .expect("a few commitments fit"); let root = block.message_hash_tree_root(); let slot_secs = config.slot_duration_ms / 1000; - let genesis = now_secs() - slot * slot_secs - 1; + let genesis = now_secs() - clock_slot * slot_secs - 1; let mut store = Store::init_beacon( Arc::new(InMemoryBackend::default()), genesis, @@ -596,6 +610,8 @@ mod tests { Checkpoint { root, slot }, slot, ); + let tick_ms = store.config().genesis_time_ms() + slot * store.config().slot_duration_ms; + store.set_time_ms(tick_ms).unwrap(); store.insert_signed_block(root, block).unwrap(); store.insert_state(root, state).unwrap(); store @@ -749,6 +765,30 @@ mod tests { assert_eq!(reply.status, StatusCode::BAD_REQUEST); } + /// The devnet case: the wall clock is in the epoch after the head's while + /// the store's own clock has not ticked past the head's, and the validator + /// client asks for the epoch after that. + #[tokio::test] + async fn the_bound_is_the_wall_clock_not_the_stores_tick() { + use ethlambda_state_transition::beacon::fork_choice::get_current_slot; + + let state = gloas_state(); + let head_epoch = compute_epoch_at_slot(state.slot()); + let clock_slot = compute_start_slot_at_epoch(head_epoch + 1); + let (store, _) = store_with_head_at_clock(state, gloas_config(), 0, clock_slot); + let tick_epoch = compute_epoch_at_slot(get_current_slot(&store, &store.config())); + assert_eq!( + tick_epoch, head_epoch, + "the store's tick is behind the clock" + ); + + let reply = request(store.clone(), duties_request(head_epoch + 2, r#"["1"]"#)).await; + assert_eq!(reply.status, StatusCode::OK); + + let reply = request(store, duties_request(head_epoch + 3, r#"["1"]"#)).await; + assert_eq!(reply.status, StatusCode::BAD_REQUEST); + } + #[tokio::test] async fn a_pre_gloas_epoch_is_answered_with_no_duties() { let state = gloas_state(); diff --git a/crates/net/rpc/src/beacon/validator.rs b/crates/net/rpc/src/beacon/validator.rs index d1bca404..ec2d9e5b 100644 --- a/crates/net/rpc/src/beacon/validator.rs +++ b/crates/net/rpc/src/beacon/validator.rs @@ -157,6 +157,50 @@ pub(crate) fn head(store: &Store) -> Result<(H256, Arc), ApiError> Ok((root, state)) } +/// The newest epoch a duties request may name: one past the later of the wall +/// clock's epoch and the head's. +/// +/// The Beacon API bounds duties by the current epoch, not the head's: at the +/// first slot of an epoch the head is still in the previous one until that +/// slot's block arrives, and a validator client asks for the next epoch's +/// duties right then. The head's epoch only matters when it is ahead of a +/// lagging clock, where refusing it would be a regression. +pub(crate) fn epoch_upper_bound(store: &Store, head_epoch: Epoch) -> Epoch { + let clock_epoch = compute_epoch_at_slot(crate::beacon::node::wall_slot(store)); + head_epoch.max(clock_epoch) + 1 +} + +/// The head's post-state advanced through empty slots to the first slot of +/// `epoch`, from fork choice's checkpoint-state cache (the one +/// `attestation_data` reads) so a repeated request, or one for the epoch whose +/// start `attestation_data` already advanced to, does no epoch processing. +/// +/// 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")) +} + +/// Turns a duties computation run on a blocking thread into a response. +fn duties_response( + computed: Result, tokio::task::JoinError>, +) -> Response { + 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 root of the latest block at or before `slot`, on the chain ending in /// `head_root`, whose post-state is `head_state`. /// @@ -194,21 +238,38 @@ struct ProposerDuty { /// /// Read from the `proposer_lookahead` fulu introduced and gloas keeps, which the state keeps for its own /// epoch and the next `MIN_SEED_LOOKAHEAD` epochs, so any epoch in that window -/// is answered without advancing a state. Any other epoch is refused. +/// is answered without advancing a state. A later epoch, up to one past the +/// wall clock's, is read from a copy of the head advanced to the first epoch +/// whose lookahead covers it (on a blocking thread). An epoch before the head's +/// or past that bound is a 400. /// /// `dependent_root` is v1's definition, the block root at /// `compute_start_slot_at_epoch(epoch) - 1` (the genesis block's at epoch 0). /// It is what `ethlambda validator` compares across fetches to notice a reorg. async fn get_proposer_duties(Path(epoch): Path, State(store): State) -> Response { - match proposer_duties(&store, &epoch) { - Ok(body) => crate::json_response(body), - Err(err) => err.into_response(), - } + let computed = tokio::task::spawn_blocking(move || proposer_duties(&store, &epoch)).await; + duties_response(computed) } fn proposer_duties(store: &Store, epoch: &str) -> Result { let epoch = parse_epoch(epoch)?; - let (head_root, state) = head(store)?; + let (head_root, head_state) = head(store)?; + let head_epoch = compute_epoch_at_slot(head_state.slot()); + if epoch < head_epoch || epoch > epoch_upper_bound(store, head_epoch) { + return Err(ApiError::BadRequest( + "epoch is outside the range the node serves duties for", + )); + } + + // The lookahead holds the state's own epoch and the next `MIN_SEED_LOOKAHEAD`. + // A later epoch the wall clock already allows (the head is behind it) is + // read from the head advanced to the first epoch whose lookahead covers it. + let state = if epoch > head_epoch + preset::MIN_SEED_LOOKAHEAD { + state_for_epoch(store, head_root, epoch - preset::MIN_SEED_LOOKAHEAD)? + } else { + head_state.clone() + }; + let state_epoch = compute_epoch_at_slot(state.slot()); // Gloas keeps fulu's lookahead as it is, and `upgrade_to_gloas` carries it // over, so a fulu head already holds the first gloas epoch's proposers. @@ -221,13 +282,10 @@ fn proposer_duties(store: &Store, epoch: &str) -> Result Result, ApiError>>()?; - let dependent_root = block_root_at_or_before(&state, head_root, first_slot.saturating_sub(1))?; + let dependent_root = + block_root_at_or_before(&head_state, head_root, first_slot.saturating_sub(1))?; Ok(serde_json::json!({ "dependent_root": dependent_root, "execution_optimistic": store.is_beacon_optimistic(head_root), @@ -278,8 +337,14 @@ struct AttesterDuty { /// The head state answers for its previous, current and next epoch as it is: /// an epoch's committees depend only on its seed, whose RANDAO mix is fixed a /// full epoch earlier, and on which validators are active in it, which the -/// registry records `MAX_SEED_LOOKAHEAD` epochs ahead. Any other epoch is -/// refused rather than computed from an advanced state. +/// registry records `MAX_SEED_LOOKAHEAD` epochs ahead. A later epoch is +/// computed from a copy of the head advanced with `process_slots` to the start +/// of the epoch before it (cached as fork choice's checkpoint state, run on a +/// blocking thread). The upper bound is the wall clock's, as the Beacon API +/// defines it: one past the current epoch (or the head's, if that is later), +/// since the head lags the clock until the boundary slot's block arrives and a +/// validator client asks for the next epoch's duties then. An epoch more than +/// one before the head's, or past that bound, is a 400. /// /// `dependent_root` is the block root at /// `compute_start_slot_at_epoch(epoch - 1) - 1` (the genesis block's on @@ -294,10 +359,9 @@ async fn post_attester_duties( State(store): State, Json(indices): Json>, ) -> Response { - match attester_duties(&store, &epoch, &indices) { - Ok(body) => crate::json_response(body), - Err(err) => err.into_response(), - } + let computed = + tokio::task::spawn_blocking(move || attester_duties(&store, &epoch, &indices)).await; + duties_response(computed) } fn attester_duties( @@ -311,15 +375,23 @@ fn attester_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)?; - let state_epoch = compute_epoch_at_slot(state.slot()); - if epoch + 1 < state_epoch || epoch > state_epoch + 1 { + let head_epoch = compute_epoch_at_slot(head_state.slot()); + if epoch + 1 < head_epoch || epoch > epoch_upper_bound(store, head_epoch) { return Err(ApiError::BadRequest( - "epoch is not within one epoch of the head state's", + "epoch is outside the range the node serves duties for", )); } + // The head answers as it is up to its next epoch (see the endpoint's doc); + // beyond that it is advanced to the start of the epoch before `epoch`. + let state = if epoch > head_epoch + 1 { + state_for_epoch(store, head_root, epoch - 1)? + } else { + head_state.clone() + }; + let committees = store.committee_cache().committees(&state, epoch); let committees_at_slot = committees.committees_per_slot(); let first_slot = compute_start_slot_at_epoch(epoch); @@ -350,7 +422,7 @@ fn attester_duties( } let dependent_slot = compute_start_slot_at_epoch(epoch.saturating_sub(1)).saturating_sub(1); - let dependent_root = block_root_at_or_before(&state, head_root, dependent_slot)?; + let dependent_root = block_root_at_or_before(&head_state, head_root, dependent_slot)?; Ok(serde_json::json!({ "dependent_root": dependent_root, "execution_optimistic": store.is_beacon_optimistic(head_root), @@ -582,6 +654,14 @@ mod tests { body: serde_json::Value, ) -> (StatusCode, serde_json::Value) { let (store, _root) = beacon_store_at(state); + post_to(store, uri, body).await + } + + async fn post_to( + store: Store, + uri: &str, + body: serde_json::Value, + ) -> (StatusCode, serde_json::Value) { let request = Request::post(uri) .header("content-type", "application/json") .body(Body::from(body.to_string())) @@ -658,19 +738,137 @@ mod tests { assert_eq!(json["dependent_root"], format!("{expected}")); } + /// The state's own epoch, with the wall clock at the first slot of the + /// epoch after it: the head is a slot behind the clock at an epoch + /// boundary, which is when validator clients ask for the next epoch. + fn boundary_store() -> (Store, BeaconState, Epoch) { + let state = fulu_state(); + let head_epoch = compute_epoch_at_slot(state.slot()); + let clock_slot = compute_start_slot_at_epoch(head_epoch + 1); + let (store, _root) = crate::test_utils::beacon_store_at_clock(state.clone(), clock_slot); + (store, state, head_epoch) + } + #[tokio::test] - async fn attester_duties_two_epochs_ahead_are_a_400() { + async fn attester_duties_past_the_clocks_next_epoch_are_a_400() { + // The old head-relative bound is gone: with the clock in the head's + // epoch, two epochs ahead is exactly one past the clock and so is + // refused, while the clock itself decides how far a request may reach. let state = fulu_state(); - let too_far = compute_epoch_at_slot(state.slot()) + 2; - let (status, _) = post( - state, - &format!("/eth/v1/validator/duties/attester/{too_far}"), + let head_epoch = compute_epoch_at_slot(state.slot()); + let clock_slot = compute_start_slot_at_epoch(head_epoch); + let (store, _root) = crate::test_utils::beacon_store_at_clock(state, clock_slot); + let (status, _) = post_to( + store, + &format!("/eth/v1/validator/duties/attester/{}", head_epoch + 2), serde_json::json!(["0"]), ) .await; assert_eq!(status, StatusCode::BAD_REQUEST); } + #[tokio::test] + async fn attester_duties_two_epochs_past_the_clock_are_a_400() { + let (store, _state, head_epoch) = boundary_store(); + let clock_epoch = head_epoch + 1; + let (status, _) = post_to( + store, + &format!("/eth/v1/validator/duties/attester/{}", clock_epoch + 2), + serde_json::json!(["0"]), + ) + .await; + assert_eq!(status, StatusCode::BAD_REQUEST); + } + + #[tokio::test] + async fn the_next_epoch_past_a_lagging_head_is_served_from_an_advanced_state() { + use ethlambda_state_transition::beacon::{ + helpers::accessors::get_beacon_committee, stf::process_slots, + }; + + let (store, state, head_epoch) = boundary_store(); + let epoch = head_epoch + 2; + let (status, json) = post_to( + store, + &format!("/eth/v1/validator/duties/attester/{epoch}"), + serde_json::json!(["3", "17"]), + ) + .await; + assert_eq!(status, StatusCode::OK, "{json}"); + + // What the spec derives from the head advanced to the first epoch + // that can derive `epoch`. + let mut advanced = state.clone(); + let start = compute_start_slot_at_epoch(epoch - 1); + process_slots(&mut advanced, start, &Config::mainnet()).unwrap(); + + let duties = json["data"].as_array().unwrap(); + assert_eq!(duties.len(), 2); + for duty in duties { + let slot: u64 = duty["slot"].as_str().unwrap().parse().unwrap(); + let index: u64 = duty["committee_index"].as_str().unwrap().parse().unwrap(); + let position: usize = duty["validator_committee_index"] + .as_str() + .unwrap() + .parse() + .unwrap(); + assert_eq!(compute_epoch_at_slot(slot), epoch); + let committee = get_beacon_committee(&advanced, slot, index).unwrap(); + assert_eq!( + committee[position].to_string(), + duty["validator_index"].as_str().unwrap() + ); + assert_eq!(duty["committee_length"], committee.len().to_string()); + } + } + + #[tokio::test] + async fn the_dependent_root_two_epochs_ahead_is_the_head_when_no_block_followed_it() { + let (store, state, head_epoch) = boundary_store(); + let head_root = store.beacon_head().unwrap().1; + let epoch = head_epoch + 2; + assert!(compute_start_slot_at_epoch(epoch - 1) > state.slot()); + let (_, json) = post_to( + store, + &format!("/eth/v1/validator/duties/attester/{epoch}"), + serde_json::json!(["0"]), + ) + .await; + assert_eq!(json["dependent_root"], format!("{head_root}")); + } + + #[tokio::test] + async fn proposer_duties_for_the_epoch_after_a_lagging_heads_next_are_served() { + use ethlambda_state_transition::beacon::stf::process_slots; + + let (store, state, head_epoch) = boundary_store(); + let epoch = head_epoch + 2; + let uri = format!("/eth/v1/validator/duties/proposer/{epoch}"); + let request = Request::get(uri).body(Body::empty()).unwrap(); + let response = routes().with_state(store).oneshot(request).await.unwrap(); + assert_eq!(response.status(), StatusCode::OK); + let body = response.into_body().collect().await.unwrap().to_bytes(); + let json: serde_json::Value = serde_json::from_slice(&body).unwrap(); + + let mut advanced = state.clone(); + process_slots( + &mut advanced, + compute_start_slot_at_epoch(head_epoch + 1), + &Config::mainnet(), + ) + .unwrap(); + let expected = get_beacon_proposer_indices(&advanced, epoch).unwrap(); + let data = json["data"].as_array().unwrap(); + assert_eq!(data.len() as u64, preset::SLOTS_PER_EPOCH); + for (duty, (slot, proposer)) in data + .iter() + .zip((compute_start_slot_at_epoch(epoch)..).zip(expected)) + { + assert_eq!(duty["slot"], slot.to_string()); + assert_eq!(duty["validator_index"], proposer.to_string()); + } + } + /// `fulu_state` moved `slots_past_boundary` slots into its epoch, with a /// current justified checkpoint distinct from the default so a source read /// from anywhere else shows up. @@ -994,12 +1192,25 @@ mod tests { } #[tokio::test] - async fn an_epoch_outside_the_lookahead_is_a_400() { + async fn an_epoch_past_the_clocks_next_is_a_400() { + // Adapted from a head-relative bound: the clock now decides, and the + // request is two past it. + let (store, _state, head_epoch) = boundary_store(); + let request = get_request(format!( + "/eth/v1/validator/duties/proposer/{}", + head_epoch + 3 + )); + let response = routes().with_state(store).oneshot(request).await.unwrap(); + assert_eq!(response.status(), StatusCode::BAD_REQUEST); + } + + #[tokio::test] + async fn an_epoch_before_the_head_is_a_400() { let state = fulu_state(); - let too_far = compute_epoch_at_slot(state.slot()) + 2; + let before = compute_epoch_at_slot(state.slot()) - 1; let (status, _) = get( state, - &format!("/eth/v1/validator/duties/proposer/{too_far}"), + &format!("/eth/v1/validator/duties/proposer/{before}"), ) .await; assert_eq!(status, StatusCode::BAD_REQUEST); diff --git a/crates/net/rpc/src/lib.rs b/crates/net/rpc/src/lib.rs index e03ae786..c4b6e79d 100644 --- a/crates/net/rpc/src/lib.rs +++ b/crates/net/rpc/src/lib.rs @@ -626,12 +626,33 @@ pub(crate) mod test_utils { /// [`beacon_store_at`], under `config` instead of mainnet's schedule. pub(crate) fn beacon_store_with_config(state: BeaconState, config: Config) -> (Store, H256) { + beacon_store_with_genesis(state, config, 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 config = Config::mainnet(); + let slot_secs = config.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, config, genesis) + } + + fn beacon_store_with_genesis( + state: BeaconState, + config: Config, + 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, root, Checkpoint { root, slot }, diff --git a/docs/rpc.md b/docs/rpc.md index 3c335723..3c9e9f68 100644 --- a/docs/rpc.md +++ b/docs/rpc.md @@ -238,7 +238,7 @@ surface rather than sitting beside it; a `/lean/v0` path on a beacon node is a | `GET` | `/eth/v1/node/version` | JSON | Client version string | | `GET` | `/eth/v1/node/identity` | JSON | Peer ID and metadata only (see below) | | `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 | +| `GET` | `/eth/v1/validator/duties/proposer/{epoch}` | JSON | Proposers for any epoch from the head's up to one past the wall clock's | | `POST` | `/eth/v1/validator/duties/attester/{epoch}` | JSON | Committee assignments for the given indices | | `POST` | `/eth/v1/validator/duties/ptc/{epoch}` | JSON | Payload timeliness committee seats for the given indices (gloas) | | `GET` | `/eth/v1/validator/attestation_data` | JSON | What to attest to at `slot` | @@ -262,10 +262,17 @@ These are what `ethlambda validator` needs to attest through this node. Every answer is computed from the fork-choice head's post-state, read off the store the chain actor writes, so no request waits on the actor. -- **Duties** answer for a window around the head, not any epoch. Proposer - duties read fulu's `proposer_lookahead`, which covers the head's epoch and the - next; attester duties cover the head's previous, current and next epoch, - which is as far as its shuffling is already fixed. Anything else is a `400`. +- **Duties** are bounded by the wall clock, as the Beacon API defines it: any + epoch up to one past the current one is served (or one past the head's, if + the head is ahead of a lagging clock), so a validator client's next-epoch + lookahead at an epoch boundary, while the head is still in the previous + epoch, is answered. Proposer duties read fulu's `proposer_lookahead`, which + covers the head's epoch and the next; attester duties read the head state as + it is for its previous, current and next epoch. A later epoch is computed from + a copy of the head advanced with `process_slots` to the first epoch that can + derive it, taken from fork choice's `checkpoint_state` cache and run on a + blocking thread. An epoch before the head's (attester: more than one before) + or past the bound 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. Gloas epochs are served like fulu ones: gloas keeps fulu's `proposer_lookahead` and