diff --git a/config.example.yml b/config.example.yml index 8447b00a..352fbae6 100644 --- a/config.example.yml +++ b/config.example.yml @@ -51,6 +51,7 @@ block_merging_config: # - builder_coinbase: "0x0000000000000000000000000000000000000000" # collateral_safe: "0x0000000000000000000000000000000000000000" target_get_payload_propagation_duration_ms: 0 +get_payload_v1_response_buffer_ms: 800 primev_config: null discord_webhook_url: null is_submission_instance: true diff --git a/crates/common/src/config.rs b/crates/common/src/config.rs index 0d86fb27..0858c35a 100644 --- a/crates/common/src/config.rs +++ b/crates/common/src/config.rs @@ -45,6 +45,9 @@ pub struct RelayConfig { pub router_config: RouterConfig, #[serde(default = "default_duration")] pub target_get_payload_propagation_duration_ms: u64, + /// Requests this close to the V1 getPayload cutoff are published but not returned. + #[serde(default = "default_get_payload_v1_response_buffer_ms")] + pub get_payload_v1_response_buffer_ms: u64, /// Configuration for block merging parameters. #[serde(default)] pub block_merging_config: BlockMergingConfig, @@ -101,6 +104,7 @@ impl RelayConfig { validator_preferences: Default::default(), router_config: Default::default(), target_get_payload_propagation_duration_ms: Default::default(), + get_payload_v1_response_buffer_ms: Default::default(), block_merging_config: Default::default(), header_stream: Default::default(), primev_config: Default::default(), @@ -720,6 +724,10 @@ fn default_duration() -> u64 { 1000 } +fn default_get_payload_v1_response_buffer_ms() -> u64 { + 800 +} + #[derive(Serialize, Deserialize, Clone)] pub struct S3Config { pub bucket: String, diff --git a/crates/common/src/metrics.rs b/crates/common/src/metrics.rs index 62ba049f..9ee5ec77 100644 --- a/crates/common/src/metrics.rs +++ b/crates/common/src/metrics.rs @@ -449,6 +449,12 @@ lazy_static! { &RELAY_METRICS_REGISTRY ) .unwrap(); + pub static ref GET_PAYLOAD_V1_RESPONSE_WITHHELD: IntCounter = register_int_counter_with_registry!( + "get_payload_v1_response_withheld_total", + "Count of V1 getPayload responses withheld due to the response-safety buffer", + &RELAY_METRICS_REGISTRY + ) + .unwrap(); //////////////// GET HEADER //////////////// pub static ref HEADER_TIMEOUT_FETCH: IntCounter = register_int_counter_with_registry!( diff --git a/crates/relay/src/api/proposer/get_payload.rs b/crates/relay/src/api/proposer/get_payload.rs index 06698fc5..32fabb66 100644 --- a/crates/relay/src/api/proposer/get_payload.rs +++ b/crates/relay/src/api/proposer/get_payload.rs @@ -8,7 +8,7 @@ use helix_common::{ beacon::types::BroadcastValidation, chain_info::ChainInfo, decoder::{Encoding, HEADER_SSZ}, - metrics::BEACON_BLOCK_PUBLISH_FAILURES, + metrics::{BEACON_BLOCK_PUBLISH_FAILURES, GET_PAYLOAD_V1_RESPONSE_WITHHELD}, spawn_tracked, utils::{extract_request_id, utcnow_ms, utcnow_ns}, }; @@ -344,22 +344,24 @@ impl ProposerApi { trace.payload_fetched = utcnow_ns(); // Handle early/late requests - if let Err(err) = - self.await_and_validate_slot_start_time(head_slot + 1, trace.receive).await - { - warn!(error = %err, "get_payload was sent too late"); - - self.db.save_too_late_get_payload( - (head_slot + 1).into(), - proposer_public_key, - block_hash, - trace.receive, - trace.payload_fetched, - ); + let slot_timing = + match self.await_and_validate_slot_start_time(head_slot + 1, trace.receive).await { + Ok(slot_timing) => slot_timing, + Err(err) => { + warn!(error = %err, "get_payload was sent too late"); + + self.db.save_too_late_get_payload( + (head_slot + 1).into(), + proposer_public_key, + block_hash, + trace.receive, + trace.payload_fetched, + ); - let _ = dedup_tx.send(Arc::new(None)); - return Err(err); - } + let _ = dedup_tx.send(Arc::new(None)); + return Err(err); + } + }; self.gossip_payload( to_publish.signed_block.slot(), @@ -462,6 +464,13 @@ impl ProposerApi { if remaining_sleep_ms > 0 { sleep(Duration::from_millis(remaining_sleep_ms)).await; } + + if let SlotTiming::WithinResponseBuffer(err) = slot_timing { + warn!(error = %err, "get_payload landed in the V1 response-safety buffer, withholding payload"); + GET_PAYLOAD_V1_RESPONSE_WITHHELD.inc(); + let _ = dedup_tx.send(Arc::new(None)); + return Err(err); + } } // Notify dedup waiters with the successful response. @@ -476,7 +485,7 @@ impl ProposerApi { &self, slot: Slot, request_time_ns: u64, - ) -> Result<(), ProposerApiError> { + ) -> Result { let Some((since_slot_start, until_slot_start)) = calculate_slot_time_info(&self.chain_info, slot, request_time_ns) else { @@ -491,15 +500,16 @@ impl ProposerApi { if let Some(until_slot_start) = until_slot_start { info!("waiting until slot start t=0: {} ms", until_slot_start.as_millis()); sleep(until_slot_start).await; - } else if let Some(since_slot_start) = since_slot_start && - since_slot_start.as_millis() > GET_PAYLOAD_REQUEST_CUTOFF_MS as u128 - { - return Err(ProposerApiError::GetPayloadRequestTooLate { - cutoff: GET_PAYLOAD_REQUEST_CUTOFF_MS as u64, - request_time: since_slot_start.as_millis() as u64, - }); + return Ok(SlotTiming::OnTime); } - Ok(()) + + let Some(since_slot_start) = since_slot_start else { return Ok(SlotTiming::OnTime) }; + + evaluate_response_buffer( + since_slot_start.as_millis() as u64, + GET_PAYLOAD_REQUEST_CUTOFF_MS as u64, + self.relay_config.get_payload_v1_response_buffer_ms, + ) } async fn save_delivered_payload_info( @@ -565,6 +575,30 @@ impl ProposerApi { } } +enum SlotTiming { + OnTime, + WithinResponseBuffer(ProposerApiError), +} + +fn evaluate_response_buffer( + since_slot_start_ms: u64, + cutoff_ms: u64, + buffer_ms: u64, +) -> Result { + let too_late = || ProposerApiError::GetPayloadRequestTooLate { + cutoff: cutoff_ms, + request_time: since_slot_start_ms, + }; + + if since_slot_start_ms > cutoff_ms { + return Err(too_late()); + } + if since_slot_start_ms > cutoff_ms.saturating_sub(buffer_ms) { + return Ok(SlotTiming::WithinResponseBuffer(too_late())); + } + Ok(SlotTiming::OnTime) +} + /// Calculates the time information for a given slot. fn calculate_slot_time_info( chain_info: &ChainInfo, @@ -585,3 +619,51 @@ pub(super) fn fork_name_from_header(headers: &HeaderMap) -> Result