From 7b681d3e6bef92b36274a4f763231c4ce1c72d91 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Tom=C3=A1s=20Gr=C3=BCner?= <47506558+MegaRedHand@users.noreply.github.com> Date: Tue, 29 Sep 2026 20:23:11 -0300 Subject: [PATCH] feat(beacon): score gossipsub peers with lighthouse's parameters The beacon wire had no way to push out peers that send invalid messages, break IWANT promises or cannot keep up: gossipsub kept exchanging with all of them, and the connection cap was the only bound on bad peers. The parameters are lighthouse's, ported formula for formula. Teku, Lodestar and Grandine run the same ones, so a score means here what it means on most of mainnet. Two deviations: - Mesh-delivery scoring (P3) is off while the head lags the wall clock by more than four slots, and back on only after an epoch within that. Lighthouse joins these topics only once synced. This node subscribes at startup and, while catching up, ignores every aggregate voting for a block it has not imported; a mesh peer credited with no aggregates scores below the graylist, so a restart would graylist the whole aggregate mesh for our own lag. SyncStatus cannot be the gate: its network-stall rule reported synced through a 98-slot catch-up on the mainnet follower. - Peers below the graylist are disconnected every 10 s. Gossipsub alone only ignores them, and they keep a connection slot. Columns, sync committee contributions and BLS changes stay unscored. This node ignores every contribution and change for lack of a consumer, so P3 there would penalize the mesh for a gap that is ours. --- CLAUDE.md | 1 + .../src/beacon/helpers/accessors.rs | 6 +- crates/common/types/src/beacon/committees.rs | 7 + crates/net/p2p/src/beacon/mod.rs | 1 + crates/net/p2p/src/beacon/scoring.rs | 903 ++++++++++++++++++ crates/net/p2p/src/beacon/swarm.rs | 11 +- crates/net/p2p/src/gossipsub/handler.rs | 2 +- crates/net/p2p/src/gossipsub/mod.rs | 1 + crates/net/p2p/src/lib.rs | 77 +- crates/net/p2p/src/metrics.rs | 45 + crates/net/p2p/src/req_resp/handlers.rs | 1 + crates/net/p2p/src/swarm_adapter.rs | 66 +- crates/storage/src/committee_cache.rs | 42 + docs/beacon_wire.md | 63 ++ docs/metrics.md | 20 + 15 files changed, 1224 insertions(+), 22 deletions(-) create mode 100644 crates/net/p2p/src/beacon/scoring.rs diff --git a/CLAUDE.md b/CLAUDE.md index 3e93d8f6..3d0af7f6 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -297,6 +297,7 @@ actual_slot = finalized_slot + 1 + relative_index - `fork_digest` is a 4-byte hex string (no `0x` prefix); currently the dummy `12345678` agreed across clients - Mesh size: 8 (6-12 bounds), heartbeat: 700ms - Beacon wire: `validate_messages()` is on, so every beacon message waits for a verdict (~4.2s before gossipsub's cache evicts it). Rules in `state_transition::beacon::gossip` (cheap half inline, stateful half on a bounded `spawn_blocking` task); plumbing in `p2p/src/beacon/verdict.rs`. Lean gossip still auto-forwards + - Beacon peer scoring: lighthouse's parameters ported in `p2p/src/beacon/scoring.rs` (block, aggregate, 64 attestation subnets, exits, slashings; columns and sync contributions unscored), refreshed every slot from the head's active validator count. Mesh-delivery scoring (P3) is gated off while the head lags more than 4 slots and for an epoch after, because a catching-up node `Ignore`s every aggregate and would otherwise graylist its own mesh. The swarm disconnects peers below the graylist (-16000) every 10 s. Lean runs unscored. See [`docs/beacon_wire.md`](docs/beacon_wire.md#peer-scoring) - Data columns: every check runs in p2p. A column gossip did not accept (`Queue`/`Overloaded`), every fetched column, and parked columns replayed after their parent imports go through `column::chain_checks` in `p2p/src/beacon/column_checks.rs`. The chain actor stores what it gets unchecked; only debug builds re-run `chain_checks` there - Beacon subscribes seven global topics plus two node-id-derived subnet families: custody columns and backbone attestation subnets. `beacon_aggregate_and_proof` and diff --git a/crates/blockchain/state_transition/src/beacon/helpers/accessors.rs b/crates/blockchain/state_transition/src/beacon/helpers/accessors.rs index 677e895e..5abb213c 100644 --- a/crates/blockchain/state_transition/src/beacon/helpers/accessors.rs +++ b/crates/blockchain/state_transition/src/beacon/helpers/accessors.rs @@ -129,8 +129,10 @@ pub fn get_seed(state: &BeaconState, epoch: Epoch, domain_type: DomainType) -> B /// split out so [`build_epoch_committees`] can share one /// [`get_active_validator_indices`] scan between this and its own committee /// derivation, rather than [`get_committee_count_per_slot`] repeating the scan -/// the caller already did to get `active_count` in the first place. -fn committee_count_per_slot(active_count: u64) -> u64 { +/// the caller already did to get `active_count` in the first place. Public +/// for p2p's gossipsub scoring, which sizes expected aggregate traffic from +/// a validator count rather than from a state. +pub fn committee_count_per_slot(active_count: u64) -> u64 { let ideal = active_count / preset::SLOTS_PER_EPOCH / preset::TARGET_COMMITTEE_SIZE; ideal.clamp(1, preset::MAX_COMMITTEES_PER_SLOT as u64) } diff --git a/crates/common/types/src/beacon/committees.rs b/crates/common/types/src/beacon/committees.rs index ba6d1606..b3845000 100644 --- a/crates/common/types/src/beacon/committees.rs +++ b/crates/common/types/src/beacon/committees.rs @@ -104,6 +104,13 @@ impl EpochCommittees { self.committees_per_slot } + /// How many validators are active in this epoch: the shuffle is a + /// permutation of exactly that set, so its length is the count without + /// another registry scan. + pub fn active_validator_count(&self) -> u64 { + self.shuffled.len() as u64 + } + /// The committee at `slot` with `index`. /// /// The same members `ethlambda-state-transition`'s `get_beacon_committee` diff --git a/crates/net/p2p/src/beacon/mod.rs b/crates/net/p2p/src/beacon/mod.rs index 3369ccba..3b349acb 100644 --- a/crates/net/p2p/src/beacon/mod.rs +++ b/crates/net/p2p/src/beacon/mod.rs @@ -19,6 +19,7 @@ pub mod encoding; pub mod handler; pub mod messages; pub mod protocols; +pub mod scoring; pub mod subnets; pub mod swarm; pub mod topics; diff --git a/crates/net/p2p/src/beacon/scoring.rs b/crates/net/p2p/src/beacon/scoring.rs new file mode 100644 index 00000000..733c3960 --- /dev/null +++ b/crates/net/p2p/src/beacon/scoring.rs @@ -0,0 +1,903 @@ +//! Gossipsub peer scoring on the beacon wire. +//! +//! Lighthouse's parameters (`lighthouse_network/src/service/gossipsub_scoring_parameters.rs`), +//! ported formula for formula. Grandine, Teku and Lodestar run the same +//! constants, so a peer is scored here the way most of mainnet scores it, and +//! the tests pin this port to the numbers lighthouse's formulas give. What is +//! scored: +//! +//! | topic | weight | mesh deliveries (P3) | +//! |---|---|---| +//! | `beacon_block` | 0.5 | scored | +//! | `beacon_aggregate_and_proof` | 0.5 | scored | +//! | `beacon_attestation_{subnet_id}`, every subnet | 1/64 each | scored | +//! | `voluntary_exit`, `proposer_slashing`, `attester_slashing` | 0.05 each | off | +//! +//! Every subnet is given parameters, not only the backbone ones, because an +//! aggregator joins other subnets at runtime and the parameters have to be in +//! place before the mesh forms. +//! +//! Nothing else this node subscribes to has topic parameters, so only the +//! topic-independent terms (IP colocation, broken IWANT promises, slow peers) +//! see that traffic. `data_column_sidecar_{subnet_id}` is left out, as every +//! client but Prysm leaves it out. `sync_committee_contribution_and_proof` and +//! `bls_to_execution_change` are left out because this node has no consumer +//! for either and `Ignore`s every message on them. A delivery is only credited +//! once it is accepted, so mesh-delivery scoring there would starve every mesh +//! peer for a gap that is this node's own. +//! +//! # Where this differs from lighthouse +//! +//! - `mesh_n` is this node's own ([`crate::MESH_N`]) rather than the 5 that +//! lighthouse's default network load gives. It only enters the +//! first-message-delivery cap, which is two mesh shares of the expected +//! traffic. +//! - Mesh-delivery scoring also waits for this node to keep up with the chain: +//! see [`MeshDeliveryGate`]. +//! +//! # Refreshing +//! +//! The expected message rates depend on the active validator count, so the +//! block, aggregate and attestation parameters are rebuilt once per slot +//! ([`refresh`]) from the head's current-epoch shuffling. Until the first +//! refresh they are built for [`PLACEHOLDER_ACTIVE_VALIDATORS`] with P3 off, +//! the same placeholder lighthouse starts from. +//! +//! The topics are named for the one fork digest the node subscribed under at +//! startup. When the digest can change at runtime, the old digest's topics +//! need their weight zeroed and the new digest's topics need parameters, as +//! lighthouse's `remove_topic_weight_except` does. + +use std::collections::HashSet; +use std::time::Duration; + +use ethlambda_state_transition::beacon::helpers::accessors::committee_count_per_slot; +use ethlambda_types::beacon::config::Config; +use ethlambda_types::beacon::constants::TARGET_AGGREGATORS_PER_COMMITTEE; +use ethlambda_types::beacon::preset::SLOTS_PER_EPOCH; +use ethlambda_types::beacon::primitives::ForkDigest; +use libp2p::gossipsub::{IdentTopic, PeerScoreParams, PeerScoreThresholds, TopicScoreParams}; +use tracing::info; + +use crate::beacon::constants::ATTESTATION_SUBNET_COUNT; +use crate::beacon::topics::{ + ATTESTER_SLASHING, BEACON_AGGREGATE_AND_PROOF, BEACON_BLOCK, PROPOSER_SLASHING, VOLUNTARY_EXIT, + attestation_topic_name, topic_name, +}; +use crate::gossipsub::beacon_wall_slot; +use crate::{P2PServer, metrics}; + +/// The most a topic's time-in-mesh term (P1) contributes, before its weight. +const MAX_IN_MESH_SCORE: f64 = 10.0; +/// The most a topic's first-message-delivery term (P2) contributes, before +/// its weight. +const MAX_FIRST_MESSAGE_DELIVERIES_SCORE: f64 = 40.0; + +const BEACON_BLOCK_WEIGHT: f64 = 0.5; +const BEACON_AGGREGATE_PROOF_WEIGHT: f64 = 0.5; +const VOLUNTARY_EXIT_WEIGHT: f64 = 0.05; +const PROPOSER_SLASHING_WEIGHT: f64 = 0.05; +const ATTESTER_SLASHING_WEIGHT: f64 = 0.05; + +/// How late after the first delivery a mesh peer's copy still counts towards +/// its mesh deliveries: the time a hostile peer would need to replay a message +/// this node just forwarded to it and be credited for it. +const MESH_MESSAGE_DELIVERIES_WINDOW: Duration = Duration::from_secs(2); + +/// A counter decayed below this is treated as zero. +const DECAY_TO_ZERO: f64 = 0.01; + +/// Below this score a peer gets no IHAVE/IWANT gossip from this node. +pub const GOSSIP_THRESHOLD: f64 = -4000.0; +/// Below this score a peer is left out of this node's publishes. +pub const PUBLISH_THRESHOLD: f64 = -8000.0; +/// Below this score gossipsub ignores every RPC the peer sends, and +/// [`crate::swarm_adapter`] disconnects it. +pub const GRAYLIST_THRESHOLD: f64 = -16000.0; + +/// The active validator count the parameters are built for before the head +/// has named one: lighthouse's `minimum_validator_count`, one per slot. +pub const PLACEHOLDER_ACTIVE_VALIDATORS: u64 = SLOTS_PER_EPOCH; + +/// Head lag, in slots, beyond which [`MeshDeliveryGate`] switches mesh-delivery +/// scoring off. +/// +/// Two slots behind is already enough to `Ignore` most of what arrives, +/// since current aggregates and attestations vote for a block this node has +/// not imported. The margin is set by how long the delivery counters can go +/// without credit before a penalty starts. The aggregate counter falls from +/// its cap to the threshold in about 19 slots, and the other scored topics +/// take longer, so closing the gate within a few slots of losing the head +/// leaves nothing to penalize. Empty slots count towards the lag as well, +/// which is harmless: those votes are for a block this node has. +pub const MESH_DELIVERY_MAX_HEAD_LAG: u64 = 4; + +/// How many slots the head has to stay within [`MESH_DELIVERY_MAX_HEAD_LAG`] +/// before [`MeshDeliveryGate`] switches mesh-delivery scoring back on. +/// +/// While the gate was closed the counters kept counting, but a catch-up +/// credits nothing, so they come out of it near zero. Reopening the moment the +/// head catches up would score the whole mesh on that empty history. Every +/// scored topic refills to its threshold within a few slots of accepted +/// traffic, so one epoch leaves ample margin. +pub const MESH_DELIVERY_WARMUP_SLOTS: u64 = SLOTS_PER_EPOCH; + +/// The score thresholds, which do not depend on the chain. +pub fn thresholds() -> PeerScoreThresholds { + PeerScoreThresholds { + gossip_threshold: GOSSIP_THRESHOLD, + publish_threshold: PUBLISH_THRESHOLD, + graylist_threshold: GRAYLIST_THRESHOLD, + // Above what any peer can reach (the topic score cap), so peer + // exchange is never accepted, as on lighthouse. + accept_px_threshold: 100.0, + opportunistic_graft_threshold: 5.0, + } +} + +/// What the dynamic topic parameters are rebuilt from. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub struct ScoreInputs { + /// Active validators in the head's current epoch. + pub active_validators: u64, + /// The wall-clock slot. Only matters on a young chain: a topic's + /// mesh-delivery scoring stays off until the chain is older than that + /// topic's decay window, as on lighthouse. + pub current_slot: u64, + /// [`MeshDeliveryGate`]'s verdict. + pub mesh_deliveries_scored: bool, +} + +impl ScoreInputs { + /// What the parameters are built from before the first refresh: a + /// placeholder validator count, and no mesh-delivery scoring. + pub fn placeholder() -> Self { + Self { + active_validators: PLACEHOLDER_ACTIVE_VALIDATORS, + current_slot: 0, + mesh_deliveries_scored: false, + } + } +} + +/// A topic's mesh-delivery (P3) setup: lighthouse's `mesh_message_info`. +struct MeshDeliveries { + /// How many slots the delivery counter takes to decay to zero. + decay_slots: u64, + /// The counter's cap, as a multiple of its threshold. + cap_factor: f64, + /// How long a peer is in the mesh before its deficit counts. + activation: Duration, +} + +/// The chain-derived constants every parameter is computed from, fixed for +/// the life of the node. +#[derive(Debug, Clone)] +pub struct ScoreSettings { + slot: Duration, + epoch: Duration, + decay_interval: Duration, + mesh_n: usize, + attestation_subnet_weight: f64, + /// The most every scored topic's P1 and P2 can add up to, the unit the + /// invalid-message penalty is sized in. + max_positive_score: f64, +} + +impl ScoreSettings { + pub fn new(config: &Config, mesh_n: usize) -> Self { + let slot = Duration::from_millis(config.slot_duration_ms); + let attestation_subnet_weight = 1.0 / ATTESTATION_SUBNET_COUNT as f64; + let max_positive_score = (MAX_IN_MESH_SCORE + MAX_FIRST_MESSAGE_DELIVERIES_SCORE) + * (BEACON_BLOCK_WEIGHT + + BEACON_AGGREGATE_PROOF_WEIGHT + + attestation_subnet_weight * ATTESTATION_SUBNET_COUNT as f64 + + VOLUNTARY_EXIT_WEIGHT + + PROPOSER_SLASHING_WEIGHT + + ATTESTER_SLASHING_WEIGHT); + Self { + slot, + epoch: slot * SLOTS_PER_EPOCH as u32, + // One decay tick per slot, which is what every expected rate + // below is stated per. + decay_interval: slot.max(Duration::from_secs(1)), + mesh_n, + attestation_subnet_weight, + max_positive_score, + } + } + + /// The whole parameter set: the global terms, and every scored topic + /// under `fork_digest`. + pub fn peer_score_params( + &self, + fork_digest: ForkDigest, + inputs: ScoreInputs, + ) -> PeerScoreParams { + let behaviour_penalty_decay = self.decay(self.epoch * 10); + // The weight puts a peer that keeps earning ten behaviour penalties + // an epoch exactly at the gossip threshold once its counter settles. + let behaviour_penalty_threshold = 6.0; + let settled_excess = + decay_convergence(behaviour_penalty_decay, 10.0 / SLOTS_PER_EPOCH as f64) + - behaviour_penalty_threshold; + let topic_score_cap = self.max_positive_score * 0.5; + + let topics = self + .fixed_topic_params(fork_digest) + .into_iter() + .chain(self.dynamic_topic_params(fork_digest, inputs)) + .map(|(topic, params)| (topic.hash(), params)) + .collect(); + + PeerScoreParams { + topics, + topic_score_cap, + // Lighthouse's value. This node never sets an application score, + // so the term is always zero. + app_specific_weight: 1.0, + ip_colocation_factor_weight: -topic_score_cap, + // Up to eight peers per IP before the penalty starts. + ip_colocation_factor_threshold: 8.0, + ip_colocation_factor_whitelist: HashSet::new(), + behaviour_penalty_weight: GOSSIP_THRESHOLD / settled_excess.powi(2), + behaviour_penalty_threshold, + behaviour_penalty_decay, + decay_interval: self.decay_interval, + decay_to_zero: DECAY_TO_ZERO, + retain_score: self.epoch * 100, + // A peer whose send queue is full, or whose queued messages time + // out, loses ten a time, mostly gone by the next tick. + slow_peer_weight: -10.0, + slow_peer_threshold: 0.0, + slow_peer_decay: 0.1, + } + } + + /// Exits and slashings: rare enough that their rates do not depend on + /// the validator count, and never mesh-delivery scored. + fn fixed_topic_params(&self, fork_digest: ForkDigest) -> Vec<(IdentTopic, TopicScoreParams)> { + let per_slot = |per_epoch: f64| per_epoch / SLOTS_PER_EPOCH as f64; + [ + (VOLUNTARY_EXIT, VOLUNTARY_EXIT_WEIGHT, per_slot(4.0)), + ( + ATTESTER_SLASHING, + ATTESTER_SLASHING_WEIGHT, + per_slot(1.0 / 5.0), + ), + ( + PROPOSER_SLASHING, + PROPOSER_SLASHING_WEIGHT, + per_slot(1.0 / 5.0), + ), + ] + .into_iter() + .map(|(kind, weight, rate)| { + let topic = IdentTopic::new(topic_name(fork_digest, kind)); + let params = self.topic_params(weight, rate, self.epoch * 100, None, false); + (topic, params) + }) + .collect() + } + + /// Blocks, aggregates and every attestation subnet: the topics whose + /// expected traffic follows the validator count, and so the ones + /// [`refresh`] rebuilds. + pub fn dynamic_topic_params( + &self, + fork_digest: ForkDigest, + inputs: ScoreInputs, + ) -> Vec<(IdentTopic, TopicScoreParams)> { + // Zero would leave every rate zero and every first-delivery cap with + // it, which gossipsub rejects. + let active = inputs.active_validators.max(PLACEHOLDER_ACTIVE_VALIDATORS); + let (aggregators_per_slot, committees_per_slot) = expected_aggregators_per_slot(active); + // Whether each subnet sees more than one committee's burst an epoch, + // which shortens its decay windows. + let multiple_bursts_per_subnet_per_epoch = + committees_per_slot >= 2 * ATTESTATION_SUBNET_COUNT / SLOTS_PER_EPOCH; + // Lighthouse's young-chain rule and this node's sync gate: P3 needs + // both a chain older than the topic's decay window and a node that + // has kept up with it. + let scored = + |decay_slots: u64| inputs.mesh_deliveries_scored && inputs.current_slot > decay_slots; + + let block_mesh = MeshDeliveries { + decay_slots: SLOTS_PER_EPOCH * 5, + cap_factor: 3.0, + activation: self.epoch, + }; + let block = self.topic_params( + BEACON_BLOCK_WEIGHT, + 1.0, + self.epoch * 20, + Some(&block_mesh), + scored(block_mesh.decay_slots), + ); + + let aggregate_mesh = MeshDeliveries { + decay_slots: SLOTS_PER_EPOCH * 2, + cap_factor: 4.0, + activation: self.epoch, + }; + let aggregate = self.topic_params( + BEACON_AGGREGATE_PROOF_WEIGHT, + aggregators_per_slot, + self.epoch, + Some(&aggregate_mesh), + scored(aggregate_mesh.decay_slots), + ); + + let attestation_mesh = if multiple_bursts_per_subnet_per_epoch { + MeshDeliveries { + decay_slots: SLOTS_PER_EPOCH * 4, + cap_factor: 16.0, + activation: self.slot * (SLOTS_PER_EPOCH as u32 / 2 + 1), + } + } else { + MeshDeliveries { + decay_slots: SLOTS_PER_EPOCH * 16, + cap_factor: 16.0, + activation: self.epoch * 3, + } + }; + let attestation_fmd_decay_time = if multiple_bursts_per_subnet_per_epoch { + self.epoch + } else { + self.epoch * 4 + }; + let attestation = self.topic_params( + self.attestation_subnet_weight, + active as f64 / ATTESTATION_SUBNET_COUNT as f64 / SLOTS_PER_EPOCH as f64, + attestation_fmd_decay_time, + Some(&attestation_mesh), + scored(attestation_mesh.decay_slots), + ); + + let mut topics = vec![ + ( + IdentTopic::new(topic_name(fork_digest, BEACON_BLOCK)), + block, + ), + ( + IdentTopic::new(topic_name(fork_digest, BEACON_AGGREGATE_AND_PROOF)), + aggregate, + ), + ]; + topics.extend((0..ATTESTATION_SUBNET_COUNT).map(|subnet_id| { + let topic = IdentTopic::new(attestation_topic_name(fork_digest, subnet_id)); + (topic, attestation.clone()) + })); + topics + } + + /// One topic's parameters: lighthouse's `get_topic_params`. + /// + /// `rate` is the topic's expected messages per slot, one decay tick. + /// With `mesh_scored` false the mesh-delivery threshold and weight are + /// zero, so no deficit can build and no P3b is charged at prune, but the + /// cap stays where the threshold would put it: the counters keep + /// counting to their real ceiling, so switching P3 on later scores the + /// history that was actually delivered. + fn topic_params( + &self, + topic_weight: f64, + rate: f64, + first_message_decay_time: Duration, + mesh: Option<&MeshDeliveries>, + mesh_scored: bool, + ) -> TopicScoreParams { + let time_in_mesh_cap = 3600.0 / self.slot.as_secs_f64(); + let first_message_deliveries_decay = self.decay(first_message_decay_time); + // Two mesh shares of the traffic that arrives over one decay window. + let first_message_deliveries_cap = decay_convergence( + first_message_deliveries_decay, + 2.0 * rate / self.mesh_n as f64, + ); + + let mut params = TopicScoreParams { + topic_weight, + time_in_mesh_weight: MAX_IN_MESH_SCORE / time_in_mesh_cap, + time_in_mesh_quantum: self.slot, + time_in_mesh_cap, + first_message_deliveries_weight: MAX_FIRST_MESSAGE_DELIVERIES_SCORE + / first_message_deliveries_cap, + first_message_deliveries_decay, + first_message_deliveries_cap, + mesh_message_deliveries_weight: 0.0, + mesh_message_deliveries_decay: 0.0, + mesh_message_deliveries_cap: 0.0, + mesh_message_deliveries_threshold: 0.0, + mesh_message_deliveries_window: Duration::ZERO, + mesh_message_deliveries_activation: Duration::ZERO, + mesh_failure_penalty_weight: 0.0, + mesh_failure_penalty_decay: 0.0, + // One invalid message costs a whole `max_positive_score`, on any + // topic, whatever its weight. + invalid_message_deliveries_weight: -self.max_positive_score / topic_weight, + invalid_message_deliveries_decay: self.decay(self.epoch * 50), + }; + + if let Some(mesh) = mesh { + let decay = self.decay(self.slot * mesh.decay_slots as u32); + // Two percent of the expected traffic, over the decay window. + let threshold = decay_convergence(decay, rate / 50.0) * decay; + params.mesh_message_deliveries_decay = decay; + params.mesh_message_deliveries_cap = (mesh.cap_factor * threshold).max(2.0); + params.mesh_message_deliveries_activation = mesh.activation; + params.mesh_message_deliveries_window = MESH_MESSAGE_DELIVERIES_WINDOW; + params.mesh_failure_penalty_decay = decay; + params.mesh_failure_penalty_weight = -topic_weight; + if mesh_scored { + params.mesh_message_deliveries_threshold = threshold; + params.mesh_message_deliveries_weight = -topic_weight; + } + } + + params + } + + /// The per-tick factor that takes a counter from one to + /// [`DECAY_TO_ZERO`] over `decay_time`. + fn decay(&self, decay_time: Duration) -> f64 { + let ticks = decay_time.as_secs_f64() / self.decay_interval.as_secs_f64(); + DECAY_TO_ZERO.powf(1.0 / ticks) + } +} + +/// Where a counter that gains `rate` a tick and decays by `decay` a tick +/// settles. +fn decay_convergence(decay: f64, rate: f64) -> f64 { + rate / (1.0 - decay) +} + +/// Aggregates expected per slot with `active` validators, and the committee +/// count per slot it was derived from. +/// +/// Every committee selects `len / max(1, len / TARGET_AGGREGATORS_PER_COMMITTEE)` +/// aggregators in expectation (phase0's `is_aggregator`), and an epoch's +/// committees differ in size by at most one member. +fn expected_aggregators_per_slot(active: u64) -> (f64, u64) { + let committees_per_slot = committee_count_per_slot(active); + let committees = committees_per_slot * SLOTS_PER_EPOCH; + let smaller = active / committees; + let larger_count = active - smaller * committees; + let smaller_modulo = (smaller / TARGET_AGGREGATORS_PER_COMMITTEE).max(1); + let larger_modulo = ((smaller + 1) / TARGET_AGGREGATORS_PER_COMMITTEE).max(1); + let per_epoch = ((committees - larger_count) * smaller) as f64 / smaller_modulo as f64 + + (larger_count * (smaller + 1)) as f64 / larger_modulo as f64; + (per_epoch / SLOTS_PER_EPOCH as f64, committees_per_slot) +} + +/// Whether mesh-delivery scoring (P3) is safe to switch on: the head has +/// kept up with the wall clock long enough that the delivery counters +/// reflect what mesh peers sent, not what this node was able to accept. +/// +/// A delivery is credited only once it is accepted. While this node catches +/// up it `Ignore`s every aggregate and attestation voting for a block it has +/// not imported, so no mesh peer is credited for anything. Aggregate traffic +/// is heavy enough that a mesh peer credited with nothing scores below the +/// graylist. Without the gate, a restart would graylist the whole aggregate +/// mesh because of this node's lag, not anything the peers did. +/// +/// Lighthouse avoids the same trap differently: it joins these topics only +/// once it has synced, so it has no mesh while it catches up. This node +/// subscribes at startup. The gate also covers the case lighthouse leaves +/// open, a node that falls behind after it has synced. +#[derive(Debug, Default)] +pub struct MeshDeliveryGate { + /// The wall-clock slot since which the head has stayed within + /// [`MESH_DELIVERY_MAX_HEAD_LAG`], or `None` while it lags further. + keeping_up_since: Option, +} + +impl MeshDeliveryGate { + /// Record where the head is, and return whether mesh deliveries should + /// be scored. + pub fn update(&mut self, wall_slot: u64, head_slot: u64) -> bool { + if wall_slot.saturating_sub(head_slot) > MESH_DELIVERY_MAX_HEAD_LAG { + self.keeping_up_since = None; + return false; + } + let since = *self.keeping_up_since.get_or_insert(wall_slot); + // Saturating: the wall clock can step backwards under an NTP correction. + wall_slot.saturating_sub(since) >= MESH_DELIVERY_WARMUP_SLOTS + } +} + +/// The scoring state [`P2PServer`] keeps on a beacon wire. +pub(crate) struct PeerScoring { + settings: ScoreSettings, + gate: MeshDeliveryGate, + /// The last active validator count the head named, kept for the ticks + /// where its shuffling is not built yet. + active_validators: u64, + /// The gate's last verdict, so a change is logged once. + mesh_deliveries_scored: bool, +} + +impl PeerScoring { + pub(crate) fn new(config: &Config) -> Self { + Self { + settings: ScoreSettings::new(config, crate::MESH_N), + gate: MeshDeliveryGate::default(), + active_validators: PLACEHOLDER_ACTIVE_VALIDATORS, + mesh_deliveries_scored: false, + } + } +} + +/// Rebuild the dynamic topic parameters from the head and hand them to the +/// swarm. Beacon only; a no-op on lean, which runs without peer scoring. +pub(crate) fn refresh(server: &mut P2PServer) { + let (Some(wire), Some(scoring)) = (server.wire.beacon(), server.peer_scoring.as_mut()) else { + return; + }; + let wall_slot = beacon_wall_slot(wire); + let head_slot = server.store.beacon_head().map_or(0, |(slot, _)| slot); + let mesh_deliveries_scored = scoring.gate.update(wall_slot, head_slot); + if mesh_deliveries_scored != scoring.mesh_deliveries_scored { + info!( + wall_slot, + head_slot, mesh_deliveries_scored, "Gossipsub mesh-delivery scoring switched" + ); + scoring.mesh_deliveries_scored = mesh_deliveries_scored; + } + metrics::set_gossipsub_mesh_delivery_scoring(mesh_deliveries_scored); + + if let Some(committees) = server.store.committee_cache().head_current_committees() { + scoring.active_validators = committees.active_validator_count(); + } + + let inputs = ScoreInputs { + active_validators: scoring.active_validators, + current_slot: wall_slot, + mesh_deliveries_scored, + }; + let topics = scoring + .settings + .dynamic_topic_params(wire.fork_digest, inputs); + server.swarm_handle.set_topic_score_params(topics); +} + +/// Which score band a peer falls in, for `lean_gossipsub_peers_by_score`. +/// +/// Mutually exclusive, so the bands add up to the peers gossipsub knows. +pub(crate) fn score_band(score: f64) -> &'static str { + if score >= 0.0 { + "non_negative" + } else if score >= GOSSIP_THRESHOLD { + "negative" + } else if score >= PUBLISH_THRESHOLD { + "below_gossip" + } else if score >= GRAYLIST_THRESHOLD { + "below_publish" + } else { + "below_graylist" + } +} + +/// Every band [`score_band`] can return, so a band that empties is published +/// as zero rather than keeping its last count. +pub(crate) const SCORE_BANDS: [&str; 5] = [ + "non_negative", + "negative", + "below_gossip", + "below_publish", + "below_graylist", +]; + +#[cfg(test)] +mod tests { + use super::*; + + const DIGEST: ForkDigest = [0x8c, 0x9f, 0x62, 0xfe]; + + /// Lighthouse's default network load gives `mesh_n` 5, the value every + /// number below was computed at on lighthouse v8.2.2's formulas. + const LIGHTHOUSE_MESH_N: usize = 5; + + fn settings(mesh_n: usize) -> ScoreSettings { + ScoreSettings::new(&Config::mainnet(), mesh_n) + } + + fn inputs(active_validators: u64, mesh_deliveries_scored: bool) -> ScoreInputs { + ScoreInputs { + active_validators, + // Far past every topic's decay window, like mainnet. + current_slot: 12_000_000, + mesh_deliveries_scored, + } + } + + fn assert_close(actual: f64, expected: f64, what: &str) { + let tolerance = expected.abs() * 1e-3; + assert!( + (actual - expected).abs() <= tolerance, + "{what}: {actual} is not {expected}" + ); + } + + fn topic<'a>(topics: &'a [(IdentTopic, TopicScoreParams)], kind: &str) -> &'a TopicScoreParams { + let name = topic_name(DIGEST, kind); + &topics + .iter() + .find(|(topic, _)| topic.to_string() == name) + .unwrap_or_else(|| panic!("{kind} has parameters")) + .1 + } + + /// The globals match lighthouse's, which is what makes the thresholds + /// mean what they mean on every other client. + #[test] + fn the_global_terms_are_lighthouses() { + let params = settings(LIGHTHOUSE_MESH_N).peer_score_params(DIGEST, inputs(1_000_000, true)); + assert_close(params.topic_score_cap, 53.75, "topic_score_cap"); + assert_close( + params.ip_colocation_factor_weight, + -53.75, + "ip colocation weight", + ); + assert_close(params.behaviour_penalty_decay, 0.985_712, "behaviour decay"); + assert_close(params.behaviour_penalty_weight, -15.879, "behaviour weight"); + assert_eq!(params.decay_interval, Duration::from_secs(12)); + assert_eq!(params.retain_score, Duration::from_secs(38_400)); + } + + /// The per-topic numbers at 1,000,000 active validators match + /// lighthouse's, computed from its own formulas. + #[test] + fn the_topic_terms_are_lighthouses() { + let topics = + settings(LIGHTHOUSE_MESH_N).dynamic_topic_params(DIGEST, inputs(1_000_000, true)); + + let block = topic(&topics, BEACON_BLOCK); + assert_close(block.first_message_deliveries_cap, 55.79, "block fmd cap"); + assert_close( + block.mesh_message_deliveries_threshold, + 0.685, + "block mmd threshold", + ); + assert_close( + block.invalid_message_deliveries_weight, + -215.0, + "block imd weight", + ); + + let aggregate = topic(&topics, BEACON_AGGREGATE_AND_PROOF); + assert_close( + aggregate.first_message_deliveries_cap, + 3108.6, + "aggregate fmd cap", + ); + assert_close( + aggregate.mesh_message_deliveries_threshold, + 279.24, + "aggregate threshold", + ); + assert_close( + aggregate.mesh_message_deliveries_cap, + 1116.95, + "aggregate mmd cap", + ); + + let attestation = topic(&topics, "beacon_attestation_0"); + assert_close( + attestation.first_message_deliveries_cap, + 1457.2, + "attestation fmd cap", + ); + assert_close( + attestation.mesh_message_deliveries_threshold, + 266.58, + "attestation threshold", + ); + assert_eq!( + attestation.mesh_message_deliveries_activation, + Duration::from_secs(204) + ); + assert_close( + attestation.invalid_message_deliveries_weight, + -6880.0, + "attestation imd", + ); + } + + /// Every combination gossipsub can be handed passes its own validation, + /// including the placeholder the swarm starts on and a registry of zero. + #[test] + fn every_parameter_set_is_valid() { + let settings = settings(crate::MESH_N); + for active in [ + 0, + 1, + PLACEHOLDER_ACTIVE_VALIDATORS, + 1_000, + 100_000, + 1_000_000, + 2_500_000, + ] { + for current_slot in [0, 100, 600, 12_000_000] { + for mesh_deliveries_scored in [false, true] { + let inputs = ScoreInputs { + active_validators: active, + current_slot, + mesh_deliveries_scored, + }; + settings + .peer_score_params(DIGEST, inputs) + .validate() + .unwrap_or_else(|err| panic!("{inputs:?}: {err}")); + } + } + } + thresholds().validate().expect("valid thresholds"); + } + + /// With the gate closed no deficit can build, but the cap stays at its + /// real value, so the counters are full when the gate opens again. + #[test] + fn a_closed_gate_zeroes_the_deficit_but_keeps_the_cap() { + let settings = settings(crate::MESH_N); + let open = settings.dynamic_topic_params(DIGEST, inputs(1_000_000, true)); + let closed = settings.dynamic_topic_params(DIGEST, inputs(1_000_000, false)); + for kind in [ + BEACON_BLOCK, + BEACON_AGGREGATE_AND_PROOF, + "beacon_attestation_7", + ] { + let (open, closed) = (topic(&open, kind), topic(&closed, kind)); + assert!( + open.mesh_message_deliveries_weight < 0.0, + "{kind} scored when open" + ); + assert_eq!(closed.mesh_message_deliveries_weight, 0.0, "{kind}"); + assert_eq!(closed.mesh_message_deliveries_threshold, 0.0, "{kind}"); + assert_eq!( + closed.mesh_message_deliveries_cap, open.mesh_message_deliveries_cap, + "{kind} keeps counting to its real cap" + ); + } + } + + /// Lighthouse's young-chain rule survives the port: an open gate still + /// leaves a topic unscored until the chain outlives its decay window. + #[test] + fn a_young_chain_leaves_mesh_deliveries_unscored() { + let young = ScoreInputs { + current_slot: SLOTS_PER_EPOCH * 2, + ..inputs(1_000_000, true) + }; + let topics = settings(crate::MESH_N).dynamic_topic_params(DIGEST, young); + assert_eq!( + topic(&topics, BEACON_BLOCK).mesh_message_deliveries_weight, + 0.0 + ); + assert_eq!( + topic(&topics, BEACON_AGGREGATE_AND_PROOF).mesh_message_deliveries_weight, + 0.0, + "the aggregate window is exactly two epochs, and the slot is not past it" + ); + } + + /// Every attestation subnet has parameters, not only the ones this node + /// backbones, since an aggregator joins others at runtime. + #[test] + fn every_attestation_subnet_is_scored() { + let params = settings(crate::MESH_N).peer_score_params(DIGEST, inputs(1_000_000, true)); + for subnet_id in 0..ATTESTATION_SUBNET_COUNT { + let topic = IdentTopic::new(attestation_topic_name(DIGEST, subnet_id)); + assert!( + params.topics.contains_key(&topic.hash()), + "subnet {subnet_id}" + ); + } + // Block, aggregate, three fixed topics and the 64 subnets. + assert_eq!( + params.topics.len(), + 2 + 3 + ATTESTATION_SUBNET_COUNT as usize + ); + } + + /// Mainnet's committees saturate at `MAX_COMMITTEES_PER_SLOT`, so about + /// sixteen aggregators each: lighthouse's figure of 1041.67 per slot. + #[test] + fn a_million_validators_expect_about_a_thousand_aggregators_a_slot() { + let (per_slot, committees_per_slot) = expected_aggregators_per_slot(1_000_000); + assert_eq!(committees_per_slot, 64); + assert_close(per_slot, 1041.67, "aggregators per slot"); + } + + #[test] + fn the_gate_opens_only_after_a_warmup_within_the_lag() { + let mut gate = MeshDeliveryGate::default(); + // Catching up: far behind. + assert!(!gate.update(1_000, 900)); + // Caught up, but the counters have not refilled yet. + assert!(!gate.update(1_001, 1_000)); + assert!(!gate.update( + 1_000 + MESH_DELIVERY_WARMUP_SLOTS, + 1_000 + MESH_DELIVERY_WARMUP_SLOTS + )); + assert!(gate.update( + 1_001 + MESH_DELIVERY_WARMUP_SLOTS, + 1_000 + MESH_DELIVERY_WARMUP_SLOTS + )); + // A lag within the bound, empty slots say, keeps it open. + let wall = 1_001 + MESH_DELIVERY_WARMUP_SLOTS + MESH_DELIVERY_MAX_HEAD_LAG; + assert!(gate.update(wall, 1_001 + MESH_DELIVERY_WARMUP_SLOTS)); + // Falling further behind closes it at once, and restarts the warmup. + assert!(!gate.update(wall + 1, 1_001 + MESH_DELIVERY_WARMUP_SLOTS)); + assert!(!gate.update(wall + 2, wall + 2)); + } + + /// The built beacon swarm scores peers and already has parameters for a + /// subnet it does not subscribe to; a lean swarm scores nothing. + #[tokio::test] + async fn only_the_beacon_swarm_scores_peers() { + use ethlambda_types::beacon::fork::ForkName; + use ethlambda_types::beacon::primitives::Root; + + use crate::beacon::swarm::BeaconWireConfig; + use crate::{LeanWireConfig, SwarmConfig, WireConfig, build_swarm}; + + let swarm_config = |wire| SwarmConfig { + node_key: vec![3u8; 32], + bootnodes: Vec::new(), + listening_socket: "127.0.0.1:0".parse().expect("valid socket"), + target_peers: crate::discovery::DEFAULT_DISCOVERY_TARGET_PEERS, + wire, + }; + + let beacon = build_swarm(swarm_config(WireConfig::Beacon(Box::new( + BeaconWireConfig { + fork_digest: DIGEST, + fork: ForkName::Fulu, + config: Config::mainnet(), + genesis_time: 1_606_824_023, + genesis_validators_root: Root::ZERO, + custody_columns: vec![3], + attestation_subnets: vec![5], + }, + )))) + .expect("beacon swarm builds"); + let gossipsub = &beacon.swarm.behaviour().gossipsub; + assert!(gossipsub.peer_score(&beacon.local_peer_id).is_some()); + let unsubscribed_subnet = IdentTopic::new(attestation_topic_name(DIGEST, 40)); + assert!(gossipsub.get_topic_params(&unsubscribed_subnet).is_some()); + let column = IdentTopic::new(crate::beacon::topics::data_column_topic_name(DIGEST, 3)); + assert!( + gossipsub.get_topic_params(&column).is_none(), + "columns are not scored" + ); + + let lean = build_swarm(swarm_config(WireConfig::Lean(LeanWireConfig { + validator_ids: Vec::new(), + attestation_committee_count: 1, + subscription_subnets: HashSet::new(), + milliseconds_per_slot: 4_000, + }))) + .expect("lean swarm builds"); + assert!( + lean.swarm + .behaviour() + .gossipsub + .peer_score(&lean.local_peer_id) + .is_none() + ); + } + + #[test] + fn the_score_bands_are_split_at_the_thresholds() { + assert_eq!(score_band(10.0), "non_negative"); + assert_eq!(score_band(0.0), "non_negative"); + assert_eq!(score_band(-1.0), "negative"); + assert_eq!(score_band(GOSSIP_THRESHOLD), "negative"); + assert_eq!(score_band(GOSSIP_THRESHOLD - 1.0), "below_gossip"); + assert_eq!(score_band(PUBLISH_THRESHOLD - 1.0), "below_publish"); + assert_eq!(score_band(GRAYLIST_THRESHOLD - 1.0), "below_graylist"); + for score in [10.0, -1.0, -5000.0, -9000.0, -20000.0] { + assert!(SCORE_BANDS.contains(&score_band(score))); + } + } +} diff --git a/crates/net/p2p/src/beacon/swarm.rs b/crates/net/p2p/src/beacon/swarm.rs index 9c264ef0..67986fc6 100644 --- a/crates/net/p2p/src/beacon/swarm.rs +++ b/crates/net/p2p/src/beacon/swarm.rs @@ -1,10 +1,11 @@ //! The beacon half of the swarm configuration. //! -//! Five things differ from lean at the swarm level: the topic set, the protocol -//! set, the `seen_ttl`, the identify protocol version and the connection limits. -//! [`crate::build_swarm`] resolves all five from the variant it is handed, and -//! every one of them that is a beacon *value* rather than a lean one lives here, -//! so tuning the mainnet numbers never means editing the crate root. +//! Six things differ from lean at the swarm level: the topic set, the protocol +//! set, the `seen_ttl`, the identify protocol version, the connection limits and +//! peer scoring. [`crate::build_swarm`] resolves all six from the variant it is +//! handed, and every one of them that is a beacon *value* rather than a lean one +//! lives here or, for scoring, in [`crate::beacon::scoring`], so tuning the +//! mainnet numbers never means editing the crate root. use std::time::Duration; diff --git a/crates/net/p2p/src/gossipsub/handler.rs b/crates/net/p2p/src/gossipsub/handler.rs index a1797a81..f66dc190 100644 --- a/crates/net/p2p/src/gossipsub/handler.rs +++ b/crates/net/p2p/src/gossipsub/handler.rs @@ -659,7 +659,7 @@ pub async fn publish_beacon_block(server: &mut P2PServer, block: SignedBeaconBlo } /// The beacon wall-clock slot, from the wire's genesis and slot duration. -fn beacon_wall_slot(wire: &BeaconWire) -> u64 { +pub(crate) fn beacon_wall_slot(wire: &BeaconWire) -> u64 { let genesis_ms = wire.genesis_time.saturating_mul(1000); unix_now_ms().saturating_sub(genesis_ms) / wire.config.slot_duration_ms.max(1) } diff --git a/crates/net/p2p/src/gossipsub/mod.rs b/crates/net/p2p/src/gossipsub/mod.rs index 608584a2..ce5cc560 100644 --- a/crates/net/p2p/src/gossipsub/mod.rs +++ b/crates/net/p2p/src/gossipsub/mod.rs @@ -3,6 +3,7 @@ mod handler; mod messages; pub use encoding::decompress_message; +pub(crate) use handler::beacon_wall_slot; pub use handler::{ handle_gossip_message, join_aggregator_subnets, leave_expired_aggregator_subnets, prune_attestation_pool, publish_aggregated_attestation, publish_attestation, diff --git a/crates/net/p2p/src/lib.rs b/crates/net/p2p/src/lib.rs index 081f1cd6..fe22edea 100644 --- a/crates/net/p2p/src/lib.rs +++ b/crates/net/p2p/src/lib.rs @@ -573,9 +573,9 @@ pub struct SwarmConfig { /// Ignored on lean, which runs [`unlimited_connections`]. pub target_peers: usize, /// Which network's wire to build. Decides the gossip topics, the req/resp - /// protocol set, the gossipsub `seen_ttl`, the identify protocol version and - /// the connection limits, and so decides the [`Wire`] the built swarm - /// carries. + /// protocol set, the gossipsub `seen_ttl`, the identify protocol version, + /// the connection limits and whether peers are scored, and so decides the + /// [`Wire`] the built swarm carries. pub wire: WireConfig, } @@ -725,6 +725,11 @@ pub enum SwarmBuildError { Subscription(#[from] libp2p::gossipsub::SubscriptionError), } +/// Gossipsub's target mesh size, the spec's `D`. Named because beacon peer +/// scoring sizes its first-delivery caps by it; see +/// [`beacon::scoring::ScoreSettings`]. +pub(crate) const MESH_N: usize = 8; + /// The gossipsub parameters both wires share. /// /// `mesh_n` 8, low 6, high 12, the 700ms heartbeat, and the 6/3 history already @@ -743,7 +748,7 @@ pub(crate) fn gossipsub_config( let mut builder = libp2p::gossipsub::ConfigBuilder::default(); builder // d - .mesh_n(8) + .mesh_n(MESH_N) // d_low .mesh_n_low(6) // d_high @@ -771,9 +776,10 @@ pub(crate) fn gossipsub_config( /// Build and configure the libp2p swarm, dial bootnodes, subscribe to topics. /// /// One builder for both networks. Four things differ at the behaviour level and -/// are resolved in the first match below; the topic set differs and is resolved -/// in the last one. Everything in between, the transport, the QUIC and TCP -/// listeners and the static bootnode dialing, is the same on either wire. +/// are resolved in the first match below, peer scoring (beacon only) just after +/// it; the topic set differs and is resolved in the last one. Everything in +/// between, the transport, the QUIC and TCP listeners and the static bootnode +/// dialing, is the same on either wire. pub fn build_swarm(config: SwarmConfig) -> Result { let SwarmConfig { node_key, @@ -807,11 +813,27 @@ pub fn build_swarm(config: SwarmConfig) -> Result { }; let validate_messages = matches!(wire, WireConfig::Beacon(_)); - let gossipsub = libp2p::gossipsub::Behaviour::new( + let mut gossipsub = libp2p::gossipsub::Behaviour::new( MessageAuthenticity::Anonymous, gossipsub_config(seen_ttl, validate_messages), ) .expect("failed to initiate behaviour"); + // Beacon only: the parameters are sized from the beacon chain's topics + // and message rates, and lean has no set of its own. Without a verdict + // on each message (lean's are auto-accepted) the invalid-message term + // would have nothing to count either. Built for a placeholder validator + // count; `P2PServer`'s first scoring refresh, at startup, replaces the + // topic parameters with the head's. + if let WireConfig::Beacon(beacon) = &wire { + let settings = beacon::scoring::ScoreSettings::new(&beacon.config, MESH_N); + let params = settings.peer_score_params( + beacon.fork_digest, + beacon::scoring::ScoreInputs::placeholder(), + ); + gossipsub + .with_peer_score(params, beacon::scoring::thresholds()) + .expect("valid peer score parameters"); + } let req_resp = ReqResp::new(codec, &wire); @@ -832,8 +854,6 @@ pub fn build_swarm(config: SwarmConfig) -> Result { connection_limits, }; - // TODO: set peer scoring params - let mut swarm = libp2p::SwarmBuilder::with_existing_identity(identity) .with_tokio() .with_tcp( @@ -1051,6 +1071,10 @@ impl P2P { .wire .beacon() .map_or(0, |beacon| beacon.attestation_subnets.len()); + let peer_scoring = built + .wire + .beacon() + .map(|beacon| beacon::scoring::PeerScoring::new(&beacon.config)); let server = P2PServer { swarm_handle, @@ -1085,8 +1109,19 @@ impl P2P { )), attestation_pool, aggregator_subnets: HashMap::new(), + peer_scoring, }; + let scores_peers = server.peer_scoring.is_some(); let handle = server.start(); + // At once rather than a slot in: the swarm starts on placeholder + // topic scores, and the head's are already known. + if scores_peers { + send_after( + Duration::ZERO, + handle.context(), + p2p_protocol::RefreshPeerScoring, + ); + } send_after( AGGREGATOR_SUBNET_SWEEP_INTERVAL, handle.context(), @@ -1197,6 +1232,10 @@ pub struct P2PServer { /// advertised in `attnets`, and left once the slot has passed. The /// backbone subnets are separate and never left. pub(crate) aggregator_subnets: HashMap, + + /// Gossipsub scoring's inputs and gate. `None` on lean, which runs + /// without peer scoring. See [`beacon::scoring`]. + pub(crate) peer_scoring: Option, } impl P2PServer { @@ -1282,6 +1321,8 @@ pub(crate) trait P2PProtocol: Send + Sync { fn leave_expired_aggregator_subnets(&self) -> Result<(), ActorError>; #[allow(dead_code)] // invoked via send_after, not called directly fn retry_beacon_range_batch(&self) -> Result<(), ActorError>; + #[allow(dead_code)] // invoked via send_after, not called directly + fn refresh_peer_scoring(&self) -> Result<(), ActorError>; } #[actor(protocol = P2PProtocol)] @@ -1392,6 +1433,21 @@ impl P2PServer { send_after(interval, ctx.clone(), p2p_protocol::DiscoverPeers); } + /// Rebuild the gossipsub topic scores from the head, once a slot. Beacon + /// only: never scheduled on lean. See [`beacon::scoring::refresh`]. + #[send_handler] + async fn handle_refresh_peer_scoring( + &mut self, + _msg: p2p_protocol::RefreshPeerScoring, + ctx: &Context, + ) { + if let Some(wire) = self.wire.beacon() { + let slot = Duration::from_millis(wire.config.slot_duration_ms); + send_after(slot, ctx.clone(), p2p_protocol::RefreshPeerScoring); + } + beacon::scoring::refresh(self); + } + /// The deadline of a range batch held back for custody. Scheduled once, /// when the batch starts waiting; see [`RANGE_BATCH_CUSTODY_WAIT`]. #[send_handler] @@ -2518,6 +2574,7 @@ pub(crate) mod test_support { )), attestation_pool: Default::default(), aggregator_subnets: HashMap::new(), + peer_scoring: None, } } diff --git a/crates/net/p2p/src/metrics.rs b/crates/net/p2p/src/metrics.rs index 3f1c060d..081d817c 100644 --- a/crates/net/p2p/src/metrics.rs +++ b/crates/net/p2p/src/metrics.rs @@ -582,6 +582,51 @@ pub fn set_swarm_established_connections(inbound: u32, outbound: u32) { .set(i64::from(outbound)); } +/// Set how many gossipsub peers fall in each score band. `bands` names every +/// band, empty ones as zero, so a band that empties does not keep its last +/// count. See `beacon::scoring::score_band`. +pub fn set_gossipsub_peers_by_score(bands: &HashMap<&'static str, i64>) { + static LEAN_GOSSIPSUB_PEERS_BY_SCORE: LazyLock = LazyLock::new(|| { + register_int_gauge_vec!( + "lean_gossipsub_peers_by_score", + "Gossipsub peers by score band: non_negative, negative, below_gossip, below_publish, below_graylist", + &["band"] + ) + .unwrap() + }); + for (band, count) in bands { + LEAN_GOSSIPSUB_PEERS_BY_SCORE + .with_label_values(&[band]) + .set(*count); + } +} + +/// Count a peer disconnected for a gossipsub score below the graylist. +pub fn inc_gossipsub_score_disconnects() { + static LEAN_GOSSIPSUB_SCORE_DISCONNECTS_TOTAL: LazyLock = LazyLock::new(|| { + register_int_counter!( + "lean_gossipsub_score_disconnects_total", + "Peers disconnected for a gossipsub score below the graylist threshold" + ) + .unwrap() + }); + LEAN_GOSSIPSUB_SCORE_DISCONNECTS_TOTAL.inc(); +} + +/// Whether gossipsub's mesh-delivery term (P3) is being scored: 0 while the +/// head lags or is still warming up after catching up. See +/// `beacon::scoring::MeshDeliveryGate`. +pub fn set_gossipsub_mesh_delivery_scoring(scored: bool) { + static LEAN_GOSSIPSUB_MESH_DELIVERY_SCORING: LazyLock = LazyLock::new(|| { + register_int_gauge!( + "lean_gossipsub_mesh_delivery_scoring", + "1 while gossipsub scores mesh message deliveries, 0 while the head lags or is warming up" + ) + .unwrap() + }); + LEAN_GOSSIPSUB_MESH_DELIVERY_SCORING.set(i64::from(scored)); +} + /// Set how many connected peers are known to custody `column`. /// /// A peer counts only once it has answered `metadata/3` or arrived with a diff --git a/crates/net/p2p/src/req_resp/handlers.rs b/crates/net/p2p/src/req_resp/handlers.rs index d324af81..0fbd6a73 100644 --- a/crates/net/p2p/src/req_resp/handlers.rs +++ b/crates/net/p2p/src/req_resp/handlers.rs @@ -2540,6 +2540,7 @@ mod tests { )), attestation_pool: Default::default(), aggregator_subnets: HashMap::new(), + peer_scoring: None, } } diff --git a/crates/net/p2p/src/swarm_adapter.rs b/crates/net/p2p/src/swarm_adapter.rs index f060c3fb..b4290991 100644 --- a/crates/net/p2p/src/swarm_adapter.rs +++ b/crates/net/p2p/src/swarm_adapter.rs @@ -4,18 +4,20 @@ use std::time::Duration; use libp2p::{ PeerId, futures::StreamExt, + gossipsub::{IdentTopic, TopicScoreParams}, request_response, swarm::{SwarmEvent, dial_opts::DialOpts}, }; use tokio::{sync::mpsc, time::MissedTickBehavior}; -use tracing::{debug, error, warn}; +use tracing::{debug, error, info, warn}; use crate::{ - Behaviour, BehaviourEvent, ReqRespProtocol, ReqRespRequestId, metrics, req_resp::Request, - req_resp::Response, + Behaviour, BehaviourEvent, ReqRespProtocol, ReqRespRequestId, beacon::scoring, metrics, + req_resp::Request, req_resp::Response, }; -/// Interval between gossipsub mesh peer metric refreshes. +/// Interval between gossipsub mesh peer metric refreshes, and between the +/// checks that disconnect peers scored below the graylist. const MESH_METRIC_REFRESH_INTERVAL: Duration = Duration::from_secs(10); pub enum SwarmCommand { @@ -49,6 +51,9 @@ pub enum SwarmCommand { channel: request_response::ResponseChannel, response: Response, }, + /// Replace these topics' score parameters (beacon only; see + /// [`scoring::refresh`]). + SetTopicScoreParams(Vec<(IdentTopic, TopicScoreParams)>), /// A verdict for a gossip message gossipsub is holding (beacon only). ReportValidation { message_id: libp2p::gossipsub::MessageId, @@ -204,6 +209,13 @@ impl SwarmHandle { .inspect_err(|_| debug!("Swarm adapter closed, cannot send response")); } + pub fn set_topic_score_params(&self, topics: Vec<(IdentTopic, TopicScoreParams)>) { + let _ = self + .cmd_tx + .send(SwarmCommand::SetTopicScoreParams(topics)) + .inspect_err(|_| debug!("Swarm adapter closed, cannot set topic score parameters")); + } + pub fn report_validation( &self, message_id: libp2p::gossipsub::MessageId, @@ -277,12 +289,50 @@ async fn swarm_loop( counters.num_established_incoming(), counters.num_established_outgoing(), ); + disconnect_graylisted_peers(&mut swarm); } } } error!("Swarm adapter loop exited — P2P networking is no longer functional"); } +/// Disconnect every peer gossipsub scores below the graylist, and publish how +/// the rest are spread across the score bands. +/// +/// Gossipsub on its own only ignores a graylisted peer's RPCs, so the peer +/// would keep a connection slot it is no use in. There is no ban list, so the +/// peer may reconnect, but its score is retained for `retain_score` after it +/// leaves: it comes back graylisted and is dropped again on the next pass. +/// +/// Does nothing on lean, which runs without scoring. +fn disconnect_graylisted_peers(swarm: &mut libp2p::Swarm) { + let gossipsub = &swarm.behaviour().gossipsub; + // `peer_score` is `Some` for any peer id, this node's own included, + // exactly when scoring is on. + if gossipsub.peer_score(swarm.local_peer_id()).is_none() { + return; + } + let mut bands: HashMap<&'static str, i64> = + scoring::SCORE_BANDS.iter().map(|&band| (band, 0)).collect(); + let mut graylisted = Vec::new(); + for (peer_id, _topics) in gossipsub.all_peers() { + let Some(score) = gossipsub.peer_score(peer_id) else { + continue; + }; + *bands.entry(scoring::score_band(score)).or_default() += 1; + if score < scoring::GRAYLIST_THRESHOLD { + graylisted.push((*peer_id, score)); + } + } + metrics::set_gossipsub_peers_by_score(&bands); + for (peer_id, score) in graylisted { + if swarm.disconnect_peer_id(peer_id).is_ok() { + metrics::inc_gossipsub_score_disconnects(); + info!(%peer_id, score, "Disconnecting peer scored below the gossipsub graylist"); + } + } +} + fn execute_command(swarm: &mut libp2p::Swarm, cmd: SwarmCommand) { match cmd { SwarmCommand::Publish { topic, data } => { @@ -386,6 +436,14 @@ fn execute_command(swarm: &mut libp2p::Swarm, cmd: SwarmCommand) { .send_response(channel, response) .inspect_err(|response| debug!(%response, "Swarm adapter: send_response failed")); } + SwarmCommand::SetTopicScoreParams(topics) => { + let gossipsub = &mut swarm.behaviour_mut().gossipsub; + for (topic, params) in topics { + let _ = gossipsub + .set_topic_params(topic, params) + .inspect_err(|err| debug!(%err, "Swarm adapter: topic score update failed")); + } + } SwarmCommand::ReportValidation { message_id, propagation_source, diff --git a/crates/storage/src/committee_cache.rs b/crates/storage/src/committee_cache.rs index 763ee235..35c079d9 100644 --- a/crates/storage/src/committee_cache.rs +++ b/crates/storage/src/committee_cache.rs @@ -83,6 +83,10 @@ const COMMITTEE_CACHE_CAPACITY: usize = 8; /// pins: its previous, current, and next epochs'. const HEAD_SHUFFLINGS: usize = 3; +/// Where the head's current epoch sits among [`CommitteeCache::pin_head`]'s +/// keys, which name its previous, current and next epochs in that order. +const HEAD_CURRENT_EPOCH: usize = 1; + const _: () = assert!( COMMITTEE_CACHE_CAPACITY > HEAD_SHUFFLINGS, "the cache must hold at least one shuffling the head does not pin" @@ -267,6 +271,22 @@ impl CommitteeCache { pub fn pin_head(&self, head_root: Root, keys: [Option; HEAD_SHUFFLINGS]) { self.state.lock().unwrap().head = Some((head_root, keys)); } + + /// The pinned head's current-epoch committees, if that shuffling has been + /// built. + /// + /// Never builds one: a caller that only wants a figure off the head (the + /// active validator count gossipsub scoring sizes its expected message + /// rates by) has no state to build from, and waiting for the chain actor + /// to need that shuffling is cheaper than deriving it here. `None` until + /// a head has been pinned and its current epoch's shuffling filled in. + pub fn head_current_committees(&self) -> Option> { + let state = self.state.lock().unwrap(); + let (_, keys) = state.head.as_ref()?; + let current = keys[HEAD_CURRENT_EPOCH]?; + let (_, slot) = state.entries.iter().find(|(key, _)| *key == current)?; + slot.get().cloned() + } } #[cfg(test)] @@ -380,6 +400,28 @@ mod tests { assert!(cache.state.lock().unwrap().entries.len() <= COMMITTEE_CACHE_CAPACITY); } + /// `head_current_committees` answers with the pinned head's + /// current-epoch entry, the middle of the three it pins, and only once + /// that entry has been built: it is a read, never a derivation. + #[test] + fn the_head_current_committees_are_the_middle_pin_once_built() { + let cache = CommitteeCache::default(); + assert!(cache.head_current_committees().is_none(), "nothing pinned"); + + let pinned_keys = [Some(key(1, 1)), Some(key(2, 2)), Some(key(3, 3))]; + cache.pin_head(Root::repeat_byte(0xAA), pinned_keys); + assert!( + cache.head_current_committees().is_none(), + "pinned but not built" + ); + + cache.get_or_init(key(2, 2), || EpochCommittees::new(2, vec![7, 8, 9], 1)); + let current = cache + .head_current_committees() + .expect("the current epoch's shuffling is built"); + assert_eq!(current.active_validator_count(), 3); + } + /// Concurrent misses on the same key must run `build` exactly once: the /// rest wait on the first call's result rather than each deriving their /// own. This is the property the module documentation calls load-bearing diff --git a/docs/beacon_wire.md b/docs/beacon_wire.md index 59912834..60791089 100644 --- a/docs/beacon_wire.md +++ b/docs/beacon_wire.md @@ -154,6 +154,66 @@ since nothing consumes them; an undecodable payload on any topic is REJECTed. Nothing is published on any topic, columns included: nothing this node can produce today would be signature-valid. +### Peer scoring + +Gossipsub peer scoring is on for this wire only +(`crate::beacon::scoring`). The parameters are lighthouse's +(`gossipsub_scoring_parameters.rs`), ported formula for formula; Grandine, +Teku and Lodestar use the same constants, so thresholds mean the same thing +here as on most of mainnet. + +| Threshold | Score | Effect | +| --- | --- | --- | +| gossip | -4000 | no IHAVE/IWANT to or from the peer | +| publish | -8000 | left out of this node's publishes | +| graylist | -16000 | every RPC ignored, and the peer is disconnected | + +| Topic | Weight | Mesh deliveries (P3) | +| --- | --- | --- | +| `beacon_block` | 0.5 | scored | +| `beacon_aggregate_and_proof` | 0.5 | scored | +| `beacon_attestation_{0..63}` | 1/64 each | scored | +| `voluntary_exit`, `proposer_slashing`, `attester_slashing` | 0.05 each | off | + +Columns, sync committee contributions and BLS changes carry no topic +parameters. Only Prysm scores columns. This node `Ignore`s every contribution +and change (no consumer), and a delivery is only credited once accepted, so +scoring P3 there would penalize every mesh peer for this node's own gap. All +64 attestation subnets get parameters, not just the backbone ones, because +aggregator duties join others at runtime. + +Two differences from lighthouse: + +- `mesh_n` is this node's `D` (8), not lighthouse's 5. It only enters the + first-message-delivery cap. +- **P3 waits for this node to keep up.** Lighthouse joins these topics only + once synced, so it has no mesh while it catches up. This node subscribes at + startup, and while catching up it `Ignore`s every aggregate and attestation + voting for a block it has not imported. A mesh peer credited with no + aggregates scores below the graylist, so a restart would graylist the whole + aggregate mesh because of this node's lag. `MeshDeliveryGate` turns P3 off + while the head lags the wall clock by more than + `MESH_DELIVERY_MAX_HEAD_LAG` slots, and back on only after an epoch within + that lag, so the counters have refilled. While P3 is off its threshold and + weight are zero but its cap is not, so the counters keep counting to their + real ceiling. `SyncStatus` cannot serve as this gate: on the mainnet + follower it reported `synced` through a 10-minute catch-up up to 98 slots + behind, since its network-stall rule reads a lagging freshest-known block + as a stalled network. + +The block, aggregate and attestation parameters depend on the active +validator count, so the p2p actor rebuilds them every slot from the head's +pinned current-epoch shuffling (`CommitteeCache::head_current_committees`) +and hands them to the swarm. The swarm starts on lighthouse's placeholder (32 +validators, P3 off) until the first refresh, which runs at startup. Every 10 +s the swarm task disconnects peers below the graylist. There is no ban list: +a peer's negative score is kept for `retain_score` (100 epochs) after it +leaves, so a peer that reconnects comes back graylisted and is dropped again. + +Topics are named for the one fork digest subscribed at startup. Once the +digest can change at runtime, the old topics need their weight zeroed and +the new ones need parameters (lighthouse's `remove_topic_weight_except`). + ## Aggregate attestations `beacon_aggregate_and_proof` reaches fork choice. It is how a follower learns @@ -552,6 +612,9 @@ as a query filter, so a `quic`-only record is invisible to it. | `lean_beacon_aggregate_end_to_end_seconds` | Wire to fork choice, for aggregates applied on arrival | | `lean_beacon_aggregate_total{outcome}` | Aggregates by `applied`, `invalid`, `known_subset` or `queue_full` | | `lean_beacon_aggregates_deferred` | Aggregates held until their own slot has passed | +| `lean_gossipsub_peers_by_score{band}` | Gossipsub peers by score band; see [Peer scoring](#peer-scoring) | +| `lean_gossipsub_score_disconnects_total` | Peers disconnected for a score below the graylist | +| `lean_gossipsub_mesh_delivery_scoring` | 1 while P3 is scored, 0 while the head lags or is warming up | The four aggregate histograms no longer cover what they used to: gossip validation (committees, all three signatures, the seen caches) runs in p2p now diff --git a/docs/metrics.md b/docs/metrics.md index 214392e8..f41d4bfa 100644 --- a/docs/metrics.md +++ b/docs/metrics.md @@ -387,6 +387,26 @@ instead. Neither pool has a metric of its own yet; a permit exhausted on either shows up as `outcome="ignore",reason="overloaded"` on `lean_beacon_gossip_validation_total`, for the topics that draw from it. +### Beacon Peer Scoring + +Gossipsub peer scoring runs on the beacon wire only; see +[beacon_wire.md](./beacon_wire.md#peer-scoring) for the parameters and the +mesh-delivery gate. These are ethlambda-specific, not part of the leanMetrics +spec. + +| Name | Type | Usage | Sample collection event | Labels | Buckets | +|------|------|-------|-------------------------|--------|---------| +| `lean_gossipsub_peers_by_score` | Gauge | Gossipsub peers by score band | Every 10 s, in the swarm task | band=non_negative,negative,below_gossip,below_publish,below_graylist | | +| `lean_gossipsub_score_disconnects_total` | Counter | Peers disconnected for a gossipsub score below the graylist (-16000) | When the 10 s check disconnects one | | | +| `lean_gossipsub_mesh_delivery_scoring` | Gauge | 1 while mesh message deliveries (P3) are scored, 0 while the head lags the wall clock or is warming up after catching up | Every slot, on the scoring refresh | | | + +The bands are mutually exclusive and split at the gossip (-4000), publish +(-8000) and graylist (-16000) thresholds, so they add up to the peers +gossipsub knows. `below_graylist` counts the peers that same pass is about to +disconnect, so it should not hold a value across passes. A persistent +`below_gossip`/`below_publish` population points at slow peers first: a full +send queue costs a peer 10 points per dropped message. + ### Beacon Committee Cache `ethlambda beacon` derives an epoch's attester committees with one whole-epoch