From 7e8d38defd68c7e46b50b24163ddcecc8b7040ca Mon Sep 17 00:00:00 2001 From: shane-moore Date: Sat, 18 Jul 2026 00:08:50 -0700 Subject: [PATCH 01/15] Eagerly sign sync committee messages on non-optimistic head events Publish sync committee messages as soon as a non-optimistic head event for the current slot arrives, instead of always sleeping to the due point, mirroring the attestation service change from #7892. Contributions remain delayed to their point in the slot. - Convert the head event channel to tokio broadcast so multiple services can subscribe via BeaconNodeFallback::subscribe_to_head_events. - Share the head-event-vs-deadline select (including the degrade to timer-only on channel close) between the attestation and sync committee services via beacon_head_monitor::head_event_or_deadline. This also fixes the attestation service busy-looping if the head event channel ever closed. - If a head event arrives before sync duties are computed, wait until the sync message deadline and check once more, reusing the event root. - Make sync message timing fork-aware via ChainSpec::get_sync_message_due_at_slot (SYNC_MESSAGE_DUE_BPS_GLOAS). - Extend MockBeaconNode with sync duty, head root, sync message pool and subscription mocks, and add sync committee service tests on the shared ValidatorClientHarness. Co-Authored-By: Claude Fable 5 Claude-Session: https://claude.ai/code/session_01MSiv1h1jtjqp8JE2AXbntQ --- book/src/help_vc.md | 10 +- consensus/types/src/core/chain_spec.rs | 62 +- .../src/mock_beacon_node.rs | 70 ++- .../src/validator_client_harness.rs | 4 +- .../src/beacon_head_monitor.rs | 133 +++- .../beacon_node_fallback/src/lib.rs | 13 +- validator_client/src/cli.rs | 4 +- validator_client/src/lib.rs | 24 +- .../src/attestation_service.rs | 58 +- .../src/sync_committee_service.rs | 573 ++++++++++++++++-- 10 files changed, 804 insertions(+), 147 deletions(-) diff --git a/book/src/help_vc.md b/book/src/help_vc.md index 719b02a5a53..8a9be7939eb 100644 --- a/book/src/help_vc.md +++ b/book/src/help_vc.md @@ -186,11 +186,11 @@ Flags: validators-dir. Validators will need to be manually added to the validator_definitions.yml file. --disable-beacon-head-monitor - Disable the beacon head monitor which tries to attest as soon as any - of the configured beacon nodes sends a head event. Leaving the service - enabled is recommended, but disabling it can lead to reduced bandwidth - and more predictable usage of the primary beacon node (rather than the - fastest BN). + Disable the beacon head monitor which triggers attestations and sync + committee messages when a configured beacon node sends a head event. + Leaving it enabled is recommended, but disabling it can lead to + reduced bandwidth and more predictable usage of the primary beacon + node (rather than the fastest BN). --disable-latency-measurement-service Disables the service that periodically attempts to measure latency to BNs. diff --git a/consensus/types/src/core/chain_spec.rs b/consensus/types/src/core/chain_spec.rs index f3d04bcff43..1c2fc0e1878 100644 --- a/consensus/types/src/core/chain_spec.rs +++ b/consensus/types/src/core/chain_spec.rs @@ -113,6 +113,7 @@ pub struct ChainSpec { pub payload_attestation_due_bps: u64, pub aggregate_due_bps: u64, pub sync_message_due_bps: u64, + pub sync_message_due_bps_gloas: u64, pub contribution_due_bps: u64, /* @@ -124,6 +125,7 @@ pub struct ChainSpec { pub payload_attestation_due: Duration, pub aggregate_attestation_due: Duration, pub sync_message_due: Duration, + pub sync_message_due_gloas: Duration, pub contribution_and_proof_due: Duration, /* @@ -933,6 +935,15 @@ impl ChainSpec { self.sync_message_due } + /// Spec: `get_sync_message_due_ms`. Returns the epoch-appropriate threshold. + pub fn get_sync_message_due_at_slot(&self, slot: Slot) -> Duration { + if self.fork_name_at_slot::(slot).gloas_enabled() { + self.sync_message_due_gloas + } else { + self.sync_message_due + } + } + /// Calculate the duration into a slot for a given slot component pub fn compute_slot_component_duration( &self, @@ -978,6 +989,11 @@ impl ChainSpec { "invalid chain spec: sync_message_due_bps ({}) exceeds slot duration", self.sync_message_due_bps ); + assert!( + self.sync_message_due_bps_gloas <= BASIS_POINTS, + "invalid chain spec: sync_message_due_bps_gloas ({}) exceeds slot duration", + self.sync_message_due_bps_gloas + ); assert!( self.contribution_due_bps <= BASIS_POINTS, "invalid chain spec: contribution_due_bps ({}) exceeds slot duration", @@ -1002,6 +1018,9 @@ impl ChainSpec { self.sync_message_due = self .compute_slot_component_duration(self.sync_message_due_bps) .expect("invalid chain spec: cannot compute sync_message_due"); + self.sync_message_due_gloas = self + .compute_slot_component_duration(self.sync_message_due_bps_gloas) + .expect("invalid chain spec: cannot compute sync_message_due_gloas"); self.contribution_and_proof_due = self .compute_slot_component_duration(self.contribution_due_bps) .expect("invalid chain spec: cannot compute contribution_and_proof_due"); @@ -1131,6 +1150,7 @@ impl ChainSpec { payload_attestation_due_bps: 7500, aggregate_due_bps: 6667, sync_message_due_bps: 3333, + sync_message_due_bps_gloas: 2500, contribution_due_bps: 6667, /* @@ -1142,6 +1162,7 @@ impl ChainSpec { payload_attestation_due: Duration::from_millis(9000), aggregate_attestation_due: Duration::from_millis(8000), sync_message_due: Duration::from_millis(3999), + sync_message_due_gloas: Duration::from_millis(3000), contribution_and_proof_due: Duration::from_millis(8000), /* @@ -1471,6 +1492,7 @@ impl ChainSpec { payload_attestation_due: Duration::from_millis(4500), aggregate_attestation_due: Duration::from_millis(4000), sync_message_due: Duration::from_millis(1999), + sync_message_due_gloas: Duration::from_millis(1500), contribution_and_proof_due: Duration::from_millis(4000), // Networking Fulu @@ -1562,6 +1584,7 @@ impl ChainSpec { payload_due_bps: 7500, payload_attestation_due_bps: 7500, aggregate_due_bps: 6667, + sync_message_due_bps_gloas: 2500, /* * Derived time values (set by `compute_derived_values()`) @@ -1573,6 +1596,7 @@ impl ChainSpec { payload_attestation_due: Duration::from_millis(3750), aggregate_attestation_due: Duration::from_millis(3333), sync_message_due: Duration::from_millis(1666), + sync_message_due_gloas: Duration::from_millis(1250), contribution_and_proof_due: Duration::from_millis(3333), /* @@ -2190,6 +2214,9 @@ pub struct Config { #[serde(default = "default_sync_message_due_bps")] #[serde(with = "serde_utils::quoted_u64")] sync_message_due_bps: u64, + #[serde(default = "default_sync_message_due_bps_gloas")] + #[serde(with = "serde_utils::quoted_u64")] + sync_message_due_bps_gloas: u64, #[serde(default = "default_contribution_due_bps")] #[serde(with = "serde_utils::quoted_u64")] contribution_due_bps: u64, @@ -2444,6 +2471,10 @@ const fn default_sync_message_due_bps() -> u64 { 3333 } +const fn default_sync_message_due_bps_gloas() -> u64 { + 2500 +} + const fn default_contribution_due_bps() -> u64 { 6667 } @@ -2723,6 +2754,7 @@ impl Config { payload_attestation_due_bps: spec.payload_attestation_due_bps, aggregate_due_bps: spec.aggregate_due_bps, sync_message_due_bps: spec.sync_message_due_bps, + sync_message_due_bps_gloas: spec.sync_message_due_bps_gloas, contribution_due_bps: spec.contribution_due_bps, min_builder_withdrawability_delay: spec.min_builder_withdrawability_delay.as_u64(), @@ -2827,6 +2859,7 @@ impl Config { payload_attestation_due_bps, aggregate_due_bps, sync_message_due_bps, + sync_message_due_bps_gloas, contribution_due_bps, confirmation_byzantine_threshold, min_builder_withdrawability_delay, @@ -2939,6 +2972,7 @@ impl Config { payload_attestation_due_bps, aggregate_due_bps, sync_message_due_bps, + sync_message_due_bps_gloas, contribution_due_bps, min_builder_withdrawability_delay: Epoch::new(min_builder_withdrawability_delay), @@ -3714,6 +3748,7 @@ mod yaml_tests { // Test sync message (3333 bps = 33.33% of 12s = 4s) let sync_msg_due = spec.get_sync_message_due(); assert_eq!(sync_msg_due, Duration::from_millis(3999)); // 12000 * 3333 / 10000 + assert_eq!(spec.sync_message_due_gloas, Duration::from_millis(3000)); // Test contribution message (6667 bps = 66.67% of 12s = 8s) let contribution_due = spec.get_contribution_message_due(); @@ -3795,6 +3830,24 @@ mod yaml_tests { ); // 12000 * 5000 / 10000 } + #[test] + fn sync_message_due_at_slot() { + let mut spec = ChainSpec::mainnet(); + let gloas_fork_epoch = Epoch::new(1); + spec.gloas_fork_epoch = Some(gloas_fork_epoch); + let spec = spec.compute_derived_values::(); + let gloas_fork_slot = gloas_fork_epoch.start_slot(MainnetEthSpec::slots_per_epoch()); + + assert_eq!( + spec.get_sync_message_due_at_slot::(gloas_fork_slot - 1), + Duration::from_millis(3999) + ); + assert_eq!( + spec.get_sync_message_due_at_slot::(gloas_fork_slot), + Duration::from_millis(3000) + ); + } + #[test] fn test_default_duration_values_without_compute_derived_values() { // Verify that mainnet, minimal, and gnosis have correct pre-computed defaults @@ -3893,6 +3946,14 @@ mod yaml_tests { spec.compute_derived_values::(); } + #[test] + #[should_panic(expected = "sync_message_due_bps_gloas")] + fn compute_derived_values_rejects_invalid_gloas_sync_message_due() { + let mut spec = ChainSpec::mainnet(); + spec.sync_message_due_bps_gloas = BASIS_POINTS + 1; + spec.compute_derived_values::(); + } + fn configs_base_path() -> PathBuf { env::var("CARGO_MANIFEST_DIR") .expect("should know manifest dir") @@ -3912,7 +3973,6 @@ mod yaml_tests { "EIP7928_FORK_EPOCH", // Gloas params not yet in Config "AGGREGATE_DUE_BPS_GLOAS", - "SYNC_MESSAGE_DUE_BPS_GLOAS", "CONTRIBUTION_DUE_BPS_GLOAS", "MAX_REQUEST_PAYLOADS", // Heze networking diff --git a/testing/validator_test_rig/src/mock_beacon_node.rs b/testing/validator_test_rig/src/mock_beacon_node.rs index 255831da502..71a34b7563f 100644 --- a/testing/validator_test_rig/src/mock_beacon_node.rs +++ b/testing/validator_test_rig/src/mock_beacon_node.rs @@ -1,4 +1,4 @@ -use eth2::types::{GenericResponse, PublishBlockRequest, SyncingData}; +use eth2::types::{GenericResponse, PublishBlockRequest, RootData, SyncingData}; use eth2::{BeaconNodeHttpClient, Timeouts}; use mockito::{Matcher, Mock, Server, ServerGuard}; use regex::Regex; @@ -11,9 +11,9 @@ use std::sync::{Arc, Mutex}; use std::time::Duration; use tracing::info; use types::{ - BeaconBlock, ChainSpec, ConfigAndPreset, EthSpec, ExecutionPayloadEnvelope, ForkName, Hash256, - PayloadAttestationData, PayloadAttestationMessage, SignedBlindedBeaconBlock, - SignedExecutionPayloadEnvelope, Slot, + BeaconBlock, ChainSpec, ConfigAndPreset, Epoch, EthSpec, ExecutionPayloadEnvelope, ForkName, + Hash256, PayloadAttestationData, PayloadAttestationMessage, SignedBlindedBeaconBlock, + SignedExecutionPayloadEnvelope, Slot, SyncCommitteeMessage, SyncDuty, }; pub struct MockBeaconNode { @@ -24,6 +24,7 @@ pub struct MockBeaconNode { pub received_full_blocks: Arc>>>, pub execution_payload_envelope: Arc>>>, pub payload_attestation_message: Arc>>, + pub sync_committee_messages: Arc>>, } impl MockBeaconNode { @@ -42,6 +43,7 @@ impl MockBeaconNode { received_full_blocks: Arc::new(Mutex::new(Vec::new())), execution_payload_envelope: Arc::new(Mutex::new(Vec::new())), payload_attestation_message: Arc::new(Mutex::new(Vec::new())), + sync_committee_messages: Arc::new(Mutex::new(Vec::new())), } } @@ -102,6 +104,66 @@ impl MockBeaconNode { .create(); } + /// Mocks `POST /eth/v1/validator/duties/sync/{epoch}` + pub fn mock_sync_duties(&mut self, epoch: Epoch, duties: Vec) -> Mock { + let path_pattern = Regex::new(&format!( + r"^/eth/v1/validator/duties/sync/{}$", + epoch.as_u64() + )) + .unwrap(); + let response = + GenericResponse::from(duties).add_execution_optimistic_finalized(false, false); + + self.server + .mock("POST", Matcher::Regex(path_pattern.to_string())) + .with_status(200) + .with_header("content-type", "application/json") + .with_body(serde_json::to_string(&response).unwrap()) + .create() + } + + /// Mocks `GET /eth/v1/beacon/blocks/head/root` + pub fn mock_get_head_block_root(&mut self, root: Hash256) -> Mock { + let path_pattern = Regex::new(r"^/eth/v1/beacon/blocks/head/root$").unwrap(); + let response = GenericResponse::from(RootData { root }) + .add_execution_optimistic_finalized(false, false); + + self.server + .mock("GET", Matcher::Regex(path_pattern.to_string())) + .with_status(200) + .with_header("content-type", "application/json") + .with_body(serde_json::to_string(&response).unwrap()) + .create() + } + + /// Mocks `POST /eth/v1/beacon/pool/sync_committees` + pub fn mock_post_sync_committee_messages(&mut self) -> Mock { + let path_pattern = Regex::new(r"^/eth/v1/beacon/pool/sync_committees$").unwrap(); + let sync_committee_messages = Arc::clone(&self.sync_committee_messages); + + self.server + .mock("POST", Matcher::Regex(path_pattern.to_string())) + .with_status(200) + .with_body_from_request(move |request| { + let body = request.body().expect("Failed to get request body"); + let messages: Vec = serde_json::from_slice(body) + .expect("Failed to deserialize sync committee messages"); + sync_committee_messages.lock().unwrap().extend(messages); + vec![] + }) + .create() + } + + /// Mocks `POST /eth/v1/validator/sync_committee_subscriptions` + pub fn mock_sync_committee_subscriptions(&mut self) -> Mock { + let path_pattern = Regex::new(r"^/eth/v1/validator/sync_committee_subscriptions$").unwrap(); + + self.server + .mock("POST", Matcher::Regex(path_pattern.to_string())) + .with_status(200) + .create() + } + /// Mocks `GET /eth/v4/validator/blocks/{slot}` pub fn mock_get_validator_blocks_v4( &mut self, diff --git a/testing/validator_test_rig/src/validator_client_harness.rs b/testing/validator_test_rig/src/validator_client_harness.rs index fe6194ffada..22ac9ca8212 100644 --- a/testing/validator_test_rig/src/validator_client_harness.rs +++ b/testing/validator_test_rig/src/validator_client_harness.rs @@ -38,8 +38,10 @@ impl ValidatorClientHarness { pub async fn new(num_validators: usize) -> Self { let mut default_spec = MainnetEthSpec::default_spec(); default_spec.gloas_fork_epoch = Some(Epoch::new(0)); - let spec = Arc::new(default_spec); + Self::new_with_spec(num_validators, Arc::new(default_spec)).await + } + pub async fn new_with_spec(num_validators: usize, spec: Arc) -> Self { let test_runtime = TestRuntime::default(); let executor = test_runtime.task_executor.clone(); let slot_duration = spec.get_slot_duration(); diff --git a/validator_client/beacon_node_fallback/src/beacon_head_monitor.rs b/validator_client/beacon_node_fallback/src/beacon_head_monitor.rs index bed107d856d..5b16fdaa86d 100644 --- a/validator_client/beacon_node_fallback/src/beacon_head_monitor.rs +++ b/validator_client/beacon_node_fallback/src/beacon_head_monitor.rs @@ -4,21 +4,78 @@ use futures::StreamExt; use slot_clock::SlotClock; use std::collections::HashMap; use std::sync::Arc; -use tokio::sync::RwLock; +use std::time::Duration; +use tokio::sync::{RwLock, broadcast}; +use tokio::time::sleep; use tracing::{debug, info, warn}; use types::EthSpec; type CacheHashMap = HashMap; -// This is used to send the index derived from `CandidateBeaconNode` to the -// `AttestationService` for further processing -#[derive(Debug)] +// This is used to send the index derived from `CandidateBeaconNode` to validator services. +#[derive(Clone, Debug)] pub struct HeadEvent { pub beacon_node_index: usize, pub slot: types::Slot, pub beacon_block_root: Hash256, } +pub async fn poll_for_current_slot_head( + receiver: &mut broadcast::Receiver, + slot_clock: &T, +) -> Option { + loop { + match receiver.recv().await { + Ok(head_event) => { + // Skip the event on a clock read failure rather than returning `None`, + // which callers treat as a terminal error that disables head monitoring. + let Some(current_slot) = slot_clock.now() else { + continue; + }; + if head_event.slot == current_slot { + return Some(head_event); + } + } + Err(broadcast::error::RecvError::Lagged(skipped)) => { + warn!(skipped, "Head monitor channel lagged"); + } + Err(broadcast::error::RecvError::Closed) => { + warn!("Head monitor channel closed unexpectedly"); + return None; + } + } + } +} + +/// Wait for a head event for the current slot, or until `deadline` has elapsed. +/// +/// Returns `None` when the deadline fires. If the head event channel closes, head monitoring +/// is disabled by clearing `head_monitor_rx` and the deadline is awaited as usual. +pub async fn head_event_or_deadline( + head_monitor_rx: &mut Option>, + slot_clock: &T, + deadline: Duration, +) -> Option { + let deadline = sleep(deadline); + tokio::pin!(deadline); + if let Some(receiver) = head_monitor_rx { + tokio::select! { + _ = &mut deadline => None, + event = poll_for_current_slot_head(receiver, slot_clock) => match event { + Some(event) => Some(event), + None => { + *head_monitor_rx = None; + deadline.await; + None + } + }, + } + } else { + deadline.await; + None + } +} + /// Cache to maintain the latest head received from each of the beacon nodes /// in the `BeaconNodeFallback`. #[derive(Debug)] @@ -68,16 +125,7 @@ impl Default for BeaconHeadCache { } } -// Runs a non-terminating loop to update the `BeaconHeadCache` with the latest head received -// from the candidate beacon_nodes. This is an attempt to stream events to beacon nodes and -// potential start attestation duties earlier as soon as latest head is receive from any of the -// beacon node in contrast to attest at the 1/3rd mark in the slot. -// -// -// The cache and the candidate BNs list are refresh/purged to avoid dangling reference conditions -// that arise due to `update_candidates_list`. -// -// Starts the service to perpetually stream head events from connected beacon_nodes +// Updates the head cache and streams the latest non-optimistic head events from connected BNs. pub async fn poll_head_event_from_beacon_nodes( beacon_nodes: Arc>, ) -> Result<(), String> { @@ -167,7 +215,6 @@ pub async fn poll_head_event_from_beacon_nodes SseHead { @@ -313,6 +362,60 @@ mod tests { assert_eq!(event.beacon_block_root, block_root); } + fn head_event(slot: u64, block_root: u64) -> HeadEvent { + HeadEvent { + beacon_node_index: 0, + slot: Slot::new(slot), + beacon_block_root: Hash256::from_low_u64_be(block_root), + } + } + + #[tokio::test] + async fn ignores_stale_head_events() { + let slot_clock = + ManualSlotClock::new(Slot::new(0), Duration::ZERO, Duration::from_secs(12)); + slot_clock.set_slot(2); + let (sender, mut receiver) = broadcast::channel(2); + + sender.send(head_event(1, 1)).unwrap(); + sender.send(head_event(2, 2)).unwrap(); + + let event = poll_for_current_slot_head(&mut receiver, &slot_clock) + .await + .unwrap(); + assert_eq!(event.beacon_block_root, Hash256::from_low_u64_be(2)); + } + + #[tokio::test] + async fn recovers_from_lagged_head_events() { + let slot_clock = + ManualSlotClock::new(Slot::new(0), Duration::ZERO, Duration::from_secs(12)); + slot_clock.set_slot(2); + let (sender, mut receiver) = broadcast::channel(1); + + sender.send(head_event(1, 1)).unwrap(); + sender.send(head_event(2, 2)).unwrap(); + + let event = poll_for_current_slot_head(&mut receiver, &slot_clock) + .await + .unwrap(); + assert_eq!(event.beacon_block_root, Hash256::from_low_u64_be(2)); + } + + #[tokio::test] + async fn returns_when_head_event_channel_closes() { + let slot_clock = + ManualSlotClock::new(Slot::new(0), Duration::ZERO, Duration::from_secs(12)); + let (sender, mut receiver) = broadcast::channel(1); + drop(sender); + + assert!( + poll_for_current_slot_head(&mut receiver, &slot_clock) + .await + .is_none() + ); + } + #[tokio::test] async fn test_cache_caches_multiple_heads_from_different_nodes() { let cache = BeaconHeadCache::new(); diff --git a/validator_client/beacon_node_fallback/src/lib.rs b/validator_client/beacon_node_fallback/src/lib.rs index 5eb94633a3f..bd4cdbab976 100644 --- a/validator_client/beacon_node_fallback/src/lib.rs +++ b/validator_client/beacon_node_fallback/src/lib.rs @@ -26,7 +26,7 @@ use std::vec::Vec; use strum::VariantNames; use task_executor::TaskExecutor; use tokio::{ - sync::{RwLock, mpsc}, + sync::{RwLock, broadcast}, time::sleep, }; use tracing::{debug, error, warn}; @@ -415,7 +415,7 @@ pub struct BeaconNodeFallback { distance_tiers: BeaconNodeSyncDistanceTiers, slot_clock: Option, beacon_head_cache: Option>, - head_monitor_send: Option>>, + head_monitor_send: Option>, broadcast_topics: Vec, spec: Arc, } @@ -452,11 +452,18 @@ impl BeaconNodeFallback { /// validator client is connected in the `BeaconNodeFallback`. This also initializes the /// beacon_head_cache under the assumption the beacon_head_cache will always be needed when /// head_monitor_send is set. - pub fn set_head_send(&mut self, head_monitor_send: Arc>) { + pub fn set_head_send(&mut self, head_monitor_send: broadcast::Sender) { self.head_monitor_send = Some(head_monitor_send); self.beacon_head_cache = Some(Arc::new(BeaconHeadCache::new())); } + /// Subscribe to the stream of head events, if the head monitor is enabled. + pub fn subscribe_to_head_events(&self) -> Option> { + self.head_monitor_send + .as_ref() + .map(|sender| sender.subscribe()) + } + /// The count of candidates, regardless of their state. pub async fn num_total(&self) -> usize { self.candidates.read().await.len() diff --git a/validator_client/src/cli.rs b/validator_client/src/cli.rs index cf21e276d7b..e789153ce05 100644 --- a/validator_client/src/cli.rs +++ b/validator_client/src/cli.rs @@ -490,8 +490,8 @@ pub struct ValidatorClient { #[clap( long, - help = "Disable the beacon head monitor which tries to attest as soon as any of the \ - configured beacon nodes sends a head event. Leaving the service enabled is \ + help = "Disable the beacon head monitor which triggers attestations and sync committee \ + messages when a configured beacon node sends a head event. Leaving it enabled is \ recommended, but disabling it can lead to reduced bandwidth and more predictable \ usage of the primary beacon node (rather than the fastest BN).", display_order = 0, diff --git a/validator_client/src/lib.rs b/validator_client/src/lib.rs index 88844918431..d81aa8e60d0 100644 --- a/validator_client/src/lib.rs +++ b/validator_client/src/lib.rs @@ -9,7 +9,6 @@ use metrics::set_gauge; use monitoring_api::{MonitoringHttpClient, ProcessType}; use sensitive_url::SensitiveUrl; use slashing_protection::{SLASHING_PROTECTION_FILENAME, SlashingDatabase}; -use tokio::sync::Mutex; use account_utils::validator_definitions::ValidatorDefinitions; use beacon_node_fallback::{ @@ -33,7 +32,7 @@ use std::path::Path; use std::sync::Arc; use std::time::{SystemTime, UNIX_EPOCH}; use tokio::{ - sync::mpsc, + sync::{broadcast, mpsc}, time::{Duration, sleep}, }; use tracing::{debug, error, info, warn}; @@ -410,16 +409,18 @@ impl ProductionValidatorClient { // Only the beacon_nodes are used for attestation duties and thus biconditionally // proposer_nodes do not need head_send ref. - let head_monitor_rx = if config.enable_beacon_head_monitor { - let (head_monitor_tx, head_receiver) = - mpsc::channel::(MAX_HEAD_EVENT_QUEUE_LEN); - beacon_nodes.set_head_send(Arc::new(head_monitor_tx)); - Some(Mutex::new(head_receiver)) - } else { - None - }; + if config.enable_beacon_head_monitor { + let (head_monitor_tx, _) = broadcast::channel::(MAX_HEAD_EVENT_QUEUE_LEN); + beacon_nodes.set_head_send(head_monitor_tx); + } let beacon_nodes = Arc::new(beacon_nodes); + + // Subscribe before the head monitor starts so that no events are sent while the + // channel has no receivers. + let attestation_head_monitor_rx = beacon_nodes.subscribe_to_head_events(); + let sync_head_monitor_rx = beacon_nodes.subscribe_to_head_events(); + start_fallback_updater_service::<_, E>(context.executor.clone(), beacon_nodes.clone())?; let proposer_nodes = Arc::new(proposer_nodes); @@ -536,7 +537,7 @@ impl ProductionValidatorClient { .validator_store(validator_store.clone()) .beacon_nodes(beacon_nodes.clone()) .executor(context.executor.clone()) - .head_monitor_rx(head_monitor_rx) + .head_monitor_rx(attestation_head_monitor_rx) .chain_spec(context.eth2_config.spec.clone()) .disable(config.disable_attesting); @@ -557,6 +558,7 @@ impl ProductionValidatorClient { slot_clock.clone(), beacon_nodes.clone(), context.executor.clone(), + sync_head_monitor_rx, ); let payload_attestation_service = PayloadAttestationService::new( diff --git a/validator_client/validator_services/src/attestation_service.rs b/validator_client/validator_services/src/attestation_service.rs index e28d0912be4..634379886fa 100644 --- a/validator_client/validator_services/src/attestation_service.rs +++ b/validator_client/validator_services/src/attestation_service.rs @@ -1,5 +1,8 @@ use crate::duties_service::{DutiesService, DutyAndProof}; -use beacon_node_fallback::{ApiTopic, BeaconNodeFallback, beacon_head_monitor::HeadEvent}; +use beacon_node_fallback::{ + ApiTopic, BeaconNodeFallback, + beacon_head_monitor::{HeadEvent, head_event_or_deadline}, +}; use futures::StreamExt; use logging::crit; use slot_clock::SlotClock; @@ -7,8 +10,7 @@ use std::collections::HashMap; use std::ops::Deref; use std::sync::Arc; use task_executor::TaskExecutor; -use tokio::sync::Mutex; -use tokio::sync::mpsc; +use tokio::sync::{Mutex, broadcast}; use tokio::time::{Duration, Instant, sleep, sleep_until}; use tracing::{Instrument, debug, error, info, info_span, instrument, warn}; use tree_hash::TreeHash; @@ -24,7 +26,7 @@ pub struct AttestationServiceBuilder beacon_nodes: Option>>, executor: Option, chain_spec: Option>, - head_monitor_rx: Option>>, + head_monitor_rx: Option>, disable: bool, } @@ -79,7 +81,7 @@ impl AttestationServiceBuil pub fn head_monitor_rx( mut self, - head_monitor_rx: Option>>, + head_monitor_rx: Option>, ) -> Self { self.head_monitor_rx = head_monitor_rx; self @@ -105,7 +107,7 @@ impl AttestationServiceBuil chain_spec: self .chain_spec .ok_or("Cannot build AttestationService without chain_spec")?, - head_monitor_rx: self.head_monitor_rx, + head_monitor_rx: Mutex::new(self.head_monitor_rx), disable: self.disable, latest_attested_slot: Mutex::new(Slot::default()), }), @@ -121,7 +123,7 @@ pub struct Inner { beacon_nodes: Arc>, executor: TaskExecutor, chain_spec: Arc, - head_monitor_rx: Option>>, + head_monitor_rx: Mutex>>, disable: bool, latest_attested_slot: Mutex, } @@ -176,6 +178,7 @@ impl AttestationService AttestationService None, - event = self.poll_for_head_events() => - event.map(|event| (event.beacon_node_index, event.beacon_block_root)), - } - } else { - sleep(duration + unaggregated_attestation_due).await; - None - }; + let beacon_node_data = head_event_or_deadline( + &mut head_monitor_rx, + &self.slot_clock, + duration + unaggregated_attestation_due, + ) + .await + .map(|event| (event.beacon_node_index, event.beacon_block_root)); let Some(current_slot) = self.slot_clock.now() else { error!("Failed to read slot clock after trigger"); @@ -221,30 +221,6 @@ impl AttestationService Option { - let Some(receiver) = &self.head_monitor_rx else { - return None; - }; - let mut receiver = receiver.lock().await; - loop { - match receiver.recv().await { - Some(head_event) => { - // Only return head events for the current slot - this ensures the - // block for this slot has been produced before triggering attestation - let current_slot = self.slot_clock.now()?; - if head_event.slot == current_slot { - return Some(head_event); - } - // Head event is for a previous slot, keep waiting - } - None => { - warn!("Head monitor channel closed unexpectedly"); - return None; - } - } - } - } - /// Spawn only one new task for attestation post-Electra /// For each required aggregates, spawn a new task that downloads, signs and uploads the /// aggregates to the beacon node. diff --git a/validator_client/validator_services/src/sync_committee_service.rs b/validator_client/validator_services/src/sync_committee_service.rs index 8ea16fc3ad0..310f3b5849a 100644 --- a/validator_client/validator_services/src/sync_committee_service.rs +++ b/validator_client/validator_services/src/sync_committee_service.rs @@ -1,5 +1,8 @@ use crate::duties_service::DutiesService; -use beacon_node_fallback::{ApiTopic, BeaconNodeFallback}; +use beacon_node_fallback::{ + ApiTopic, BeaconNodeFallback, + beacon_head_monitor::{HeadEvent, head_event_or_deadline}, +}; use bls::PublicKeyBytes; use eth2::types::BlockId; use futures::StreamExt; @@ -11,6 +14,7 @@ use std::ops::Deref; use std::sync::Arc; use std::sync::atomic::{AtomicBool, Ordering}; use task_executor::TaskExecutor; +use tokio::sync::{Mutex, broadcast}; use tokio::time::{Duration, Instant, sleep, sleep_until}; use tracing::{Instrument, debug, error, info, info_span, instrument, trace, warn}; use types::{ @@ -21,6 +25,21 @@ use validator_store::{ContributionToSign, SyncMessageToSign, ValidatorStore}; pub const SUBSCRIPTION_LOOKAHEAD_EPOCHS: u64 = 4; +/// Duration from now until `due` past the start of `slot`, or `None` if `slot` is not the +/// current slot. +fn delay_until_slot_offset( + slot_clock: &T, + slot: Slot, + due: Duration, +) -> Option { + let now = slot_clock.now_duration()?; + if slot_clock.slot_of(now)? != slot { + return None; + } + let due_at = slot_clock.start_of(slot)?.checked_add(due)?; + Some(due_at.saturating_sub(now)) +} + pub struct SyncCommitteeService { inner: Arc>, } @@ -47,6 +66,7 @@ pub struct Inner { slot_clock: T, beacon_nodes: Arc>, executor: TaskExecutor, + head_monitor_rx: Mutex>>, /// Boolean to track whether the service has posted subscriptions to the BN at least once. /// /// This acts as a latch that fires once upon start-up, and then never again. @@ -60,6 +80,7 @@ impl SyncCommitteeService>, executor: TaskExecutor, + head_monitor_rx: Option>, ) -> Self { Self { inner: Arc::new(Inner { @@ -68,6 +89,7 @@ impl SyncCommitteeService SyncCommitteeService = None; loop { - if let Some(duration_to_next_slot) = self.slot_clock.duration_to_next_slot() { - // Wait for contribution broadcast interval 1/3 of the way through the slot. - sleep(duration_to_next_slot + sync_message_slot_component).await; + let Some((duration_to_next_slot, next_slot)) = self + .slot_clock + .duration_to_next_slot() + .zip(self.slot_clock.now().map(|slot| slot + 1)) + else { + error!("Failed to read slot clock"); + sleep(slot_duration).await; + continue; + }; - // Do nothing if the Altair fork has not yet occurred. - if !self.altair_fork_activated() { - continue; - } + // Wait for the sync message due point of the next slot, or a head event for the + // current slot, whichever comes first. + let sync_message_due = self + .duties_service + .spec + .get_sync_message_due_at_slot::(next_slot); + let head_event_root = head_event_or_deadline( + &mut head_monitor_rx, + &self.slot_clock, + duration_to_next_slot + sync_message_due, + ) + .await + .map(|event| event.beacon_block_root); + + let Some(current_slot) = self.slot_clock.now() else { + error!("Failed to read slot clock after trigger"); + continue; + }; + + if last_sync_message_slot.is_some_and(|last_slot| current_slot <= last_slot) { + debug!(%current_slot, "Sync messages already initiated for the slot"); + continue; + } - if let Err(e) = self.spawn_contribution_tasks().await { + // Do nothing if the Altair fork has not yet occurred. + if !self.altair_fork_activated() { + continue; + } + + match self.spawn_contribution_tasks(head_event_root).await { + Ok(()) => { + last_sync_message_slot = Some(current_slot); + trace!("Spawned sync contribution tasks"); + } + Err(e) => { crit!( error = ?e, "Failed to spawn sync contribution tasks" ); - } else { - trace!("Spawned sync contribution tasks"); } - - // Do subscriptions for future slots/epochs. - self.spawn_subscription_tasks(); - } else { - error!("Failed to read slot clock"); - // If we can't read the slot clock, just wait another slot. - sleep(slot_duration).await; } + + // Do subscriptions for future slots/epochs. + self.spawn_subscription_tasks(); } }; @@ -142,27 +193,36 @@ impl SyncCommitteeService Result<(), String> { + async fn spawn_contribution_tasks( + &self, + head_event_root: Option, + ) -> Result<(), String> { let spec = &self.duties_service.spec; let slot = self.slot_clock.now().ok_or("Failed to read slot clock")?; - let duration_to_next_slot = self - .slot_clock - .duration_to_next_slot() - .ok_or("Unable to determine duration to next slot")?; - // If a validator needs to publish a sync aggregate, they must do so at 2/3 - // through the slot. This delay triggers at this time - let aggregate_production_instant = Instant::now() - + duration_to_next_slot - .checked_add(spec.get_contribution_message_due()) - .and_then(|offset| offset.checked_sub(spec.get_slot_duration())) - .unwrap_or_else(|| Duration::from_secs(0)); - - let Some(slot_duties) = self + let mut slot_duties = self .duties_service .sync_duties - .get_duties_for_slot::(slot, &self.duties_service.spec) - else { + .get_duties_for_slot::(slot, spec); + + // If a head event triggered us before the duties were computed, wait until the sync + // message deadline and check for duties once more. + if slot_duties.is_none() && head_event_root.is_some() { + let duration_to_deadline = delay_until_slot_offset( + &self.slot_clock, + slot, + spec.get_sync_message_due_at_slot::(slot), + ) + .unwrap_or(Duration::ZERO); + sleep(duration_to_deadline).await; + + slot_duties = self + .duties_service + .sync_duties + .get_duties_for_slot::(slot, spec); + } + + let Some(slot_duties) = slot_duties else { debug!("No duties known for slot {}", slot); return Ok(()); }; @@ -172,35 +232,57 @@ impl SyncCommitteeService { - Ok(block) - } - Ok(Some(_)) => { - Err(format!("To sign sync committee messages for slot {slot} a non-optimistic head block is required")) - } - Ok(None) => Err(format!("No block root found for slot {}", slot)), - Err(e) => Err(e.to_string()), - } - }, - ) - .await; + // If a validator needs to publish a sync aggregate, they must do so at 2/3 + // through the slot. This delay triggers at this time + let Some(contribution_delay) = + delay_until_slot_offset(&self.slot_clock, slot, spec.get_contribution_message_due()) + else { + debug!(%slot, "Skipping sync committee tasks for expired slot"); + return Ok(()); + }; + let aggregate_production_instant = Instant::now() + contribution_delay; - let block_root = match response { - Ok(block) => block.data.root, - Err(errs) => { - warn!( - errors = errs.to_string(), - %slot, - "Refusing to sign sync committee messages for an optimistic head block or \ - a block head with unknown optimistic status" - ); - return Ok(()); + debug!( + %slot, + from_head_monitor = head_event_root.is_some(), + "Starting sync committee message production" + ); + + let block_root = if let Some(block_root) = head_event_root { + // The head monitor only forwards non-optimistic heads, so the event root can be + // used directly. + block_root + } else { + // Fetch `block_root` with non optimistic execution for `SyncCommitteeContribution`. + let response = self + .beacon_nodes + .first_success( + |beacon_node| async move { + match beacon_node.get_beacon_blocks_root(BlockId::Head).await { + Ok(Some(block)) if block.execution_optimistic == Some(false) => { + Ok(block) + } + Ok(Some(_)) => { + Err(format!("To sign sync committee messages for slot {slot} a non-optimistic head block is required")) + } + Ok(None) => Err(format!("No block root found for slot {}", slot)), + Err(e) => Err(e.to_string()), + } + }, + ) + .await; + + match response { + Ok(block) => block.data.root, + Err(errs) => { + warn!( + errors = errs.to_string(), + %slot, + "Refusing to sign sync committee messages for an optimistic head block or \ + a block head with unknown optimistic status" + ); + return Ok(()); + } } }; @@ -568,3 +650,366 @@ fn subscriptions_from_sync_duties( until_epoch, }) } + +#[cfg(test)] +mod tests { + use super::*; + use crate::{ + duties_service::DutiesServiceBuilder, sync::poll_sync_committee_duties_for_period, + }; + use bls::FixedBytesExtended; + use slot_clock::ManualSlotClock; + use types::{Epoch, MainnetEthSpec}; + use validator_test_rig::validator_client_harness::{S, ValidatorClientHarness}; + + type E = MainnetEthSpec; + + struct TestHarness { + harness: ValidatorClientHarness, + service: SyncCommitteeService, + head_sender: Option>, + } + + impl TestHarness { + async fn new(head_monitoring: bool) -> Self { + let mut spec = E::default_spec(); + spec.altair_fork_epoch = Some(Epoch::new(0)); + Self::new_with_spec(head_monitoring, Arc::new(spec)).await + } + + async fn new_with_spec(head_monitoring: bool, spec: Arc) -> Self { + let mut harness = ValidatorClientHarness::new_with_spec(1, spec).await; + harness + .mock_beacon_node_1 + .mock_sync_committee_subscriptions(); + + let duties_service = Arc::new( + DutiesServiceBuilder::new() + .validator_store(harness.validator_store.clone()) + .slot_clock(harness.slot_clock.clone()) + .beacon_nodes(harness.beacon_nodes.clone()) + .executor(harness.test_runtime.task_executor.clone()) + .spec(harness.spec.clone()) + .build() + .unwrap(), + ); + + let (head_sender, head_monitor_rx) = if head_monitoring { + let (sender, receiver) = broadcast::channel(8); + (Some(sender), Some(receiver)) + } else { + (None, None) + }; + + let service = SyncCommitteeService::new( + duties_service, + harness.validator_store.clone(), + harness.slot_clock.clone(), + harness.beacon_nodes.clone(), + harness.test_runtime.task_executor.clone(), + head_monitor_rx, + ); + + Self { + harness, + service, + head_sender, + } + } + + async fn insert_duties(&mut self) { + tokio::time::resume(); + let duty = SyncDuty { + pubkey: self.harness.pubkeys[0], + validator_index: 0, + validator_sync_committee_indices: vec![0], + }; + let mock = self + .harness + .mock_beacon_node_1 + .mock_sync_duties(Epoch::new(0), vec![duty]); + poll_sync_committee_duties_for_period(&self.service.duties_service, &[0], 0) + .await + .unwrap(); + mock.assert(); + tokio::time::pause(); + } + + fn start(&self) { + self.service + .clone() + .start_update_service(&self.harness.spec) + .unwrap(); + } + + fn send_head(&self, slot: u64, block_root: u64) { + self.head_sender + .as_ref() + .unwrap() + .send(head_event(slot, block_root)) + .unwrap(); + } + + async fn advance_time(&self, duration: Duration) { + self.harness.slot_clock.advance_time(duration); + tokio::time::advance(duration).await; + yield_to_service().await; + } + + fn messages(&self) -> Vec { + self.harness + .mock_beacon_node_1 + .sync_committee_messages + .lock() + .unwrap() + .clone() + } + } + + async fn yield_to_service() { + for _ in 0..20 { + tokio::task::yield_now().await; + } + } + + /// Resume real time so the service can complete signing and HTTP requests, then pause again. + async fn wait_for_message_count(harness: &TestHarness, count: usize) { + tokio::time::resume(); + for _ in 0..100 { + if harness.messages().len() >= count { + tokio::time::pause(); + return; + } + tokio::time::sleep(Duration::from_millis(1)).await; + } + tokio::time::pause(); + assert_eq!(harness.messages().len(), count); + } + + fn head_event(slot: u64, block_root: u64) -> HeadEvent { + HeadEvent { + beacon_node_index: 0, + slot: Slot::new(slot), + beacon_block_root: Hash256::from_low_u64_be(block_root), + } + } + + #[test] + fn delay_stays_attached_to_requested_slot() { + let spec = E::default_spec(); + let slot_clock = + ManualSlotClock::new(Slot::new(0), Duration::ZERO, spec.get_slot_duration()); + let contribution_due = spec.get_contribution_message_due(); + + slot_clock.set_current_time(Duration::from_secs(5)); + assert_eq!( + delay_until_slot_offset(&slot_clock, Slot::new(0), contribution_due), + Some(Duration::from_secs(3)) + ); + + slot_clock.set_current_time(Duration::from_secs(9)); + assert_eq!( + delay_until_slot_offset(&slot_clock, Slot::new(0), contribution_due), + Some(Duration::ZERO) + ); + + slot_clock.set_slot(1); + assert_eq!( + delay_until_slot_offset(&slot_clock, Slot::new(0), contribution_due), + None + ); + } + + #[tokio::test(start_paused = true)] + async fn eager_publish_on_current_slot_head_event() { + let mut harness = TestHarness::new(true).await; + harness.insert_duties().await; + let root_mock = harness + .harness + .mock_beacon_node_1 + .mock_get_head_block_root(Hash256::from_low_u64_be(22)); + let post_mock = harness + .harness + .mock_beacon_node_1 + .mock_post_sync_committee_messages(); + harness.start(); + yield_to_service().await; + + harness.send_head(0, 11); + wait_for_message_count(&harness, 1).await; + + let messages = harness.messages(); + assert_eq!(messages[0].slot, Slot::new(0)); + assert_eq!(messages[0].beacon_block_root, Hash256::from_low_u64_be(11)); + // The head event root is used directly, without fetching a root from the BN. + root_mock.expect(0).assert(); + post_mock.expect(1).assert(); + } + + #[tokio::test(start_paused = true)] + async fn duplicate_head_events_launch_once() { + let mut harness = TestHarness::new(true).await; + harness.insert_duties().await; + let post_mock = harness + .harness + .mock_beacon_node_1 + .mock_post_sync_committee_messages(); + harness.start(); + yield_to_service().await; + + harness.send_head(0, 11); + harness.send_head(0, 12); + wait_for_message_count(&harness, 1).await; + yield_to_service().await; + + let messages = harness.messages(); + assert_eq!(messages.len(), 1); + assert_eq!(messages[0].beacon_block_root, Hash256::from_low_u64_be(11)); + post_mock.expect(1).assert(); + } + + #[tokio::test(start_paused = true)] + async fn duties_available_before_deadline_are_retried_once() { + let mut harness = TestHarness::new(true).await; + let post_mock = harness + .harness + .mock_beacon_node_1 + .mock_post_sync_committee_messages(); + harness.start(); + yield_to_service().await; + + harness.send_head(0, 11); + yield_to_service().await; + harness.insert_duties().await; + assert!(harness.messages().is_empty()); + + harness + .advance_time(Duration::from_secs(4) + Duration::from_millis(1)) + .await; + wait_for_message_count(&harness, 1).await; + + let messages = harness.messages(); + assert_eq!(messages[0].slot, Slot::new(0)); + assert_eq!(messages[0].beacon_block_root, Hash256::from_low_u64_be(11)); + post_mock.expect(1).assert(); + } + + #[tokio::test(start_paused = true)] + async fn unavailable_duties_at_deadline_advance_without_spinning() { + let mut harness = TestHarness::new(true).await; + let expected_root = Hash256::from_low_u64_be(22); + let _root_mock = harness + .harness + .mock_beacon_node_1 + .mock_get_head_block_root(expected_root); + let post_mock = harness + .harness + .mock_beacon_node_1 + .mock_post_sync_committee_messages(); + harness.start(); + yield_to_service().await; + + harness.send_head(0, 11); + yield_to_service().await; + harness + .advance_time(Duration::from_secs(4) + Duration::from_millis(1)) + .await; + assert!(harness.messages().is_empty()); + + harness.insert_duties().await; + harness.advance_time(Duration::from_secs(12)).await; + wait_for_message_count(&harness, 1).await; + + let messages = harness.messages(); + assert_eq!(messages[0].slot, Slot::new(1)); + assert_eq!(messages[0].beacon_block_root, expected_root); + post_mock.expect(1).assert(); + } + + #[tokio::test(start_paused = true)] + async fn timer_fallback_works_without_head_monitoring() { + let mut harness = TestHarness::new(false).await; + harness.insert_duties().await; + let expected_root = Hash256::from_low_u64_be(22); + let _root_mock = harness + .harness + .mock_beacon_node_1 + .mock_get_head_block_root(expected_root); + let post_mock = harness + .harness + .mock_beacon_node_1 + .mock_post_sync_committee_messages(); + harness.start(); + yield_to_service().await; + + harness + .advance_time(Duration::from_secs(16) + Duration::from_millis(1)) + .await; + wait_for_message_count(&harness, 1).await; + + let messages = harness.messages(); + assert_eq!(messages[0].slot, Slot::new(1)); + assert_eq!(messages[0].beacon_block_root, expected_root); + post_mock.expect(1).assert(); + } + + #[tokio::test(start_paused = true)] + async fn closed_head_channel_preserves_timer_fallback() { + let mut harness = TestHarness::new(true).await; + harness.insert_duties().await; + let expected_root = Hash256::from_low_u64_be(22); + let _root_mock = harness + .harness + .mock_beacon_node_1 + .mock_get_head_block_root(expected_root); + let post_mock = harness + .harness + .mock_beacon_node_1 + .mock_post_sync_committee_messages(); + // Drop the sender before the service starts so the channel is closed. + harness.head_sender = None; + harness.start(); + yield_to_service().await; + + harness + .advance_time(Duration::from_secs(16) + Duration::from_millis(1)) + .await; + wait_for_message_count(&harness, 1).await; + + let messages = harness.messages(); + assert_eq!(messages[0].slot, Slot::new(1)); + assert_eq!(messages[0].beacon_block_root, expected_root); + post_mock.expect(1).assert(); + } + + #[tokio::test(start_paused = true)] + async fn timer_deadline_is_fork_aware_at_gloas() { + let mut spec = E::default_spec(); + spec.altair_fork_epoch = Some(Epoch::new(0)); + spec.gloas_fork_epoch = Some(Epoch::new(0)); + let mut harness = TestHarness::new_with_spec(false, Arc::new(spec)).await; + harness.insert_duties().await; + let _root_mock = harness + .harness + .mock_beacon_node_1 + .mock_get_head_block_root(Hash256::from_low_u64_be(22)); + let post_mock = harness + .harness + .mock_beacon_node_1 + .mock_post_sync_committee_messages(); + harness.start(); + yield_to_service().await; + + // The Gloas deadline for slot 1 is 12s + 3s. The pre-Gloas deadline would be + // 12s + 3.999s. + harness.advance_time(Duration::from_millis(14_900)).await; + assert!(harness.messages().is_empty()); + + harness.advance_time(Duration::from_millis(200)).await; + wait_for_message_count(&harness, 1).await; + + let messages = harness.messages(); + assert_eq!(messages[0].slot, Slot::new(1)); + post_mock.expect(1).assert(); + } +} From d698f59bd9221ada73df1c1a21c7d10548cd6407 Mon Sep 17 00:00:00 2001 From: shane-moore Date: Sat, 18 Jul 2026 07:23:13 -0700 Subject: [PATCH 02/15] Harden sync service loop against slot-boundary clock races Three fixes from adversarial review of the eager sync message loop: - Derive next slot and its start offset from a single clock read (next_slot_with_duration), so a slot boundary passing between reads cannot anchor the timer to one slot's start with another slot's due point (previously ~1s early for one slot at the Gloas transition). - Pass the triggering slot into spawn_contribution_tasks instead of re-reading the clock, keeping the signed slot, dedupe marking, and duty lookup consistent with a single slot value. - Discard the head event root when the missing-duties retry sleeps to the deadline, falling back to a fresh non-optimistic head lookup, as the attestation service's deadline retry does. The event root may be stale after sleeping most of a slot. Co-Authored-By: Claude Fable 5 Claude-Session: https://claude.ai/code/session_01MSiv1h1jtjqp8JE2AXbntQ --- .../src/sync_committee_service.rs | 39 +++++++++++++++---- 1 file changed, 31 insertions(+), 8 deletions(-) diff --git a/validator_client/validator_services/src/sync_committee_service.rs b/validator_client/validator_services/src/sync_committee_service.rs index 310f3b5849a..bf9b9a9dddf 100644 --- a/validator_client/validator_services/src/sync_committee_service.rs +++ b/validator_client/validator_services/src/sync_committee_service.rs @@ -40,6 +40,16 @@ fn delay_until_slot_offset( Some(due_at.saturating_sub(now)) } +/// The next slot and the duration until it starts, derived from a single clock read so that a +/// slot boundary passing between reads cannot pair one slot's start with another slot's due +/// point. +fn next_slot_with_duration(slot_clock: &T) -> Option<(Slot, Duration)> { + let now = slot_clock.now_duration()?; + let next_slot = slot_clock.slot_of(now)? + 1; + let duration_to_next_slot = slot_clock.start_of(next_slot)?.saturating_sub(now); + Some((next_slot, duration_to_next_slot)) +} + pub struct SyncCommitteeService { inner: Arc>, } @@ -132,10 +142,8 @@ impl SyncCommitteeService = None; loop { - let Some((duration_to_next_slot, next_slot)) = self - .slot_clock - .duration_to_next_slot() - .zip(self.slot_clock.now().map(|slot| slot + 1)) + let Some((next_slot, duration_to_next_slot)) = + next_slot_with_duration(&self.slot_clock) else { error!("Failed to read slot clock"); sleep(slot_duration).await; @@ -171,7 +179,10 @@ impl SyncCommitteeService { last_sync_message_slot = Some(current_slot); trace!("Spawned sync contribution tasks"); @@ -195,10 +206,10 @@ impl SyncCommitteeService, + slot: Slot, + mut head_event_root: Option, ) -> Result<(), String> { let spec = &self.duties_service.spec; - let slot = self.slot_clock.now().ok_or("Failed to read slot clock")?; let mut slot_duties = self .duties_service @@ -220,6 +231,10 @@ impl SyncCommitteeService(slot, spec); + + // The head may have changed while sleeping, so discard the event root and fall + // back to a fresh head lookup below. + head_event_root = None; } let Some(slot_duties) = slot_duties else { @@ -871,6 +886,13 @@ mod tests { #[tokio::test(start_paused = true)] async fn duties_available_before_deadline_are_retried_once() { let mut harness = TestHarness::new(true).await; + // The head may have changed during the retry sleep, so the retry discards the event + // root (11) and fetches the current head from the beacon node instead. + let expected_root = Hash256::from_low_u64_be(33); + let root_mock = harness + .harness + .mock_beacon_node_1 + .mock_get_head_block_root(expected_root); let post_mock = harness .harness .mock_beacon_node_1 @@ -890,7 +912,8 @@ mod tests { let messages = harness.messages(); assert_eq!(messages[0].slot, Slot::new(0)); - assert_eq!(messages[0].beacon_block_root, Hash256::from_low_u64_be(11)); + assert_eq!(messages[0].beacon_block_root, expected_root); + root_mock.expect(1).assert(); post_mock.expect(1).assert(); } From f3716dee6ef2b0117435baa2af13a58532086a20 Mon Sep 17 00:00:00 2001 From: shane-moore Date: Sat, 18 Jul 2026 07:41:59 -0700 Subject: [PATCH 03/15] Derive the triggered slot from the trigger, not the clock Take the slot for sync message signing from the trigger itself: the validated event slot for the head event arm, or the armed next slot for the timer arm. Previously the loop re-read the clock after triggering, so a slot boundary (or backwards clock movement, which SlotClock permits) between the event's validation and the re-read could sign slot N+1 with slot N's root and, worse, mark N+1 as handled so the real N+1 head event was suppressed by the dedupe check. With the trigger-derived slot there are no post-trigger clock reads: a late event signs for its own slot (or is skipped by the expired-slot guard), and the next slot's messages are never suppressed. Co-Authored-By: Claude Fable 5 Claude-Session: https://claude.ai/code/session_01MSiv1h1jtjqp8JE2AXbntQ --- .../src/sync_committee_service.rs | 16 ++++++++++------ 1 file changed, 10 insertions(+), 6 deletions(-) diff --git a/validator_client/validator_services/src/sync_committee_service.rs b/validator_client/validator_services/src/sync_committee_service.rs index bf9b9a9dddf..43117372b27 100644 --- a/validator_client/validator_services/src/sync_committee_service.rs +++ b/validator_client/validator_services/src/sync_committee_service.rs @@ -156,17 +156,21 @@ impl SyncCommitteeService(next_slot); - let head_event_root = head_event_or_deadline( + let head_event = head_event_or_deadline( &mut head_monitor_rx, &self.slot_clock, duration_to_next_slot + sync_message_due, ) - .await - .map(|event| event.beacon_block_root); + .await; - let Some(current_slot) = self.slot_clock.now() else { - error!("Failed to read slot clock after trigger"); - continue; + // Take the slot from the trigger itself rather than re-reading the clock: the + // validated event slot for the head event arm, or the armed slot for the timer + // arm. If a slot boundary passes during triggering, the event still signs for + // its own slot and can never suppress the next slot's messages via the dedupe + // below. + let (current_slot, head_event_root) = match head_event { + Some(event) => (event.slot, Some(event.beacon_block_root)), + None => (next_slot, None), }; if last_sync_message_slot.is_some_and(|last_slot| current_slot <= last_slot) { From 9f37ed0de173ab14a46c388150783d35a03c477d Mon Sep 17 00:00:00 2001 From: shane-moore Date: Tue, 21 Jul 2026 11:32:54 -0700 Subject: [PATCH 04/15] Group gnosis sync message due keys together Move sync_message_due_bps_gloas next to sync_message_due_bps in the gnosis constructor, which places both pairs of due keys adjacent to their base values. Co-Authored-By: Claude Fable 5 Claude-Session: https://claude.ai/code/session_01MSiv1h1jtjqp8JE2AXbntQ --- consensus/types/src/core/chain_spec.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/consensus/types/src/core/chain_spec.rs b/consensus/types/src/core/chain_spec.rs index 1c2fc0e1878..4ae2b6e20e5 100644 --- a/consensus/types/src/core/chain_spec.rs +++ b/consensus/types/src/core/chain_spec.rs @@ -1584,7 +1584,6 @@ impl ChainSpec { payload_due_bps: 7500, payload_attestation_due_bps: 7500, aggregate_due_bps: 6667, - sync_message_due_bps_gloas: 2500, /* * Derived time values (set by `compute_derived_values()`) @@ -1674,6 +1673,7 @@ impl ChainSpec { altair_fork_version: [0x01, 0x00, 0x00, 0x64], altair_fork_epoch: Some(Epoch::new(512)), sync_message_due_bps: 3333, + sync_message_due_bps_gloas: 2500, contribution_due_bps: 6667, /* From 40db2899e663e1318cdefe1623fbc652e12fb42d Mon Sep 17 00:00:00 2001 From: shane-moore Date: Tue, 21 Jul 2026 12:38:09 -0700 Subject: [PATCH 05/15] Trim comments to constraint statements Co-Authored-By: Claude Fable 5 Claude-Session: https://claude.ai/code/session_01MSiv1h1jtjqp8JE2AXbntQ --- .../validator_services/src/sync_committee_service.rs | 12 ++++-------- 1 file changed, 4 insertions(+), 8 deletions(-) diff --git a/validator_client/validator_services/src/sync_committee_service.rs b/validator_client/validator_services/src/sync_committee_service.rs index 43117372b27..1e48ddb2e07 100644 --- a/validator_client/validator_services/src/sync_committee_service.rs +++ b/validator_client/validator_services/src/sync_committee_service.rs @@ -40,9 +40,8 @@ fn delay_until_slot_offset( Some(due_at.saturating_sub(now)) } -/// The next slot and the duration until it starts, derived from a single clock read so that a -/// slot boundary passing between reads cannot pair one slot's start with another slot's due -/// point. +/// The next slot and the duration until it starts, derived from a single clock read so the +/// two values always describe the same slot. fn next_slot_with_duration(slot_clock: &T) -> Option<(Slot, Duration)> { let now = slot_clock.now_duration()?; let next_slot = slot_clock.slot_of(now)? + 1; @@ -163,11 +162,8 @@ impl SyncCommitteeService (event.slot, Some(event.beacon_block_root)), None => (next_slot, None), From 4c836b9098bb5faf14d89341207a854c992e5411 Mon Sep 17 00:00:00 2001 From: shane-moore Date: Tue, 21 Jul 2026 13:09:41 -0700 Subject: [PATCH 06/15] Rename head monitor tests to match module convention Co-Authored-By: Claude Fable 5 Claude-Session: https://claude.ai/code/session_01MSiv1h1jtjqp8JE2AXbntQ --- .../beacon_node_fallback/src/beacon_head_monitor.rs | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/validator_client/beacon_node_fallback/src/beacon_head_monitor.rs b/validator_client/beacon_node_fallback/src/beacon_head_monitor.rs index 5b16fdaa86d..66f955347da 100644 --- a/validator_client/beacon_node_fallback/src/beacon_head_monitor.rs +++ b/validator_client/beacon_node_fallback/src/beacon_head_monitor.rs @@ -371,7 +371,7 @@ mod tests { } #[tokio::test] - async fn ignores_stale_head_events() { + async fn test_ignores_stale_head_events() { let slot_clock = ManualSlotClock::new(Slot::new(0), Duration::ZERO, Duration::from_secs(12)); slot_clock.set_slot(2); @@ -387,7 +387,7 @@ mod tests { } #[tokio::test] - async fn recovers_from_lagged_head_events() { + async fn test_recovers_from_lagged_head_events() { let slot_clock = ManualSlotClock::new(Slot::new(0), Duration::ZERO, Duration::from_secs(12)); slot_clock.set_slot(2); @@ -403,7 +403,7 @@ mod tests { } #[tokio::test] - async fn returns_when_head_event_channel_closes() { + async fn test_returns_when_head_event_channel_closes() { let slot_clock = ManualSlotClock::new(Slot::new(0), Duration::ZERO, Duration::from_secs(12)); let (sender, mut receiver) = broadcast::channel(1); From 909373e0ba3490b03801134daedf180ef2d8ba74 Mon Sep 17 00:00:00 2001 From: shane-moore Date: Tue, 21 Jul 2026 13:36:40 -0700 Subject: [PATCH 07/15] Simplify head event deadline helper Co-Authored-By: Claude Fable 5 Claude-Session: https://claude.ai/code/session_01MSiv1h1jtjqp8JE2AXbntQ --- .../src/beacon_head_monitor.rs | 19 ++++++++----------- 1 file changed, 8 insertions(+), 11 deletions(-) diff --git a/validator_client/beacon_node_fallback/src/beacon_head_monitor.rs b/validator_client/beacon_node_fallback/src/beacon_head_monitor.rs index 66f955347da..dcb3c86d64c 100644 --- a/validator_client/beacon_node_fallback/src/beacon_head_monitor.rs +++ b/validator_client/beacon_node_fallback/src/beacon_head_monitor.rs @@ -20,7 +20,7 @@ pub struct HeadEvent { pub beacon_block_root: Hash256, } -pub async fn poll_for_current_slot_head( +async fn poll_for_current_slot_head( receiver: &mut broadcast::Receiver, slot_clock: &T, ) -> Option { @@ -60,20 +60,17 @@ pub async fn head_event_or_deadline( tokio::pin!(deadline); if let Some(receiver) = head_monitor_rx { tokio::select! { - _ = &mut deadline => None, - event = poll_for_current_slot_head(receiver, slot_clock) => match event { - Some(event) => Some(event), - None => { - *head_monitor_rx = None; - deadline.await; - None + _ = &mut deadline => return None, + event = poll_for_current_slot_head(receiver, slot_clock) => { + if event.is_some() { + return event; } + *head_monitor_rx = None; }, } - } else { - deadline.await; - None } + deadline.await; + None } /// Cache to maintain the latest head received from each of the beacon nodes From e1b9bb4fa29daacf2e52155262162098a478dc72 Mon Sep 17 00:00:00 2001 From: shane-moore Date: Tue, 4 Aug 2026 13:29:45 -0700 Subject: [PATCH 08/15] Remove dead sync committee task result --- .../src/sync_committee_service.rs | 40 ++++++------------- 1 file changed, 12 insertions(+), 28 deletions(-) diff --git a/validator_client/validator_services/src/sync_committee_service.rs b/validator_client/validator_services/src/sync_committee_service.rs index 62ba1467cb2..502d2f2fb1f 100644 --- a/validator_client/validator_services/src/sync_committee_service.rs +++ b/validator_client/validator_services/src/sync_committee_service.rs @@ -139,7 +139,7 @@ impl SyncCommitteeService = None; + let mut last_processed_slot: Option = None; loop { let Some((next_slot, duration_to_next_slot)) = next_slot_with_duration(&self.slot_clock) @@ -169,8 +169,8 @@ impl SyncCommitteeService (next_slot, None), }; - if last_sync_message_slot.is_some_and(|last_slot| current_slot <= last_slot) { - debug!(%current_slot, "Sync messages already initiated for the slot"); + if last_processed_slot.is_some_and(|last_slot| current_slot <= last_slot) { + debug!(%current_slot, "Sync message slot already processed"); continue; } @@ -179,21 +179,9 @@ impl SyncCommitteeService { - last_sync_message_slot = Some(current_slot); - trace!("Spawned sync contribution tasks"); - } - Err(e) => { - crit!( - error = ?e, - "Failed to spawn sync contribution tasks" - ); - } - } + self.spawn_contribution_tasks(current_slot, head_event_root) + .await; + last_processed_slot = Some(current_slot); // Do subscriptions for future slots/epochs. self.spawn_subscription_tasks(); @@ -204,11 +192,7 @@ impl SyncCommitteeService, - ) -> Result<(), String> { + async fn spawn_contribution_tasks(&self, slot: Slot, mut head_event_root: Option) { let spec = &self.duties_service.spec; let mut slot_duties = self @@ -239,12 +223,12 @@ impl SyncCommitteeService SyncCommitteeService SyncCommitteeService SyncCommitteeService Date: Tue, 4 Aug 2026 15:08:20 -0700 Subject: [PATCH 09/15] Relax sync committee test timeout --- .../validator_services/src/sync_committee_service.rs | 11 +++++------ 1 file changed, 5 insertions(+), 6 deletions(-) diff --git a/validator_client/validator_services/src/sync_committee_service.rs b/validator_client/validator_services/src/sync_committee_service.rs index 502d2f2fb1f..82e0920494c 100644 --- a/validator_client/validator_services/src/sync_committee_service.rs +++ b/validator_client/validator_services/src/sync_committee_service.rs @@ -775,13 +775,12 @@ mod tests { /// Resume real time so the service can complete signing and HTTP requests, then pause again. async fn wait_for_message_count(harness: &TestHarness, count: usize) { tokio::time::resume(); - for _ in 0..100 { - if harness.messages().len() >= count { - tokio::time::pause(); - return; - } - tokio::time::sleep(Duration::from_millis(1)).await; + let deadline = Instant::now() + Duration::from_secs(5); + + while harness.messages().len() < count && Instant::now() < deadline { + tokio::time::sleep(Duration::from_millis(10)).await; } + tokio::time::pause(); assert_eq!(harness.messages().len(), count); } From df5a158ce8d1258fefdecdb544b7ce3d39d4ffb4 Mon Sep 17 00:00:00 2001 From: shane-moore Date: Mon, 10 Aug 2026 11:24:48 -0700 Subject: [PATCH 10/15] Skip sync committee tasks past the contribution deadline --- .../src/sync_committee_service.rs | 41 +++++++++++++++++++ 1 file changed, 41 insertions(+) diff --git a/validator_client/validator_services/src/sync_committee_service.rs b/validator_client/validator_services/src/sync_committee_service.rs index 82e0920494c..27d445c2847 100644 --- a/validator_client/validator_services/src/sync_committee_service.rs +++ b/validator_client/validator_services/src/sync_committee_service.rs @@ -239,6 +239,12 @@ impl SyncCommitteeService Date: Mon, 10 Aug 2026 12:01:08 -0700 Subject: [PATCH 11/15] Tidy sync committee skip paths --- .../validator_services/src/sync_committee_service.rs | 10 ++++++---- 1 file changed, 6 insertions(+), 4 deletions(-) diff --git a/validator_client/validator_services/src/sync_committee_service.rs b/validator_client/validator_services/src/sync_committee_service.rs index 27d445c2847..5bee30d523e 100644 --- a/validator_client/validator_services/src/sync_committee_service.rs +++ b/validator_client/validator_services/src/sync_committee_service.rs @@ -203,12 +203,14 @@ impl SyncCommitteeService(slot), - ) - .unwrap_or(Duration::ZERO); + ) else { + debug!(%slot, "Skipping sync committee tasks for expired slot"); + return; + }; sleep(duration_to_deadline).await; slot_duties = self @@ -222,7 +224,7 @@ impl SyncCommitteeService Date: Thu, 13 Aug 2026 13:06:30 -0700 Subject: [PATCH 12/15] Treat clock read failure as terminal in head monitor --- .../src/beacon_head_monitor.rs | 26 ++++++++++++++++--- 1 file changed, 23 insertions(+), 3 deletions(-) diff --git a/validator_client/beacon_node_fallback/src/beacon_head_monitor.rs b/validator_client/beacon_node_fallback/src/beacon_head_monitor.rs index dcb3c86d64c..d70f998b357 100644 --- a/validator_client/beacon_node_fallback/src/beacon_head_monitor.rs +++ b/validator_client/beacon_node_fallback/src/beacon_head_monitor.rs @@ -27,10 +27,11 @@ async fn poll_for_current_slot_head( loop { match receiver.recv().await { Ok(head_event) => { - // Skip the event on a clock read failure rather than returning `None`, - // which callers treat as a terminal error that disables head monitoring. + // A clock read only fails when the system clock is broken, so treat it like + // a closed channel and disable head monitoring. let Some(current_slot) = slot_clock.now() else { - continue; + warn!("Failed to read slot clock while polling head events"); + return None; }; if head_event.slot == current_slot { return Some(head_event); @@ -399,6 +400,25 @@ mod tests { assert_eq!(event.beacon_block_root, Hash256::from_low_u64_be(2)); } + #[tokio::test] + async fn test_clock_read_failure_is_terminal() { + let slot_clock = ManualSlotClock::new( + Slot::new(0), + Duration::from_secs(10), + Duration::from_secs(12), + ); + // A time before genesis makes clock reads fail. + slot_clock.set_current_time(Duration::from_secs(5)); + let (sender, mut receiver) = broadcast::channel(1); + sender.send(head_event(0, 1)).unwrap(); + + assert!( + poll_for_current_slot_head(&mut receiver, &slot_clock) + .await + .is_none() + ); + } + #[tokio::test] async fn test_returns_when_head_event_channel_closes() { let slot_clock = From 367bb6fa561302d43225efd6d8be88959eae7c52 Mon Sep 17 00:00:00 2001 From: shane-moore Date: Mon, 17 Aug 2026 11:33:51 -0700 Subject: [PATCH 13/15] Clarify clock failure comment in head monitor --- .../beacon_node_fallback/src/beacon_head_monitor.rs | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/validator_client/beacon_node_fallback/src/beacon_head_monitor.rs b/validator_client/beacon_node_fallback/src/beacon_head_monitor.rs index d70f998b357..82b5e978588 100644 --- a/validator_client/beacon_node_fallback/src/beacon_head_monitor.rs +++ b/validator_client/beacon_node_fallback/src/beacon_head_monitor.rs @@ -27,8 +27,9 @@ async fn poll_for_current_slot_head( loop { match receiver.recv().await { Ok(head_event) => { - // A clock read only fails when the system clock is broken, so treat it like - // a closed channel and disable head monitoring. + // A clock read only fails pre-genesis (services start after the genesis + // wait) or when the system clock is broken, so treat it like a closed + // channel and disable head monitoring. let Some(current_slot) = slot_clock.now() else { warn!("Failed to read slot clock while polling head events"); return None; From cf939e139313ee888aad16781f9d80de6326d2e8 Mon Sep 17 00:00:00 2001 From: shane-moore Date: Mon, 24 Aug 2026 08:01:15 -0700 Subject: [PATCH 14/15] Use upstream sync_message_deadline helper in service loop Replaces next_slot_with_duration with the sync_message_deadline helper that landed on unstable, restoring its unit test and the defensive pre-genesis fallback, and mirroring attestation_deadline in the attestation service. Co-Authored-By: Claude Fable 5 Claude-Session: https://claude.ai/code/session_012YyYbMaMeNuxVq3Bv4a4FA --- .../src/sync_committee_service.rs | 85 +++++++++++++++---- 1 file changed, 68 insertions(+), 17 deletions(-) diff --git a/validator_client/validator_services/src/sync_committee_service.rs b/validator_client/validator_services/src/sync_committee_service.rs index 8f024eb863a..e66b3577220 100644 --- a/validator_client/validator_services/src/sync_committee_service.rs +++ b/validator_client/validator_services/src/sync_committee_service.rs @@ -40,15 +40,6 @@ fn delay_until_slot_offset( Some(due_at.saturating_sub(now)) } -/// The next slot and the duration until it starts, derived from a single clock read so the -/// two values always describe the same slot. -fn next_slot_with_duration(slot_clock: &T) -> Option<(Slot, Duration)> { - let now = slot_clock.now_duration()?; - let next_slot = slot_clock.slot_of(now)? + 1; - let duration_to_next_slot = slot_clock.start_of(next_slot)?.saturating_sub(now); - Some((next_slot, duration_to_next_slot)) -} - pub struct SyncCommitteeService { inner: Arc>, } @@ -141,24 +132,25 @@ impl SyncCommitteeService = None; loop { - let Some((next_slot, duration_to_next_slot)) = - next_slot_with_duration(&self.slot_clock) - else { + let Some(now) = self.slot_clock.now_duration() else { error!("Failed to read slot clock"); sleep(slot_duration).await; continue; }; + let (next_slot, Some(duration_to_sync_message_deadline)) = + sync_message_deadline::(&self.slot_clock, &self.duties_service.spec, now) + else { + error!("Failed to determine sync message deadline"); + sleep(slot_duration).await; + continue; + }; // Wait for the sync message due point of the next slot, or a head event for the // current slot, whichever comes first. - let sync_message_due = self - .duties_service - .spec - .get_sync_message_due::(next_slot); let head_event = head_event_or_deadline( &mut head_monitor_rx, &self.slot_clock, - duration_to_next_slot + sync_message_due, + duration_to_sync_message_deadline, ) .await; @@ -639,6 +631,23 @@ impl SyncCommitteeService( + slot_clock: &impl SlotClock, + chain_spec: &ChainSpec, + now: Duration, +) -> (Slot, Option) { + let sync_message_slot = slot_clock + .slot_of(now) + .map_or_else(|| slot_clock.genesis_slot(), |slot| slot + 1); + let duration_to_sync_message_deadline = slot_clock + .start_of(sync_message_slot) + .and_then(|slot_start| { + slot_start.checked_add(chain_spec.get_sync_message_due::(sync_message_slot)) + }) + .and_then(|deadline| deadline.checked_sub(now)); + (sync_message_slot, duration_to_sync_message_deadline) +} + fn sync_period_of_slot(slot: Slot, spec: &ChainSpec) -> Result { slot.epoch(E::slots_per_epoch()) .sync_committee_period(spec) @@ -803,6 +812,48 @@ mod tests { } } + #[test] + fn duration_to_sync_message_deadline_is_fork_aware() { + let mut spec = E::default_spec(); + let gloas_fork_epoch = Epoch::new(1); + spec.gloas_fork_epoch = Some(gloas_fork_epoch); + + let slot_duration = spec.get_slot_duration(); + let genesis_time = slot_duration; + let slot_clock = ManualSlotClock::new(Slot::new(0), genesis_time, slot_duration); + let first_gloas_slot = gloas_fork_epoch.start_slot(E::slots_per_epoch()); + let last_pre_gloas_slot = first_gloas_slot - 1; + + let test_cases = [ + ( + "pre-genesis", + genesis_time - Duration::from_secs(1), + slot_clock.genesis_slot(), + Duration::from_millis(4999), + ), + ( + "pre-Gloas", + slot_clock.start_of(last_pre_gloas_slot - 1).unwrap(), + last_pre_gloas_slot, + Duration::from_millis(15999), + ), + ( + "post-Gloas", + slot_clock.start_of(last_pre_gloas_slot).unwrap(), + first_gloas_slot, + Duration::from_millis(15000), + ), + ]; + + for (case, now, expected_slot, expected_duration) in test_cases { + assert_eq!( + sync_message_deadline::(&slot_clock, &spec, now), + (expected_slot, Some(expected_duration)), + "{case}" + ); + } + } + #[test] fn delay_stays_attached_to_requested_slot() { let spec = E::default_spec(); From 5e6c06c2d3a58457809282a7203eb6d56b444b27 Mon Sep 17 00:00:00 2001 From: shane-moore Date: Mon, 28 Sep 2026 09:34:04 -0700 Subject: [PATCH 15/15] Read sync aggregators at the contribution deadline Selection proofs can be stored after a head event triggers the slot, for example with `--distributed`, where the proof for a slot is only computed once that slot has started. Reading the aggregators at trigger time then skipped the contribution. Read them at the contribution deadline instead. Add tests for the aggregate path, a head event arriving after the timer handled the slot, and `head_event_or_deadline`. Co-Authored-By: Claude Fable 5.1 --- .../src/mock_beacon_node.rs | 55 ++++- .../beacon_node_fallback/Cargo.toml | 1 + .../src/beacon_head_monitor.rs | 53 +++++ .../validator_services/src/sync.rs | 26 +++ .../src/sync_committee_service.rs | 196 ++++++++++++++++-- 5 files changed, 315 insertions(+), 16 deletions(-) diff --git a/testing/validator_test_rig/src/mock_beacon_node.rs b/testing/validator_test_rig/src/mock_beacon_node.rs index 135c4fdfafe..b3fefc161e6 100644 --- a/testing/validator_test_rig/src/mock_beacon_node.rs +++ b/testing/validator_test_rig/src/mock_beacon_node.rs @@ -16,7 +16,8 @@ use tracing::info; use types::{ ChainSpec, ConfigAndPreset, Epoch, EthSpec, ExecutionPayloadEnvelope, ForkName, Hash256, PayloadAttestationData, PayloadAttestationMessage, SignedBlindedBeaconBlock, - SignedExecutionPayloadEnvelope, Slot, SyncCommitteeMessage, SyncDuty, + SignedContributionAndProof, SignedExecutionPayloadEnvelope, Slot, SyncCommitteeContribution, + SyncCommitteeMessage, SyncDuty, }; pub struct MockBeaconNode { @@ -31,6 +32,7 @@ pub struct MockBeaconNode { pub payload_attestation_message: Arc>>, pub builder_preferences: Arc>>, pub sync_committee_messages: Arc>>, + pub sync_committee_contributions: Arc>>>, } impl MockBeaconNode { @@ -52,6 +54,7 @@ impl MockBeaconNode { payload_attestation_message: Arc::new(Mutex::new(Vec::new())), builder_preferences: Arc::new(Mutex::new(Vec::new())), sync_committee_messages: Arc::new(Mutex::new(Vec::new())), + sync_committee_contributions: Arc::new(Mutex::new(Vec::new())), } } @@ -162,6 +165,56 @@ impl MockBeaconNode { .create() } + /// Mocks `GET /eth/v1/validator/sync_committee_contribution`, matching the slot, block root + /// and subcommittee index of `contribution`. + pub fn mock_get_sync_committee_contribution( + &mut self, + contribution: &SyncCommitteeContribution, + ) -> Mock { + let path_pattern = Regex::new(r"^/eth/v1/validator/sync_committee_contribution$").unwrap(); + let response = GenericResponse::from(contribution.clone()); + + self.server + .mock("GET", Matcher::Regex(path_pattern.to_string())) + .match_query(Matcher::AllOf(vec![ + Matcher::UrlEncoded("slot".into(), contribution.slot.to_string()), + Matcher::UrlEncoded( + "beacon_block_root".into(), + format!("{:?}", contribution.beacon_block_root), + ), + Matcher::UrlEncoded( + "subcommittee_index".into(), + contribution.subcommittee_index.to_string(), + ), + ])) + .with_status(200) + .with_header("content-type", "application/json") + .with_body(serde_json::to_string(&response).unwrap()) + .create() + } + + /// Mocks `POST /eth/v1/validator/contribution_and_proofs` + pub fn mock_post_contribution_and_proofs(&mut self) -> Mock { + let path_pattern = Regex::new(r"^/eth/v1/validator/contribution_and_proofs$").unwrap(); + let sync_committee_contributions = Arc::clone(&self.sync_committee_contributions); + + self.server + .mock("POST", Matcher::Regex(path_pattern.to_string())) + .with_status(200) + .with_body_from_request(move |request| { + let body = request.body().expect("Failed to get request body"); + let contributions: Vec> = + serde_json::from_slice(body) + .expect("Failed to deserialize sync committee contributions"); + sync_committee_contributions + .lock() + .unwrap() + .extend(contributions); + vec![] + }) + .create() + } + /// Mocks `POST /eth/v1/validator/sync_committee_subscriptions` pub fn mock_sync_committee_subscriptions(&mut self) -> Mock { let path_pattern = Regex::new(r"^/eth/v1/validator/sync_committee_subscriptions$").unwrap(); diff --git a/validator_client/beacon_node_fallback/Cargo.toml b/validator_client/beacon_node_fallback/Cargo.toml index bc1ac20d44c..9410edac5cf 100644 --- a/validator_client/beacon_node_fallback/Cargo.toml +++ b/validator_client/beacon_node_fallback/Cargo.toml @@ -25,4 +25,5 @@ types = { workspace = true } validator_metrics = { workspace = true } [dev-dependencies] +tokio = { workspace = true, features = ["test-util"] } validator_test_rig = { workspace = true } diff --git a/validator_client/beacon_node_fallback/src/beacon_head_monitor.rs b/validator_client/beacon_node_fallback/src/beacon_head_monitor.rs index 1eeeb90befe..d093813964f 100644 --- a/validator_client/beacon_node_fallback/src/beacon_head_monitor.rs +++ b/validator_client/beacon_node_fallback/src/beacon_head_monitor.rs @@ -326,6 +326,7 @@ mod tests { use bls::FixedBytesExtended; use slot_clock::ManualSlotClock; use std::time::Duration; + use tokio::time::Instant; use types::{Hash256, Slot}; fn create_sse_head(slot: u64, block_root: u8) -> SseHead { @@ -517,6 +518,58 @@ mod tests { ); } + #[tokio::test(start_paused = true)] + async fn test_head_event_returned_before_deadline() { + let slot_clock = + ManualSlotClock::new(Slot::new(0), Duration::ZERO, Duration::from_secs(12)); + let (sender, receiver) = broadcast::channel(1); + let mut head_monitor_rx = Some(receiver); + sender.send(head_event(0, 1)).unwrap(); + + let start = Instant::now(); + let event = + head_event_or_deadline(&mut head_monitor_rx, &slot_clock, Duration::from_secs(4)) + .await + .unwrap(); + + assert_eq!(event.beacon_block_root, Hash256::from_low_u64_be(1)); + assert_eq!(start.elapsed(), Duration::ZERO); + assert!(head_monitor_rx.is_some()); + } + + #[tokio::test(start_paused = true)] + async fn test_deadline_keeps_head_monitoring_enabled() { + let slot_clock = + ManualSlotClock::new(Slot::new(0), Duration::ZERO, Duration::from_secs(12)); + let (_sender, receiver) = broadcast::channel::(1); + let mut head_monitor_rx = Some(receiver); + + let start = Instant::now(); + let event = + head_event_or_deadline(&mut head_monitor_rx, &slot_clock, Duration::from_secs(4)).await; + + assert!(event.is_none()); + assert_eq!(start.elapsed(), Duration::from_secs(4)); + assert!(head_monitor_rx.is_some()); + } + + #[tokio::test(start_paused = true)] + async fn test_closed_channel_disables_head_monitoring_and_awaits_deadline() { + let slot_clock = + ManualSlotClock::new(Slot::new(0), Duration::ZERO, Duration::from_secs(12)); + let (sender, receiver) = broadcast::channel::(1); + let mut head_monitor_rx = Some(receiver); + drop(sender); + + let start = Instant::now(); + let event = + head_event_or_deadline(&mut head_monitor_rx, &slot_clock, Duration::from_secs(4)).await; + + assert!(event.is_none()); + assert_eq!(start.elapsed(), Duration::from_secs(4)); + assert!(head_monitor_rx.is_none()); + } + #[tokio::test] async fn test_cache_caches_multiple_heads_from_different_nodes() { let cache = BeaconHeadCache::new(); diff --git a/validator_client/validator_services/src/sync.rs b/validator_client/validator_services/src/sync.rs index 0f456a70507..4ab5523bb7d 100644 --- a/validator_client/validator_services/src/sync.rs +++ b/validator_client/validator_services/src/sync.rs @@ -225,6 +225,32 @@ impl SyncDutiesMap { }) } + /// Store a selection proof for `validator_index`, as `fill_in_aggregation_proofs` does. + #[cfg(test)] + pub(crate) fn insert_proof( + &self, + committee_period: u64, + validator_index: u64, + slot: Slot, + subnet_id: SyncSubnetId, + proof: SyncSelectionProof, + ) { + let committees = self.committees.read(); + let validators = committees + .get(&committee_period) + .expect("duties should exist for period") + .validators + .read(); + validators + .get(&validator_index) + .and_then(Option::as_ref) + .expect("validator should have a sync duty") + .aggregation_duties + .proofs + .write() + .insert((slot, subnet_id), proof); + } + /// Prune duties for past sync committee periods from the map. fn prune(&self, current_sync_committee_period: u64) { self.committees diff --git a/validator_client/validator_services/src/sync_committee_service.rs b/validator_client/validator_services/src/sync_committee_service.rs index af536b79daf..e5d1d8f19e1 100644 --- a/validator_client/validator_services/src/sync_committee_service.rs +++ b/validator_client/validator_services/src/sync_committee_service.rs @@ -9,7 +9,6 @@ use futures::StreamExt; use futures::future::FutureExt; use logging::crit; use slot_clock::SlotClock; -use std::collections::HashMap; use std::ops::Deref; use std::sync::Arc; use std::sync::atomic::{AtomicBool, Ordering}; @@ -145,7 +144,7 @@ impl SyncCommitteeService SyncCommitteeService SyncCommitteeService SyncCommitteeService>, aggregate_instant: Instant, ) { - for (subnet_id, subnet_aggregators) in aggregators { + sleep_until(aggregate_instant).await; + + // Read the aggregators now rather than when the slot was triggered, since selection + // proofs can be stored after the trigger (e.g. with `--distributed`, where the proof + // for a slot is only computed once that slot has started). + let Some(slot_duties) = self + .duties_service + .sync_duties + .get_duties_for_slot::(slot, &self.duties_service.spec) + else { + debug!(%slot, "No duties known for slot at contribution deadline"); + return; + }; + + for (subnet_id, subnet_aggregators) in slot_duties.aggregators { let service = self.clone(); self.inner.executor.spawn( async move { @@ -402,7 +412,6 @@ impl SyncCommitteeService SyncCommitteeService, - aggregate_instant: Instant, ) -> Result<(), ()> { - sleep_until(aggregate_instant).await; - let contribution = &self .beacon_nodes .first_success(|beacon_node| async move { @@ -677,7 +683,7 @@ mod tests { }; use bls::FixedBytesExtended; use slot_clock::ManualSlotClock; - use types::{Epoch, MainnetEthSpec}; + use types::{Epoch, MainnetEthSpec, SignedContributionAndProof, SyncCommitteeContribution}; use validator_test_rig::validator_client_harness::{S, ValidatorClientHarness}; type E = MainnetEthSpec; @@ -755,6 +761,23 @@ mod tests { tokio::time::pause(); } + /// Store a selection proof making the validator an aggregator for `slot`. + async fn insert_selection_proof(&self, slot: Slot) { + tokio::time::resume(); + let subnet_id = SyncSubnetId::new(0); + let proof = self + .harness + .validator_store + .produce_sync_selection_proof(&self.harness.pubkeys[0], slot, subnet_id) + .await + .unwrap(); + tokio::time::pause(); + self.service + .duties_service + .sync_duties + .insert_proof(0, 0, slot, subnet_id, proof); + } + fn start(&self) { self.service .clone() @@ -784,6 +807,20 @@ mod tests { .unwrap() .clone() } + + fn contributions(&self) -> Vec> { + self.harness + .mock_beacon_node_1 + .sync_committee_contributions + .lock() + .unwrap() + .clone() + } + + /// The contribution a beacon node would build from the first published message. + fn contribution(&self) -> SyncCommitteeContribution { + SyncCommitteeContribution::from_message(&self.messages()[0], 0, 0).unwrap() + } } async fn yield_to_service() { @@ -793,16 +830,24 @@ mod tests { } /// Resume real time so the service can complete signing and HTTP requests, then pause again. - async fn wait_for_message_count(harness: &TestHarness, count: usize) { + async fn wait_for_count(len: impl Fn() -> usize, count: usize) { tokio::time::resume(); let deadline = Instant::now() + Duration::from_secs(5); - while harness.messages().len() < count && Instant::now() < deadline { + while len() < count && Instant::now() < deadline { tokio::time::sleep(Duration::from_millis(10)).await; } tokio::time::pause(); - assert_eq!(harness.messages().len(), count); + assert_eq!(len(), count); + } + + async fn wait_for_message_count(harness: &TestHarness, count: usize) { + wait_for_count(|| harness.messages().len(), count).await; + } + + async fn wait_for_contribution_count(harness: &TestHarness, count: usize) { + wait_for_count(|| harness.contributions().len(), count).await; } fn head_event(slot: u64, block_root: u64) -> HeadEvent { @@ -1018,7 +1063,7 @@ mod tests { yield_to_service().await; assert!(harness.messages().is_empty()); - // The skip is latched and the timer still covers the next slot at its due point. + // The skip is latched and the timer still covers the next slot at its deadline. harness .advance_time(Duration::from_secs(7) + Duration::from_millis(1)) .await; @@ -1086,6 +1131,127 @@ mod tests { post_mock.expect(1).assert(); } + #[tokio::test(start_paused = true)] + async fn head_event_after_timer_handled_slot_is_ignored() { + let mut harness = TestHarness::new(true).await; + harness.insert_duties().await; + let expected_root = Hash256::from_low_u64_be(22); + let _root_mock = harness + .harness + .mock_beacon_node_1 + .mock_get_head_block_root(expected_root); + let post_mock = harness + .harness + .mock_beacon_node_1 + .mock_post_sync_committee_messages(); + harness.start(); + yield_to_service().await; + + // The timer handles slot 1 at its sync message deadline (12s + 4s). + harness + .advance_time(Duration::from_secs(16) + Duration::from_millis(1)) + .await; + wait_for_message_count(&harness, 1).await; + + // A head event for slot 1 arriving after the deadline must not sign again. + harness.advance_time(Duration::from_secs(1)).await; + harness.send_head(1, 11); + yield_to_service().await; + + // The next messages are those of the timer for slot 2. + harness.advance_time(Duration::from_secs(11)).await; + wait_for_message_count(&harness, 2).await; + + let messages = harness.messages(); + assert_eq!(messages[0].slot, Slot::new(1)); + assert_eq!(messages[1].slot, Slot::new(2)); + assert_eq!(messages[0].beacon_block_root, expected_root); + assert_eq!(messages[1].beacon_block_root, expected_root); + post_mock.expect(2).assert(); + } + + #[tokio::test(start_paused = true)] + async fn contribution_published_at_contribution_deadline() { + let mut harness = TestHarness::new(true).await; + harness.insert_duties().await; + harness.insert_selection_proof(Slot::new(0)).await; + let _post_mock = harness + .harness + .mock_beacon_node_1 + .mock_post_sync_committee_messages(); + harness.start(); + yield_to_service().await; + + harness.send_head(0, 11); + wait_for_message_count(&harness, 1).await; + let contribution = harness.contribution(); + let _get_contribution_mock = harness + .harness + .mock_beacon_node_1 + .mock_get_sync_committee_contribution(&contribution); + let post_contribution_mock = harness + .harness + .mock_beacon_node_1 + .mock_post_contribution_and_proofs(); + + // The pre-Gloas contribution deadline for slot 0 is 8s. + harness.advance_time(Duration::from_secs(6)).await; + assert!(harness.contributions().is_empty()); + + harness + .advance_time(Duration::from_secs(2) + Duration::from_millis(1)) + .await; + wait_for_contribution_count(&harness, 1).await; + + let contribution = &harness.contributions()[0].message; + assert_eq!(contribution.aggregator_index, 0); + assert_eq!(contribution.contribution.slot, Slot::new(0)); + assert_eq!( + contribution.contribution.beacon_block_root, + Hash256::from_low_u64_be(11) + ); + post_contribution_mock.expect(1).assert(); + } + + #[tokio::test(start_paused = true)] + async fn contribution_uses_selection_proof_stored_after_head_event() { + let mut harness = TestHarness::new(true).await; + harness.insert_duties().await; + let _post_mock = harness + .harness + .mock_beacon_node_1 + .mock_post_sync_committee_messages(); + harness.start(); + yield_to_service().await; + + // The head event arrives before the selection proof for the slot is stored, as happens + // when proofs are only computed once the slot has started (e.g. `--distributed`). + harness.send_head(0, 11); + wait_for_message_count(&harness, 1).await; + let contribution = harness.contribution(); + let _get_contribution_mock = harness + .harness + .mock_beacon_node_1 + .mock_get_sync_committee_contribution(&contribution); + let post_contribution_mock = harness + .harness + .mock_beacon_node_1 + .mock_post_contribution_and_proofs(); + + harness.advance_time(Duration::from_secs(2)).await; + harness.insert_selection_proof(Slot::new(0)).await; + + harness + .advance_time(Duration::from_secs(6) + Duration::from_millis(1)) + .await; + wait_for_contribution_count(&harness, 1).await; + + let contribution = &harness.contributions()[0].message; + assert_eq!(contribution.aggregator_index, 0); + assert_eq!(contribution.contribution.slot, Slot::new(0)); + post_contribution_mock.expect(1).assert(); + } + #[tokio::test(start_paused = true)] async fn timer_deadline_is_fork_aware_at_gloas() { let mut spec = E::default_spec();