From ee05f3692cbb9080a08a8047915769ceda869fa0 Mon Sep 17 00:00:00 2001 From: Mark Mackey Date: Thu, 13 Aug 2026 14:44:20 -0500 Subject: [PATCH] Add Gloas Builder API client and adapt execution layer (Gloas builder API 2/5) Second PR of the Gloas builder API stack (builder-specs #165): - builder_client: relocate the legacy MEV-boost client to `PreGloasBuilderHttpClient` and add the Gloas Builder API HTTP client (request-auth, bid requests, preference submission, signed-block forwarding) plus the stateless `Builders` service used by block production - execution_layer: adopt the relocated pre-Gloas client and drop dead error variants The new client is not yet wired into the beacon chain; that lands in the next PR of this stack. Change-Id: I5bcab56ff499b183f26cc90697bc7c7f2baddfe7 --- Cargo.lock | 5 + Cargo.toml | 1 + beacon_node/builder_client/Cargo.toml | 5 + .../builder_client/src/builder_http_client.rs | 504 ++++++++++++ beacon_node/builder_client/src/builders.rs | 410 ++++++++++ beacon_node/builder_client/src/error.rs | 94 +++ beacon_node/builder_client/src/lib.rs | 716 +----------------- .../src/pre_gloas_builder_http_client.rs | 676 +++++++++++++++++ beacon_node/execution_layer/Cargo.toml | 2 +- beacon_node/execution_layer/src/engine_api.rs | 7 - beacon_node/execution_layer/src/engines.rs | 1 - beacon_node/execution_layer/src/lib.rs | 15 +- 12 files changed, 1735 insertions(+), 701 deletions(-) create mode 100644 beacon_node/builder_client/src/builder_http_client.rs create mode 100644 beacon_node/builder_client/src/builders.rs create mode 100644 beacon_node/builder_client/src/error.rs create mode 100644 beacon_node/builder_client/src/pre_gloas_builder_http_client.rs diff --git a/Cargo.lock b/Cargo.lock index 480f4f9b49d..76426650d4a 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1694,16 +1694,21 @@ version = "0.1.0" dependencies = [ "arbitrary", "bls", + "builder_types", "context_deserialize", "eth2", "ethereum_ssz", + "futures", "lighthouse_version", "mockito", + "parking_lot", + "pretty_reqwest_error", "reqwest", "sensitive_url", "serde", "serde_json", "tokio", + "tracing", "types", ] diff --git a/Cargo.toml b/Cargo.toml index 6e0867a692d..48dbefc8b1a 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -119,6 +119,7 @@ beacon_processor = { path = "beacon_node/beacon_processor" } bincode = "1" bitvec = "1" bls = { path = "crypto/bls" } +builder_client = { path = "beacon_node/builder_client" } builder_types = { path = "common/builder_types" } byteorder = "1" bytes = "1.11.1" diff --git a/beacon_node/builder_client/Cargo.toml b/beacon_node/builder_client/Cargo.toml index a329379160f..e6facfb8c2d 100644 --- a/beacon_node/builder_client/Cargo.toml +++ b/beacon_node/builder_client/Cargo.toml @@ -9,14 +9,19 @@ bls = { workspace = true } context_deserialize = { workspace = true } eth2 = { workspace = true } ethereum_ssz = { workspace = true } +futures = { workspace = true } lighthouse_version = { workspace = true } +parking_lot = { workspace = true } +pretty_reqwest_error = { workspace = true } reqwest = { workspace = true } sensitive_url = { workspace = true } serde = { workspace = true } serde_json = { workspace = true } +tracing = { workspace = true } [dev-dependencies] arbitrary = { workspace = true } +builder_types = { workspace = true, features = ["arbitrary"] } mockito = { workspace = true } tokio = { workspace = true } types = { workspace = true, features = ["arbitrary"] } diff --git a/beacon_node/builder_client/src/builder_http_client.rs b/beacon_node/builder_client/src/builder_http_client.rs new file mode 100644 index 00000000000..49e84b4adf0 --- /dev/null +++ b/beacon_node/builder_client/src/builder_http_client.rs @@ -0,0 +1,504 @@ +use crate::{ + DEFAULT_USER_AGENT, Error, JSON_ACCEPT_VALUE, PREFERENCE_ACCEPT_VALUE, + content_type_from_header, ok_or_error, success_or_error, +}; +use bls::PublicKeyBytes; +use eth2::types::{ + BuilderPreferencesRequest, ContentType, EthSpec, ExecutionBlockHash, ForkName, + ForkVersionedResponse, Hash256, SignedBeaconBlock, SignedExecutionPayloadBid, + SignedRequestAuth, Slot, +}; +use eth2::{ + CONSENSUS_VERSION_HEADER, CONTENT_TYPE_HEADER, JSON_CONTENT_TYPE_HEADER, + SSZ_CONTENT_TYPE_HEADER, +}; +use reqwest::StatusCode; +use reqwest::header::{ACCEPT, HeaderMap, HeaderName, HeaderValue}; +use sensitive_url::SensitiveUrl; +use ssz::{Decode, Encode}; +use std::time::Duration; +use tracing::warn; + +/// This is a whole rabbithole.. see discussion: +/// https://discord.com/channels/595666850260713488/874767108809031740/1529125867484348577 +pub const DEFAULT_GET_EXECUTION_PAYLOAD_BID_TIMEOUT_MILLIS: u64 = 400; + +/// Default timeout for builder submit requests (preferences and signed block). +pub const DEFAULT_SUBMIT_TIMEOUT_MILLIS: u64 = 1000; + +/// Header advertising the proposer's request timeout (in milliseconds) to the builder. +const X_TIMEOUT_MS: HeaderName = HeaderName::from_static("x-timeout-ms"); +/// Header carrying the Unix send-time (in milliseconds) so the builder can measure latency. +const DATE_MILLISECONDS: HeaderName = HeaderName::from_static("date-milliseconds"); + +/// A client for the Gloas (ePBS) Builder API. +/// +/// This client is **not** bound to a single builder URL and holds **no** per-connection state: +/// every request takes the target `builder_url` as a parameter, so one instance can fan out to any +/// number of builders. SSZ negotiation is done per-request rather than cached, because in Gloas the +/// bid request and the signed-block submission are separated by a full VC round-trip +/// (produce -> sign -> publish) and so cannot share instance state. +#[derive(Clone)] +pub struct BuilderHttpClient { + client: reqwest::Client, + user_agent: String, + /// Only use json for all request/response types. + disable_ssz: bool, +} + +impl BuilderHttpClient { + pub fn new(user_agent: Option, disable_ssz: bool) -> Result { + let user_agent = user_agent.unwrap_or_else(|| DEFAULT_USER_AGENT.to_string()); + let client = reqwest::Client::builder().user_agent(&user_agent).build()?; + Ok(Self { + client, + user_agent, + disable_ssz, + }) + } + + pub fn get_user_agent(&self) -> &str { + &self.user_agent + } + + /// Build the HTTP headers sent with a `getExecutionPayloadBid` request. + /// + /// Sets three headers: + /// - `Accept`: requests SSZ (with JSON fallback) for the response, or JSON only when + /// `disable_ssz` is set. This governs the (larger) bid response encoding only. + /// - `X-Timeout-Ms`: the proposer's request timeout, measured from `Date-Milliseconds`. The + /// builder must respond within this window; required by the builder spec. + /// - `Date-Milliseconds`: the Unix ms send time, letting the builder estimate transit delay; + /// required by the builder spec. + /// + /// The `Accept` header is best-effort (logged and skipped if it cannot be constructed). The two + /// required timing headers are built from a static timeout and the system clock, so their + /// construction cannot realistically fail. + fn compute_get_execution_payload_bid_headers(&self) -> HeaderMap { + let mut headers = HeaderMap::new(); + + let accept_value = if self.disable_ssz { + JSON_ACCEPT_VALUE + } else { + PREFERENCE_ACCEPT_VALUE + }; + + match HeaderValue::from_str(accept_value) { + Ok(accept_header) => { + headers.insert(ACCEPT, accept_header); + } + Err(e) => { + warn!("Invalid accept value: {}", e); + } + } + + // Advertise our timeout to the builder so it can bound its own work. + headers.insert( + X_TIMEOUT_MS, + HeaderValue::from(DEFAULT_GET_EXECUTION_PAYLOAD_BID_TIMEOUT_MILLIS), + ); + + // Timestamp the request (Unix ms) so the builder can measure one-way latency. + match std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .map(|d| d.as_millis() as u64) + { + Ok(now_millis) => { + headers.insert(DATE_MILLISECONDS, HeaderValue::from(now_millis)); + } + Err(e) => { + warn!("Failed to compute date header: {}", e); + } + } + + headers + } + + /// `POST /eth/v1/builder/execution_payload_bid/{slot}/{parent_hash}/{parent_root}/{proposer_pubkey}` + /// + /// Request a bid from a single builder. Returns `Ok(None)` if the builder has no bid available + /// (HTTP 204). + /// + /// The `SignedRequestAuth` body is required by the builder spec (a builder returns 400 if it + /// is missing), and `RequestAuth` is fork-versioned, so the `Eth-Consensus-Version` header is + /// required too. The body is small and always sent as JSON; SSZ is only negotiated for the + /// (larger) response via the `Accept` header, and the response is decoded according to its + /// `Content-Type`. + #[allow(clippy::too_many_arguments)] + pub async fn get_execution_payload_bid( + &self, + builder_url: &SensitiveUrl, + slot: Slot, + parent_hash: ExecutionBlockHash, + parent_root: Hash256, + proposer_pubkey: &PublicKeyBytes, + signed_request_auth: &SignedRequestAuth, + fork_name: ForkName, + ) -> Result>, Error> { + let mut path = builder_url.expose_full().clone(); + + path.path_segments_mut() + .map_err(|()| Error::InvalidUrl(builder_url.clone()))? + .push("eth") + .push("v1") + .push("builder") + .push("execution_payload_bid") + .push(slot.to_string().as_str()) + .push(format!("{parent_hash:?}").as_str()) + .push(format!("{parent_root:?}").as_str()) + .push(proposer_pubkey.as_hex_string().as_str()); + + let timeout = Duration::from_millis(DEFAULT_GET_EXECUTION_PAYLOAD_BID_TIMEOUT_MILLIS); + let headers = self.compute_get_execution_payload_bid_headers(); + // The auth body is tiny; always send it as JSON. SSZ-encoding it buys nothing and avoids + // having to probe the builder's SSZ request-ingest support. + let request = self + .client + .post(path) + .timeout(timeout) + .headers(headers) + .header(CONSENSUS_VERSION_HEADER, fork_name.to_string()) + .json(signed_request_auth); + + let response = ok_or_error(request.send().await.map_err(Error::from)?).await?; + + if response.status() == StatusCode::NO_CONTENT { + return Ok(None); + } + + let response_headers = response.headers().clone(); + let response_bytes = response.bytes().await?; + + match content_type_from_header(&response_headers) { + ContentType::Ssz => { + let bid = SignedExecutionPayloadBid::::from_ssz_bytes(&response_bytes) + .map_err(Error::InvalidSsz)?; + Ok(Some(bid)) + } + ContentType::Json => { + let versioned: ForkVersionedResponse> = + serde_json::from_slice(&response_bytes).map_err(Error::InvalidJson)?; + Ok(Some(versioned.data)) + } + } + } + + /// `POST /eth/v1/builder/builder_preferences/{validator_pubkey}` + /// + /// Submit a validator's builder preferences to a builder ahead of the bid request (typically in + /// the epoch before the proposal, so the builder has them before `getExecutionPayloadBid` + /// arrives). Success is HTTP 202. + /// + /// `BuilderPreferencesRequest` is fork-versioned, so `fork_name` is sent as the required + /// `Eth-Consensus-Version` header (builder-specs #165); the body is small and sent as JSON. + pub async fn submit_builder_preferences( + &self, + builder_url: &SensitiveUrl, + proposer_pubkey: &PublicKeyBytes, + preferences: &BuilderPreferencesRequest, + fork_name: ForkName, + ) -> Result<(), Error> { + let mut path = builder_url.expose_full().clone(); + + path.path_segments_mut() + .map_err(|()| Error::InvalidUrl(builder_url.clone()))? + .push("eth") + .push("v1") + .push("builder") + .push("builder_preferences") + .push(proposer_pubkey.as_hex_string().as_str()); + + let timeout = Duration::from_millis(DEFAULT_SUBMIT_TIMEOUT_MILLIS); + let request = self + .client + .post(path) + .timeout(timeout) + .header(CONSENSUS_VERSION_HEADER, fork_name.to_string()) + .json(preferences); + + let response = success_or_error(request.send().await.map_err(Error::from)?).await?; + + if response.status() == StatusCode::ACCEPTED { + Ok(()) + } else { + // ACCEPTED is the only valid status code response + Err(Error::StatusCode(response.status())) + } + } + + /// `POST /eth/v1/builder/beacon_blocks` + /// + /// Submit the signed Gloas beacon block to the builder that won selection. On success (HTTP + /// 202) the builder becomes responsible for publishing the execution payload envelope. + /// + /// `ssz_request` selects the request-body encoding: SSZ when `true` and the client has SSZ + /// enabled, otherwise JSON. + pub async fn submit_signed_beacon_block( + &self, + builder_url: &SensitiveUrl, + block: &SignedBeaconBlock, + ssz_request: bool, + ) -> Result<(), Error> { + let mut path = builder_url.expose_full().clone(); + + path.path_segments_mut() + .map_err(|()| Error::InvalidUrl(builder_url.clone()))? + .push("eth") + .push("v1") + .push("builder") + .push("beacon_blocks"); + + let mut headers = HeaderMap::new(); + headers.insert( + CONSENSUS_VERSION_HEADER, + HeaderValue::from_str(&block.fork_name_unchecked().to_string()) + .map_err(|e| Error::InvalidHeaders(format!("{}", e)))?, + ); + + let timeout = Duration::from_millis(DEFAULT_SUBMIT_TIMEOUT_MILLIS); + let request = if ssz_request && !self.disable_ssz { + headers.insert( + CONTENT_TYPE_HEADER, + HeaderValue::from_str(SSZ_CONTENT_TYPE_HEADER) + .map_err(|e| Error::InvalidHeaders(format!("{}", e)))?, + ); + self.client + .post(path) + .timeout(timeout) + .headers(headers) + .body(block.as_ssz_bytes()) + } else { + headers.insert( + CONTENT_TYPE_HEADER, + HeaderValue::from_str(JSON_CONTENT_TYPE_HEADER) + .map_err(|e| Error::InvalidHeaders(format!("{}", e)))?, + ); + self.client + .post(path) + .timeout(timeout) + .headers(headers) + .json(block) + }; + + let response = success_or_error(request.send().await.map_err(Error::from)?).await?; + + if response.status() == StatusCode::ACCEPTED { + Ok(()) + } else { + // ACCEPTED is the only valid status code response + Err(Error::StatusCode(response.status())) + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + use arbitrary::Arbitrary; + use eth2::types::beacon_response::EmptyMetadata; + use eth2::types::{ForkName, MainnetEthSpec}; + use mockito::{Matcher, Server, ServerGuard}; + use std::str::FromStr; + + type E = MainnetEthSpec; + + fn client_for() -> BuilderHttpClient { + BuilderHttpClient::new(None, false).unwrap() + } + + fn builder_url(server: &ServerGuard) -> SensitiveUrl { + SensitiveUrl::from_str(&server.url()).unwrap() + } + + fn signed_request_auth() -> SignedRequestAuth { + let mut u = types::test_utils::test_unstructured(); + SignedRequestAuth::arbitrary(&mut u).unwrap() + } + + fn empty_bid_response() -> ForkVersionedResponse> { + ForkVersionedResponse { + version: ForkName::Gloas, + metadata: EmptyMetadata {}, + data: SignedExecutionPayloadBid::empty(), + } + } + + fn mock_bid(server: &mut ServerGuard, content_type: ContentType) { + let body = empty_bid_response(); + let mut mock = server.mock( + "POST", + Matcher::Regex(r"^/eth/v1/builder/execution_payload_bid/.+$".to_string()), + ); + mock = match content_type { + ContentType::Json => mock + .with_header(CONTENT_TYPE_HEADER, JSON_CONTENT_TYPE_HEADER) + .with_header(CONSENSUS_VERSION_HEADER, "gloas") + .with_body(serde_json::to_string(&body).unwrap()), + ContentType::Ssz => mock + .with_header(CONTENT_TYPE_HEADER, SSZ_CONTENT_TYPE_HEADER) + .with_header(CONSENSUS_VERSION_HEADER, "gloas") + .with_body(body.data.as_ssz_bytes()), + }; + mock.with_status(200).create(); + } + + async fn request_bid(server: &ServerGuard) -> Option> { + client_for() + .get_execution_payload_bid::( + &builder_url(server), + Slot::new(1), + ExecutionBlockHash::repeat_byte(1), + Hash256::repeat_byte(2), + &PublicKeyBytes::empty(), + &signed_request_auth(), + ForkName::Gloas, + ) + .await + .expect("bid request should succeed") + } + + #[tokio::test] + async fn get_execution_payload_bid_json() { + let mut server = Server::new_async().await; + mock_bid(&mut server, ContentType::Json); + let bid = request_bid(&server).await.expect("should have a bid"); + assert_eq!(bid, SignedExecutionPayloadBid::empty()); + } + + #[tokio::test] + async fn get_execution_payload_bid_ssz() { + let mut server = Server::new_async().await; + mock_bid(&mut server, ContentType::Ssz); + let bid = request_bid(&server).await.expect("should have a bid"); + assert_eq!(bid, SignedExecutionPayloadBid::empty()); + } + + #[tokio::test] + async fn submit_builder_preferences_accepted() { + use arbitrary::Arbitrary; + let mut server = Server::new_async().await; + server + .mock( + "POST", + Matcher::Regex(r"^/eth/v1/builder/builder_preferences/.+$".to_string()), + ) + .with_status(202) + .create(); + + let mut u = types::test_utils::test_unstructured(); + let preferences = BuilderPreferencesRequest::arbitrary(&mut u).unwrap(); + + client_for() + .submit_builder_preferences( + &builder_url(&server), + &PublicKeyBytes::empty(), + &preferences, + ForkName::Gloas, + ) + .await + .expect("preferences should be accepted"); + } + + /// The bid request must carry the spec-required headers: `Eth-Consensus-Version` (fork of the + /// auth body), `Date-Milliseconds` + `X-Timeout-Ms` (timing), a JSON `Content-Type` for the + /// auth body, and an `Accept` preferring SSZ for the response. + #[tokio::test] + async fn bid_request_sends_expected_headers() { + let mut server = Server::new_async().await; + let mock = server + .mock( + "POST", + Matcher::Regex(r"^/eth/v1/builder/execution_payload_bid/.+$".to_string()), + ) + .match_header("accept", PREFERENCE_ACCEPT_VALUE) + .match_header(CONSENSUS_VERSION_HEADER, "gloas") + .match_header(CONTENT_TYPE_HEADER, JSON_CONTENT_TYPE_HEADER) + .match_header( + X_TIMEOUT_MS.as_str(), + DEFAULT_GET_EXECUTION_PAYLOAD_BID_TIMEOUT_MILLIS + .to_string() + .as_str(), + ) + .match_header( + DATE_MILLISECONDS.as_str(), + Matcher::Regex(r"^\d+$".to_string()), + ) + .match_header("user-agent", DEFAULT_USER_AGENT) + .with_status(204) + .create(); + + assert!(request_bid(&server).await.is_none()); + mock.assert_async().await; + } + + /// With SSZ responses disabled, the bid request's `Accept` asks for JSON only. + #[tokio::test] + async fn bid_request_accepts_json_only_when_ssz_disabled() { + let mut server = Server::new_async().await; + let mock = server + .mock( + "POST", + Matcher::Regex(r"^/eth/v1/builder/execution_payload_bid/.+$".to_string()), + ) + .match_header("accept", JSON_ACCEPT_VALUE) + .with_status(204) + .create(); + + BuilderHttpClient::new(None, true) + .unwrap() + .get_execution_payload_bid::( + &builder_url(&server), + Slot::new(1), + ExecutionBlockHash::repeat_byte(1), + Hash256::repeat_byte(2), + &PublicKeyBytes::empty(), + &signed_request_auth(), + ForkName::Gloas, + ) + .await + .expect("bid request should succeed"); + mock.assert_async().await; + } + + /// The preferences submission must carry `Eth-Consensus-Version` (the body is fork-versioned) + /// and a JSON `Content-Type`. + #[tokio::test] + async fn submit_builder_preferences_sends_expected_headers() { + let mut server = Server::new_async().await; + let mock = server + .mock( + "POST", + Matcher::Regex(r"^/eth/v1/builder/builder_preferences/.+$".to_string()), + ) + .match_header(CONSENSUS_VERSION_HEADER, "gloas") + .match_header(CONTENT_TYPE_HEADER, JSON_CONTENT_TYPE_HEADER) + .with_status(202) + .create(); + + let mut u = types::test_utils::test_unstructured(); + let preferences = BuilderPreferencesRequest::arbitrary(&mut u).unwrap(); + client_for() + .submit_builder_preferences( + &builder_url(&server), + &PublicKeyBytes::empty(), + &preferences, + ForkName::Gloas, + ) + .await + .expect("preferences should be accepted"); + mock.assert_async().await; + } + + #[tokio::test] + async fn get_execution_payload_bid_no_content() { + let mut server = Server::new_async().await; + server + .mock( + "POST", + Matcher::Regex(r"^/eth/v1/builder/execution_payload_bid/.+$".to_string()), + ) + .with_status(204) + .create(); + assert!(request_bid(&server).await.is_none()); + } +} diff --git a/beacon_node/builder_client/src/builders.rs b/beacon_node/builder_client/src/builders.rs new file mode 100644 index 00000000000..1342d0c3d51 --- /dev/null +++ b/beacon_node/builder_client/src/builders.rs @@ -0,0 +1,410 @@ +use crate::{BuilderHttpClient, Error as BuilderClientError}; +use bls::PublicKeyBytes; +use eth2::types::{ + BuilderEntry, BuilderPreferenceEntry, BuilderPreferences, BuilderPreferencesRequest, + BuilderPubkeys, EthSpec, ExecutionBlockHash, ForkName, Hash256, SignedBeaconBlock, + SignedExecutionPayloadBid, Slot, +}; +use futures::future::join_all; +use sensitive_url::SensitiveUrl; +use std::fmt::Display; +use std::future::Future; +use std::sync::Arc; +use tracing::{debug, warn}; + +/// A validated direct builder bid, with the provenance and per-builder policy needed to turn it into +/// a selection candidate. +/// +/// The per-builder `min_bid` / `max_execution_payment` / `builder_boost_factor` are carried up as-is; +/// this crate applies no bid math (the `min_bid` floor and the boost are proposer policy resolved on +/// the beacon-chain side). +#[derive(Clone)] +pub struct DirectBid { + /// The signed bid returned by the builder. + pub signed_bid: Arc>, + /// URL of the builder that returned this bid, so a winning block can be forwarded to it via + /// `submitSignedBeaconBlock` (echoed to the beacon node as `Eth-Builder-Url`). + pub builder_url: SensitiveUrl, + /// The proposer's `max_execution_payment` cap for this builder, from its `BuilderEntry`. + pub max_execution_payment: u64, + /// The proposer's `builder_boost_factor` for this builder, from its `BuilderEntry`. + pub builder_boost_factor: u64, + /// The proposer's `min_bid` acceptance floor (gwei) for this builder, from its `BuilderEntry`. + pub min_bid: u64, +} + +/// The per-proposal parameters used to address each `getExecutionPayloadBid` request. +/// +/// Validation of returned bids is performed entirely by the caller's `validate` callback (which has +/// the beacon-chain state), so this only carries what's needed to build the request. +#[derive(Clone)] +pub struct BidRequestContext { + pub slot: Slot, + pub parent_hash: ExecutionBlockHash, + pub parent_root: Hash256, + pub proposer_pubkey: PublicKeyBytes, + /// The active consensus version at `slot`, sent as the required `Eth-Consensus-Version` + /// header on each bid request. + pub fork_name: ForkName, +} + +/// Orchestrates direct builder bid requests. +/// +/// Fans `getExecutionPayloadBid` out to the builders a proposer configured and returns the validated +/// bids for the block producer to rank against the local and gossip payloads. Stateless — it holds +/// no bids between requests. +pub struct Builders { + client: Arc, +} + +/// A single failed builder-preference submission, identified by its position in the submitted list. +pub struct SubmissionFailure { + /// Index of the failing entry in the submitted list. + pub index: usize, + /// Why the submission failed. + pub error: BuilderClientError, +} + +impl Builders { + pub fn new(client: Arc) -> Self { + Self { client } + } + + /// Forward a signed beacon block to the builder that won this slot's bid, via + /// `submitSignedBeaconBlock`. + /// + /// Submitted as JSON: the builder's SSZ preference from bid time isn't carried across the + /// `Eth-Builder-Url` header round-trip, and builders must accept JSON. + pub async fn forward_signed_block( + &self, + builder_url: &SensitiveUrl, + block: &SignedBeaconBlock, + ) -> Result<(), BuilderClientError> { + self.client + .submit_signed_beacon_block(builder_url, block, false) + .await + } + + /// Request bids from every builder in `entries` concurrently, validate them, and return the + /// valid ones. + /// + /// Every entry is a bid request to its `url`, which beacon-APIs #630 requires (a zero-length url + /// is invalid); an entry whose `url` is empty, malformed, or not http(s) can't be requested and + /// is skipped. One request is made **per entry** — several entries MAY share a `url` with + /// different `auth`, so requests are not de-duplicated by URL (#630 forbids two entries sharing + /// both a `url` and their `auth`'s `data`). + /// + /// Each builder runs in its own pipeline — request, then the producer-supplied `validate` + /// callback, which performs *all* bid validation against the block producer's advanced beacon + /// state (consensus consistency, builder eligibility, collateral, and the BLS signature). The + /// entry's `builder_pubkeys` filter (empty accepts any builder) is passed to `validate` so it + /// can enforce that the bid is signed by one of the expected builders — the state and signing + /// domain that check needs live on the producer side, not here. The per-builder `min_bid` floor + /// is likewise a proposer policy applied by the caller (see `DirectBid::min_bid`), not here. These + /// pipelines run + /// **concurrently across builders**, so a slow builder or an expensive validation for one bid + /// does not hold up the others. A failure, timeout, empty (204) response, or validation error + /// for one builder is isolated: it is logged and that bid is skipped. + /// + /// Returns every bid that passed validation; the block producer turns each into a selection + /// candidate and ranks them. + pub async fn request_and_validate_bids( + &self, + ctx: &BidRequestContext, + entries: &[BuilderEntry], + validate: F, + ) -> Vec> + where + F: Fn(Arc>, BuilderPubkeys) -> Fut, + Fut: Future>, + Err: Display, + { + // Resolve each entry to a `(resolved_url, entry)` target. Every entry must carry a valid url + // (#630); one that's empty, malformed, or non-http(s) can't be requested and is skipped. One + // request is made per entry (no URL de-duplication). + let mut targets = Vec::new(); + for entry in entries { + let url = match entry.url.to_sensitive_url() { + Ok(url) => url, + Err(e) => { + warn!(error = ?e, "Skipping builder entry with a malformed URL"); + continue; + } + }; + if !matches!(url.expose_full().scheme(), "http" | "https") { + warn!(url = ?url, "Skipping builder entry with an unsupported URL scheme"); + continue; + } + targets.push((url, entry)); + } + + // Run one pipeline per builder — request, then the producer's `validate` callback — and let + // them run concurrently across builders. Each request carries its own timeout, so a slow + // builder cannot delay the others. + let client = &self.client; + let validate = &validate; + let pipelines = targets.iter().map(|(url, entry)| async move { + let response = client + .get_execution_payload_bid::( + url, + ctx.slot, + ctx.parent_hash, + ctx.parent_root, + &ctx.proposer_pubkey, + &entry.auth, + ctx.fork_name, + ) + .await; + + match response { + Ok(Some(bid)) => { + let direct_bid = DirectBid { + signed_bid: Arc::new(bid), + builder_url: url.clone(), + max_execution_payment: entry.max_execution_payment, + builder_boost_factor: entry.builder_boost_factor, + min_bid: entry.min_bid, + }; + + if let Err(error) = + validate(direct_bid.signed_bid.clone(), entry.builder_pubkeys.clone()).await + { + warn!(url = ?url, %error, "Builder bid failed validation"); + return None; + } + Some(direct_bid) + } + Ok(None) => { + debug!(url = ?url, "Builder returned no bid"); + None + } + Err(error) => { + warn!(url = ?url, error = %error, "Builder bid request failed"); + None + } + } + }); + + join_all(pipelines).await.into_iter().flatten().collect() + } + + /// Submit a proposer's builder preferences to each entry's builder, concurrently and + /// best-effort. + /// + /// One submission is made per entry — entries are **not** de-duplicated by URL, since + /// beacon-APIs #630 allows several entries to share a `url`. Each submission is isolated: a + /// malformed URL or a failed request is recorded against that entry's index and never aborts the + /// others. The submissions run **concurrently**, so a slow builder cannot delay the rest. + /// + /// Returns `Ok(())` when every entry was submitted, or the per-entry [`SubmissionFailure`]s by + /// index. + pub async fn submit_builder_preferences( + &self, + entries: Vec, + fork_name: ForkName, + ) -> Result<(), Vec> { + let client = &self.client; + let submissions = entries + .into_iter() + .enumerate() + .map(|(index, entry)| async move { + let url = entry + .url + .to_sensitive_url() + .map_err(|e| SubmissionFailure { + index, + error: e.into(), + })?; + let request = BuilderPreferencesRequest::new( + BuilderPreferences { + max_execution_payment: entry.max_execution_payment, + }, + entry.auth, + ); + client + .submit_builder_preferences(&url, &entry.proposer_pubkey, &request, fork_name) + .await + .map_err(|error| SubmissionFailure { index, error }) + }); + + let failures: Vec = join_all(submissions) + .await + .into_iter() + .filter_map(Result::err) + .collect(); + + if failures.is_empty() { + Ok(()) + } else { + Err(failures) + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + use bls::Signature; + use eth2::types::beacon_response::EmptyMetadata; + use eth2::types::{ + ExecutionPayloadBid, ForkName, ForkVersionedResponse, MainnetEthSpec, RequestAuth, + RequestAuthData, SignedExecutionPayloadBid, SignedRequestAuth, + }; + use eth2::{CONSENSUS_VERSION_HEADER, CONTENT_TYPE_HEADER, JSON_CONTENT_TYPE_HEADER}; + use mockito::{Matcher, Mock, Server, ServerGuard}; + + type E = MainnetEthSpec; + + const BID_PATH: &str = r"^/eth/v1/builder/execution_payload_bid/.+$"; + + fn entry(url: &str, max_execution_payment: u64) -> BuilderEntry { + BuilderEntry { + url: url.parse().unwrap(), + auth: SignedRequestAuth { + message: RequestAuth { + data: RequestAuthData::new(url.as_bytes().to_vec()).unwrap(), + slot: Slot::new(1), + }, + signature: Signature::empty(), + }, + builder_pubkeys: BuilderPubkeys::default(), + max_execution_payment, + min_bid: 0, + builder_boost_factor: 100, + } + } + + fn bid_body(value: u64) -> String { + let body = ForkVersionedResponse { + version: ForkName::Gloas, + metadata: EmptyMetadata {}, + data: SignedExecutionPayloadBid:: { + message: ExecutionPayloadBid { + slot: Slot::new(1), + parent_block_hash: ExecutionBlockHash::zero(), + parent_block_root: Hash256::ZERO, + value, + ..ExecutionPayloadBid::default() + }, + signature: Signature::empty(), + }, + }; + serde_json::to_string(&body).unwrap() + } + + fn mock_bid(server: &mut ServerGuard, value: u64) -> Mock { + server + .mock("POST", Matcher::Regex(BID_PATH.to_string())) + .with_header(CONTENT_TYPE_HEADER, JSON_CONTENT_TYPE_HEADER) + .with_header(CONSENSUS_VERSION_HEADER, "gloas") + .with_body(bid_body(value)) + .with_status(200) + .create() + } + + fn context() -> BidRequestContext { + BidRequestContext { + slot: Slot::new(1), + parent_hash: ExecutionBlockHash::zero(), + parent_root: Hash256::ZERO, + proposer_pubkey: PublicKeyBytes::empty(), + fork_name: ForkName::Gloas, + } + } + + fn builders() -> Builders { + Builders::new(Arc::new(BuilderHttpClient::new(None, false).unwrap())) + } + + #[tokio::test] + async fn fans_out_and_returns_all_valid_bids() { + let mut server_a = Server::new_async().await; + let mut server_b = Server::new_async().await; + mock_bid(&mut server_a, 100); + mock_bid(&mut server_b, 200); + + let builders = builders(); + let entries = vec![entry(&server_a.url(), 1000), entry(&server_b.url(), 1000)]; + + let bids: Vec> = builders + .request_and_validate_bids(&context(), &entries, |_bid, _expected| async { + Ok::<(), String>(()) + }) + .await; + let mut values: Vec = bids.iter().map(|b| b.signed_bid.message.value).collect(); + values.sort_unstable(); + assert_eq!(values, vec![100, 200]); + } + + #[tokio::test] + async fn skips_invalid_url_entry() { + let builders = builders(); + // #630 requires a url; an empty one is invalid and can't be requested, so it is skipped. + let entries = vec![entry("", 1000)]; + + let bids: Vec> = builders + .request_and_validate_bids(&context(), &entries, |_bid, _expected| async { + Ok::<(), String>(()) + }) + .await; + assert!(bids.is_empty()); + } + + #[tokio::test] + async fn requests_each_entry_even_when_url_is_shared() { + let mut server = Server::new_async().await; + // Two entries share a URL but carry different `auth`, so both are requested (one per entry). + let mock = mock_bid(&mut server, 100).expect(2); + + let builders = builders(); + let entry_a = entry(&server.url(), 1000); + let mut entry_b = entry(&server.url(), 1000); + entry_b.auth.message.slot = Slot::new(2); + let entries = vec![entry_a, entry_b]; + + let bids: Vec> = builders + .request_and_validate_bids(&context(), &entries, |_bid, _expected| async { + Ok::<(), String>(()) + }) + .await; + assert_eq!(bids.len(), 2); + mock.assert(); + } + + #[tokio::test] + async fn returns_bid_carrying_min_bid_for_the_caller() { + // The transport layer does not enforce the `min_bid` floor: it returns the bid carrying its + // entry's `min_bid` for the beacon-chain-side caller to enforce. + let mut server = Server::new_async().await; + mock_bid(&mut server, 100); + + let builders = builders(); + let mut entry = entry(&server.url(), 1000); + entry.min_bid = 500; + let entries = vec![entry]; + + let bids: Vec> = builders + .request_and_validate_bids(&context(), &entries, |_bid, _expected| async { + Ok::<(), String>(()) + }) + .await; + assert_eq!(bids.len(), 1); + assert_eq!(bids[0].min_bid, 500); + } + + #[tokio::test] + async fn rejects_bid_failing_producer_validation() { + let mut server = Server::new_async().await; + mock_bid(&mut server, 100); + + let builders = builders(); + let entries = vec![entry(&server.url(), 1000)]; + // The producer callback rejects the bid (e.g. a failed signature or ineligible builder). + let bids: Vec> = builders + .request_and_validate_bids(&context(), &entries, |_bid, _expected| async { + Err::<(), String>("rejected by producer".to_string()) + }) + .await; + assert!(bids.is_empty()); + } +} diff --git a/beacon_node/builder_client/src/error.rs b/beacon_node/builder_client/src/error.rs new file mode 100644 index 00000000000..1fe4af0fa43 --- /dev/null +++ b/beacon_node/builder_client/src/error.rs @@ -0,0 +1,94 @@ +//! The error type for the Gloas [`BuilderHttpClient`](crate::BuilderHttpClient), aligned with the +//! Builder API spec (`builder-specs`). +//! +//! This is deliberately separate from `eth2::Error` (the beacon-node API client's error), which +//! carries beacon-node concerns irrelevant to a builder — API tokens, impostor-signature headers, +//! server-sent events — and collapses every builder-spec status (204 no-bid, 401 auth failed, +//! 406/415 negotiation) into an opaque status code. The pre-Gloas builder client still uses +//! `eth2::Error`. + +use eth2::types::{BuilderUrlError, ErrorMessage}; +use pretty_reqwest_error::PrettyReqwestError; +use reqwest::{Response, StatusCode}; +use sensitive_url::SensitiveUrl; +use std::fmt; + +#[derive(Debug)] +pub enum Error { + /// A transport-level failure sending the request or reading the response. + Reqwest(PrettyReqwestError), + /// A builder URL could not be turned into a request URL. + InvalidUrl(SensitiveUrl), + /// A `BuilderUrl` supplied in config did not parse as a URL. + InvalidBuilderUrl(BuilderUrlError), + /// The builder returned an error response with a parseable `{code, message}` body (the + /// builder-specs `ErrorMessage`), e.g. 400 invalid request or 401 authentication failed. + ServerMessage(ErrorMessage), + /// The builder returned a non-success status whose body could not be parsed. + StatusCode(StatusCode), + /// The builder's JSON response could not be decoded. + InvalidJson(serde_json::Error), + /// The builder's SSZ response could not be decoded. + InvalidSsz(ssz::DecodeError), + /// Request headers could not be constructed. + InvalidHeaders(String), +} + +impl From for Error { + fn from(error: reqwest::Error) -> Self { + Error::Reqwest(error.into()) + } +} + +impl From for Error { + fn from(error: BuilderUrlError) -> Self { + Error::InvalidBuilderUrl(error) + } +} + +impl fmt::Display for Error { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + match self { + Error::Reqwest(e) => write!(f, "HTTP transport error: {e}"), + Error::InvalidUrl(url) => write!(f, "invalid builder URL: {url:?}"), + Error::InvalidBuilderUrl(e) => write!(f, "invalid builder URL: {e:?}"), + Error::ServerMessage(m) => { + write!(f, "builder returned error {}: {}", m.code, m.message) + } + Error::StatusCode(status) => write!(f, "builder returned unexpected status {status}"), + Error::InvalidJson(e) => write!(f, "invalid JSON response: {e}"), + Error::InvalidSsz(e) => write!(f, "invalid SSZ response: {e:?}"), + Error::InvalidHeaders(e) => write!(f, "invalid response headers: {e}"), + } + } +} + +impl std::error::Error for Error {} + +/// Returns `Ok(response)` for a builder success status (200/202/204), otherwise parses the body +/// into an [`Error`]. Mirrors `eth2::ok_or_error` but produces the builder-spec [`Error`]. +pub async fn ok_or_error(response: Response) -> Result { + let status = response.status(); + if matches!( + status, + StatusCode::OK | StatusCode::ACCEPTED | StatusCode::NO_CONTENT + ) { + Ok(response) + } else if let Ok(message) = response.json::().await { + Err(Error::ServerMessage(message)) + } else { + Err(Error::StatusCode(status)) + } +} + +/// Like [`ok_or_error`] but accepts any 2xx status as success. +pub async fn success_or_error(response: Response) -> Result { + let status = response.status(); + if status.is_success() { + Ok(response) + } else if let Ok(message) = response.json::().await { + Err(Error::ServerMessage(message)) + } else { + Err(Error::StatusCode(status)) + } +} diff --git a/beacon_node/builder_client/src/lib.rs b/beacon_node/builder_client/src/lib.rs index bd064ca8bf9..d0f7dafde80 100644 --- a/beacon_node/builder_client/src/lib.rs +++ b/beacon_node/builder_client/src/lib.rs @@ -1,32 +1,17 @@ -use bls::PublicKeyBytes; -use context_deserialize::ContextDeserialize; -pub use eth2::Error; -use eth2::types::beacon_response::EmptyMetadata; -use eth2::types::builder::SignedBuilderBid; -use eth2::types::{ - ContentType, EthSpec, ExecutionBlockHash, ForkName, ForkVersionDecode, ForkVersionedResponse, - SignedValidatorRegistrationData, Slot, -}; -use eth2::types::{FullPayloadContents, SignedBlindedBeaconBlock}; -use eth2::{ - CONSENSUS_VERSION_HEADER, CONTENT_TYPE_HEADER, JSON_CONTENT_TYPE_HEADER, - SSZ_CONTENT_TYPE_HEADER, ok_or_error, success_or_error, -}; -use reqwest::header::{ACCEPT, HeaderMap, HeaderValue}; -use reqwest::{IntoUrl, Response, StatusCode}; -use sensitive_url::SensitiveUrl; -use serde::Serialize; -use serde::de::DeserializeOwned; -use ssz::Encode; +use eth2::types::{ContentType, ForkName}; +use eth2::{CONSENSUS_VERSION_HEADER, CONTENT_TYPE_HEADER, SSZ_CONTENT_TYPE_HEADER}; +use reqwest::header::HeaderMap; use std::str::FromStr; -use std::sync::Arc; -use std::sync::atomic::{AtomicBool, Ordering}; -use std::time::Duration; -pub const DEFAULT_TIMEOUT_MILLIS: u64 = 15000; +pub mod builder_http_client; +pub mod builders; +pub mod error; +pub mod pre_gloas_builder_http_client; -/// This timeout is in accordance with v0.2.0 of the [builder specs](https://github.com/flashbots/mev-boost/pull/20). -pub const DEFAULT_GET_HEADER_TIMEOUT_MILLIS: u64 = 1000; +pub use builder_http_client::BuilderHttpClient; +pub use builders::{BidRequestContext, Builders, DirectBid, SubmissionFailure}; +pub use error::{Error, ok_or_error, success_or_error}; +pub use pre_gloas_builder_http_client::PreGloasBuilderHttpClient; /// Default user agent for HTTP requests. pub const DEFAULT_USER_AGENT: &str = lighthouse_version::VERSION; @@ -36,667 +21,28 @@ pub const PREFERENCE_ACCEPT_VALUE: &str = "application/octet-stream;q=1.0,applic /// Only accept json responses. pub const JSON_ACCEPT_VALUE: &str = "application/json"; -#[derive(Clone)] -pub struct Timeouts { - get_header: Duration, - post_validators: Duration, - post_blinded_blocks: Duration, - get_builder_status: Duration, -} - -impl Timeouts { - fn new(get_header_timeout: Option) -> Self { - let get_header = - get_header_timeout.unwrap_or(Duration::from_millis(DEFAULT_GET_HEADER_TIMEOUT_MILLIS)); - - Self { - get_header, - post_validators: Duration::from_millis(DEFAULT_TIMEOUT_MILLIS), - post_blinded_blocks: Duration::from_millis(DEFAULT_TIMEOUT_MILLIS), - get_builder_status: Duration::from_millis(DEFAULT_TIMEOUT_MILLIS), - } - } -} - -#[derive(Clone)] -pub struct BuilderHttpClient { - client: reqwest::Client, - server: SensitiveUrl, - timeouts: Timeouts, - user_agent: String, - /// Only use json for all requests/responses types. - disable_ssz: bool, - /// Indicates that the `get_header` response had content-type ssz - /// so we can set content-type header to ssz to make the `submit_blinded_blocks` - /// request. - ssz_available: Arc, -} - -impl BuilderHttpClient { - pub fn new( - server: SensitiveUrl, - user_agent: Option, - builder_header_timeout: Option, - disable_ssz: bool, - ) -> Result { - let user_agent = user_agent.unwrap_or(DEFAULT_USER_AGENT.to_string()); - let client = reqwest::Client::builder().user_agent(&user_agent).build()?; - Ok(Self { - client, - server, - timeouts: Timeouts::new(builder_header_timeout), - user_agent, - disable_ssz, - ssz_available: Arc::new(false.into()), +/// Parse the `Eth-Consensus-Version` response header into a `ForkName`, if present. +pub fn fork_name_from_header(headers: &HeaderMap) -> Result, String> { + headers + .get(CONSENSUS_VERSION_HEADER) + .map(|fork_name| { + fork_name + .to_str() + .map_err(|e| e.to_string()) + .and_then(ForkName::from_str) }) - } - - pub fn get_user_agent(&self) -> &str { - &self.user_agent - } - - fn fork_name_from_header(&self, headers: &HeaderMap) -> Result, String> { - headers - .get(CONSENSUS_VERSION_HEADER) - .map(|fork_name| { - fork_name - .to_str() - .map_err(|e| e.to_string()) - .and_then(ForkName::from_str) - }) - .transpose() - } - - fn content_type_from_header(&self, headers: &HeaderMap) -> ContentType { - let Some(content_type) = headers.get(CONTENT_TYPE_HEADER).map(|content_type| { - let content_type = content_type.to_str(); - match content_type { - Ok(SSZ_CONTENT_TYPE_HEADER) => ContentType::Ssz, - _ => ContentType::Json, - } - }) else { - return ContentType::Json; - }; - content_type - } - - async fn get_with_header< - T: DeserializeOwned + ForkVersionDecode + for<'de> ContextDeserialize<'de, ForkName>, - U: IntoUrl, - >( - &self, - url: U, - timeout: Duration, - headers: HeaderMap, - ) -> Result, Error> { - let response = self - .get_response_with_header(url, Some(timeout), headers) - .await?; - - let headers = response.headers().clone(); - let response_bytes = response.bytes().await?; - - let Ok(Some(fork_name)) = self.fork_name_from_header(&headers) else { - // if no fork version specified, attempt to fallback to JSON - self.ssz_available.store(false, Ordering::SeqCst); - return serde_json::from_slice(&response_bytes).map_err(Error::InvalidJson); - }; - - let content_type = self.content_type_from_header(&headers); - - match content_type { - ContentType::Ssz => { - self.ssz_available.store(true, Ordering::SeqCst); - T::from_ssz_bytes_by_fork(&response_bytes, fork_name) - .map(|data| ForkVersionedResponse { - version: fork_name, - metadata: EmptyMetadata {}, - data, - }) - .map_err(Error::InvalidSsz) - } - ContentType::Json => { - self.ssz_available.store(false, Ordering::SeqCst); - serde_json::from_slice(&response_bytes).map_err(Error::InvalidJson) - } - } - } - - /// Return `true` if the most recently received response from the builder had SSZ Content-Type. - /// Return `false` otherwise. - /// Also returns `false` if we have explicitly disabled ssz. - pub fn is_ssz_available(&self) -> bool { - !self.disable_ssz && self.ssz_available.load(Ordering::SeqCst) - } - - async fn get_with_timeout( - &self, - url: U, - timeout: Duration, - ) -> Result { - self.get_response_with_timeout(url, Some(timeout)) - .await? - .json() - .await - .map_err(Into::into) - } - - /// Perform a HTTP GET request, returning the `Response` for further processing. - async fn get_response_with_header( - &self, - url: U, - timeout: Option, - headers: HeaderMap, - ) -> Result { - let mut builder = self.client.get(url); - if let Some(timeout) = timeout { - builder = builder.timeout(timeout); - } - let response = builder.headers(headers).send().await.map_err(Error::from)?; - ok_or_error(response).await - } - - /// Perform a HTTP GET request, returning the `Response` for further processing. - async fn get_response_with_timeout( - &self, - url: U, - timeout: Option, - ) -> Result { - let mut builder = self.client.get(url); - if let Some(timeout) = timeout { - builder = builder.timeout(timeout); - } - let response = builder.send().await.map_err(Error::from)?; - ok_or_error(response).await - } - - /// Generic POST function supporting arbitrary responses and timeouts. - async fn post_generic( - &self, - url: U, - body: &T, - timeout: Option, - ) -> Result { - let mut builder = self.client.post(url); - if let Some(timeout) = timeout { - builder = builder.timeout(timeout); - } - let response = builder.json(body).send().await?; - ok_or_error(response).await - } - - async fn post_ssz_with_raw_response( - &self, - url: U, - ssz_body: Vec, - headers: HeaderMap, - timeout: Option, - ) -> Result { - let mut builder = self.client.post(url); - if let Some(timeout) = timeout { - builder = builder.timeout(timeout); - } - - let response = builder - .headers(headers) - .body(ssz_body) - .send() - .await - .map_err(Error::from)?; - success_or_error(response).await - } - - async fn post_with_raw_response( - &self, - url: U, - body: &T, - headers: HeaderMap, - timeout: Option, - ) -> Result { - let mut builder = self.client.post(url); - if let Some(timeout) = timeout { - builder = builder.timeout(timeout); - } - - let response = builder - .headers(headers) - .json(body) - .send() - .await - .map_err(Error::from)?; - success_or_error(response).await - } - - /// `POST /eth/v1/builder/validators` - pub async fn post_builder_validators( - &self, - validator: &[SignedValidatorRegistrationData], - ) -> Result<(), Error> { - let mut path = self.server.expose_full().clone(); - - path.path_segments_mut() - .map_err(|()| Error::InvalidUrl(self.server.clone()))? - .push("eth") - .push("v1") - .push("builder") - .push("validators"); - - self.post_generic(path, &validator, Some(self.timeouts.post_validators)) - .await?; - Ok(()) - } - - /// `POST /eth/v1/builder/blinded_blocks` with SSZ serialized request body - pub async fn post_builder_blinded_blocks_v1_ssz( - &self, - blinded_block: &SignedBlindedBeaconBlock, - ) -> Result, Error> { - let mut path = self.server.expose_full().clone(); - - let body = blinded_block.as_ssz_bytes(); - - path.path_segments_mut() - .map_err(|()| Error::InvalidUrl(self.server.clone()))? - .push("eth") - .push("v1") - .push("builder") - .push("blinded_blocks"); - - let mut headers = HeaderMap::new(); - headers.insert( - CONSENSUS_VERSION_HEADER, - HeaderValue::from_str(&blinded_block.fork_name_unchecked().to_string()) - .map_err(|e| Error::InvalidHeaders(format!("{}", e)))?, - ); - headers.insert( - CONTENT_TYPE_HEADER, - HeaderValue::from_str(SSZ_CONTENT_TYPE_HEADER) - .map_err(|e| Error::InvalidHeaders(format!("{}", e)))?, - ); - headers.insert( - ACCEPT, - HeaderValue::from_str(PREFERENCE_ACCEPT_VALUE) - .map_err(|e| Error::InvalidHeaders(format!("{}", e)))?, - ); - - let result = self - .post_ssz_with_raw_response( - path, - body, - headers, - Some(self.timeouts.post_blinded_blocks), - ) - .await? - .bytes() - .await?; - - FullPayloadContents::from_ssz_bytes_by_fork(&result, blinded_block.fork_name_unchecked()) - .map_err(Error::InvalidSsz) - } - - /// `POST /eth/v2/builder/blinded_blocks` with SSZ serialized request body - pub async fn post_builder_blinded_blocks_v2_ssz( - &self, - blinded_block: &SignedBlindedBeaconBlock, - ) -> Result<(), Error> { - let mut path = self.server.expose_full().clone(); - - let body = blinded_block.as_ssz_bytes(); - - path.path_segments_mut() - .map_err(|()| Error::InvalidUrl(self.server.clone()))? - .push("eth") - .push("v2") - .push("builder") - .push("blinded_blocks"); - - let mut headers = HeaderMap::new(); - headers.insert( - CONSENSUS_VERSION_HEADER, - HeaderValue::from_str(&blinded_block.fork_name_unchecked().to_string()) - .map_err(|e| Error::InvalidHeaders(format!("{}", e)))?, - ); - headers.insert( - CONTENT_TYPE_HEADER, - HeaderValue::from_str(SSZ_CONTENT_TYPE_HEADER) - .map_err(|e| Error::InvalidHeaders(format!("{}", e)))?, - ); - headers.insert( - ACCEPT, - HeaderValue::from_str(PREFERENCE_ACCEPT_VALUE) - .map_err(|e| Error::InvalidHeaders(format!("{}", e)))?, - ); - - let result = self - .post_ssz_with_raw_response( - path, - body, - headers, - Some(self.timeouts.post_blinded_blocks), - ) - .await?; - - if result.status() == StatusCode::ACCEPTED { - Ok(()) - } else { - // ACCEPTED is the only valid status code response - Err(Error::StatusCode(result.status())) - } - } - - /// `POST /eth/v1/builder/blinded_blocks` - pub async fn post_builder_blinded_blocks_v1( - &self, - blinded_block: &SignedBlindedBeaconBlock, - ) -> Result>, Error> { - let mut path = self.server.expose_full().clone(); - - path.path_segments_mut() - .map_err(|()| Error::InvalidUrl(self.server.clone()))? - .push("eth") - .push("v1") - .push("builder") - .push("blinded_blocks"); - - let mut headers = HeaderMap::new(); - headers.insert( - CONSENSUS_VERSION_HEADER, - HeaderValue::from_str(&blinded_block.fork_name_unchecked().to_string()) - .map_err(|e| Error::InvalidHeaders(format!("{}", e)))?, - ); - headers.insert( - CONTENT_TYPE_HEADER, - HeaderValue::from_str(JSON_CONTENT_TYPE_HEADER) - .map_err(|e| Error::InvalidHeaders(format!("{}", e)))?, - ); - headers.insert( - ACCEPT, - HeaderValue::from_str(JSON_ACCEPT_VALUE) - .map_err(|e| Error::InvalidHeaders(format!("{}", e)))?, - ); - - Ok(self - .post_with_raw_response( - path, - &blinded_block, - headers, - Some(self.timeouts.post_blinded_blocks), - ) - .await? - .json() - .await?) - } - - /// `POST /eth/v2/builder/blinded_blocks` - pub async fn post_builder_blinded_blocks_v2( - &self, - blinded_block: &SignedBlindedBeaconBlock, - ) -> Result<(), Error> { - let mut path = self.server.expose_full().clone(); - - path.path_segments_mut() - .map_err(|()| Error::InvalidUrl(self.server.clone()))? - .push("eth") - .push("v2") - .push("builder") - .push("blinded_blocks"); - - let mut headers = HeaderMap::new(); - headers.insert( - CONSENSUS_VERSION_HEADER, - HeaderValue::from_str(&blinded_block.fork_name_unchecked().to_string()) - .map_err(|e| Error::InvalidHeaders(format!("{}", e)))?, - ); - headers.insert( - CONTENT_TYPE_HEADER, - HeaderValue::from_str(JSON_CONTENT_TYPE_HEADER) - .map_err(|e| Error::InvalidHeaders(format!("{}", e)))?, - ); - headers.insert( - ACCEPT, - HeaderValue::from_str(JSON_ACCEPT_VALUE) - .map_err(|e| Error::InvalidHeaders(format!("{}", e)))?, - ); - - let result = self - .post_with_raw_response( - path, - &blinded_block, - headers, - Some(self.timeouts.post_blinded_blocks), - ) - .await?; - - if result.status() == StatusCode::ACCEPTED { - Ok(()) - } else { - // ACCEPTED is the only valid status code response - Err(Error::StatusCode(result.status())) - } - } - - /// `GET /eth/v1/builder/header` - pub async fn get_builder_header( - &self, - slot: Slot, - parent_hash: ExecutionBlockHash, - pubkey: &PublicKeyBytes, - ) -> Result>>, Error> { - let mut path = self.server.expose_full().clone(); - - path.path_segments_mut() - .map_err(|()| Error::InvalidUrl(self.server.clone()))? - .push("eth") - .push("v1") - .push("builder") - .push("header") - .push(slot.to_string().as_str()) - .push(format!("{parent_hash:?}").as_str()) - .push(pubkey.as_hex_string().as_str()); - - let mut headers = HeaderMap::new(); - if self.disable_ssz { - headers.insert( - ACCEPT, - HeaderValue::from_str(JSON_CONTENT_TYPE_HEADER) - .map_err(|e| Error::InvalidHeaders(format!("{}", e)))?, - ); - } else { - // Indicate preference for ssz response in the accept header - headers.insert( - ACCEPT, - HeaderValue::from_str(PREFERENCE_ACCEPT_VALUE) - .map_err(|e| Error::InvalidHeaders(format!("{}", e)))?, - ); - } - - let resp = self - .get_with_header(path, self.timeouts.get_header, headers) - .await; - - if matches!(resp, Err(Error::StatusCode(StatusCode::NO_CONTENT))) { - Ok(None) - } else { - resp.map(Some) - } - } - - /// `GET /eth/v1/builder/status` - pub async fn get_builder_status(&self) -> Result<(), Error> { - let mut path = self.server.expose_full().clone(); - - path.path_segments_mut() - .map_err(|()| Error::InvalidUrl(self.server.clone()))? - .push("eth") - .push("v1") - .push("builder") - .push("status"); - - self.get_with_timeout(path, self.timeouts.get_builder_status) - .await - } + .transpose() } -#[cfg(test)] -mod tests { - use super::*; - use arbitrary::Arbitrary; - use bls::Signature; - use eth2::types::MainnetEthSpec; - use eth2::types::builder::{BuilderBid, BuilderBidFulu}; - use mockito::{Matcher, Server, ServerGuard}; - - type E = MainnetEthSpec; - - #[test] - fn test_headers_no_panic() { - for fork in ForkName::list_all() { - assert!(HeaderValue::from_str(&fork.to_string()).is_ok()); - } - assert!(HeaderValue::from_str(PREFERENCE_ACCEPT_VALUE).is_ok()); - assert!(HeaderValue::from_str(JSON_ACCEPT_VALUE).is_ok()); - assert!(HeaderValue::from_str(JSON_CONTENT_TYPE_HEADER).is_ok()); - } - - #[tokio::test] - async fn test_get_builder_header_ssz_response() { - // Set up mock server - let mut server = Server::new_async().await; - let mock_response_body = fulu_signed_builder_bid(); - mock_get_header_response( - &mut server, - Some("fulu"), - ContentType::Ssz, - mock_response_body.clone(), - ); - - let builder_client = BuilderHttpClient::new( - SensitiveUrl::from_str(&server.url()).unwrap(), - None, - None, - false, - ) - .unwrap(); - - let response = builder_client - .get_builder_header( - Slot::new(1), - ExecutionBlockHash::repeat_byte(1), - &PublicKeyBytes::empty(), - ) - .await - .expect("should succeed in get_builder_header") - .expect("should have response body"); - - assert_eq!(response, mock_response_body); - } - - #[tokio::test] - async fn test_get_builder_header_json_response() { - // Set up mock server - let mut server = Server::new_async().await; - let mock_response_body = fulu_signed_builder_bid(); - mock_get_header_response( - &mut server, - None, - ContentType::Json, - mock_response_body.clone(), - ); - - let builder_client = BuilderHttpClient::new( - SensitiveUrl::from_str(&server.url()).unwrap(), - None, - None, - false, - ) - .unwrap(); - - let response = builder_client - .get_builder_header( - Slot::new(1), - ExecutionBlockHash::repeat_byte(1), - &PublicKeyBytes::empty(), - ) - .await - .expect("should succeed in get_builder_header") - .expect("should have response body"); - - assert_eq!(response, mock_response_body); - } - - #[tokio::test] - async fn test_get_builder_header_no_version_header_fallback_json() { - // Set up mock server - let mut server = Server::new_async().await; - let mock_response_body = fulu_signed_builder_bid(); - mock_get_header_response( - &mut server, - Some("fulu"), - ContentType::Json, - mock_response_body.clone(), - ); - - let builder_client = BuilderHttpClient::new( - SensitiveUrl::from_str(&server.url()).unwrap(), - None, - None, - false, - ) - .unwrap(); - - let response = builder_client - .get_builder_header( - Slot::new(1), - ExecutionBlockHash::repeat_byte(1), - &PublicKeyBytes::empty(), - ) - .await - .expect("should succeed in get_builder_header") - .expect("should have response body"); - - assert_eq!(response, mock_response_body); - } - - fn mock_get_header_response( - server: &mut ServerGuard, - header_version_opt: Option<&str>, - content_type: ContentType, - response_body: ForkVersionedResponse>, - ) { - let mut mock = server.mock( - "GET", - Matcher::Regex(r"^/eth/v1/builder/header/\d+/.+/.+$".to_string()), - ); - - if let Some(version) = header_version_opt { - mock = mock.with_header(CONSENSUS_VERSION_HEADER, version); - } - - match content_type { - ContentType::Json => { - mock = mock - .with_header(CONTENT_TYPE_HEADER, JSON_CONTENT_TYPE_HEADER) - .with_body(serde_json::to_string(&response_body).unwrap()); - } - ContentType::Ssz => { - mock = mock - .with_header(CONTENT_TYPE_HEADER, SSZ_CONTENT_TYPE_HEADER) - .with_body(response_body.data.as_ssz_bytes()); - } - } - - mock.with_status(200).create(); - } - - fn fulu_signed_builder_bid() -> ForkVersionedResponse> { - let mut u = types::test_utils::test_unstructured(); - ForkVersionedResponse { - version: ForkName::Fulu, - metadata: EmptyMetadata {}, - data: SignedBuilderBid { - message: BuilderBid::Fulu(BuilderBidFulu::arbitrary(&mut u).unwrap()), - signature: Signature::empty(), - }, - } +/// Determine the `ContentType` of a response from its `Content-Type` header. +/// +/// Defaults to JSON when the header is absent or unrecognized. +pub fn content_type_from_header(headers: &HeaderMap) -> ContentType { + match headers + .get(CONTENT_TYPE_HEADER) + .and_then(|content_type| content_type.to_str().ok()) + { + Some(SSZ_CONTENT_TYPE_HEADER) => ContentType::Ssz, + _ => ContentType::Json, } } diff --git a/beacon_node/builder_client/src/pre_gloas_builder_http_client.rs b/beacon_node/builder_client/src/pre_gloas_builder_http_client.rs new file mode 100644 index 00000000000..75428a0247b --- /dev/null +++ b/beacon_node/builder_client/src/pre_gloas_builder_http_client.rs @@ -0,0 +1,676 @@ +use crate::{ + DEFAULT_USER_AGENT, JSON_ACCEPT_VALUE, PREFERENCE_ACCEPT_VALUE, content_type_from_header, + fork_name_from_header, +}; +use bls::PublicKeyBytes; +// The pre-Gloas builder client keeps the beacon-node API client's error type, unlike the Gloas +// `BuilderHttpClient` which has its own builder-spec-aligned `crate::Error`. +use context_deserialize::ContextDeserialize; +use eth2::Error; +use eth2::types::beacon_response::EmptyMetadata; +use eth2::types::builder::SignedBuilderBid; +use eth2::types::{ + ContentType, EthSpec, ExecutionBlockHash, ForkName, ForkVersionDecode, ForkVersionedResponse, + SignedValidatorRegistrationData, Slot, +}; +use eth2::types::{FullPayloadContents, SignedBlindedBeaconBlock}; +use eth2::{ + CONSENSUS_VERSION_HEADER, CONTENT_TYPE_HEADER, JSON_CONTENT_TYPE_HEADER, + SSZ_CONTENT_TYPE_HEADER, ok_or_error, success_or_error, +}; +use reqwest::header::{ACCEPT, HeaderMap, HeaderValue}; +use reqwest::{IntoUrl, Response, StatusCode}; +use sensitive_url::SensitiveUrl; +use serde::Serialize; +use serde::de::DeserializeOwned; +use ssz::Encode; +use std::sync::Arc; +use std::sync::atomic::{AtomicBool, Ordering}; +use std::time::Duration; + +/// Default timeout for builder requests without a more specific timeout. +pub const DEFAULT_TIMEOUT_MILLIS: u64 = 15000; + +/// This timeout is in accordance with v0.2.0 of the [builder specs](https://github.com/flashbots/mev-boost/pull/20). +pub const DEFAULT_GET_HEADER_TIMEOUT_MILLIS: u64 = 1000; + +#[derive(Clone)] +pub struct Timeouts { + get_header: Duration, + post_validators: Duration, + post_blinded_blocks: Duration, + get_builder_status: Duration, +} + +impl Timeouts { + fn new(get_header_timeout: Option) -> Self { + let get_header = + get_header_timeout.unwrap_or(Duration::from_millis(DEFAULT_GET_HEADER_TIMEOUT_MILLIS)); + + Self { + get_header, + post_validators: Duration::from_millis(DEFAULT_TIMEOUT_MILLIS), + post_blinded_blocks: Duration::from_millis(DEFAULT_TIMEOUT_MILLIS), + get_builder_status: Duration::from_millis(DEFAULT_TIMEOUT_MILLIS), + } + } +} + +#[derive(Clone)] +pub struct PreGloasBuilderHttpClient { + client: reqwest::Client, + server: SensitiveUrl, + timeouts: Timeouts, + user_agent: String, + /// Only use json for all requests/responses types. + disable_ssz: bool, + /// Indicates that the `get_header` response had content-type ssz + /// so we can set content-type header to ssz to make the `submit_blinded_blocks` + /// request. + ssz_available: Arc, +} + +impl PreGloasBuilderHttpClient { + pub fn new( + server: SensitiveUrl, + user_agent: Option, + builder_header_timeout: Option, + disable_ssz: bool, + ) -> Result { + let user_agent = user_agent.unwrap_or(DEFAULT_USER_AGENT.to_string()); + let client = reqwest::Client::builder().user_agent(&user_agent).build()?; + Ok(Self { + client, + server, + timeouts: Timeouts::new(builder_header_timeout), + user_agent, + disable_ssz, + ssz_available: Arc::new(false.into()), + }) + } + + pub fn get_user_agent(&self) -> &str { + &self.user_agent + } + + async fn get_with_header< + T: DeserializeOwned + ForkVersionDecode + for<'de> ContextDeserialize<'de, ForkName>, + U: IntoUrl, + >( + &self, + url: U, + timeout: Duration, + headers: HeaderMap, + ) -> Result, Error> { + let response = self + .get_response_with_header(url, Some(timeout), headers) + .await?; + + let headers = response.headers().clone(); + let response_bytes = response.bytes().await?; + + let Ok(Some(fork_name)) = fork_name_from_header(&headers) else { + // if no fork version specified, attempt to fallback to JSON + self.ssz_available.store(false, Ordering::SeqCst); + return serde_json::from_slice(&response_bytes).map_err(Error::InvalidJson); + }; + + let content_type = content_type_from_header(&headers); + + match content_type { + ContentType::Ssz => { + self.ssz_available.store(true, Ordering::SeqCst); + T::from_ssz_bytes_by_fork(&response_bytes, fork_name) + .map(|data| ForkVersionedResponse { + version: fork_name, + metadata: EmptyMetadata {}, + data, + }) + .map_err(Error::InvalidSsz) + } + ContentType::Json => { + self.ssz_available.store(false, Ordering::SeqCst); + serde_json::from_slice(&response_bytes).map_err(Error::InvalidJson) + } + } + } + + /// Return `true` if the most recently received response from the builder had SSZ Content-Type. + /// Return `false` otherwise. + /// Also returns `false` if we have explicitly disabled ssz. + pub fn is_ssz_available(&self) -> bool { + !self.disable_ssz && self.ssz_available.load(Ordering::SeqCst) + } + + async fn get_with_timeout( + &self, + url: U, + timeout: Duration, + ) -> Result { + self.get_response_with_timeout(url, Some(timeout)) + .await? + .json() + .await + .map_err(Into::into) + } + + /// Perform a HTTP GET request, returning the `Response` for further processing. + async fn get_response_with_header( + &self, + url: U, + timeout: Option, + headers: HeaderMap, + ) -> Result { + let mut builder = self.client.get(url); + if let Some(timeout) = timeout { + builder = builder.timeout(timeout); + } + let response = builder.headers(headers).send().await.map_err(Error::from)?; + ok_or_error(response).await + } + + /// Perform a HTTP GET request, returning the `Response` for further processing. + async fn get_response_with_timeout( + &self, + url: U, + timeout: Option, + ) -> Result { + let mut builder = self.client.get(url); + if let Some(timeout) = timeout { + builder = builder.timeout(timeout); + } + let response = builder.send().await.map_err(Error::from)?; + ok_or_error(response).await + } + + /// Generic POST function supporting arbitrary responses and timeouts. + async fn post_generic( + &self, + url: U, + body: &T, + timeout: Option, + ) -> Result { + let mut builder = self.client.post(url); + if let Some(timeout) = timeout { + builder = builder.timeout(timeout); + } + let response = builder.json(body).send().await?; + ok_or_error(response).await + } + + async fn post_ssz_with_raw_response( + &self, + url: U, + ssz_body: Vec, + headers: HeaderMap, + timeout: Option, + ) -> Result { + let mut builder = self.client.post(url); + if let Some(timeout) = timeout { + builder = builder.timeout(timeout); + } + + let response = builder + .headers(headers) + .body(ssz_body) + .send() + .await + .map_err(Error::from)?; + success_or_error(response).await + } + + async fn post_with_raw_response( + &self, + url: U, + body: &T, + headers: HeaderMap, + timeout: Option, + ) -> Result { + let mut builder = self.client.post(url); + if let Some(timeout) = timeout { + builder = builder.timeout(timeout); + } + + let response = builder + .headers(headers) + .json(body) + .send() + .await + .map_err(Error::from)?; + success_or_error(response).await + } + + /// `POST /eth/v1/builder/validators` + pub async fn post_builder_validators( + &self, + validator: &[SignedValidatorRegistrationData], + ) -> Result<(), Error> { + let mut path = self.server.expose_full().clone(); + + path.path_segments_mut() + .map_err(|()| Error::InvalidUrl(self.server.clone()))? + .push("eth") + .push("v1") + .push("builder") + .push("validators"); + + self.post_generic(path, &validator, Some(self.timeouts.post_validators)) + .await?; + Ok(()) + } + + /// `POST /eth/v1/builder/blinded_blocks` with SSZ serialized request body + pub async fn post_builder_blinded_blocks_v1_ssz( + &self, + blinded_block: &SignedBlindedBeaconBlock, + ) -> Result, Error> { + let mut path = self.server.expose_full().clone(); + + let body = blinded_block.as_ssz_bytes(); + + path.path_segments_mut() + .map_err(|()| Error::InvalidUrl(self.server.clone()))? + .push("eth") + .push("v1") + .push("builder") + .push("blinded_blocks"); + + let mut headers = HeaderMap::new(); + headers.insert( + CONSENSUS_VERSION_HEADER, + HeaderValue::from_str(&blinded_block.fork_name_unchecked().to_string()) + .map_err(|e| Error::InvalidHeaders(format!("{}", e)))?, + ); + headers.insert( + CONTENT_TYPE_HEADER, + HeaderValue::from_str(SSZ_CONTENT_TYPE_HEADER) + .map_err(|e| Error::InvalidHeaders(format!("{}", e)))?, + ); + headers.insert( + ACCEPT, + HeaderValue::from_str(PREFERENCE_ACCEPT_VALUE) + .map_err(|e| Error::InvalidHeaders(format!("{}", e)))?, + ); + + let result = self + .post_ssz_with_raw_response( + path, + body, + headers, + Some(self.timeouts.post_blinded_blocks), + ) + .await? + .bytes() + .await?; + + FullPayloadContents::from_ssz_bytes_by_fork(&result, blinded_block.fork_name_unchecked()) + .map_err(Error::InvalidSsz) + } + + /// `POST /eth/v2/builder/blinded_blocks` with SSZ serialized request body + pub async fn post_builder_blinded_blocks_v2_ssz( + &self, + blinded_block: &SignedBlindedBeaconBlock, + ) -> Result<(), Error> { + let mut path = self.server.expose_full().clone(); + + let body = blinded_block.as_ssz_bytes(); + + path.path_segments_mut() + .map_err(|()| Error::InvalidUrl(self.server.clone()))? + .push("eth") + .push("v2") + .push("builder") + .push("blinded_blocks"); + + let mut headers = HeaderMap::new(); + headers.insert( + CONSENSUS_VERSION_HEADER, + HeaderValue::from_str(&blinded_block.fork_name_unchecked().to_string()) + .map_err(|e| Error::InvalidHeaders(format!("{}", e)))?, + ); + headers.insert( + CONTENT_TYPE_HEADER, + HeaderValue::from_str(SSZ_CONTENT_TYPE_HEADER) + .map_err(|e| Error::InvalidHeaders(format!("{}", e)))?, + ); + headers.insert( + ACCEPT, + HeaderValue::from_str(PREFERENCE_ACCEPT_VALUE) + .map_err(|e| Error::InvalidHeaders(format!("{}", e)))?, + ); + + let result = self + .post_ssz_with_raw_response( + path, + body, + headers, + Some(self.timeouts.post_blinded_blocks), + ) + .await?; + + if result.status() == StatusCode::ACCEPTED { + Ok(()) + } else { + // ACCEPTED is the only valid status code response + Err(Error::StatusCode(result.status())) + } + } + + /// `POST /eth/v1/builder/blinded_blocks` + pub async fn post_builder_blinded_blocks_v1( + &self, + blinded_block: &SignedBlindedBeaconBlock, + ) -> Result>, Error> { + let mut path = self.server.expose_full().clone(); + + path.path_segments_mut() + .map_err(|()| Error::InvalidUrl(self.server.clone()))? + .push("eth") + .push("v1") + .push("builder") + .push("blinded_blocks"); + + let mut headers = HeaderMap::new(); + headers.insert( + CONSENSUS_VERSION_HEADER, + HeaderValue::from_str(&blinded_block.fork_name_unchecked().to_string()) + .map_err(|e| Error::InvalidHeaders(format!("{}", e)))?, + ); + headers.insert( + CONTENT_TYPE_HEADER, + HeaderValue::from_str(JSON_CONTENT_TYPE_HEADER) + .map_err(|e| Error::InvalidHeaders(format!("{}", e)))?, + ); + headers.insert( + ACCEPT, + HeaderValue::from_str(JSON_ACCEPT_VALUE) + .map_err(|e| Error::InvalidHeaders(format!("{}", e)))?, + ); + + Ok(self + .post_with_raw_response( + path, + &blinded_block, + headers, + Some(self.timeouts.post_blinded_blocks), + ) + .await? + .json() + .await?) + } + + /// `POST /eth/v2/builder/blinded_blocks` + pub async fn post_builder_blinded_blocks_v2( + &self, + blinded_block: &SignedBlindedBeaconBlock, + ) -> Result<(), Error> { + let mut path = self.server.expose_full().clone(); + + path.path_segments_mut() + .map_err(|()| Error::InvalidUrl(self.server.clone()))? + .push("eth") + .push("v2") + .push("builder") + .push("blinded_blocks"); + + let mut headers = HeaderMap::new(); + headers.insert( + CONSENSUS_VERSION_HEADER, + HeaderValue::from_str(&blinded_block.fork_name_unchecked().to_string()) + .map_err(|e| Error::InvalidHeaders(format!("{}", e)))?, + ); + headers.insert( + CONTENT_TYPE_HEADER, + HeaderValue::from_str(JSON_CONTENT_TYPE_HEADER) + .map_err(|e| Error::InvalidHeaders(format!("{}", e)))?, + ); + headers.insert( + ACCEPT, + HeaderValue::from_str(JSON_ACCEPT_VALUE) + .map_err(|e| Error::InvalidHeaders(format!("{}", e)))?, + ); + + let result = self + .post_with_raw_response( + path, + &blinded_block, + headers, + Some(self.timeouts.post_blinded_blocks), + ) + .await?; + + if result.status() == StatusCode::ACCEPTED { + Ok(()) + } else { + // ACCEPTED is the only valid status code response + Err(Error::StatusCode(result.status())) + } + } + + /// `GET /eth/v1/builder/header` + pub async fn get_builder_header( + &self, + slot: Slot, + parent_hash: ExecutionBlockHash, + pubkey: &PublicKeyBytes, + ) -> Result>>, Error> { + let mut path = self.server.expose_full().clone(); + + path.path_segments_mut() + .map_err(|()| Error::InvalidUrl(self.server.clone()))? + .push("eth") + .push("v1") + .push("builder") + .push("header") + .push(slot.to_string().as_str()) + .push(format!("{parent_hash:?}").as_str()) + .push(pubkey.as_hex_string().as_str()); + + let mut headers = HeaderMap::new(); + if self.disable_ssz { + headers.insert( + ACCEPT, + HeaderValue::from_str(JSON_CONTENT_TYPE_HEADER) + .map_err(|e| Error::InvalidHeaders(format!("{}", e)))?, + ); + } else { + // Indicate preference for ssz response in the accept header + headers.insert( + ACCEPT, + HeaderValue::from_str(PREFERENCE_ACCEPT_VALUE) + .map_err(|e| Error::InvalidHeaders(format!("{}", e)))?, + ); + } + + let resp = self + .get_with_header(path, self.timeouts.get_header, headers) + .await; + + if matches!(resp, Err(Error::StatusCode(StatusCode::NO_CONTENT))) { + Ok(None) + } else { + resp.map(Some) + } + } + + /// `GET /eth/v1/builder/status` + pub async fn get_builder_status(&self) -> Result<(), Error> { + let mut path = self.server.expose_full().clone(); + + path.path_segments_mut() + .map_err(|()| Error::InvalidUrl(self.server.clone()))? + .push("eth") + .push("v1") + .push("builder") + .push("status"); + + self.get_with_timeout(path, self.timeouts.get_builder_status) + .await + } +} + +#[cfg(test)] +mod tests { + use super::*; + use arbitrary::Arbitrary; + use bls::Signature; + use eth2::types::MainnetEthSpec; + use eth2::types::builder::{BuilderBid, BuilderBidFulu}; + use mockito::{Matcher, Server, ServerGuard}; + use std::str::FromStr; + + type E = MainnetEthSpec; + + #[test] + fn test_headers_no_panic() { + for fork in ForkName::list_all() { + assert!(HeaderValue::from_str(&fork.to_string()).is_ok()); + } + assert!(HeaderValue::from_str(PREFERENCE_ACCEPT_VALUE).is_ok()); + assert!(HeaderValue::from_str(JSON_ACCEPT_VALUE).is_ok()); + assert!(HeaderValue::from_str(JSON_CONTENT_TYPE_HEADER).is_ok()); + } + + #[tokio::test] + async fn test_get_builder_header_ssz_response() { + // Set up mock server + let mut server = Server::new_async().await; + let mock_response_body = fulu_signed_builder_bid(); + mock_get_header_response( + &mut server, + Some("fulu"), + ContentType::Ssz, + mock_response_body.clone(), + ); + + let builder_client = PreGloasBuilderHttpClient::new( + SensitiveUrl::from_str(&server.url()).unwrap(), + None, + None, + false, + ) + .unwrap(); + + let response = builder_client + .get_builder_header( + Slot::new(1), + ExecutionBlockHash::repeat_byte(1), + &PublicKeyBytes::empty(), + ) + .await + .expect("should succeed in get_builder_header") + .expect("should have response body"); + + assert_eq!(response, mock_response_body); + } + + #[tokio::test] + async fn test_get_builder_header_json_response() { + // Set up mock server + let mut server = Server::new_async().await; + let mock_response_body = fulu_signed_builder_bid(); + mock_get_header_response( + &mut server, + None, + ContentType::Json, + mock_response_body.clone(), + ); + + let builder_client = PreGloasBuilderHttpClient::new( + SensitiveUrl::from_str(&server.url()).unwrap(), + None, + None, + false, + ) + .unwrap(); + + let response = builder_client + .get_builder_header( + Slot::new(1), + ExecutionBlockHash::repeat_byte(1), + &PublicKeyBytes::empty(), + ) + .await + .expect("should succeed in get_builder_header") + .expect("should have response body"); + + assert_eq!(response, mock_response_body); + } + + #[tokio::test] + async fn test_get_builder_header_no_version_header_fallback_json() { + // Set up mock server + let mut server = Server::new_async().await; + let mock_response_body = fulu_signed_builder_bid(); + mock_get_header_response( + &mut server, + Some("fulu"), + ContentType::Json, + mock_response_body.clone(), + ); + + let builder_client = PreGloasBuilderHttpClient::new( + SensitiveUrl::from_str(&server.url()).unwrap(), + None, + None, + false, + ) + .unwrap(); + + let response = builder_client + .get_builder_header( + Slot::new(1), + ExecutionBlockHash::repeat_byte(1), + &PublicKeyBytes::empty(), + ) + .await + .expect("should succeed in get_builder_header") + .expect("should have response body"); + + assert_eq!(response, mock_response_body); + } + + fn mock_get_header_response( + server: &mut ServerGuard, + header_version_opt: Option<&str>, + content_type: ContentType, + response_body: ForkVersionedResponse>, + ) { + let mut mock = server.mock( + "GET", + Matcher::Regex(r"^/eth/v1/builder/header/\d+/.+/.+$".to_string()), + ); + + if let Some(version) = header_version_opt { + mock = mock.with_header(CONSENSUS_VERSION_HEADER, version); + } + + match content_type { + ContentType::Json => { + mock = mock + .with_header(CONTENT_TYPE_HEADER, JSON_CONTENT_TYPE_HEADER) + .with_body(serde_json::to_string(&response_body).unwrap()); + } + ContentType::Ssz => { + mock = mock + .with_header(CONTENT_TYPE_HEADER, SSZ_CONTENT_TYPE_HEADER) + .with_body(response_body.data.as_ssz_bytes()); + } + } + + mock.with_status(200).create(); + } + + fn fulu_signed_builder_bid() -> ForkVersionedResponse> { + let mut u = types::test_utils::test_unstructured(); + ForkVersionedResponse { + version: ForkName::Fulu, + metadata: EmptyMetadata {}, + data: SignedBuilderBid { + message: BuilderBid::Fulu(BuilderBidFulu::arbitrary(&mut u).unwrap()), + signature: Signature::empty(), + }, + } + } +} diff --git a/beacon_node/execution_layer/Cargo.toml b/beacon_node/execution_layer/Cargo.toml index 0d90cdaf2f2..a04e065aa9f 100644 --- a/beacon_node/execution_layer/Cargo.toml +++ b/beacon_node/execution_layer/Cargo.toml @@ -11,7 +11,7 @@ alloy-rlp = { workspace = true } alloy-rpc-types-eth = { workspace = true } arc-swap = "1.6.0" bls = { workspace = true } -builder_client = { path = "../builder_client" } +builder_client = { workspace = true } bytes = { workspace = true } eth2 = { workspace = true, features = ["events", "lighthouse", "network"] } ethereum_serde_utils = { workspace = true } diff --git a/beacon_node/execution_layer/src/engine_api.rs b/beacon_node/execution_layer/src/engine_api.rs index 3aff96c9b15..048e232d567 100644 --- a/beacon_node/execution_layer/src/engine_api.rs +++ b/beacon_node/execution_layer/src/engine_api.rs @@ -65,7 +65,6 @@ pub enum Error { DeserializeWithdrawals(ssz_types::Error), DeserializeDepositRequests(ssz_types::Error), DeserializeWithdrawalRequests(ssz_types::Error), - BuilderApi(builder_client::Error), IncorrectStateVariant, RequiredMethodUnsupported(&'static str), UnsupportedForkVariant(String), @@ -98,12 +97,6 @@ impl From for Error { } } -impl From for Error { - fn from(e: builder_client::Error) -> Self { - Error::BuilderApi(e) - } -} - impl From for Error { fn from(e: ssz_types::Error) -> Self { Error::SszError(e) diff --git a/beacon_node/execution_layer/src/engines.rs b/beacon_node/execution_layer/src/engines.rs index aac170d48c1..bc1516a4b89 100644 --- a/beacon_node/execution_layer/src/engines.rs +++ b/beacon_node/execution_layer/src/engines.rs @@ -115,7 +115,6 @@ struct PayloadIdCacheKey { pub enum EngineError { Offline, Api { error: EngineApiError }, - BuilderApi { error: EngineApiError }, Auth, } diff --git a/beacon_node/execution_layer/src/lib.rs b/beacon_node/execution_layer/src/lib.rs index 239ffdb4e60..5c94a5fd65a 100644 --- a/beacon_node/execution_layer/src/lib.rs +++ b/beacon_node/execution_layer/src/lib.rs @@ -10,7 +10,7 @@ use arc_swap::ArcSwapOption; use auth::{Auth, JwtKey, strip_prefix}; pub use block_hash::calculate_execution_block_hash; use bls::{PublicKeyBytes, Signature}; -use builder_client::BuilderHttpClient; +use builder_client::PreGloasBuilderHttpClient; pub use engine_api::EngineCapabilities; use engine_api::Error as ApiError; pub use engine_api::*; @@ -138,7 +138,8 @@ pub enum Error { NoEngine, NoPayloadBuilder, ApiError(ApiError), - Builder(builder_client::Error), + // The pre-Gloas builder client uses the beacon-node API client's error type. + Builder(eth2::Error), NoHeaderFromBuilder, CannotProduceHeader, EngineError(Box), @@ -464,7 +465,7 @@ type PayloadContentsRefTuple<'a, E> = (ExecutionPayloadRef<'a, E>, Option<&'a Bl struct Inner { engine: Arc, - builder: ArcSwapOption, + builder: ArcSwapOption, execution_engine_forkchoice_lock: Mutex<()>, suggested_fee_recipient: Option
, proposer_preparation_data: Mutex>, @@ -603,7 +604,7 @@ impl ExecutionLayer { &self.inner.engine } - pub fn builder(&self) -> Option> { + pub fn builder(&self) -> Option> { self.inner.builder.load_full() } @@ -618,7 +619,7 @@ impl ExecutionLayer { builder_header_timeout: Option, disable_ssz: bool, ) -> Result<(), Error> { - let builder_client = BuilderHttpClient::new( + let builder_client = PreGloasBuilderHttpClient::new( builder_url.clone(), builder_user_agent, builder_header_timeout, @@ -1045,11 +1046,11 @@ impl ExecutionLayer { /// Fetches local and builder paylaods concurrently, Logs and returns results. async fn fetch_builder_and_local_payloads( &self, - builder: &BuilderHttpClient, + builder: &PreGloasBuilderHttpClient, builder_params: &BuilderParams, payload_parameters: PayloadParameters<'_>, ) -> ( - Result>>, builder_client::Error>, + Result>>, eth2::Error>, Result, Error>, ) { let slot = builder_params.slot;