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
130 changes: 129 additions & 1 deletion crates/common/src/beacon/beacon_client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,10 @@ use std::{sync::Arc, task::Poll, time::Duration};

use ::ssz::Encode;
use alloy_primitives::B256;
use helix_types::{ForkName, LhConfig, VersionedSignedProposal, spec_from_config};
use helix_types::{
ForkName, LhConfig, SignedExecutionPayloadEnvelopeContents, VersionedSignedProposal,
spec_from_config,
};
use http::{Request, header::CONTENT_TYPE};
use http_body_util::Full;
use hyper::body::Bytes;
Expand All @@ -20,6 +23,8 @@ use crate::{
};

const CONSENSUS_VERSION_HEADER: &str = "eth-consensus-version";
// Always "false": helix always has blobs cached from the builder's own submission.
const BLOB_DATA_INCLUDED_HEADER: &str = "eth-blob-data-included";
const PUBLISH_BLOCK_TIMEOUT: Duration = Duration::from_secs(4);
const GET_TIMEOUT: Duration = Duration::from_secs(5);

Expand Down Expand Up @@ -112,6 +117,48 @@ impl BeaconClient {
}
}

/// Publishes a signed execution payload envelope SSZ-encoded, so a connected beacon node
/// broadcasts it to the `execution_payload` gossip topic on helix's behalf.
/// <https://github.com/ethereum/beacon-APIs/blob/master/apis/beacon/execution_payload/envelope_post.yaml>
pub async fn publish_execution_payload_envelope(
&self,
envelope: Arc<SignedExecutionPayloadEnvelopeContents>,
fork: ForkName,
) -> Result<u16, BeaconClientError> {
let target = self.config.url.join("eth/v1/beacon/execution_payload_envelopes")?;
let body_bytes = Bytes::from(envelope.as_ssz_bytes());
Comment thread
0w3n-d marked this conversation as resolved.
let req = Request::builder()
.method("POST")
.uri(target.as_str())
.header(CONSENSUS_VERSION_HEADER, fork.to_string())
.header(BLOB_DATA_INCLUDED_HEADER, "true")
.header(CONTENT_TYPE, "application/octet-stream")
.body(Full::new(body_bytes))?;
let mut pending = self.http.send(&target, req)?.with_timeout(PUBLISH_BLOCK_TIMEOUT);

let (status, body) = loop {
match pending.poll_bytes() {
Poll::Pending => {}
Poll::Ready(Ok(r)) => break r,
Poll::Ready(Err(e)) => return Err(e.into()),
}
tokio::task::yield_now().await;
};

match status {
200 => Ok(200),
202 => {
let body_str = String::from_utf8_lossy(&body);
warn!("Envelope broadcast but not integrated: {body_str}");
Ok(202)
}
_ => {
let api_err: ApiError = serde_json::from_slice(&body)?;
Err(BeaconClientError::Api(api_err))
}
}
}

pub async fn get_chain_info(&self) -> Result<ChainInfo, BeaconClientError> {
let spec: BeaconResponse<LhConfig> = self.get("eth/v1/config/spec").await?;
let spec = spec_from_config(spec.data);
Expand All @@ -130,3 +177,84 @@ impl BeaconClient {
Ok(chain_info)
}
}

#[cfg(test)]
mod tests {
use helix_types::{BlsSignature, ExecutionPayloadEnvelope, SignedExecutionPayloadEnvelope};
use httpmock::{Method::POST, MockServer};
use reqwest::Url;

use super::*;

fn test_client(url: Url) -> BeaconClient {
crate::utils::install_default_crypto_provider();
BeaconClient::new(BeaconClientConfig { url })
}

fn empty_envelope() -> Arc<SignedExecutionPayloadEnvelopeContents> {
Arc::new(SignedExecutionPayloadEnvelopeContents {
signed_execution_payload_envelope: SignedExecutionPayloadEnvelope {
message: ExecutionPayloadEnvelope::empty(),
signature: BlsSignature::empty(),
},
kzg_proofs: Default::default(),
blobs: Default::default(),
})
}

#[tokio::test]
async fn publish_execution_payload_envelope_sends_ssz_with_fork_and_blob_headers() {
let server = MockServer::start();
let mock = server.mock(|when, then| {
when.method(POST)
.path("/eth/v1/beacon/execution_payload_envelopes")
.header("eth-consensus-version", "gloas")
.header("eth-blob-data-included", "true")
.header("content-type", "application/octet-stream");
then.status(200);
});

let client = test_client(Url::parse(&server.url("/")).unwrap());
let result =
client.publish_execution_payload_envelope(empty_envelope(), ForkName::Gloas).await;

mock.assert();
assert_eq!(result.unwrap(), 200);
}

#[tokio::test]
async fn publish_execution_payload_envelope_202_is_ok() {
let server = MockServer::start();
server.mock(|when, then| {
when.method(POST).path("/eth/v1/beacon/execution_payload_envelopes");
then.status(202).body("envelope failed integration but was broadcast");
});

let client = test_client(Url::parse(&server.url("/")).unwrap());
let result =
client.publish_execution_payload_envelope(empty_envelope(), ForkName::Gloas).await;

assert_eq!(result.unwrap(), 202);
}

#[tokio::test]
async fn publish_execution_payload_envelope_error_response_parses_api_error() {
let server = MockServer::start();
server.mock(|when, then| {
when.method(POST).path("/eth/v1/beacon/execution_payload_envelopes");
then.status(400).json_body(serde_json::json!({
"code": 400,
"message": "Invalid signed execution payload envelope"
}));
});

let client = test_client(Url::parse(&server.url("/")).unwrap());
let result =
client.publish_execution_payload_envelope(empty_envelope(), ForkName::Gloas).await;

match result {
Err(BeaconClientError::Api(ApiError::ErrorMessage { code: 400, .. })) => {}
other => panic!("expected a 400 ApiError, got {other:?}"),
}
}
}
92 changes: 91 additions & 1 deletion crates/common/src/beacon/multi_beacon_client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@ use std::sync::{
};

