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
102 changes: 102 additions & 0 deletions crates/net/rpc/src/beacon/states.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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(

Copy link
Copy Markdown
Collaborator

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) repeats validators_response. Running the list selection with one id and returning its single match or a 404 would keep one copy.

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

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The 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 spawn_blocking) would help if this runs on a large network; fine for kurtosis.

.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::*;
Expand Down Expand Up @@ -520,6 +579,49 @@ mod tests {
}

/// What `ethlambda validator` sends: its keys, to learn their indices.

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Nit: get_one was inserted under the existing doc comment, so "What ethlambda validator sends: its keys, to learn their indices." now documents the GET helper (which is what Lighthouse uses), and a_pubkey_resolves_to_its_index lost its comment.

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();
Expand Down
129 changes: 126 additions & 3 deletions crates/net/rpc/src/beacon/validator.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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}",

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Nit: the module docs still say this file serves the endpoints under /eth/v1/validator/; it now serves /eth/v2/validator/duties/proposer too. Same for the header list in states.rs, which omits /fork (from #642) and /validators/{validator_id}.

get(get_proposer_duties_v2),
)
.route(
"/eth/v1/validator/duties/attester/{epoch}",
post(post_attester_duties),
Expand Down Expand Up @@ -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 {

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The 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 ethlambda validator uses) needs it too; if it doesn't, v2 can drop it. Also, this if Syncing { 503 } block is now in three handlers (sync duties, liveness, v2 proposer duties); a reject_if_syncing helper or a route_layer on the gated routes would keep them identical.

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),

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The 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 proposer_shuffling_decision_slot (consensus/types/src/core/chain_spec.rs) is fork-aware: it uses start(epoch - 1) - 1 only when Fulu is active at epoch - 1, and the v1 rule otherwise. At the Fulu fork epoch F itself, the proposers were fixed by the upgrade from the end of epoch F-1, so v2 there equals v1.

With fulu_fork_epoch: 0 (the kurtosis configs) the two agree. With a later fork epoch, at epoch F this answers a different root than Lighthouse's BN, so a VC failing over between the two sees a spurious reorg, and a reorg of F-1's tail goes unnoticed. Checking the fork at epoch - 1, and using MIN_SEED_LOOKAHEAD instead of the literal 1, would match.

(This rule is fork-dependent and attester_duties' identical-looking formula is not, so they shouldn't be merged into one helper.)

};
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)?;

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The 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 offset <= MIN_SEED_LOOKAHEAD check below), not the clock. Lighthouse's VC asks v2 for current_epoch and current_epoch + 1 on every poll (poll_beacon_proposers in duties_service.rs). From the start of epoch E until a block lands in E, the head state is still in E-1, so the E+1 request is a 400 and Lighthouse logs an error each poll; with several missed slots at the start of E, next-epoch duties stay unavailable that long.

Advancing the head state to the clock's epoch (as attestation_data does through fork_choice::checkpoint_state) would fix it. v1 has the same bound, so it isn't new, but v2 is the path Lighthouse uses.


Expand Down Expand Up @@ -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),
Expand Down Expand Up @@ -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
Expand Down
7 changes: 6 additions & 1 deletion docs/rpc.md
Original file line number Diff line number Diff line change
Expand Up @@ -239,7 +239,9 @@ 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/beacon/states/{state_id}/validators/{validator_id}` | JSON | One registry entry, by index or pubkey; `404` if there is none |
| `GET` | `/eth/v1/validator/duties/proposer/{epoch}` | JSON | Proposers for the head's epoch or the next |
| `GET` | `/eth/v2/validator/duties/proposer/{epoch}` | JSON | The same, with v2's `dependent_root` (what Lighthouse asks for) |
| `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/liveness/{epoch}` | JSON | Whether this node saw each validator act in the epoch (doppelganger protection) |
Expand All @@ -262,7 +264,10 @@ the chain actor writes, so no request waits on the actor.
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`.
`dependent_root` follows each endpoint's v1 definition. Attester duties walk
`dependent_root` follows each endpoint's v1 definition, except proposer
duties v2, whose root is the block before the *previous* epoch (fulu fixes
proposers an epoch ahead), or genesis where that would underflow; v2 is a
`503` while syncing. 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
Expand Down
69 changes: 69 additions & 0 deletions tooling/kurtosis-validator/main.star
Original file line number Diff line number Diff line change
Expand Up @@ -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")

Expand Down Expand Up @@ -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":

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Nit: any client value other than exactly "lighthouse" (Lighthouse, lodestar, ...) silently runs ethlambda's VC and drops the Lighthouse-only keys. A fail() on unknown values would make a typo visible.

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",
Expand Down Expand Up @@ -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"),

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Nit: sigp/lighthouse:latest floats. Pinning the version this was checked with (v8.2.3) keeps a future Lighthouse release from changing what this config tests without a change here.

entrypoint=["sh", "-c"],
cmd=[
"mkdir -p /data && cp -r {}/raw /data/ && lighthouse vc {}".format(

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Minor: the flags are joined into one sh -c string without quoting, so any value with a space breaks the command (graffiti: "ethlambda vc" passes vc as a stray positional argument and lighthouse vc exits), and ; or $(...) runs as shell. Shell-quoting each flag before the join would avoid both.

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.

Expand Down
Loading
Loading