Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions config.example.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
8 changes: 8 additions & 0 deletions crates/common/src/config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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(),
Expand Down Expand Up @@ -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,
Expand Down
6 changes: 6 additions & 0 deletions crates/common/src/metrics.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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!(
Expand Down
132 changes: 107 additions & 25 deletions crates/relay/src/api/proposer/get_payload.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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},
};
Expand Down Expand Up @@ -344,22 +344,24 @@ impl<A: Api> ProposerApi<A> {
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(),
Expand Down Expand Up @@ -462,6 +464,13 @@ impl<A: Api> ProposerApi<A> {
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.
Expand All @@ -476,7 +485,7 @@ impl<A: Api> ProposerApi<A> {
&self,
slot: Slot,
request_time_ns: u64,
) -> Result<(), ProposerApiError> {
) -> Result<SlotTiming, ProposerApiError> {
let Some((since_slot_start, until_slot_start)) =
calculate_slot_time_info(&self.chain_info, slot, request_time_ns)
else {
Expand All @@ -491,15 +500,16 @@ impl<A: Api> ProposerApi<A> {
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(
Expand Down Expand Up @@ -565,6 +575,30 @@ impl<A: Api> ProposerApi<A> {
}
}

enum SlotTiming {
OnTime,
WithinResponseBuffer(ProposerApiError),
}

fn evaluate_response_buffer(
since_slot_start_ms: u64,
cutoff_ms: u64,
buffer_ms: u64,
) -> Result<SlotTiming, ProposerApiError> {
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,
Expand All @@ -585,3 +619,51 @@ pub(super) fn fork_name_from_header(headers: &HeaderMap) -> Result<Option<ForkNa
.map(|fork_name| fork_name.to_str().map_err(|e| e.to_string()).and_then(ForkName::from_str))
.transpose()
}

#[cfg(test)]
mod tests {
use super::*;

const CUTOFF: u64 = GET_PAYLOAD_REQUEST_CUTOFF_MS as u64;
const BUFFER: u64 = 800;

#[test]
fn safely_inside_window_returns_payload() {
assert!(matches!(evaluate_response_buffer(0, CUTOFF, BUFFER), Ok(SlotTiming::OnTime)));
assert!(matches!(
evaluate_response_buffer(CUTOFF - BUFFER - 1, CUTOFF, BUFFER),
Ok(SlotTiming::OnTime)
));
}

#[test]
fn buffer_boundary_withholds_payload() {
// just inside the buffer
assert!(matches!(
evaluate_response_buffer(CUTOFF - BUFFER + 1, CUTOFF, BUFFER),
Ok(SlotTiming::WithinResponseBuffer(ProposerApiError::GetPayloadRequestTooLate { .. }))
));
// exactly at the outer cutoff — still withheld, not outright rejected
assert!(matches!(
evaluate_response_buffer(CUTOFF, CUTOFF, BUFFER),
Ok(SlotTiming::WithinResponseBuffer(ProposerApiError::GetPayloadRequestTooLate { .. }))
));
}

#[test]
fn past_cutoff_is_rejected_outright() {
assert!(matches!(
evaluate_response_buffer(CUTOFF + 1, CUTOFF, BUFFER),
Err(ProposerApiError::GetPayloadRequestTooLate { .. })
));
}

#[test]
fn zero_buffer_is_a_no_op() {
assert!(matches!(evaluate_response_buffer(CUTOFF, CUTOFF, 0), Ok(SlotTiming::OnTime)));
assert!(matches!(
evaluate_response_buffer(CUTOFF + 1, CUTOFF, 0),
Err(ProposerApiError::GetPayloadRequestTooLate { .. })
));
}
}
Loading