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
9 changes: 9 additions & 0 deletions crates/blockchain/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2772,6 +2772,15 @@ impl BlockChainServer {
"Block imported successfully"
);

// The proposer acted in this block's epoch, for the Beacon
// API's liveness endpoint. Recorded here rather than where
// gossip accepts a block, so a block that arrived by range
// sync or through this node's own API counts too.
if self.store.chain() == Chain::Beacon {
let epoch = ethlambda_types::beacon::signing::compute_epoch_at_slot(slot);
self.store.observed_liveness().record(epoch, proposer);
}

// Recover per-attestation single-message aggregates from the
// block's merged multi-message aggregate and fold them into
// the local pool. `Some` only for a lean block imported while
Expand Down
19 changes: 19 additions & 0 deletions crates/blockchain/state_transition/src/beacon/helpers/altair.rs
Original file line number Diff line number Diff line change
Expand Up @@ -170,6 +170,16 @@ pub fn get_next_sync_committee(state: &BeaconState) -> Result<altair::SyncCommit
})
}

/// The sync committee period `epoch` falls in.
///
/// Altair's `validator.md` `compute_sync_committee_period`. A state's
/// `current_sync_committee` serves its own period and `next_sync_committee`
/// the one after, so this is what tells a caller which of the two answers for
/// a given epoch.
pub fn compute_sync_committee_period(epoch: Epoch) -> u64 {
epoch / preset::EPOCHS_PER_SYNC_COMMITTEE_PERIOD
}

/// The reward every one-increment slice of a validator's effective balance
/// earns for a single timely, correct component of its attestation.
///
Expand Down Expand Up @@ -592,4 +602,13 @@ mod tests {
let elapsed = start.elapsed() / ITERATIONS;
println!("get_flag_index_deltas, {VALIDATOR_COUNT} validators -> {elapsed:?}/call");
}

