diff --git a/book/src/help_vc.md b/book/src/help_vc.md index fed2a1c48e5..e4d8efd9173 100644 --- a/book/src/help_vc.md +++ b/book/src/help_vc.md @@ -187,11 +187,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/testing/validator_test_rig/src/mock_beacon_node.rs b/testing/validator_test_rig/src/mock_beacon_node.rs index c1be2f73940..b3fefc161e6 100644 --- a/testing/validator_test_rig/src/mock_beacon_node.rs +++ b/testing/validator_test_rig/src/mock_beacon_node.rs @@ -1,5 +1,5 @@ use eth2::types::{ - GenericResponse, ProduceBlockV4Response, PublishBlockRequest, + GenericResponse, ProduceBlockV4Response, PublishBlockRequest, RootData, SignedExecutionPayloadEnvelopeContents, SubmittedBuilderPreferences, SyncingData, }; use eth2::{BLOB_DATA_INCLUDED_HEADER, BeaconNodeHttpClient, CONSENSUS_VERSION_HEADER, Timeouts}; @@ -14,9 +14,10 @@ use std::sync::{Arc, Mutex}; use std::time::Duration; use tracing::info; use types::{ - ChainSpec, ConfigAndPreset, EthSpec, ExecutionPayloadEnvelope, ForkName, Hash256, + ChainSpec, ConfigAndPreset, Epoch, EthSpec, ExecutionPayloadEnvelope, ForkName, Hash256, PayloadAttestationData, PayloadAttestationMessage, SignedBlindedBeaconBlock, - SignedExecutionPayloadEnvelope, Slot, + SignedContributionAndProof, SignedExecutionPayloadEnvelope, Slot, SyncCommitteeContribution, + SyncCommitteeMessage, SyncDuty, }; pub struct MockBeaconNode { @@ -30,6 +31,8 @@ pub struct MockBeaconNode { Arc>>>, pub payload_attestation_message: Arc>>, pub builder_preferences: Arc>>, + pub sync_committee_messages: Arc>>, + pub sync_committee_contributions: Arc>>>, } impl MockBeaconNode { @@ -50,6 +53,8 @@ impl MockBeaconNode { execution_payload_envelope_contents: Arc::new(Mutex::new(Vec::new())), 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())), } } @@ -110,6 +115,116 @@ 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 `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(); + + self.server + .mock("POST", Matcher::Regex(path_pattern.to_string())) + .with_status(200) + .create() + } + /// Mocks `POST /eth/v4/validator/blocks/{slot}`, matching the given `include_payload` query /// value and answering with `response`. pub fn mock_post_validator_blocks_v4( diff --git a/testing/validator_test_rig/src/validator_client_harness.rs b/testing/validator_test_rig/src/validator_client_harness.rs index 90f202edd2c..64e87c9447e 100644 --- a/testing/validator_test_rig/src/validator_client_harness.rs +++ b/testing/validator_test_rig/src/validator_client_harness.rs @@ -51,7 +51,6 @@ impl ValidatorClientHarness { config: &ValidatorStoreConfig, ) -> Self { let spec = Arc::new(spec); - 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/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 be05a7612b1..d093813964f 100644 --- a/validator_client/beacon_node_fallback/src/beacon_head_monitor.rs +++ b/validator_client/beacon_node_fallback/src/beacon_head_monitor.rs @@ -4,15 +4,16 @@ 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, @@ -27,6 +28,61 @@ pub struct PayloadAvailableEvent { pub block_root: Hash256, } +async fn poll_for_current_slot_head( + receiver: &mut broadcast::Receiver, + slot_clock: &T, +) -> Option { + loop { + match receiver.recv().await { + Ok(head_event) => { + // 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; + }; + 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 => return None, + event = poll_for_current_slot_head(receiver, slot_clock) => { + if event.is_some() { + return event; + } + *head_monitor_rx = None; + }, + } + } + deadline.await; + None +} + /// Cache to maintain the latest head received from each of the beacon nodes /// in the `BeaconNodeFallback`. #[derive(Debug)] @@ -76,16 +132,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> { @@ -175,7 +222,6 @@ pub async fn poll_head_event_from_beacon_nodes SseHead { @@ -396,6 +445,131 @@ 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 test_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 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); + 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 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 = + 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(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/beacon_node_fallback/src/lib.rs b/validator_client/beacon_node_fallback/src/lib.rs index 0866a3bc02a..a9405decdd1 100644 --- a/validator_client/beacon_node_fallback/src/lib.rs +++ b/validator_client/beacon_node_fallback/src/lib.rs @@ -29,7 +29,7 @@ use std::vec::Vec; use strum::VariantNames; use task_executor::TaskExecutor; use tokio::{ - sync::{RwLock, mpsc}, + sync::{RwLock, broadcast, mpsc}, time::sleep, }; use tracing::{debug, error, warn}; @@ -462,7 +462,7 @@ pub struct BeaconNodeFallback { distance_tiers: BeaconNodeSyncDistanceTiers, slot_clock: Option, beacon_head_cache: Option>, - head_monitor_send: Option>>, + head_monitor_send: Option>, payload_available_send: Option>>, broadcast_topics: Vec, spec: Arc, @@ -501,11 +501,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()) + } + /// This is the payload monitor channel that streams events from all the beacon nodes that the /// validator client is connected to in the `BeaconNodeFallback`. pub fn set_payload_available_send( diff --git a/validator_client/src/cli.rs b/validator_client/src/cli.rs index a721f2e679d..47709886c43 100644 --- a/validator_client/src/cli.rs +++ b/validator_client/src/cli.rs @@ -501,8 +501,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 f76584bb1d4..6cf519bc313 100644 --- a/validator_client/src/lib.rs +++ b/validator_client/src/lib.rs @@ -35,7 +35,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}; @@ -418,14 +418,10 @@ 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 payload_available_rx = if config.enable_payload_available_monitor { let (payload_available_tx, payload_available_receiver) = @@ -437,6 +433,12 @@ impl ProductionValidatorClient { }; 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); @@ -560,7 +562,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); @@ -581,6 +583,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 6cfeb301e1c..e18c940a1c6 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, } @@ -191,6 +193,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_to_attestation_deadline).await; - None - }; + let beacon_node_data = head_event_or_deadline( + &mut head_monitor_rx, + &self.slot_clock, + duration_to_attestation_deadline, + ) + .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"); @@ -244,30 +244,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.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 595fc42a3de..e5d1d8f19e1 100644 --- a/validator_client/validator_services/src/sync_committee_service.rs +++ b/validator_client/validator_services/src/sync_committee_service.rs @@ -1,16 +1,19 @@ 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; 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}; 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 +24,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 +65,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 +79,7 @@ impl SyncCommitteeService>, executor: TaskExecutor, + head_monitor_rx: Option>, ) -> Self { Self { inner: Arc::new(Inner { @@ -68,6 +88,7 @@ impl SyncCommitteeService SyncCommitteeService = None; loop { - if let Some(now_duration) = self.slot_clock.now_duration() { - let (_, Some(duration_to_sync_message_deadline)) = sync_message_deadline::( - &self.slot_clock, - &self.duties_service.spec, - now_duration, - ) else { - error!("Failed to determine sync message deadline"); - sleep(slot_duration).await; - continue; - }; - - // Wait for the fork-appropriate sync message due time. - sleep(duration_to_sync_message_deadline).await; - - // Do nothing if the Altair fork has not yet occurred. - if !self.altair_fork_activated() { - continue; - } - - if let Err(e) = self.spawn_contribution_tasks().await { - 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 { + let Some(now) = self.slot_clock.now_duration() else { error!("Failed to read slot clock"); - // If we can't read the slot clock, just wait another slot. 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 deadline for the next slot, or a head event for the + // current slot, whichever comes first. + let head_event = head_event_or_deadline( + &mut head_monitor_rx, + &self.slot_clock, + duration_to_sync_message_deadline, + ) + .await; + + // Take the slot from the trigger itself rather than re-reading the clock, so a + // head event arriving at the end of a slot is never attributed to the next slot. + let (current_slot, head_event_root) = match head_event { + Some(event) => (event.slot, Some(event.beacon_block_root)), + None => (next_slot, None), + }; + + if last_processed_slot.is_some_and(|last_slot| current_slot <= last_slot) { + debug!(%current_slot, "Sync message slot already processed"); + continue; + } + + // Do nothing if the Altair fork has not yet occurred. + if !self.altair_fork_activated() { + continue; } + + 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(); } }; @@ -150,65 +183,106 @@ impl SyncCommitteeService Result<(), String> { + async fn spawn_contribution_tasks(&self, slot: Slot, mut head_event_root: Option) { 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, trigger it at the - // fork-appropriate contribution due time. - let aggregate_production_instant = Instant::now() - + duration_to_next_slot - .checked_add(spec.get_contribution_message_due::(slot)) - .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 { - debug!("No duties known for slot {}", slot); - return Ok(()); + .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 Some(duration_to_deadline) = delay_until_slot_offset( + &self.slot_clock, + slot, + spec.get_sync_message_due::(slot), + ) else { + debug!(%slot, "Skipping sync committee tasks for expired slot"); + return; + }; + sleep(duration_to_deadline).await; + + slot_duties = self + .duties_service + .sync_duties + .get_duties_for_slot::(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 { + debug!(%slot, "No duties known for slot"); + return; }; if slot_duties.duties.is_empty() { debug!(%slot, "No local validators in current sync committee"); - return Ok(()); + return; } - // 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; + // If a validator needs to publish a sync aggregate, trigger it at the + // fork-appropriate contribution due time. + let Some(contribution_delay) = delay_until_slot_offset( + &self.slot_clock, + slot, + spec.get_contribution_message_due::(slot), + ) else { + debug!(%slot, "Skipping sync committee tasks for expired slot"); + return; + }; + // Messages past the contribution deadline can no longer be aggregated, so a trigger + // this late (a head event for the slot already in progress at startup) is skipped. + if contribution_delay.is_zero() { + debug!(%slot, "Skipping sync committee tasks, contribution deadline passed"); + return; + } + 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; + } } }; @@ -226,7 +300,6 @@ 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 { @@ -327,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 { @@ -597,13 +678,188 @@ fn subscriptions_from_sync_duties( #[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 types::{Epoch, MainnetEthSpec, SignedContributionAndProof, SyncCommitteeContribution}; + 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, spec).await + } + + async fn new_with_spec(head_monitoring: bool, spec: ChainSpec) -> Self { + let mut harness = + ValidatorClientHarness::new_with_spec_and_config(1, spec, &Default::default()) + .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(); + } + + /// 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() + .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() + } + + 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() { + 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_count(len: impl Fn() -> usize, count: usize) { + tokio::time::resume(); + let deadline = Instant::now() + Duration::from_secs(5); + + while len() < count && Instant::now() < deadline { + tokio::time::sleep(Duration::from_millis(10)).await; + } + + tokio::time::pause(); + 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 { + HeadEvent { + beacon_node_index: 0, + slot: Slot::new(slot), + beacon_block_root: Hash256::from_low_u64_be(block_root), + } + } #[test] fn duration_to_sync_message_deadline_is_fork_aware() { - type E = MainnetEthSpec; - let mut spec = E::default_spec(); let gloas_fork_epoch = Epoch::new(1); spec.gloas_fork_epoch = Some(gloas_fork_epoch); @@ -643,4 +899,387 @@ mod tests { ); } } + + #[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::new(0)); + + 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; + // 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 + .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, expected_root); + root_mock.expect(1).assert(); + 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 late_head_event_past_contribution_deadline_is_skipped() { + 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; + + // A head event arriving after the contribution deadline (8s) is too late for its + // messages to be aggregated, so the slot is skipped. + harness.advance_time(Duration::from_secs(9)).await; + harness.send_head(0, 11); + yield_to_service().await; + assert!(harness.messages().is_empty()); + + // 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; + 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 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(); + spec.altair_fork_epoch = Some(Epoch::new(0)); + spec.gloas_fork_epoch = Some(Epoch::new(0)); + let mut harness = TestHarness::new_with_spec(false, 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(); + } }