-
Notifications
You must be signed in to change notification settings - Fork 30
feat(rpc): serve states/{id}/validators/{validator_id} and v2 proposer duties, for Lighthouse's validator client #653
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. Weβll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from all commits
40b0729
922af02
5d667cb
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -46,6 +46,10 @@ pub(crate) fn routes() -> Router<Store> { | |
| "/eth/v1/beacon/states/{state_id}/validators", | ||
| get(get_validators).post(post_validators), | ||
| ) | ||
| .route( | ||
| "/eth/v1/beacon/states/{state_id}/validators/{validator_id}", | ||
| get(get_validator), | ||
| ) | ||
| } | ||
|
|
||
| /// Resolve a `state_id` to the block root its state is stored under. | ||
|
|
@@ -346,6 +350,61 @@ fn validators_response(store: &Store, state_id: &str, request: ValidatorsRequest | |
| })) | ||
| } | ||
|
|
||
| /// `GET .../validators/{validator_id}`: one registry entry, by index or by | ||
| /// public key, as the `data` object itself rather than a one-element list. | ||
| /// | ||
| /// What a validator client resolves its keys' indices with: Lighthouse calls | ||
| /// it once per key and knows no index for a key it gets a 404 for, which | ||
| /// leaves that validator unable to take any duty. A key looked up this way is | ||
| /// found by one scan of the registry. | ||
| async fn get_validator( | ||
| Path((state_id, validator_id)): Path<(String, String)>, | ||
| State(store): State<Store>, | ||
| ) -> Response { | ||
| match validator_response(&store, &state_id, &validator_id) { | ||
| Ok(body) => crate::json_response(body), | ||
| Err(err) => err.into_response(), | ||
| } | ||
| } | ||
|
|
||
| fn validator_response( | ||
| store: &Store, | ||
| state_id: &str, | ||
| validator_id: &str, | ||
| ) -> Result<serde_json::Value, ApiError> { | ||
| let id = ValidatorId::parse(validator_id)?; | ||
| let (root, state) = load(store, state_id)?; | ||
| let index = match id { | ||
| ValidatorId::Index(index) => index, | ||
| ValidatorId::Pubkey(pubkey) => state | ||
|
Collaborator
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Nit (perf): each lookup by pubkey is a linear scan of the registry, run on a tokio worker. Lighthouse's VC calls this once per key at startup and again for keys it has no index for, so on a mainnet-sized registry that's one full scan per key. A pubkey-to-index map (or |
||
| .validators() | ||
| .iter() | ||
| .position(|validator| validator.pubkey == pubkey) | ||
| .ok_or(ApiError::NotFound("validator not found"))? | ||
| as ValidatorIndex, | ||
| }; | ||
| let validator = state | ||
| .validator(index) | ||
| .map_err(|_| ApiError::NotFound("validator not found"))?; | ||
| let balance = *state | ||
| .balances() | ||
| .get(index as usize) | ||
| .ok_or(ApiError::Internal( | ||
| "the registry and balances disagree in length", | ||
| ))?; | ||
| let status = ValidatorStatus::of(validator, balance, compute_epoch_at_slot(state.slot())); | ||
| Ok(serde_json::json!({ | ||
| "execution_optimistic": store.is_beacon_optimistic(root), | ||
| "finalized": is_finalized(store, state.slot()), | ||
| "data": ValidatorEntry { | ||
| index, | ||
| balance, | ||
| status: status.name(), | ||
| validator, | ||
| }, | ||
| })) | ||
| } | ||
|
|
||
| #[cfg(test)] | ||
| mod tests { | ||
| use super::*; | ||
|
|
@@ -520,6 +579,49 @@ mod tests { | |
| } | ||
|
|
||
| /// What `ethlambda validator` sends: its keys, to learn their indices. | ||
|
Collaborator
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Nit: |
||
| async fn get_one(validator_id: &str) -> axum::response::Response { | ||
| let (app, _) = app(); | ||
| let uri = format!("/eth/v1/beacon/states/head/validators/{validator_id}"); | ||
| app.oneshot(Request::get(uri).body(Body::empty()).unwrap()) | ||
| .await | ||
| .unwrap() | ||
| } | ||
|
|
||
| /// What Lighthouse's validator client resolves each key's index with. | ||
| /// The answer is the same entry the list endpoint gives, as an object. | ||
| #[tokio::test] | ||
| async fn one_validator_by_pubkey_or_index_is_the_list_entry() { | ||
| let (_, state) = app(); | ||
| let listed = body_json(post(serde_json::json!({ "ids": ["5"] })).await).await; | ||
| let expected = &listed["data"][0]; | ||
|
|
||
| for id in [pubkey_hex(&state, 5), "5".to_owned()] { | ||
| let response = get_one(&id).await; | ||
| assert_eq!(response.status(), StatusCode::OK, "{id}"); | ||
| let json = body_json(response).await; | ||
| assert_eq!(&json["data"], expected, "{id}"); | ||
| assert!(json["execution_optimistic"].is_boolean()); | ||
| assert!(json["finalized"].is_boolean()); | ||
| } | ||
| } | ||
|
|
||
| /// Unlike the list endpoint, which omits an id naming no validator, | ||
| /// this one has nothing to answer with. | ||
| #[tokio::test] | ||
| async fn one_unknown_validator_is_a_404() { | ||
| let unknown_key = format!("0x{}", "ab".repeat(48)); | ||
| for id in [unknown_key, COUNT.to_string()] { | ||
| assert_eq!(get_one(&id).await.status(), StatusCode::NOT_FOUND, "{id}"); | ||
| } | ||
| } | ||
|
|
||
| #[tokio::test] | ||
| async fn one_malformed_validator_id_is_a_400() { | ||
| for id in ["0x1234", "not-an-index"] { | ||
| assert_eq!(get_one(id).await.status(), StatusCode::BAD_REQUEST, "{id}"); | ||
| } | ||
| } | ||
|
|
||
| #[tokio::test] | ||
| async fn a_pubkey_resolves_to_its_index() { | ||
| let (_, state) = app(); | ||
|
|
||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -46,6 +46,10 @@ pub(crate) fn routes() -> Router<Store> { | |
| "/eth/v1/validator/duties/proposer/{epoch}", | ||
| get(get_proposer_duties), | ||
| ) | ||
| .route( | ||
| "/eth/v2/validator/duties/proposer/{epoch}", | ||
|
Collaborator
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Nit: the module docs still say this file serves the endpoints under |
||
| get(get_proposer_duties_v2), | ||
| ) | ||
| .route( | ||
| "/eth/v1/validator/duties/attester/{epoch}", | ||
| post(post_attester_duties), | ||
|
|
@@ -388,13 +392,62 @@ struct ProposerDuty { | |
| /// `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<String>, State(store): State<Store>) -> Response { | ||
| match proposer_duties(&store, &epoch) { | ||
| match proposer_duties(&store, &epoch, DependentRoot::V1) { | ||
| Ok(body) => crate::json_response(body), | ||
| Err(err) => err.into_response(), | ||
| } | ||
| } | ||
|
|
||
| fn proposer_duties(store: &Store, epoch: &str) -> Result<serde_json::Value, ApiError> { | ||
| /// `GET /eth/v2/validator/duties/proposer/{epoch}`: the same duties as v1, | ||
| /// with v2's `dependent_root`. | ||
| /// | ||
| /// Fulu fixes an epoch's proposers one epoch ahead (EIP-7917's lookahead), so | ||
| /// they depend on the chain as of the end of epoch `epoch - 2`, not | ||
| /// `epoch - 1` as v1 has it. A validator client that compares v1's root | ||
| /// across fetches sees a spurious reorg every epoch; Lighthouse asks for v2 | ||
| /// by default and does not fall back to v1. `503` while syncing, as the | ||
| /// Beacon API lists. | ||
| async fn get_proposer_duties_v2( | ||
| Path(epoch): Path<String>, | ||
| State(store): State<Store>, | ||
| Extension(sync_status): Extension<SyncStatusController>, | ||
| ) -> Response { | ||
| if sync_status.get() == SyncStatus::Syncing { | ||
|
Collaborator
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Minor: v2 is gated on sync status but v1 isn't, although both answer from the same head-state lookahead. If the gate matters, v1 (what |
||
| return ApiError::ServiceUnavailable("the node is syncing").into_response(); | ||
| } | ||
| match proposer_duties(&store, &epoch, DependentRoot::V2) { | ||
| Ok(body) => crate::json_response(body), | ||
| Err(err) => err.into_response(), | ||
| } | ||
| } | ||
|
|
||
| /// Which definition of a proposer duty's `dependent_root` to answer with. | ||
| #[derive(Debug, Clone, Copy)] | ||
| enum DependentRoot { | ||
| /// `get_block_root_at_slot(state, compute_start_slot_at_epoch(epoch) - 1)`. | ||
| V1, | ||
| /// `get_block_root_at_slot(state, compute_start_slot_at_epoch(epoch - 1) - 1)`. | ||
| V2, | ||
| } | ||
|
|
||
| impl DependentRoot { | ||
| /// The slot whose block root this definition names for `epoch`. Either | ||
| /// definition names the genesis block where its subtraction would | ||
| /// underflow, which is slot 0. | ||
| fn slot(self, epoch: Epoch) -> Slot { | ||
| let epoch = match self { | ||
| Self::V1 => epoch, | ||
| Self::V2 => epoch.saturating_sub(1), | ||
|
Collaborator
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Minor (one epoch per network): v2 always uses the post-Fulu rule. Lighthouse's With (This rule is fork-dependent and |
||
| }; | ||
| compute_start_slot_at_epoch(epoch).saturating_sub(1) | ||
| } | ||
| } | ||
|
|
||
| fn proposer_duties( | ||
| store: &Store, | ||
| epoch: &str, | ||
| dependent: DependentRoot, | ||
| ) -> Result<serde_json::Value, ApiError> { | ||
| let epoch = parse_epoch(epoch)?; | ||
| let (head_root, state) = head(store)?; | ||
|
Collaborator
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Medium: the requested epoch is bounded by the head state's epoch (the Advancing the head state to the clock's epoch (as |
||
|
|
||
|
|
@@ -429,7 +482,7 @@ fn proposer_duties(store: &Store, epoch: &str) -> Result<serde_json::Value, ApiE | |
| }) | ||
| .collect::<Result<Vec<_>, 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(&state, head_root, dependent.slot(epoch))?; | ||
| Ok(serde_json::json!({ | ||
| "dependent_root": dependent_root, | ||
| "execution_optimistic": store.is_beacon_optimistic(head_root), | ||
|
|
@@ -699,6 +752,76 @@ mod tests { | |
| assert_eq!(json["dependent_root"], format!("{expected}")); | ||
| } | ||
|
|
||
| async fn get_v2( | ||
| state: BeaconState, | ||
| epoch: u64, | ||
| sync_status: SyncStatusController, | ||
| ) -> (StatusCode, serde_json::Value) { | ||
| let (store, _root) = beacon_store_at(state); | ||
| let request = Request::get(format!("/eth/v2/validator/duties/proposer/{epoch}")) | ||
| .body(Body::empty()) | ||
| .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()) | ||
| } | ||
|
|
||
| /// v2 changes only the dependent root: the duties are the same ones. | ||
| #[tokio::test] | ||
| async fn v2_proposer_duties_are_v1_s() { | ||
| let state = fulu_state(); | ||
| let state_epoch = compute_epoch_at_slot(state.slot()); | ||
| for epoch in [state_epoch, state_epoch + 1] { | ||
| let (_, v1) = get( | ||
| state.clone(), | ||
| &format!("/eth/v1/validator/duties/proposer/{epoch}"), | ||
| ) | ||
| .await; | ||
| let (status, v2) = get_v2(state.clone(), epoch, Default::default()).await; | ||
| assert_eq!(status, StatusCode::OK); | ||
| assert_eq!(v2["data"], v1["data"], "epoch {epoch}"); | ||
| } | ||
| } | ||
|
|
||
| /// v2's root is the block before the *previous* epoch: for the head's next | ||
| /// epoch that is the block before the head's own epoch, v1's root for the | ||
| /// head's epoch. | ||
| #[tokio::test] | ||
| async fn the_v2_dependent_root_is_the_block_before_the_previous_epoch() { | ||
| let state = fulu_state(); | ||
| let next = compute_epoch_at_slot(state.slot()) + 1; | ||
| let (_, json) = get_v2(state.clone(), next, Default::default()).await; | ||
| let before = compute_start_slot_at_epoch(next - 1) - 1; | ||
| let expected = get_block_root_at_slot(&state, before).unwrap(); | ||
| assert_eq!(json["dependent_root"], format!("{expected}")); | ||
| } | ||
|
|
||
| /// `compute_start_slot_at_epoch(epoch - 1) - 1` underflows for epochs 0 | ||
| /// and 1, where the spec names the genesis block instead. | ||
| #[tokio::test] | ||
| async fn the_v2_dependent_root_is_genesis_where_it_would_underflow() { | ||
| let state = fulu_state(); | ||
| assert_eq!( | ||
| compute_epoch_at_slot(state.slot()), | ||
| 1, | ||
| "the fixture's epoch" | ||
| ); | ||
| let (_, json) = get_v2(state.clone(), 1, Default::default()).await; | ||
| let genesis = get_block_root_at_slot(&state, 0).unwrap(); | ||
| assert_eq!(json["dependent_root"], format!("{genesis}")); | ||
| } | ||
|
|
||
| #[tokio::test] | ||
| async fn v2_proposer_duties_answer_503_while_syncing() { | ||
| let state = fulu_state(); | ||
| let epoch = compute_epoch_at_slot(state.slot()); | ||
| let syncing = SyncStatusController::new(SyncStatus::Syncing); | ||
| let (status, _) = get_v2(state, epoch, syncing).await; | ||
| assert_eq!(status, StatusCode::SERVICE_UNAVAILABLE); | ||
| } | ||
|
|
||
| #[tokio::test] | ||
| async fn the_next_epochs_dependent_block_is_the_head_itself() { | ||
| // The last slot of the head's epoch has not happened yet, so the latest | ||
|
|
||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -25,6 +25,15 @@ | |
| # client fails over per call, so everything ethlambda beacon serves goes through | ||
| # it. With `ethlambda_beacon.fallback: false` the participant's node is left | ||
| # out of the list entirely, and ethlambda beacon serves every duty alone. | ||
| # | ||
| # # Which validator client | ||
| # | ||
| # `ethlambda validator` by default. With `ethlambda_validator.client: | ||
| # lighthouse`, Lighthouse's validator client signs for the same keys instead, | ||
| # against the same beacon node list: that is how the Beacon API ethlambda | ||
| # beacon serves is checked against a client other than its own. | ||
| # `ethlambda_validator.doppelganger: true` turns on Lighthouse's doppelganger | ||
| # protection, which calls `/eth/v1/validator/liveness`. | ||
|
|
||
| ethereum_package = import_module("github.com/ethpandaops/ethereum-package/main.star") | ||
|
|
||
|
|
@@ -111,6 +120,15 @@ def run(plan, args={}): | |
| mnemonic = network_params.get("preregistered_validator_keys_mnemonic", DEFAULT_MNEMONIC) | ||
| derive_keys(plan, mnemonic, first, last) | ||
|
|
||
| if vc.get("client", "ethlambda") == "lighthouse": | ||
|
Collaborator
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Nit: any |
||
| launch_lighthouse_vc(plan, vc, beacon_nodes, name) | ||
| plan.print( | ||
| "Lighthouse's validator client signs for validators [{}, {}) via {}".format( | ||
| first, last, ", ".join(beacon_nodes) | ||
| ) | ||
| ) | ||
| return output | ||
|
|
||
| cmd = [ | ||
| "validator", | ||
| "--beacon-nodes", | ||
|
|
@@ -152,6 +170,57 @@ def run(plan, args={}): | |
| return output | ||
|
|
||
|
|
||
| def launch_lighthouse_vc(plan, vc, beacon_nodes, name): | ||
| """Run Lighthouse's validator client on the derived keys. | ||
|
|
||
| eth2-val-tools writes keys in Lighthouse's own layout | ||
| (`keys/<pubkey>/voting-keystore.json` and `secrets/<pubkey>`), and with no | ||
| `validator_definitions.yml` in the validators directory Lighthouse | ||
| discovers them and writes one. It writes that file into the directory, so | ||
| the keys are copied out of the artifact first rather than used in place. | ||
| No `--datadir`: Lighthouse refuses it alongside `--validators-dir`, and | ||
| keeps its slashing-protection database in the validators directory, which | ||
| is the writable copy. `--init-slashing-protection` because that database | ||
| starts empty. | ||
| """ | ||
| flags = [ | ||
| "--testnet-dir={}".format(GENESIS_MOUNT), | ||
| "--beacon-nodes={}".format(",".join(beacon_nodes)), | ||
| "--validators-dir=/data/raw/keys", | ||
| "--secrets-dir=/data/raw/secrets", | ||
| "--init-slashing-protection", | ||
| "--graffiti={}".format(vc.get("graffiti", name)), | ||
| "--metrics", | ||
| "--metrics-address=0.0.0.0", | ||
| "--metrics-port={}".format(METRICS_PORT), | ||
| ] | ||
| fee_recipient = vc.get("suggested_fee_recipient", "") | ||
| if fee_recipient: | ||
| flags.append("--suggested-fee-recipient={}".format(fee_recipient)) | ||
| if vc.get("doppelganger", False): | ||
| flags.append("--enable-doppelganger-protection") | ||
| flags += vc.get("lighthouse_extra_params", []) | ||
|
|
||
| plan.add_service( | ||
| name="vc-lighthouse", | ||
| config=ServiceConfig( | ||
| image=vc.get("lighthouse_image", "sigp/lighthouse:latest"), | ||
|
Collaborator
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Nit: |
||
| entrypoint=["sh", "-c"], | ||
| cmd=[ | ||
| "mkdir -p /data && cp -r {}/raw /data/ && lighthouse vc {}".format( | ||
|
Collaborator
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Minor: the flags are joined into one |
||
| KEYS_MOUNT, " ".join(flags) | ||
| ) | ||
| ], | ||
| files={KEYS_MOUNT: KEYS_ARTIFACT, GENESIS_MOUNT: GENESIS_ARTIFACT}, | ||
| ports={ | ||
| "metrics": PortSpec( | ||
| number=METRICS_PORT, transport_protocol="TCP", application_protocol="http" | ||
| ), | ||
| }, | ||
| ), | ||
| ) | ||
|
|
||
|
|
||
| def derive_keys(plan, mnemonic, first, last): | ||
| """Derive keystores for [first, last) and the definitions file the client reads. | ||
|
|
||
|
|
||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Nit: the entry building here (id match,
ValidatorStatus::of, balance lookup,ValidatorEntry) repeatsvalidators_response. Running the list selection with one id and returning its single match or a 404 would keep one copy.