Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
187 changes: 171 additions & 16 deletions crates/net/rpc/src/beacon/validator.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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<Arc<BeaconState>, 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,
Expand All @@ -116,21 +140,40 @@ fn sync_duties(
.map(|index| index.parse::<ValidatorIndex>())
.collect::<Result<Vec<_>, _>>()
.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:
Expand Down Expand Up @@ -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,
&ethlambda_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<String> = (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<String> = 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<String>)> = 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]
Expand Down
18 changes: 17 additions & 1 deletion crates/net/rpc/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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 },
Expand Down
11 changes: 8 additions & 3 deletions docs/rpc.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 |
Expand All @@ -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
Expand Down
9 changes: 5 additions & 4 deletions docs/spec_deviations.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
Loading