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 {

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 proposer is recorded only when the import ends in Imported. A gossip block whose proposer signature already verified (Accept), but that is then held for columns or the engine, or fails the STF, marks nobody, and record_liveness in verdict.rs skips Validated::Block explicitly. For doppelganger purposes a verified proposer signature is already evidence the key is in use; Lighthouse records observed_block_producers at gossip verification. Recording on Validated::Block Accept as well, and keeping this write for range-synced and API-published blocks, would cover both.

Also, this write has no test (the PR notes it). Moving it into a small helper with a unit test would keep a later refactor of this arm from dropping it unnoticed.

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
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 {

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: forward now checks outcome == Outcome::Accept in three separate places (here, aggregate pooling, aggregator-subnet pooling). One if outcome == Outcome::Accept { ... } block holding all three side effects would keep the Accept-only rule in one place.

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
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();

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: this is the same "an aggregate marks its aggregator and its attesters" rule as record_liveness in verdict.rs. A helper on ObservedLiveness, something like record_aggregate(aggregate, attesting_indices), would keep the two writers from drifting.

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
Loading
Loading