#[test]
fn a_sync_committee_period_spans_its_epochs_and_no_more() {
let period = preset::EPOCHS_PER_SYNC_COMMITTEE_PERIOD;
assert_eq!(compute_sync_committee_period(0), 0);
assert_eq!(compute_sync_committee_period(period - 1), 0);
assert_eq!(compute_sync_committee_period(period), 1);
assert_eq!(compute_sync_committee_period(3 * period + 1), 3);
}
}
84 changes: 84 additions & 0 deletions crates/net/p2p/src/beacon/verdict.rs
Original file line number Diff line number Diff line change
Expand Up @@ -147,7 +147,14 @@ impl Validated {
/// [`pool_aggregator_attestation`]. An accepted aggregate goes into the
/// pool too, whatever else happens to it, so block production can pack
/// other nodes' votes; see [`pool_gossip_aggregate`].
///
/// Every accepted aggregate or subnet attestation also marks the
/// validators it names as live for its target epoch; see
/// [`record_liveness`].
fn forward(self, server: &P2PServer, received_at: Instant, outcome: Outcome) {
if outcome == Outcome::Accept {
record_liveness(server, &self);
}
if let Self::Aggregate { aggregate, .. } = &self
&& outcome == Outcome::Accept
{
Expand Down Expand Up @@ -201,6 +208,33 @@ impl Validated {
}
}

/// Mark the validators an accepted aggregate or subnet attestation names as
/// live for its target epoch, for `/eth/v1/validator/liveness`.
///
/// An aggregate names its aggregator and every attester its aggregate
/// signature verified, which is what makes an aggregate worth more here than
/// the attestation subnets alone: this node joins only a few of those. A
/// block's proposer is recorded where the chain actor imports it instead, so
/// a block that reached it by range sync or through this node's own API
/// counts too.
fn record_liveness(server: &P2PServer, object: &Validated) {
let observed = server.store.observed_liveness();
match object {
Validated::Aggregate {
aggregate,
attesting_indices,
} => {
let (epoch, _root) = aggregate.target();
let aggregator = std::iter::once(aggregate.aggregator_index());
observed.record_all(epoch, aggregator.chain(attesting_indices.iter().copied()));
}
Validated::Attestation { attestation, .. } => {
observed.record(attestation.data.target.epoch, attestation.attester_index);
}
Validated::Block { .. } | Validated::Column(_) => {}
}
}

/// Pool an accepted gossip aggregate for block production to pack.
///
/// Pooled here, on `Accept`, because this is where all three of its
Expand Down Expand Up @@ -737,6 +771,56 @@ mod tests {
);
}

/// An accepted aggregate names its aggregator and every attester its
/// signature verified; all of them count as live for its target epoch.
#[tokio::test]
async fn an_accepted_aggregate_marks_its_aggregator_and_attesters_live() {
let server = unconnected_beacon_server(Config::mainnet(), 0).await;
let slot = 40;
Validated::Aggregate {
aggregate: Box::new(electra_aggregate(slot, 7)),
attesting_indices: vec![11, 12],
}
.forward(&server, Instant::now(), Outcome::Accept);

let observed = server.store.observed_liveness();
let epoch = slot / 32;
for validator in [7, 11, 12] {
assert!(observed.is_live(epoch, validator), "{validator}");
}
assert!(!observed.is_live(epoch, 13));
}

#[tokio::test]
async fn an_accepted_subnet_attestation_marks_its_attester_live() {
let server = unconnected_beacon_server(Config::mainnet(), 0).await;
Validated::Attestation {
attestation: Box::new(electra_attestation(40, 5)),
subnet_id: 0,
}
.forward(&server, Instant::now(), Outcome::Accept);
assert!(server.store.observed_liveness().is_live(40 / 32, 5));
}

/// Liveness rests on the same verdict the pool does: an object that was
/// not accepted proves nothing about the validators it names.
#[tokio::test]
async fn an_object_that_was_not_accepted_marks_nobody_live() {
let server = unconnected_beacon_server(Config::mainnet(), 0).await;
Validated::Aggregate {
aggregate: Box::new(electra_aggregate(40, 7)),
attesting_indices: vec![11],
}
.forward(
&server,
Instant::now(),
Outcome::Ignore(IgnoreReason::Overloaded),
);
let observed = server.store.observed_liveness();
assert!(!observed.is_live(1, 7));
assert!(!observed.is_live(1, 11));
}

/// Only `Accept` means the signatures were verified; anything else must
/// stay out of the pool, since one unverified attestation fails the
/// whole block it is packed into.
Expand Down
42 changes: 34 additions & 8 deletions crates/net/rpc/src/beacon/config.rs
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
//! `/eth/v1/config/spec`.
//! `/eth/v1/config/spec` and `/eth/v1/config/deposit_contract`.
//!
//! The Beacon API asks for three things in one flat object: the network's
//! configuration, the preset the node was built against, and the
Expand Down Expand Up @@ -29,7 +29,22 @@ use ethlambda_types::beacon::{config::Config, constants, preset, serde_helpers::
use serde_json::{Map, Value};

pub(crate) fn routes() -> Router<Store> {
Router::new().route("/eth/v1/config/spec", get(get_spec))
Router::new()
.route("/eth/v1/config/spec", get(get_spec))
.route("/eth/v1/config/deposit_contract", get(get_deposit_contract))
}

/// `GET /eth/v1/config/deposit_contract`: the network's deposit contract, off
/// the same `Config` `/eth/v1/config/spec` reports it from (as
/// `DEPOSIT_CHAIN_ID` and `DEPOSIT_CONTRACT_ADDRESS`).
async fn get_deposit_contract(State(store): State<Store>) -> Response {
let config = store.config();
crate::json_response(serde_json::json!({
"data": {
"chain_id": config.deposit_chain_id.to_string(),
"address": HexPrefixed(&config.deposit_contract_address).to_string(),
}
}))
}

async fn get_spec(State(store): State<Store>) -> Response {
Expand Down Expand Up @@ -202,15 +217,14 @@ mod tests {
use tower::ServiceExt as _;

async fn get_spec_json() -> serde_json::Value {
get_json("/eth/v1/config/spec").await
}

async fn get_json(uri: &str) -> serde_json::Value {
let fixture = beacon_fixture(64);
let app = routes().with_state(fixture.store);
let response = app
.oneshot(
Request::builder()
.uri("/eth/v1/config/spec")
.body(Body::empty())
.unwrap(),
)
.oneshot(Request::builder().uri(uri).body(Body::empty()).unwrap())
.await
.unwrap();

Expand All @@ -219,6 +233,18 @@ mod tests {
serde_json::from_slice(&body).unwrap()
}

/// The fixture's store is bootstrapped with `Config::mainnet()`, whose
/// deposit contract is the one on Ethereum mainnet.
#[tokio::test]
async fn the_deposit_contract_is_mainnets() {
let json = get_json("/eth/v1/config/deposit_contract").await;
assert_eq!(json["data"]["chain_id"], "1");
assert_eq!(
json["data"]["address"],
"0x00000000219ab540356cbb839cbe05303d7705fa"
);
}

#[tokio::test]
async fn the_spec_is_screaming_snake_case_with_quoted_values() {
let json = get_spec_json().await;
Expand Down
72 changes: 70 additions & 2 deletions crates/net/rpc/src/beacon/pool.rs
Original file line number Diff line number Diff line change
Expand Up @@ -114,6 +114,11 @@ async fn post_pool_attestations(
// messages, so without this an aggregator served by this node would
// be missing its own validator client's votes.
let published = checked.and_then(|checked| {
// Live for liveness too: gossip never delivers this node its own
// validator clients' messages, so this is where it sees them.
store
.observed_liveness()
.record(attestation.data.target.epoch, validator);
pool.lock().expect("attestation pool lock poisoned").insert(
&attestation,
checked.committee_position,
Expand Down Expand Up @@ -301,13 +306,19 @@ async fn post_aggregate_and_proofs(
let aggregator = aggregate.aggregator_index();
let seen = aggregate::SeenAggregates::new(capacity, capacity);
let checked = aggregate::cheap_checks(&seen, &store, &aggregate, now_ms)
.and_then(|()| aggregate::stateful_checks(&store, &aggregate).map(|_| ()));
.and_then(|()| aggregate::stateful_checks(&store, &aggregate));
let published = checked
.map_err(|outcome: Outcome| {
warn!(%slot, aggregator, ?outcome, "Refused a submitted aggregate");
"aggregate failed validation"
})
.and_then(|()| {
.and_then(|attesting_indices| {
// The aggregator and every attester its signature verified are
// live, as when P2P accepts a gossip aggregate.
let (epoch, _root) = aggregate.target();
store
.observed_liveness()
.record_all(epoch, std::iter::once(aggregator).chain(attesting_indices));
// Recorded for block production, which packs the aggregates
// this node has validated.
pool.lock()
Expand Down Expand Up @@ -728,6 +739,63 @@ mod tests {
assert!(fixture.network.aggregates.lock().unwrap().is_empty());
}

/// Gossip never delivers a node its own validator clients' messages, so a
/// submission through this API is where the node sees them act.
#[tokio::test]
async fn a_published_attestation_marks_its_attester_live() {
let fixture = fixture();
let attestation = attestation(&fixture, 0, 0);
let epoch = attestation.data.target.epoch;
submit(&fixture, std::slice::from_ref(&attestation)).await;
let observed = fixture.store.observed_liveness();
assert!(observed.is_live(epoch, attestation.attester_index));
}

#[tokio::test]
async fn a_refused_attestation_marks_nobody_live() {
let fixture = fixture();
let mut forged = attestation(&fixture, 0, 0);
forged.signature = attestation(&fixture, 0, 1).signature;
let epoch = forged.data.target.epoch;
submit(&fixture, std::slice::from_ref(&forged)).await;
assert!(
!fixture
.store
.observed_liveness()
.is_live(epoch, forged.attester_index)
);
}

/// The aggregate is built on one node and submitted to a fresh one, so its
/// attesters can only have been marked live by the aggregate itself, not
/// by their own votes.
#[tokio::test]
async fn a_published_aggregate_marks_its_aggregator_and_attesters_live() {
let builder = fixture();
let slot = builder.state.slot();
let committee = get_beacon_committee(&builder.state, slot, 0).unwrap();
let votes: Vec<SingleAttestation> = (0..committee.len())
.map(|position| attestation(&builder, 0, position))
.collect();
submit(&builder, &votes).await;
let aggregate = builder
.pool
.lock()
.unwrap()
.aggregate(votes[0].data.hash_tree_root(), slot, 0)
.unwrap();

let fresh = fixture();
let signed = signed_aggregate(&fresh, committee[0], aggregate);
let (status, json) = submit_aggregates(&fresh, std::slice::from_ref(&signed)).await;
assert_eq!(status, StatusCode::OK, "{json}");
let epoch = compute_epoch_at_slot(slot);
let observed = fresh.store.observed_liveness();
for validator in &committee {
assert!(observed.is_live(epoch, *validator), "{validator}");
}
}

#[tokio::test]
async fn a_pre_electra_fork_header_is_refused() {
let fixture = fixture();
Expand Down
53 changes: 53 additions & 0 deletions crates/net/rpc/src/beacon/states.rs
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,7 @@ use crate::{
pub(crate) fn routes() -> Router<Store> {
Router::new()
.route("/eth/v2/debug/beacon/states/{state_id}", get(get_state))
.route("/eth/v1/beacon/states/{state_id}/fork", get(get_fork))
.route(
"/eth/v1/beacon/states/{state_id}/finality_checkpoints",
get(get_finality_checkpoints),
Expand Down Expand Up @@ -100,6 +101,20 @@ async fn get_state(
with_consensus_version(response, fork)
}

/// `GET /eth/v1/beacon/states/{state_id}/fork`: the `Fork` the state carries,
/// which is what a validator client builds its signing domains from.
async fn get_fork(Path(state_id): Path<String>, State(store): State<Store>) -> Response {
let (root, state) = match load(&store, &state_id) {
Ok(found) => found,
Err(err) => return err.into_response(),
};
crate::json_response(serde_json::json!({
"execution_optimistic": store.is_beacon_optimistic(root),
"finalized": is_finalized(&store, state.slot()),
"data": state.fork(),
}))
}

async fn get_finality_checkpoints(
Path(state_id): Path<String>,
State(store): State<Store>,
Expand Down Expand Up @@ -439,6 +454,44 @@ mod tests {
}
}

#[tokio::test]
async fn the_fork_is_the_one_the_state_carries() {
let fixture = beacon_fixture(ANCHOR_SLOT);
let head_state = fixture
.store
.get_state(&fixture.head_root)
.unwrap()
.unwrap();
let expected = serde_json::to_value(head_state.fork()).unwrap();

let response = get("/eth/v1/beacon/states/head/fork", None).await;
assert_eq!(response.status(), StatusCode::OK);
let json = body_json(response).await;
assert_eq!(json["data"], expected);
assert!(
json["data"]["epoch"].is_string(),
"the epoch must be quoted"
);
assert!(json["execution_optimistic"].is_boolean());
assert!(json["finalized"].is_boolean());
}

#[tokio::test]
async fn the_finalized_states_fork_is_marked_finalized() {
let response = get("/eth/v1/beacon/states/finalized/fork", None).await;
assert_eq!(response.status(), StatusCode::OK);
assert_eq!(body_json(response).await["finalized"], true);
}

/// The same refusal every other state endpoint gives: state roots are not
/// indexed, so a `0x` id is a 404 rather than a guess.
#[tokio::test]
async fn a_fork_by_state_root_is_a_404() {
let root = format!("0x{}", "ab".repeat(32));
let response = get(&format!("/eth/v1/beacon/states/{root}/fork"), None).await;
assert_eq!(response.status(), StatusCode::NOT_FOUND);
}

mod validators {
use super::*;
use crate::test_utils::beacon_store_at;
Expand Down
Loading
Loading