use futures::future::join_all;
use helix_types::{ForkName, VersionedSignedProposal};
use helix_types::{ForkName, SignedExecutionPayloadEnvelopeContents, VersionedSignedProposal};

use crate::{
beacon::{beacon_client::BeaconClient, error::BeaconClientError, types::BroadcastValidation},
Expand Down Expand Up @@ -83,4 +83,94 @@ impl MultiBeaconClient {

Err(last_error.unwrap_or(BeaconClientError::BeaconNodeUnavailable))
}

/// Publishes the signed execution payload envelope to all beacon clients; returns on first
Comment thread
0w3n-d marked this conversation as resolved.
/// success. Unlike `publish_block`, fans out via plain concurrent futures, not
/// `spawn_tracked!`.
Comment on lines +88 to +89

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The fan-out uses plain futures inside the axum handler: if the proposer disconnects, or the 5s TimeoutLayer fires, the handler future is dropped and every in-flight publish is cancelled. This is exactly why publish_block uses spawn_tracked!, and here the payload is already gone when it happens.

pub async fn publish_execution_payload_envelope(
&self,
envelope: Arc<SignedExecutionPayloadEnvelopeContents>,
fork: ForkName,
) -> Result<(), BeaconClientError> {
let futures = self
.beacon_clients
.iter()
.map(|client| client.publish_execution_payload_envelope(envelope.clone(), fork));

let mut last_error: Option<BeaconClientError> = None;
for res in join_all(futures).await {
match res {
Ok(_) => return Ok(()),
Err(err) => last_error = Some(err),
}
}

Err(last_error.unwrap_or(BeaconClientError::BeaconNodeUnavailable))
}
}

#[cfg(test)]
mod tests {
use helix_types::{BlsSignature, ExecutionPayloadEnvelope, SignedExecutionPayloadEnvelope};
use httpmock::{Method::POST, MockServer};
use reqwest::Url;

use super::*;
use crate::BeaconClientConfig;

fn envelope() -> Arc<SignedExecutionPayloadEnvelopeContents> {
Arc::new(SignedExecutionPayloadEnvelopeContents {
signed_execution_payload_envelope: SignedExecutionPayloadEnvelope {
message: ExecutionPayloadEnvelope::empty(),
signature: BlsSignature::empty(),
},
kzg_proofs: Default::default(),
blobs: Default::default(),
})
}

fn client_for(server: &MockServer) -> Arc<BeaconClient> {
let url = Url::parse(&server.url("/")).unwrap();
Arc::new(BeaconClient::new(BeaconClientConfig { url }))
}

#[tokio::test]
async fn publish_execution_payload_envelope_returns_ok_on_first_success() {
crate::utils::install_default_crypto_provider();
let failing = MockServer::start();
failing.mock(|when, then| {
when.method(POST).path("/eth/v1/beacon/execution_payload_envelopes");
then.status(500);
});
let succeeding = MockServer::start();
succeeding.mock(|when, then| {
when.method(POST).path("/eth/v1/beacon/execution_payload_envelopes");
then.status(200);
});

let multi = MultiBeaconClient::new(vec![client_for(&failing), client_for(&succeeding)]);
let result = multi.publish_execution_payload_envelope(envelope(), ForkName::Gloas).await;

assert!(result.is_ok(), "expected Ok, got {result:?}");
}

#[tokio::test]
async fn publish_execution_payload_envelope_returns_err_when_all_clients_fail() {
crate::utils::install_default_crypto_provider();
let a = MockServer::start();
a.mock(|when, then| {
when.method(POST).path("/eth/v1/beacon/execution_payload_envelopes");
then.status(500);
});
let b = MockServer::start();
b.mock(|when, then| {
when.method(POST).path("/eth/v1/beacon/execution_payload_envelopes");
then.status(500);
});

let multi = MultiBeaconClient::new(vec![client_for(&a), client_for(&b)]);
let result = multi.publish_execution_payload_envelope(envelope(), ForkName::Gloas).await;

assert!(result.is_err(), "expected Err, got {result:?}");
}
}
5 changes: 5 additions & 0 deletions crates/common/src/config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -87,6 +87,10 @@ pub struct RelayConfig {
#[serde(default)]
pub operator_config: Option<OperatorConfig>,
pub blacklist_provider: Option<Url>,
/// This relay's on-chain Gloas (ePBS) builder_index. Placeholder until helix has a real
/// on-chain builder registration; signs under the relay's own key in the meantime.
#[serde(default)]

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

should this really have a default?

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This will be addressed by #635

pub gloas_builder_index: u64,
Comment thread
0w3n-d marked this conversation as resolved.
}

#[derive(Serialize, Deserialize, Clone, Default)]
Expand Down Expand Up @@ -160,6 +164,7 @@ impl RelayConfig {
enable_flux_profiler: false,
operator_config: None,
blacklist_provider: None,
gloas_builder_index: 0,
}
}
}
Expand Down
19 changes: 18 additions & 1 deletion crates/relay/src/api/proposer/error.rs
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
use alloy_primitives::B256;
use axum::{
self,
response::{IntoResponse, Response},
Expand Down Expand Up @@ -136,6 +137,19 @@ pub enum ProposerApiError {

#[error("invalid request: Date-Milliseconds and X-Timeout-Ms headers are required")]
MissingTimingHeaders,

#[error("no held execution payload for bid block hash {0:?}")]
NoHeldPayloadForBlock(B256),

#[error(
"held payload block hash {held:?} does not match the bid's committed block hash {bid:?}"
)]
HeldPayloadBlockHashMismatch { held: B256, bid: B256 },

#[error(
"bid builder_index {bid} does not match this relay's configured builder_index {configured}"
)]
BuilderIndexMismatch { bid: u64, configured: u64 },
}

impl ProposerApiError {
Expand Down Expand Up @@ -197,7 +211,10 @@ impl IntoResponse for ProposerApiError {
ProposerApiError::GetPayloadAlreadyReceived |
ProposerApiError::RequestForPastSlot { .. } |
ProposerApiError::RequestAuthSlotMismatch { .. } |
ProposerApiError::MissingTimingHeaders => StatusCode::BAD_REQUEST,
ProposerApiError::MissingTimingHeaders |
ProposerApiError::NoHeldPayloadForBlock(_) |
ProposerApiError::HeldPayloadBlockHashMismatch { .. } |
ProposerApiError::BuilderIndexMismatch { .. } => StatusCode::BAD_REQUEST,

// All authentication failures, kept indistinguishable by status
ProposerApiError::InvalidApiKey |
Expand Down
10 changes: 10 additions & 0 deletions crates/relay/src/api/proposer/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@ use helix_common::{
use helix_database::handle::DbHandle;
use helix_operator::OperatorPubSub;
use hyper::StatusCode;
pub use submit_signed_beacon_block::{GloasBuilderIdentity, GloasPayloadStore};

use crate::{
api::{Api, proposer::ip_tracker::IpTracker, router::Terminating},
Expand All @@ -46,6 +47,8 @@ pub struct ProposerApi<A: Api> {
pub reg_handle: RegWorkerHandle,
pub operator_api: Option<Arc<OperatorPubSub>>,
pub ip_tracker: IpTracker,
pub gloas_builder_identity: Arc<GloasBuilderIdentity>,

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

don't need an arc here because it's not shared?

pub gloas_payload_store: Arc<GloasPayloadStore>,
}

impl<A: Api> ProposerApi<A> {
Expand All @@ -64,7 +67,12 @@ impl<A: Api> ProposerApi<A> {
reg_handle: RegWorkerHandle,
alert_manager: Arc<AlertManager>,
operator_api: Option<Arc<OperatorPubSub>>,
gloas_payload_store: Arc<GloasPayloadStore>,
) -> Self {
let gloas_builder_identity = Arc::new(GloasBuilderIdentity {
builder_index: relay_config.gloas_builder_index,
keypair: signing_context.keypair.clone(),
});
Self {
local_cache,
db,
Expand All @@ -81,6 +89,8 @@ impl<A: Api> ProposerApi<A> {
reg_handle,
operator_api,
ip_tracker: IpTracker::default(),
gloas_builder_identity,
gloas_payload_store,
}
}
}
Expand Down
Loading
Loading