diff --git a/Cargo.lock b/Cargo.lock index 3e9744e92c5..05964c2bacc 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -3392,9 +3392,9 @@ dependencies = [ [[package]] name = "ethereum_serde_utils" -version = "0.8.0" +version = "0.8.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "3dc1355dbb41fbbd34ec28d4fb2a57d9a70c67ac3c19f6a5ca4d4a176b9e997a" +checksum = "38df44a7a271ab43835678f9215b53cc2523e4714a215da6643d83dc110245da" dependencies = [ "alloy-primitives", "hex", diff --git a/beacon_node/beacon_chain/src/fetch_blobs/fetch_blobs_beacon_adapter.rs b/beacon_node/beacon_chain/src/fetch_blobs/fetch_blobs_beacon_adapter.rs index 6a4d94fc5d4..db4d956e3f7 100644 --- a/beacon_node/beacon_chain/src/fetch_blobs/fetch_blobs_beacon_adapter.rs +++ b/beacon_node/beacon_chain/src/fetch_blobs/fetch_blobs_beacon_adapter.rs @@ -4,7 +4,9 @@ use crate::observed_data_sidecars::ObservationKey; use crate::partial_data_column_assembler::PartialDataColumnAssembler; use crate::pending_payload_cache::{Availability, PendingPayloadCache}; use crate::{AvailabilityProcessingStatus, BeaconChain, BeaconChainTypes}; -use execution_layer::json_structures::{BlobAndProofV2, BlobAndProofV3}; +use execution_layer::json_structures::{ + BlobAndProofV2, BlobAndProofV3, CustodyColumnsBitArray, GetBlobsV4List, +}; use kzg::Kzg; #[cfg(test)] use mockall::automock; @@ -81,6 +83,23 @@ impl FetchBlobsBeaconAdapter { .map_err(FetchEngineBlobError::RequestFailed) } + pub(crate) async fn get_blobs_v4( + &self, + versioned_hashes: Vec, + custody_columns: CustodyColumnsBitArray, + ) -> Result>, FetchEngineBlobError> { + let execution_layer = self + .chain + .execution_layer + .as_ref() + .ok_or(FetchEngineBlobError::ExecutionLayerMissing)?; + + execution_layer + .get_blobs_v4(versioned_hashes, custody_columns) + .await + .map_err(FetchEngineBlobError::RequestFailed) + } + pub(crate) fn data_column_known_for_observation_key( &self, observation_key: ObservationKey, @@ -144,4 +163,18 @@ impl FetchBlobsBeaconAdapter { .map_err(FetchEngineBlobError::RequestFailed) .map(|caps| caps.get_blobs_v3) } + + pub(crate) async fn supports_get_blobs_v4(&self) -> Result { + let execution_layer = self + .chain + .execution_layer + .as_ref() + .ok_or(FetchEngineBlobError::ExecutionLayerMissing)?; + + execution_layer + .get_engine_capabilities(None) + .await + .map_err(FetchEngineBlobError::RequestFailed) + .map(|caps| caps.get_blobs_v4) + } } diff --git a/beacon_node/beacon_chain/src/fetch_blobs/mod.rs b/beacon_node/beacon_chain/src/fetch_blobs/mod.rs index 93fb4115771..7c97553baba 100644 --- a/beacon_node/beacon_chain/src/fetch_blobs/mod.rs +++ b/beacon_node/beacon_chain/src/fetch_blobs/mod.rs @@ -24,17 +24,25 @@ use crate::{ metrics, }; use execution_layer::Error as ExecutionLayerError; -use execution_layer::json_structures::{BlobAndProofV2, BlobAndProofV3}; +use execution_layer::json_structures::{ + BlobAndProofV2, BlobAndProofV3, BlobCellsAndProofsV1, ColumnIndexTooHighError, + CustodyColumnsBitArray, +}; use metrics::{TryExt, inc_counter}; #[cfg(test)] use mockall_double::double; +use ssz_types::{ProgressiveVariableList, VariableList}; use state_processing::per_block_processing::deneb::kzg_commitment_to_versioned_hash; use std::sync::Arc; use tracing::{debug, instrument, warn}; -use types::data::{BlobSidecarError, ColumnIndex, DataColumnSidecarError}; +use types::data::{ + BlobSidecarError, CellBitmap, ColumnIndex, DataColumnSidecarError, PartialDataColumn, + PartialDataColumnFulu, PartialDataColumnGloas, PartialDataColumnHeader, + PartialDataColumnSidecarFulu, PartialDataColumnSidecarGloas, +}; use types::{ AbstractExecPayload, BeaconStateError, EthSpec, Hash256, KzgCommitment, ListRef, - PartialDataColumnHeader, SignedBeaconBlock, SignedExecutionPayloadBid, Slot, VersionedHash, + SignedBeaconBlock, SignedExecutionPayloadBid, Slot, VersionedHash, }; /// The source of the KZG commitments for a block's partial data columns. @@ -143,15 +151,30 @@ async fn fetch_and_process_engine_blobs_inner( .spec() .is_peer_das_enabled_for_epoch(header_or_bid.slot().epoch(T::EthSpec::slots_per_epoch())) { - fetch_and_process_blobs_v2_or_v3( - chain_adapter, - block_root, - header_or_bid, - versioned_hashes, - custody_columns, - publish_fn, - ) - .await + // `engine_getBlobsV4` lets us request only the columns we custody and assemble partial + // columns directly from the cells the EL returns. It supports both the Fulu partial-header + // path and the Gloas bid path; we fall back to V2/V3 only when the EL lacks the capability. + if chain_adapter.supports_get_blobs_v4().await? { + fetch_and_process_blobs_v4( + chain_adapter, + block_root, + header_or_bid, + versioned_hashes, + custody_columns, + publish_fn, + ) + .await + } else { + fetch_and_process_blobs_v2_or_v3( + chain_adapter, + block_root, + header_or_bid, + versioned_hashes, + custody_columns, + publish_fn, + ) + .await + } } else { Err(FetchEngineBlobError::InternalError( "fetch blobs v1 no longer supported".to_owned(), @@ -169,7 +192,6 @@ async fn fetch_and_process_blobs_v2_or_v3( publish_fn: impl Fn(Vec>) + Send + 'static, ) -> Result, FetchEngineBlobError> { let num_expected_blobs = versioned_hashes.len(); - let slot = header_or_bid.slot(); metrics::observe(&metrics::BLOBS_FROM_EL_EXPECTED, num_expected_blobs as f64); @@ -270,7 +292,33 @@ async fn fetch_and_process_blobs_v2_or_v3( return Ok(None); } - let availability_processing_status = match &header_or_bid { + let availability_processing_status = import_custody_partial_columns( + &chain_adapter, + block_root, + &header_or_bid, + custody_columns_to_import, + publish_fn, + ) + .await?; + + Ok(Some(availability_processing_status)) +} + +/// Merge the deduplicated custody partial columns into the appropriate fork-specific cache, publish +/// any newly-completed columns, and return the resulting availability status. +/// +/// This is the shared tail for both the V2/V3 and V4 fetch-blobs paths. The Fulu (partial-header) +/// path goes through the [`PartialDataColumnAssembler`](crate::partial_data_column_assembler), while +/// the Gloas (bid) path goes through the [`PendingPayloadCache`](crate::pending_payload_cache). +async fn import_custody_partial_columns( + chain_adapter: &Arc>, + block_root: Hash256, + header_or_bid: &PartialHeaderOrBid, + custody_columns_to_import: Vec>, + publish_fn: impl Fn(Vec>) + Send + 'static, +) -> Result { + let slot = header_or_bid.slot(); + match header_or_bid { PartialHeaderOrBid::PartialHeader(header) => { let custody_columns_to_import: Vec<_> = custody_columns_to_import .into_iter() @@ -307,9 +355,11 @@ async fn fetch_and_process_blobs_v2_or_v3( chain_adapter .process_engine_blobs_fulu(slot, block_root, full_columns) - .await? + .await } else { - AvailabilityProcessingStatus::MissingComponents(slot, block_root) + Ok(AvailabilityProcessingStatus::MissingComponents( + slot, block_root, + )) } } PartialHeaderOrBid::Bid(bid) => { @@ -340,13 +390,265 @@ async fn fetch_and_process_blobs_v2_or_v3( chain_adapter .process_payload_envelope_availability(slot, availability) - .await? + .await } + } +} + +/// EIP-8070 `engine_getBlobsV4` path: request only the columns we custody from +/// the EL and assemble `PartialDataColumn`s directly from the cells it returns, +/// skipping the local KZG-cell derivation that V2/V3 require. +/// +/// Works for both the Fulu (partial-header) and Gloas (bid) paths; the fork-specific assembly and +/// import happen in [`build_partial_columns_from_v4_response`] and [`import_custody_partial_columns`]. +#[instrument(skip_all, level = "debug")] +async fn fetch_and_process_blobs_v4( + chain_adapter: FetchBlobsBeaconAdapter, + block_root: Hash256, + header_or_bid: PartialHeaderOrBid, + versioned_hashes: Vec, + custody_columns_indices: &[ColumnIndex], + publish_fn: impl Fn(Vec>) + Send + 'static, +) -> Result, FetchEngineBlobError> { + let num_expected_blobs = versioned_hashes.len(); + + metrics::observe(&metrics::BLOBS_FROM_EL_EXPECTED, num_expected_blobs as f64); + inc_counter(&metrics::BEACON_ENGINE_GET_BLOBS_V4_REQUESTS_TOTAL); + let _timer = + metrics::start_timer(&metrics::BEACON_ENGINE_GET_BLOBS_V4_REQUEST_DURATION_SECONDS); + + let bitarray = CustodyColumnsBitArray::try_from(custody_columns_indices).map_err( + |ColumnIndexTooHighError(idx)| { + FetchEngineBlobError::InternalError(format!( + "Column index {} is too high for the getBlobsV4 bitmap", + idx + )) + }, + )?; + + debug!( + num_expected_blobs, + num_columns = custody_columns_indices.len(), + "Fetching blob cells from the EL via V4" + ); + + let response = chain_adapter + .get_blobs_v4(versioned_hashes, bitarray) + .await + .inspect_err(|_| { + inc_counter(&metrics::BLOBS_FROM_EL_ERROR_TOTAL); + })?; + + let Some(response) = response else { + warn!( + num_expected_blobs, + "engine_getBlobsV4 returned null, EL might be syncing" + ); + inc_counter(&metrics::BLOBS_FROM_EL_ERROR_TOTAL); + return Ok(None); }; + if response.len() != num_expected_blobs { + warn!( + response_len = response.len(), + num_expected_blobs, "engine_getBlobsV4 returned the wrong number of blob entries" + ); + inc_counter(&metrics::BLOBS_FROM_EL_ERROR_TOTAL); + return Ok(None); + } + + // Count present (non-null) cells across the response. + let total_cells_expected = num_expected_blobs.saturating_mul(custody_columns_indices.len()); + let total_cells_present: usize = response + .iter() + .flatten() + .flat_map(|cells_and_proofs| cells_and_proofs.blob_cells.iter()) + .filter(|c| c.is_some()) + .count(); + + if total_cells_present == 0 { + debug!(num_expected_blobs, "No cells fetched from the EL"); + inc_counter(&metrics::BLOBS_FROM_EL_MISS_TOTAL); + return Ok(None); + } else if total_cells_present == total_cells_expected { + debug!(total_cells_present, "All requested cells received from EL"); + inc_counter(&metrics::BLOBS_FROM_EL_HIT_TOTAL); + inc_counter(&metrics::BEACON_ENGINE_GET_BLOBS_V4_COMPLETE_RESPONSES_TOTAL); + } else { + debug!( + total_cells_present, + total_cells_expected, "Cells partially received from the EL" + ); + inc_counter(&metrics::BEACON_ENGINE_GET_BLOBS_V4_PARTIAL_RESPONSES_TOTAL); + } + + if chain_adapter.fork_choice_contains_block(&block_root) { + debug!( + info = "block has already been imported", + "Ignoring EL blobs response" + ); + return Ok(None); + } + + let chain_adapter = Arc::new(chain_adapter); + let custody_columns_to_import = build_partial_columns_from_v4_response( + &chain_adapter, + block_root, + &header_or_bid, + response, + custody_columns_indices, + ) + .await?; + + if custody_columns_to_import.is_empty() { + debug!( + info = "No new data columns to import", + "Ignoring EL blobs response" + ); + return Ok(None); + } + + let availability_processing_status = import_custody_partial_columns( + &chain_adapter, + block_root, + &header_or_bid, + custody_columns_to_import, + publish_fn, + ) + .await?; + Ok(Some(availability_processing_status)) } +/// Group the per-blob cells/proofs returned by `engine_getBlobsV4` into one `PartialDataColumn` per +/// custody column index. +/// +/// Builds the Fulu or Gloas partial-column variant depending on `header_or_bid`, then dedupes +/// against columns we've already observed on gossip or cached in the data availability checker. The +/// returned columns are known to be custody columns (we only request custody indices from the EL), +/// but are wrapped in [`KzgVerifiedPartialDataColumn`] so the shared +/// [`import_custody_partial_columns`] tail can assert custody and split by fork. +async fn build_partial_columns_from_v4_response( + chain_adapter: &Arc>, + block_root: Hash256, + header_or_bid: &PartialHeaderOrBid, + response: Vec>>, + custody_columns_indices: &[ColumnIndex], +) -> Result>, FetchEngineBlobError> { + let num_blobs = response.len(); + let num_columns = custody_columns_indices.len(); + let slot = header_or_bid.slot(); + let mut sorted_column_indices = custody_columns_indices.to_vec(); + sorted_column_indices.sort_unstable(); + + let mut custody_columns: Vec> = + Vec::with_capacity(num_columns); + for (col_pos, &column_index) in sorted_column_indices.iter().enumerate() { + let mut bitmap = CellBitmap::::with_capacity(num_blobs).map_err(|_| { + FetchEngineBlobError::InternalError("failed to allocate cell bitmap".to_string()) + })?; + let mut cells = Vec::with_capacity(num_blobs); + let mut proofs = Vec::with_capacity(num_blobs); + + for (blob_idx, blob) in response.iter().enumerate() { + let Some(blob) = blob else { + continue; + }; + let cell = blob.blob_cells.get(col_pos).and_then(|c| c.as_ref()); + let proof = blob.proofs.get(col_pos).and_then(|p| p.as_ref()); + + match (cell, proof) { + (Some(cell), Some(proof)) => { + bitmap.set(blob_idx, true).map_err(|_| { + FetchEngineBlobError::InternalError("unexpected oob bitmap set".to_string()) + })?; + cells.push(cell.0.clone()); + proofs.push(*proof); + } + (Some(_), None) => { + return Err(FetchEngineBlobError::InternalError( + "engine_getBlobsV4 entry has a cell but no proof".to_string(), + )); + } + (None, Some(_)) => { + return Err(FetchEngineBlobError::InternalError( + "engine_getBlobsV4 entry has a proof but no cell".to_string(), + )); + } + (None, None) => {} + } + } + if cells.is_empty() { + continue; + } + + // The cells and proofs are fork-independent; only the wrapping partial-column variant + // differs (Gloas carries the slot in the partial and omits the inline header). + let partial = match header_or_bid { + PartialHeaderOrBid::PartialHeader(_) => { + let column = VariableList::try_from(cells).map_err(|_| { + FetchEngineBlobError::InternalError("unexpectedly many cells".to_string()) + })?; + let kzg_proofs = VariableList::try_from(proofs).map_err(|_| { + FetchEngineBlobError::InternalError("unexpectedly many proofs".to_string()) + })?; + + PartialDataColumn::Fulu(PartialDataColumnFulu { + block_root, + index: column_index, + sidecar: PartialDataColumnSidecarFulu:: { + cells_present_bitmap: bitmap, + column, + kzg_proofs, + header: None.into(), + }, + }) + } + PartialHeaderOrBid::Bid(_) => { + let column = ProgressiveVariableList::new(cells); + let kzg_proofs = ProgressiveVariableList::new(proofs); + + PartialDataColumn::Gloas(PartialDataColumnGloas { + block_root, + slot, + index: column_index, + sidecar: PartialDataColumnSidecarGloas:: { + cells_present_bitmap: bitmap, + column, + kzg_proofs, + }, + }) + } + }; + custody_columns.push(KzgVerifiedPartialDataColumn::from_execution_verified( + partial, + )); + } + + // Dedupe against gossip-observed columns. + let observation_key = match header_or_bid { + PartialHeaderOrBid::PartialHeader(header) => ObservationKey::new_proposer_key( + header.signed_block_header.message.proposer_index, + header.slot(), + ), + PartialHeaderOrBid::Bid(bid) => { + ObservationKey::new_block_root_key(block_root, bid.message.slot) + } + }; + if let Some(observed_columns) = + chain_adapter.data_column_known_for_observation_key(observation_key) + { + custody_columns.retain(|col| !observed_columns.contains(&col.index())); + } + + // Dedupe against DA-checker-cached columns. + if let Some(known_columns) = chain_adapter.cached_data_column_indexes(&block_root, slot) { + custody_columns.retain(|col| !known_columns.contains(&col.index())); + } + + Ok(custody_columns) +} + /// Offload the data column computation to a blocking task to avoid holding up the async runtime. async fn compute_custody_columns_to_import( chain_adapter: &Arc>, diff --git a/beacon_node/beacon_chain/src/fetch_blobs/tests.rs b/beacon_node/beacon_chain/src/fetch_blobs/tests.rs index f5e4a33c929..10795562159 100644 --- a/beacon_node/beacon_chain/src/fetch_blobs/tests.rs +++ b/beacon_node/beacon_chain/src/fetch_blobs/tests.rs @@ -27,7 +27,7 @@ mod get_blobs_v2 { #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn test_fetch_blobs_v2_no_blobs_in_block() { - let mut mock_adapter = mock_beacon_adapter(ForkName::Fulu, false); + let mut mock_adapter = mock_beacon_adapter(ForkName::Fulu); let (publish_fn, _s) = mock_publish_fn(); let block = SignedBeaconBlock::::Fulu(SignedBeaconBlockFulu { message: BeaconBlockFulu::empty(mock_adapter.spec()), @@ -55,7 +55,7 @@ mod get_blobs_v2 { #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn test_fetch_blobs_v2_no_blobs_returned() { - let mut mock_adapter = mock_beacon_adapter(ForkName::Fulu, false); + let mut mock_adapter = mock_beacon_adapter(ForkName::Fulu); let (publish_fn, _) = mock_publish_fn(); let (block, _blobs_and_proofs) = create_test_block_and_blobs(&mock_adapter, 2); let block_root = block.canonical_root(); @@ -82,7 +82,7 @@ mod get_blobs_v2 { #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn test_fetch_blobs_v2_partial_blobs_returned() { - let mut mock_adapter = mock_beacon_adapter(ForkName::Fulu, false); + let mut mock_adapter = mock_beacon_adapter(ForkName::Fulu); let (publish_fn, publish_fn_args) = mock_publish_fn(); let (block, mut blobs_and_proofs) = create_test_block_and_blobs(&mock_adapter, 2); let block_root = block.canonical_root(); @@ -117,7 +117,7 @@ mod get_blobs_v2 { #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn test_fetch_blobs_v2_block_imported_after_el_response() { - let mut mock_adapter = mock_beacon_adapter(ForkName::Fulu, false); + let mut mock_adapter = mock_beacon_adapter(ForkName::Fulu); let (publish_fn, publish_fn_args) = mock_publish_fn(); let (block, blobs_and_proofs) = create_test_block_and_blobs(&mock_adapter, 2); let block_root = block.canonical_root(); @@ -152,7 +152,7 @@ mod get_blobs_v2 { #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn test_fetch_blobs_v2_no_new_columns_to_import() { - let mut mock_adapter = mock_beacon_adapter(ForkName::Fulu, false); + let mut mock_adapter = mock_beacon_adapter(ForkName::Fulu); let (publish_fn, publish_fn_args) = mock_publish_fn(); let (block, blobs_and_proofs) = create_test_block_and_blobs(&mock_adapter, 2); let block_root = block.canonical_root(); @@ -194,7 +194,7 @@ mod get_blobs_v2 { #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn test_fetch_blobs_v2_success() { - let mut mock_adapter = mock_beacon_adapter(ForkName::Fulu, false); + let mut mock_adapter = mock_beacon_adapter(ForkName::Fulu); let (publish_fn, publish_fn_args) = mock_publish_fn(); let (block, blobs_and_proofs) = create_test_block_and_blobs(&mock_adapter, 2); let block_root = block.canonical_root(); @@ -267,6 +267,330 @@ mod get_blobs_v2 { } } +mod get_blobs_v4 { + use super::*; + use crate::custody_context::{CustodyContext, NodeCustodyType}; + use crate::pending_payload_cache::PendingPayloadCache; + use crate::test_utils::{ + NumBlobs, generate_data_column_indices_rand_order, generate_rand_block_and_blobs, + }; + use execution_layer::json_structures::{BlobCellsAndProofsV1, GetBlobsV4List, JsonCell}; + use kzg::KzgProof; + use slot_clock::{SlotClock, TestingSlotClock}; + use std::time::Duration; + use types::test_utils::test_unstructured; + use types::{Cell, ColumnIndex, PartialDataColumnHeader, Slot}; + + const CUSTODY_COLUMNS: [ColumnIndex; 3] = [0, 1, 2]; + + /// Build an `engine_getBlobsV4` response for `num_blobs` blobs, with one entry per requested + /// custody column (positionally, lowest column index first). `present(blob_idx, col_pos)` + /// decides whether that cell/proof is present. Cell contents are irrelevant: the V4 path trusts + /// the EL and never re-verifies the returned cells. + fn make_v4_response( + num_blobs: usize, + num_custody_cols: usize, + present: impl Fn(usize, usize) -> bool, + ) -> Option> { + let list = (0..num_blobs) + .map(|blob_idx| { + let blob_cells = (0..num_custody_cols) + .map(|col_pos| { + present(blob_idx, col_pos).then(|| JsonCell(Cell::::default())) + }) + .collect(); + let proofs = (0..num_custody_cols) + .map(|col_pos| present(blob_idx, col_pos).then(KzgProof::empty)) + .collect(); + Some(BlobCellsAndProofsV1 { blob_cells, proofs }) + }) + .collect(); + Some(list) + } + + fn fulu_header(block: &SignedBeaconBlock) -> PartialHeaderOrBid { + PartialHeaderOrBid::PartialHeader(Arc::new( + PartialDataColumnHeader::try_from(block).unwrap(), + )) + } + + /// A complete V4 response assembles into full columns that are published and imported. + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn test_fetch_blobs_v4_fulu_success() { + let mut mock_adapter = mock_beacon_adapter_v4(ForkName::Fulu); + let (publish_fn, publish_fn_args) = mock_publish_fn(); + let (block, _blobs) = create_test_block_and_blobs(&mock_adapter, 2); + let block_root = block.canonical_root(); + + let response = make_v4_response(2, CUSTODY_COLUMNS.len(), |_, _| true); + mock_adapter + .expect_get_blobs_v4() + .return_once(move |_, _| Ok(response)); + mock_fork_choice_contains_block(&mut mock_adapter, vec![]); + mock_adapter + .expect_data_column_known_for_observation_key() + .returning(|_| None); + mock_adapter + .expect_cached_data_column_indexes() + .returning(|_, _| None); + mock_process_engine_blobs_result( + &mut mock_adapter, + Ok(AvailabilityProcessingStatus::Imported( + block.slot(), + block_root, + )), + ); + + let processing_status = fetch_and_process_engine_blobs_inner( + mock_adapter, + block_root, + fulu_header(block.as_ref()), + &CUSTODY_COLUMNS, + publish_fn, + ) + .await + .expect("fetch blobs should succeed"); + + assert_eq!( + processing_status, + Some(AvailabilityProcessingStatus::Imported( + block.slot(), + block_root + )) + ); + assert_eq!( + extract_published_blobs(publish_fn_args).len(), + CUSTODY_COLUMNS.len(), + "all custody columns should be published" + ); + } + + /// A partial V4 response that leaves every column short of a cell yields no full columns: nothing + /// is published or imported, and we report missing components. + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn test_fetch_blobs_v4_fulu_partial_response_completes_nothing() { + let mut mock_adapter = mock_beacon_adapter_v4(ForkName::Fulu); + let (publish_fn, publish_fn_args) = mock_publish_fn(); + let (block, _blobs) = create_test_block_and_blobs(&mock_adapter, 2); + let block_root = block.canonical_root(); + + // Only the first blob's cells are present, so each column is missing blob 1's cell. + let response = make_v4_response(2, CUSTODY_COLUMNS.len(), |blob_idx, _| blob_idx == 0); + mock_adapter + .expect_get_blobs_v4() + .return_once(move |_, _| Ok(response)); + mock_fork_choice_contains_block(&mut mock_adapter, vec![]); + mock_adapter + .expect_data_column_known_for_observation_key() + .returning(|_| None); + mock_adapter + .expect_cached_data_column_indexes() + .returning(|_, _| None); + mock_adapter.expect_process_engine_blobs_fulu().times(0); + + let processing_status = fetch_and_process_engine_blobs_inner( + mock_adapter, + block_root, + fulu_header(block.as_ref()), + &CUSTODY_COLUMNS, + publish_fn, + ) + .await + .expect("fetch blobs should succeed"); + + assert_eq!( + processing_status, + Some(AvailabilityProcessingStatus::MissingComponents( + block.slot(), + block_root + )) + ); + assert_eq!( + publish_fn_args.lock().unwrap().len(), + 0, + "no columns should be published" + ); + } + + /// An empty (all-null) V4 response is treated as a miss: nothing imported or published. + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn test_fetch_blobs_v4_no_cells_returned() { + let mut mock_adapter = mock_beacon_adapter_v4(ForkName::Fulu); + let (publish_fn, publish_fn_args) = mock_publish_fn(); + let (block, _blobs) = create_test_block_and_blobs(&mock_adapter, 2); + let block_root = block.canonical_root(); + + let response = make_v4_response(2, CUSTODY_COLUMNS.len(), |_, _| false); + mock_adapter + .expect_get_blobs_v4() + .return_once(move |_, _| Ok(response)); + mock_adapter.expect_process_engine_blobs_fulu().times(0); + + let processing_status = fetch_and_process_engine_blobs_inner( + mock_adapter, + block_root, + fulu_header(block.as_ref()), + &CUSTODY_COLUMNS, + publish_fn, + ) + .await + .expect("fetch blobs should succeed"); + + assert_eq!(processing_status, None); + assert_eq!(publish_fn_args.lock().unwrap().len(), 0); + } + + /// A response whose blob count disagrees with the block's commitment count is rejected. + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn test_fetch_blobs_v4_wrong_response_length() { + let mut mock_adapter = mock_beacon_adapter_v4(ForkName::Fulu); + let (publish_fn, publish_fn_args) = mock_publish_fn(); + let (block, _blobs) = create_test_block_and_blobs(&mock_adapter, 2); + let block_root = block.canonical_root(); + + // The block expects 2 blobs, but the EL returns only 1 entry. + let response = make_v4_response(1, CUSTODY_COLUMNS.len(), |_, _| true); + mock_adapter + .expect_get_blobs_v4() + .return_once(move |_, _| Ok(response)); + mock_adapter.expect_process_engine_blobs_fulu().times(0); + + let processing_status = fetch_and_process_engine_blobs_inner( + mock_adapter, + block_root, + fulu_header(block.as_ref()), + &CUSTODY_COLUMNS, + publish_fn, + ) + .await + .expect("fetch blobs should succeed"); + + assert_eq!(processing_status, None); + assert_eq!(publish_fn_args.lock().unwrap().len(), 0); + } + + /// Columns already observed on gossip are deduped away, leaving nothing to import. + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn test_fetch_blobs_v4_all_columns_already_observed() { + let mut mock_adapter = mock_beacon_adapter_v4(ForkName::Fulu); + let (publish_fn, publish_fn_args) = mock_publish_fn(); + let (block, _blobs) = create_test_block_and_blobs(&mock_adapter, 2); + let block_root = block.canonical_root(); + + let response = make_v4_response(2, CUSTODY_COLUMNS.len(), |_, _| true); + mock_adapter + .expect_get_blobs_v4() + .return_once(move |_, _| Ok(response)); + mock_fork_choice_contains_block(&mut mock_adapter, vec![]); + mock_adapter + .expect_data_column_known_for_observation_key() + .returning(|_| Some(hashset![0, 1, 2])); + mock_adapter + .expect_cached_data_column_indexes() + .returning(|_, _| None); + mock_adapter.expect_process_engine_blobs_fulu().times(0); + + let processing_status = fetch_and_process_engine_blobs_inner( + mock_adapter, + block_root, + fulu_header(block.as_ref()), + &CUSTODY_COLUMNS, + publish_fn, + ) + .await + .expect("fetch blobs should succeed"); + + assert_eq!(processing_status, None); + assert_eq!(publish_fn_args.lock().unwrap().len(), 0); + } + + /// The Gloas (bid) path also flows through V4: a complete response merges into the pending + /// payload cache, publishes the completed columns, and processes envelope availability. This is + /// the regression test for the previously-missing Gloas support in the V4 path. + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn test_fetch_blobs_v4_gloas_success() { + let mut mock_adapter = mock_beacon_adapter_v4(ForkName::Gloas); + let spec = mock_adapter.spec().clone(); + let (publish_fn, publish_fn_args) = mock_publish_fn(); + + // A Gloas block carries its blob commitments in the execution payload bid. + let mut u = test_unstructured(); + let (block, _blobs) = + generate_rand_block_and_blobs::(ForkName::Gloas, NumBlobs::Number(2), &mut u) + .expect("generate gloas block"); + let block_root = block.canonical_root(); + let bid = Arc::new( + block + .message() + .body() + .signed_execution_payload_bid() + .expect("gloas block has a bid") + .clone(), + ); + + // Real pending payload cache: the Gloas path inserts the bid and merges partial columns into + // it, so a mock would not exercise the actual merge/availability logic. + let slot_clock = TestingSlotClock::new( + Slot::new(0), + Duration::from_secs(0), + spec.get_slot_duration(), + ); + let custody_context = Arc::new(CustodyContext::::new( + NodeCustodyType::Supernode, + generate_data_column_indices_rand_order::(), + slot_clock, + false, + spec.clone(), + )); + let cache = Arc::new( + PendingPayloadCache::::new(get_kzg(&spec), custody_context, false, 0, spec.clone()) + .expect("create pending payload cache"), + ); + mock_adapter + .expect_pending_payload_cache() + .return_const(cache); + + let response = make_v4_response(2, CUSTODY_COLUMNS.len(), |_, _| true); + mock_adapter + .expect_get_blobs_v4() + .return_once(move |_, _| Ok(response)); + mock_fork_choice_contains_block(&mut mock_adapter, vec![]); + mock_adapter + .expect_data_column_known_for_observation_key() + .returning(|_| None); + mock_adapter + .expect_cached_data_column_indexes() + .returning(|_, _| None); + mock_adapter + .expect_process_payload_envelope_availability() + .return_once(move |slot, _availability| { + Ok(AvailabilityProcessingStatus::MissingComponents( + slot, block_root, + )) + }); + + let processing_status = fetch_and_process_engine_blobs_inner( + mock_adapter, + block_root, + PartialHeaderOrBid::Bid(bid), + &CUSTODY_COLUMNS, + publish_fn, + ) + .await + .expect("fetch blobs should succeed"); + + assert!( + processing_status.is_some(), + "the gloas path should return an availability status" + ); + assert_eq!( + extract_published_blobs(publish_fn_args).len(), + CUSTODY_COLUMNS.len(), + "all three custody columns should complete and be published" + ); + } +} + /// Extract the `Vec>` passed to the `publish_fn`. fn extract_published_blobs( publish_fn_args: Arc>>>>, @@ -352,7 +676,20 @@ fn mock_publish_fn() -> ( (publish_fn, captured_args) } -fn mock_beacon_adapter(fork_name: ForkName, get_blobs_v3: bool) -> MockFetchBlobsBeaconAdapter { +fn mock_beacon_adapter(fork_name: ForkName) -> MockFetchBlobsBeaconAdapter { + mock_beacon_adapter_with_capabilities(fork_name, false) +} + +/// Like [`mock_beacon_adapter`], but advertises `engine_getBlobsV4` support so the V4 fetch path is +/// exercised instead of V2/V3. +fn mock_beacon_adapter_v4(fork_name: ForkName) -> MockFetchBlobsBeaconAdapter { + mock_beacon_adapter_with_capabilities(fork_name, true) +} + +fn mock_beacon_adapter_with_capabilities( + fork_name: ForkName, + supports_get_blobs_v4: bool, +) -> MockFetchBlobsBeaconAdapter { let test_runtime = TestRuntime::default(); let spec = Arc::new(fork_name.make_genesis_spec(E::default_spec())); let kzg = get_kzg(&spec); @@ -366,7 +703,10 @@ fn mock_beacon_adapter(fork_name: ForkName, get_blobs_v3: bool) -> MockFetchBlob .return_const(test_runtime.task_executor.clone()); mock_adapter .expect_supports_get_blobs_v3() - .returning(move || Ok(get_blobs_v3)); + .returning(move || Ok(false)); + mock_adapter + .expect_supports_get_blobs_v4() + .returning(move || Ok(supports_get_blobs_v4)); mock_adapter .expect_partial_assembler() .return_const(Some(Arc::new(partial_assembler))); diff --git a/beacon_node/beacon_chain/src/metrics.rs b/beacon_node/beacon_chain/src/metrics.rs index 04fc668cd71..6a3acfd635c 100644 --- a/beacon_node/beacon_chain/src/metrics.rs +++ b/beacon_node/beacon_chain/src/metrics.rs @@ -1896,6 +1896,38 @@ pub static BEACON_ENGINE_GET_BLOBS_V3_REQUEST_DURATION_SECONDS: LazyLock> = + LazyLock::new(|| { + try_create_int_counter( + "beacon_engine_getBlobsV4_requests_total", + "Total number of engine_getBlobsV4 requests made to the execution layer", + ) + }); + +pub static BEACON_ENGINE_GET_BLOBS_V4_COMPLETE_RESPONSES_TOTAL: LazyLock> = + LazyLock::new(|| { + try_create_int_counter( + "beacon_engine_getBlobsV4_complete_responses_total", + "Total number of engine_getBlobsV4 responses with all requested cells present", + ) + }); + +pub static BEACON_ENGINE_GET_BLOBS_V4_PARTIAL_RESPONSES_TOTAL: LazyLock> = + LazyLock::new(|| { + try_create_int_counter( + "beacon_engine_getBlobsV4_partial_responses_total", + "Total number of engine_getBlobsV4 responses with at least one missing cell", + ) + }); + +pub static BEACON_ENGINE_GET_BLOBS_V4_REQUEST_DURATION_SECONDS: LazyLock> = + LazyLock::new(|| { + try_create_histogram( + "beacon_engine_getBlobsV4_request_duration_seconds", + "Duration of engine_getBlobsV4 requests to the execution layer in seconds", + ) + }); + /* * Standardized metrics for partial column efficiency */ diff --git a/beacon_node/beacon_chain/src/payload_bid_verification/direct_verified_bid.rs b/beacon_node/beacon_chain/src/payload_bid_verification/direct_verified_bid.rs index 9983df131d8..f57d5628695 100644 --- a/beacon_node/beacon_chain/src/payload_bid_verification/direct_verified_bid.rs +++ b/beacon_node/beacon_chain/src/payload_bid_verification/direct_verified_bid.rs @@ -1,6 +1,6 @@ use crate::payload_bid_verification::{ PayloadBidError, - gossip_verified_bid::{is_gas_limit_target_compatible, verify_bid_consistency}, + gossip_verified_bid::{is_gas_limit_target_compatible, verify_direct_bid_consistency}, }; use eth2::types::BuilderPubkeys; use state_processing::signature_sets::{ @@ -14,7 +14,8 @@ use types::{ /// Fully validate a bid fetched directly from a builder, for inclusion in a block being produced. /// /// This performs all validation a direct builder bid must pass before it can be selected: -/// - the consensus-consistency checks shared with the gossip verifier via [`verify_bid_consistency`] +/// - the consensus-consistency checks shared with the gossip verifier, bundled for this path in +/// [`verify_direct_bid_consistency`] /// (fee recipient, blob count, builder eligibility/version, and that the builder's collateral /// covers the bid value), /// - that the bid matches the block being produced — the exact `proposal_slot`, the selected @@ -81,7 +82,7 @@ pub fn verify_direct_bid( } // Consensus-consistency checks shared with the gossip verifier. - verify_bid_consistency(bid, proposal_slot, proposer_preferences, state, spec)?; + verify_direct_bid_consistency(bid, proposal_slot, proposer_preferences, state, spec)?; // If the requesting `BuilderEntry` named builder pubkeys, the bid must come from one of them: // the builder at `bid.builder_index` must have one of those pubkeys (the `builder_pubkeys` @@ -267,6 +268,38 @@ mod tests { )); } + #[test] + fn rejects_block_hash_equal_to_parent_block_hash() { + let (state, spec) = state_and_spec(); + // Passes every earlier check (slot, ancestor hash, parent root, RANDAO, gas limit), then + // claims a `block_hash` equal to its `parent_block_hash` — the consensus assert from + // `process_execution_payload_bid` that must be front-run before selection. + let executed_ancestor = ExecutionBlockHash::repeat_byte(7); + let mut bid = signed_bid( + Slot::new(1), + executed_ancestor, + Hash256::ZERO, + Hash256::ZERO, + ); + bid.message.block_hash = executed_ancestor; + bid.message.gas_limit = EXECUTED_ANCESTOR_GAS_LIMIT; + let result = verify_direct_bid( + &bid, + Slot::new(1), + executed_ancestor, + Hash256::ZERO, + EXECUTED_ANCESTOR_GAS_LIMIT, + &BuilderPubkeys::default(), + &preferences(), + &state, + &spec, + ); + assert!(matches!( + result, + Err(PayloadBidError::BlockHashEqualsParentBlockHash { .. }) + )); + } + #[test] fn rejects_gas_limit_incompatible_with_parent() { let (state, spec) = state_and_spec(); @@ -305,6 +338,9 @@ mod tests { Hash256::ZERO, ); bid.message.gas_limit = EXECUTED_ANCESTOR_GAS_LIMIT; + // A default (zero) `block_hash` would equal the zero parent hash and trip the + // block-hash-equals-parent rejection before the checks this test targets. + bid.message.block_hash = ExecutionBlockHash::repeat_byte(1); let result = verify_direct_bid( &bid, Slot::new(1), diff --git a/beacon_node/beacon_chain/src/payload_bid_verification/gossip_verified_bid.rs b/beacon_node/beacon_chain/src/payload_bid_verification/gossip_verified_bid.rs index 25d82ccc971..851bf821de7 100644 --- a/beacon_node/beacon_chain/src/payload_bid_verification/gossip_verified_bid.rs +++ b/beacon_node/beacon_chain/src/payload_bid_verification/gossip_verified_bid.rs @@ -44,14 +44,26 @@ fn verify_bid_payment_and_blobs( }); } + verify_bid_block_hash_not_parent(bid)?; + + verify_bid_blobs(bid, spec) +} + +/// Reject a bid whose `block_hash` equals its `parent_block_hash`. +/// +/// `process_execution_payload_bid` enforces this in `per_block_processing`, so every bid intake — +/// gossip *and* direct (builder-API) — must front-run it: a bid that fails only at block +/// processing has already won selection and costs the proposer the slot. +pub(crate) fn verify_bid_block_hash_not_parent( + bid: &ExecutionPayloadBid, +) -> Result<(), PayloadBidError> { if bid.block_hash == bid.parent_block_hash { return Err(PayloadBidError::BlockHashEqualsParentBlockHash { slot: bid.slot, block_hash: bid.block_hash, }); } - - verify_bid_blobs(bid, spec) + Ok(()) } fn verify_bid_blobs( @@ -71,12 +83,17 @@ fn verify_bid_blobs( Ok(()) } -/// Verify that an execution payload bid is consistent with the current chain state -/// and proposer preferences. +/// Verify that a direct (builder-API) bid is consistent with the current chain state +/// and proposer preferences: the direct path's bundle of the shared bid checks. /// -/// These checks are shared by gossip and direct bids. Source-specific checks (e.g. the gossip-only -/// requirement that `execution_payment == 0`) are applied by the caller. -pub(crate) fn verify_bid_consistency( +/// The individual checks are shared with gossip, but this bundle's only caller is +/// [`verify_direct_bid`](crate::payload_bid_verification::direct_verified_bid::verify_direct_bid): +/// the gossip verifier applies the same helpers (`verify_bid_slot`, `verify_bid_blobs`, +/// `verify_bid_block_hash_not_parent`, `verify_bid_state_conditions`) piecewise, in gossip-spec +/// order, interleaved with gossip-only work (cache checks, the preferences lookup, fork-choice +/// rules). A check that must cover both intakes belongs in one of those shared helpers — adding +/// it only here leaves gossip uncovered. +pub(crate) fn verify_direct_bid_consistency( bid: &ExecutionPayloadBid, current_slot: Slot, proposer_preferences: &SignedProposerPreferences, @@ -89,6 +106,11 @@ pub(crate) fn verify_bid_consistency( return Err(PayloadBidError::InvalidFeeRecipient); } + // Mirrors the consensus assert in `process_execution_payload_bid`. The gossip path applies + // this earlier (via `verify_bid_payment_and_blobs`); repeating it here keeps the direct path + // covered without depending on the gossip caller's composition. + verify_bid_block_hash_not_parent(bid)?; + verify_bid_blobs(bid, spec)?; verify_bid_state_conditions(bid, head_state, spec) diff --git a/beacon_node/builder_client/src/builder_http_client.rs b/beacon_node/builder_client/src/builder_http_client.rs index a857d13226c..12241036b3f 100644 --- a/beacon_node/builder_client/src/builder_http_client.rs +++ b/beacon_node/builder_client/src/builder_http_client.rs @@ -41,6 +41,11 @@ const DATE_MILLISECONDS: HeaderName = HeaderName::from_static("date-milliseconds #[derive(Clone)] pub struct BuilderHttpClient { client: reqwest::Client, + /// Client for `submitSignedBeaconBlock` only. The target URL arrives over the wire (the + /// `Eth-Builder-Url` request header echoed by the VC) and is an SSRF risk, so beacon-APIs + /// `publishBlock` requires that the forwarding request "MUST NOT follow redirects" — reqwest's + /// redirect policy is client-wide, hence a dedicated client with redirects disabled. + no_redirect_client: reqwest::Client, user_agent: String, /// Only use json for all request/response types. disable_ssz: bool, @@ -50,8 +55,13 @@ 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()?; + let no_redirect_client = reqwest::Client::builder() + .user_agent(&user_agent) + .redirect(reqwest::redirect::Policy::none()) + .build()?; Ok(Self { client, + no_redirect_client, user_agent, disable_ssz, }) @@ -234,6 +244,10 @@ impl BuilderHttpClient { /// /// `ssz_request` selects the request-body encoding: SSZ when `true` and the client has SSZ /// enabled, otherwise JSON. + /// + /// Sent via [`Self::no_redirect_client`]: `builder_url` is wire input (`Eth-Builder-Url`), and + /// the spec forbids following redirects on this request. A redirect response surfaces as + /// [`Error::StatusCode`] like any other non-202. pub async fn submit_signed_beacon_block( &self, builder_url: &SensitiveUrl, @@ -263,7 +277,7 @@ impl BuilderHttpClient { HeaderValue::from_str(SSZ_CONTENT_TYPE_HEADER) .map_err(|e| Error::InvalidHeaders(format!("{}", e)))?, ); - self.client + self.no_redirect_client .post(path) .timeout(timeout) .headers(headers) @@ -274,7 +288,7 @@ impl BuilderHttpClient { HeaderValue::from_str(JSON_CONTENT_TYPE_HEADER) .map_err(|e| Error::InvalidHeaders(format!("{}", e)))?, ); - self.client + self.no_redirect_client .post(path) .timeout(timeout) .headers(headers) diff --git a/beacon_node/execution_layer/src/engine_api.rs b/beacon_node/execution_layer/src/engine_api.rs index ceb5176974c..ccad2e617f3 100644 --- a/beacon_node/execution_layer/src/engine_api.rs +++ b/beacon_node/execution_layer/src/engine_api.rs @@ -1,12 +1,12 @@ use crate::engines::ForkchoiceState; use crate::http::{ ENGINE_FORKCHOICE_UPDATED_V1, ENGINE_FORKCHOICE_UPDATED_V2, ENGINE_FORKCHOICE_UPDATED_V3, - ENGINE_FORKCHOICE_UPDATED_V4, ENGINE_GET_BLOBS_V2, ENGINE_GET_CLIENT_VERSION_V1, - ENGINE_GET_INCLUSION_LIST_V1, ENGINE_GET_PAYLOAD_BODIES_BY_HASH_V1, - ENGINE_GET_PAYLOAD_BODIES_BY_HASH_V2, ENGINE_GET_PAYLOAD_V1, ENGINE_GET_PAYLOAD_V2, - ENGINE_GET_PAYLOAD_V3, ENGINE_GET_PAYLOAD_V4, ENGINE_GET_PAYLOAD_V5, ENGINE_GET_PAYLOAD_V6, - ENGINE_NEW_PAYLOAD_V1, ENGINE_NEW_PAYLOAD_V2, ENGINE_NEW_PAYLOAD_V3, ENGINE_NEW_PAYLOAD_V4, - ENGINE_NEW_PAYLOAD_V5, + ENGINE_FORKCHOICE_UPDATED_V4, ENGINE_GET_BLOBS_V2, ENGINE_GET_BLOBS_V3, ENGINE_GET_BLOBS_V4, + ENGINE_GET_CLIENT_VERSION_V1, ENGINE_GET_INCLUSION_LIST_V1, + ENGINE_GET_PAYLOAD_BODIES_BY_HASH_V1, ENGINE_GET_PAYLOAD_BODIES_BY_HASH_V2, + ENGINE_GET_PAYLOAD_V1, ENGINE_GET_PAYLOAD_V2, ENGINE_GET_PAYLOAD_V3, ENGINE_GET_PAYLOAD_V4, + ENGINE_GET_PAYLOAD_V5, ENGINE_GET_PAYLOAD_V6, ENGINE_NEW_PAYLOAD_V1, ENGINE_NEW_PAYLOAD_V2, + ENGINE_NEW_PAYLOAD_V3, ENGINE_NEW_PAYLOAD_V4, ENGINE_NEW_PAYLOAD_V5, }; use eth2::types::{ BlobsBundle, SsePayloadAttributes, SsePayloadAttributesV1, SsePayloadAttributesV2, @@ -623,6 +623,7 @@ pub struct EngineCapabilities { pub get_client_version_v1: bool, pub get_blobs_v2: bool, pub get_blobs_v3: bool, + pub get_blobs_v4: bool, pub get_inclusion_list_v1: bool, } @@ -686,6 +687,12 @@ impl EngineCapabilities { if self.get_blobs_v2 { response.push(ENGINE_GET_BLOBS_V2); } + if self.get_blobs_v3 { + response.push(ENGINE_GET_BLOBS_V3); + } + if self.get_blobs_v4 { + response.push(ENGINE_GET_BLOBS_V4); + } if self.get_inclusion_list_v1 { response.push(ENGINE_GET_INCLUSION_LIST_V1); } diff --git a/beacon_node/execution_layer/src/engine_api/http.rs b/beacon_node/execution_layer/src/engine_api/http.rs index 93ed7278aa7..cc7de356eb3 100644 --- a/beacon_node/execution_layer/src/engine_api/http.rs +++ b/beacon_node/execution_layer/src/engine_api/http.rs @@ -66,6 +66,7 @@ pub const ENGINE_GET_CLIENT_VERSION_TIMEOUT: Duration = Duration::from_secs(1); pub const ENGINE_GET_BLOBS_V2: &str = "engine_getBlobsV2"; pub const ENGINE_GET_BLOBS_V3: &str = "engine_getBlobsV3"; +pub const ENGINE_GET_BLOBS_V4: &str = "engine_getBlobsV4"; pub const ENGINE_GET_BLOBS_TIMEOUT: Duration = Duration::from_secs(1); pub const ENGINE_GET_INCLUSION_LIST_V1: &str = "engine_getInclusionListV1"; @@ -98,6 +99,7 @@ pub static LIGHTHOUSE_CAPABILITIES: &[&str] = &[ ENGINE_GET_CLIENT_VERSION_V1, ENGINE_GET_BLOBS_V2, ENGINE_GET_BLOBS_V3, + ENGINE_GET_BLOBS_V4, ENGINE_GET_INCLUSION_LIST_V1, ]; @@ -757,6 +759,21 @@ impl HttpJsonRpc { .await } + pub async fn get_blobs_v4( + &self, + versioned_hashes: Vec, + indices_bitarray: CustodyColumnsBitArray, + ) -> Result>, Error> { + let params = json!([versioned_hashes, indices_bitarray]); + + self.rpc_request( + ENGINE_GET_BLOBS_V4, + params, + ENGINE_GET_BLOBS_TIMEOUT * self.execution_timeout_multiplier, + ) + .await + } + pub async fn get_inclusion_list_v1(&self) -> Result { self.rpc_request::( ENGINE_GET_INCLUSION_LIST_V1, @@ -1258,6 +1275,7 @@ impl HttpJsonRpc { get_client_version_v1: capabilities.contains(ENGINE_GET_CLIENT_VERSION_V1), get_blobs_v2: capabilities.contains(ENGINE_GET_BLOBS_V2), get_blobs_v3: capabilities.contains(ENGINE_GET_BLOBS_V3), + get_blobs_v4: capabilities.contains(ENGINE_GET_BLOBS_V4), get_inclusion_list_v1: capabilities.contains(ENGINE_GET_INCLUSION_LIST_V1), }) } diff --git a/beacon_node/execution_layer/src/engine_api/json_structures.rs b/beacon_node/execution_layer/src/engine_api/json_structures.rs index 20bbe8a6434..25b18dc17f8 100644 --- a/beacon_node/execution_layer/src/engine_api/json_structures.rs +++ b/beacon_node/execution_layer/src/engine_api/json_structures.rs @@ -5,7 +5,7 @@ use ssz::{Decode, TryFromIter}; use ssz_types::{FixedVector, ProgressiveVariableList, VariableList, typenum::Unsigned}; use strum::EnumString; use superstruct::superstruct; -use types::data::BlobsList; +use types::data::{BlobsList, Cell, ColumnIndex}; use types::execution::{ BlockAccessList, BuilderDepositRequests, BuilderExitRequests, ConsolidationRequests, DepositRequests, ExecutionRequestsElectra, ExecutionRequestsGloas, ProgressiveTransactions, @@ -1051,6 +1051,72 @@ pub struct BlobAndProof { /// A BlobAndProofV3 is just a BlobAndProofV2 that may also be `null` if unknown by the EL. pub type BlobAndProofV3 = Option>; +/// CELLS_PER_EXT_BLOB per EIP-7594; the `custodyColumns` and `indices_bitarray` +/// EIP-8070 parameters are 128-bit bitarrays (=16 bytes). +pub const CUSTODY_COLUMNS_BITARRAY_BYTES: usize = 16; + +/// EIP-8070 - bitarray of length `CELLS_PER_EXT_BLOB` (=128). Bit `i` of +/// byte `i / 8` (LSB-first within each byte) indicates column `i`. Used as +/// the `indices_bitarray` parameter of `engine_getBlobsV4` and the +/// `custodyColumns` parameter of `engine_forkchoiceUpdatedV4`. +/// The TryFrom impl safeguards against invalid input. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] +#[serde(transparent)] +pub struct CustodyColumnsBitArray( + #[serde(with = "serde_utils::fixed_bytes_hex::bytes_16_hex")] + [u8; CUSTODY_COLUMNS_BITARRAY_BYTES], +); + +impl CustodyColumnsBitArray { + pub fn iter_set_bits(&self) -> impl Iterator + '_ { + (0..CUSTODY_COLUMNS_BITARRAY_BYTES * 8).filter_map(move |i| { + let byte = self.0[i / 8]; + ((byte >> (i % 8)) & 1 == 1).then_some(i as ColumnIndex) + }) + } +} + +#[derive(Debug, Clone, Copy)] +pub struct ColumnIndexTooHighError(pub ColumnIndex); + +impl TryFrom<&[ColumnIndex]> for CustodyColumnsBitArray { + type Error = ColumnIndexTooHighError; + + fn try_from(indices: &[ColumnIndex]) -> Result { + let mut buf = [0u8; CUSTODY_COLUMNS_BITARRAY_BYTES]; + for i in indices { + let byte_idx = *i as usize / 8; + let bit_idx = i % 8; + if byte_idx < CUSTODY_COLUMNS_BITARRAY_BYTES { + buf[byte_idx] |= 1u8 << bit_idx; + } else { + return Err(ColumnIndexTooHighError(*i)); + } + } + Ok(Self(buf)) + } +} + +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +#[serde(bound = "E: EthSpec", transparent)] +pub struct JsonCell( + #[serde(with = "ssz_types::serde_utils::hex_fixed_vec")] pub Cell, +); + +/// `blob_cells` is the partial column matrix slice for one blob, indexed +/// positionally over the bits set in the request's `indices_bitarray` +/// (lowest set bit first). An entry is `null` when the EL doesn't have +/// that cell. `proofs[i]` is the KZG cell proof for `blob_cells[i]` and +/// is only meaningful when the matching cell is `Some`. +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +#[serde(bound = "E: EthSpec")] +pub struct BlobCellsAndProofsV1 { + pub blob_cells: Vec>>, + pub proofs: Vec>, +} + +pub type GetBlobsV4List = Vec>>; + #[derive(Debug, PartialEq, Clone, Serialize, Deserialize)] #[serde(rename_all = "camelCase")] pub struct JsonForkchoiceStateV1 { diff --git a/beacon_node/execution_layer/src/lib.rs b/beacon_node/execution_layer/src/lib.rs index b829dfc97d2..0722edef4ed 100644 --- a/beacon_node/execution_layer/src/lib.rs +++ b/beacon_node/execution_layer/src/lib.rs @@ -4,7 +4,9 @@ //! This crate only provides useful functionality for "The Merge", it does not provide any of the //! deposit-contract functionality that the `beacon_node/eth1` crate already provides. -use crate::json_structures::{BlobAndProofV2, BlobAndProofV3}; +use crate::json_structures::{ + BlobAndProofV2, BlobAndProofV3, CustodyColumnsBitArray, GetBlobsV4List, +}; use crate::payload_cache::PayloadCache; use arc_swap::ArcSwapOption; use auth::{Auth, JwtKey, strip_prefix}; @@ -1779,6 +1781,26 @@ impl ExecutionLayer { } } + pub async fn get_blobs_v4( + &self, + query: Vec, + custody_columns: CustodyColumnsBitArray, + ) -> Result>, Error> { + let capabilities = self.get_engine_capabilities(None).await?; + + if capabilities.get_blobs_v4 { + self.engine() + .request( + |engine| async move { engine.api.get_blobs_v4(query, custody_columns).await }, + ) + .await + .map_err(Box::new) + .map_err(Error::EngineError) + } else { + Err(Error::GetBlobsNotSupported) + } + } + pub async fn get_inclusion_list_v1(&self) -> Result { let capabilities = self.get_engine_capabilities(None).await?; diff --git a/beacon_node/execution_layer/src/test_utils/execution_block_generator.rs b/beacon_node/execution_layer/src/test_utils/execution_block_generator.rs index aef5d5a4407..41f871a206f 100644 --- a/beacon_node/execution_layer/src/test_utils/execution_block_generator.rs +++ b/beacon_node/execution_layer/src/test_utils/execution_block_generator.rs @@ -15,7 +15,7 @@ use parking_lot::Mutex; use rand::{Rng, SeedableRng, rngs::StdRng}; use serde::{Deserialize, Serialize}; use ssz::Decode; -use ssz_types::{ProgressiveVariableList, VariableList}; +use ssz_types::{FixedVector, ProgressiveVariableList, VariableList}; use state_processing::per_block_processing::deneb::kzg_commitment_to_versioned_hash; use std::cmp::max; use std::collections::HashMap; @@ -23,6 +23,7 @@ use std::sync::Arc; use tracing::warn; use tree_hash::TreeHash; use tree_hash_derive::TreeHash; +use types::data::Cell; use types::{ Blob, ChainSpec, EthSpec, ExecutionBlockHash, ExecutionPayload, ExecutionPayloadBellatrix, ExecutionPayloadCapella, ExecutionPayloadDeneb, ExecutionPayloadElectra, ExecutionPayloadFulu, @@ -32,6 +33,9 @@ use types::{ const TEST_BLOB_BUNDLE: &[u8] = include_bytes!("fixtures/mainnet/test_blobs_bundle.ssz"); const TEST_BLOB_BUNDLE_V2: &[u8] = include_bytes!("fixtures/mainnet/test_blobs_bundle_v2.ssz"); +/// The cells of the (single) blob in `TEST_BLOB_BUNDLE_V2`, as an SSZ-encoded +/// `FixedVector, E::CellsPerExtBlob>`. +const TEST_BLOB_CELLS: &[u8] = include_bytes!("fixtures/mainnet/test_blobs_bundle_v2_cells.ssz"); pub const DEFAULT_GAS_LIMIT: u64 = 60_000_000; const GAS_USED: u64 = DEFAULT_GAS_LIMIT - 1; @@ -167,6 +171,9 @@ pub struct ExecutionBlockGenerator { * deneb stuff */ pub blobs_bundles: HashMap>, + /// The cells for each blob in `blobs_bundles`, keyed by payload id and indexed by the blob's + /// position within the bundle. Only populated for Fulu-enabled payloads. + pub blob_cells: HashMap>>>, pub kzg: Option>, rng: Arc>, /* @@ -221,6 +228,7 @@ impl ExecutionBlockGenerator { amsterdam_time, heze_time, blobs_bundles: <_>::default(), + blob_cells: <_>::default(), kzg, rng: make_rng(), execution_requests: <_>::default(), @@ -505,8 +513,8 @@ impl ExecutionBlockGenerator { self.inclusion_list = transactions; } - /// Look up a blob and proof by versioned hash across all stored bundles. - pub fn get_blob_and_proof(&self, versioned_hash: &Hash256) -> Option> { + /// Find the payload id, bundle, and blob index for a blob by versioned hash. + fn find_blob(&self, versioned_hash: &Hash256) -> Option<(&PayloadId, &BlobsBundle, usize)> { self.blobs_bundles .iter() .find_map(|(payload_id, blobs_bundle)| { @@ -518,27 +526,51 @@ impl ExecutionBlockGenerator { .find(|(_, commitment)| { &kzg_commitment_to_versioned_hash(commitment) == versioned_hash })?; - let is_fulu = self.payload_ids.get(payload_id)?.fork_name().fulu_enabled(); - let blob = blobs_bundle.blobs.get(blob_idx)?.clone(); - if is_fulu { - let start = blob_idx * E::cells_per_ext_blob(); - let end = start + E::cells_per_ext_blob(); - let proofs = blobs_bundle - .proofs - .get(start..end)? - .to_vec() - .try_into() - .ok()?; - Some(BlobAndProof::V2(BlobAndProofV2 { blob, proofs })) - } else { - Some(BlobAndProof::V1(BlobAndProofV1 { - blob, - proof: *blobs_bundle.proofs.get(blob_idx)?, - })) - } + Some((payload_id, blobs_bundle, blob_idx)) }) } + /// Slice the cell proofs for the blob at `blob_idx` out of a Fulu-style bundle. + fn cell_proofs(blobs_bundle: &BlobsBundle, blob_idx: usize) -> Option> { + let start = blob_idx * E::cells_per_ext_blob(); + let end = start + E::cells_per_ext_blob(); + blobs_bundle + .proofs + .get(start..end)? + .to_vec() + .try_into() + .ok() + } + + /// Look up a blob and proof by versioned hash across all stored bundles. + pub fn get_blob_and_proof(&self, versioned_hash: &Hash256) -> Option> { + let (payload_id, blobs_bundle, blob_idx) = self.find_blob(versioned_hash)?; + let is_fulu = self.payload_ids.get(payload_id)?.fork_name().fulu_enabled(); + let blob = blobs_bundle.blobs.get(blob_idx)?.clone(); + if is_fulu { + let proofs = Self::cell_proofs(blobs_bundle, blob_idx)?; + Some(BlobAndProof::V2(BlobAndProofV2 { blob, proofs })) + } else { + Some(BlobAndProof::V1(BlobAndProofV1 { + blob, + proof: *blobs_bundle.proofs.get(blob_idx)?, + })) + } + } + + /// Look up the cells and cell proofs for a blob by versioned hash across all stored bundles. + /// + /// Returns `None` for pre-Fulu blobs, which have no cells. + pub fn get_blob_cells_and_proofs( + &self, + versioned_hash: &Hash256, + ) -> Option<(Vec>, KzgProofs)> { + let (payload_id, blobs_bundle, blob_idx) = self.find_blob(versioned_hash)?; + let cells = self.blob_cells.get(payload_id)?.get(blob_idx)?.clone(); + let proofs = Self::cell_proofs(blobs_bundle, blob_idx)?; + Some((cells, proofs)) + } + pub fn new_payload(&mut self, payload: ExecutionPayload) -> PayloadStatusV1 { let Some(parent) = self.blocks.get(&payload.parent_hash()) else { return PayloadStatusV1 { @@ -876,6 +908,14 @@ impl ExecutionBlockGenerator { } } } + + if fork_name.fulu_enabled() { + // `generate_blobs` only produces copies of the static test blob, so the + // precomputed cells fixture applies to every blob in the bundle. + let cells = load_test_blob_cells::()?; + self.blob_cells.insert(id, vec![cells; bundle.blobs.len()]); + } + bundle } else { BlobsBundle::default() @@ -938,6 +978,15 @@ pub fn load_test_blobs_bundle_v2() )) } +/// Load the precomputed cells for the single blob in `TEST_BLOB_BUNDLE_V2`. +/// +/// Every blob served by the mock EL is a copy of that blob, so these cells apply to all of them. +pub fn load_test_blob_cells() -> Result>, String> { + FixedVector::, E::CellsPerExtBlob>::from_ssz_bytes(TEST_BLOB_CELLS) + .map(|cells| cells.to_vec()) + .map_err(|e| format!("Unable to decode ssz: {:?}", e)) +} + pub fn generate_blobs( n_blobs: usize, fork_name: ForkName, @@ -1170,4 +1219,36 @@ mod test { Kzg::new_from_trusted_setup(&get_trusted_setup()) .map_err(|e| format!("Failed to load trusted setup: {e:?}")) } + + #[test] + fn valid_test_blob_cells() { + validate_test_blob_cells::() + .expect("Mainnet preset cells fixture should match the test blob"); + validate_test_blob_cells::() + .expect("Minimal preset cells fixture should match the test blob"); + } + + fn validate_test_blob_cells() -> Result<(), String> { + let kzg = load_kzg()?; + let (_, _, blob) = load_test_blobs_bundle_v2::()?; + let kzg_blob: KzgBlobRef = blob + .as_ref() + .try_into() + .map_err(|e| format!("Error converting blob to kzg blob ref: {e:?}"))?; + let computed = kzg + .compute_cells(kzg_blob) + .map_err(|e| format!("Failed to compute cells: {e:?}"))?; + let embedded = load_test_blob_cells::()?; + if embedded.len() != computed.len() + || embedded + .iter() + .zip(computed.iter()) + .any(|(embedded_cell, computed_cell)| { + embedded_cell[..] != computed_cell.as_ref()[..] + }) + { + return Err("cells fixture does not match cells computed from the test blob".into()); + } + Ok(()) + } } diff --git a/beacon_node/execution_layer/src/test_utils/fixtures/mainnet/test_blobs_bundle_v2_cells.ssz b/beacon_node/execution_layer/src/test_utils/fixtures/mainnet/test_blobs_bundle_v2_cells.ssz new file mode 100644 index 00000000000..819cddfe472 Binary files /dev/null and b/beacon_node/execution_layer/src/test_utils/fixtures/mainnet/test_blobs_bundle_v2_cells.ssz differ diff --git a/beacon_node/execution_layer/src/test_utils/handle_rpc.rs b/beacon_node/execution_layer/src/test_utils/handle_rpc.rs index 62cf415b77d..5609f6639b7 100644 --- a/beacon_node/execution_layer/src/test_utils/handle_rpc.rs +++ b/beacon_node/execution_layer/src/test_utils/handle_rpc.rs @@ -529,6 +529,41 @@ pub async fn handle_rpc( let response: Option>> = results.into_iter().collect(); Ok(serde_json::to_value(response).unwrap()) } + ENGINE_GET_BLOBS_V4 => { + let versioned_hashes = + get_param::>(params, 0).map_err(|s| (s, BAD_PARAMS_ERROR_CODE))?; + let indices_bitarray = get_param::(params, 1) + .map_err(|s| (s, BAD_PARAMS_ERROR_CODE))?; + let requested_indices = indices_bitarray + .iter_set_bits() + .map(|i| i as usize) + .collect::>(); + + let generator = ctx.execution_block_generator.read(); + let response = versioned_hashes + .iter() + .map(|hash| { + // Pre-Fulu blobs have no cells and cannot be served over V4. + let Some((cells, cell_proofs)) = generator.get_blob_cells_and_proofs(hash) + else { + return Ok(None); + }; + let mut blob_cells = Vec::with_capacity(requested_indices.len()); + let mut proofs = Vec::with_capacity(requested_indices.len()); + for &index in &requested_indices { + let (cell, proof) = + cells.get(index).zip(cell_proofs.get(index)).ok_or(( + format!("cell index {index} out of range"), + BAD_PARAMS_ERROR_CODE, + ))?; + blob_cells.push(Some(JsonCell(cell.clone()))); + proofs.push(Some(*proof)); + } + Ok(Some(BlobCellsAndProofsV1:: { blob_cells, proofs })) + }) + .collect::, (String, i64)>>()?; + Ok(serde_json::to_value(response).unwrap()) + } ENGINE_GET_INCLUSION_LIST_V1 => { let transactions = ctx.execution_block_generator.read().get_inclusion_list(); diff --git a/beacon_node/execution_layer/src/test_utils/mod.rs b/beacon_node/execution_layer/src/test_utils/mod.rs index 67a9cbba9f4..9e553c17863 100644 --- a/beacon_node/execution_layer/src/test_utils/mod.rs +++ b/beacon_node/execution_layer/src/test_utils/mod.rs @@ -58,7 +58,10 @@ pub const DEFAULT_ENGINE_CAPABILITIES: EngineCapabilities = EngineCapabilities { get_payload_v6: true, get_client_version_v1: true, get_blobs_v2: true, - get_blobs_v3: true, + // The mock server has no `engine_getBlobsV3` handler, so it must not advertise it: nodes + // prefer the advertised version and get method-not-found errors instead of blobs. + get_blobs_v3: false, + get_blobs_v4: true, get_inclusion_list_v1: true, }; diff --git a/beacon_node/http_api/src/lib.rs b/beacon_node/http_api/src/lib.rs index 4f41c2a1a1c..e73cf06c00a 100644 --- a/beacon_node/http_api/src/lib.rs +++ b/beacon_node/http_api/src/lib.rs @@ -68,7 +68,9 @@ use eth2::types::{ self as api_types, BroadcastValidation, EndpointVersion, ForkChoice, ForkChoiceExtraData, ForkChoiceNode, LightClientUpdatesQuery, PublishBlockRequest, ValidatorId, }; -use eth2::{CONSENSUS_VERSION_HEADER, CONTENT_TYPE_HEADER, SSZ_CONTENT_TYPE_HEADER}; +use eth2::{ + BUILDER_URL_HEADER, CONSENSUS_VERSION_HEADER, CONTENT_TYPE_HEADER, SSZ_CONTENT_TYPE_HEADER, +}; use health_metrics::observe::Observe; use lighthouse_network::Enr; use lighthouse_network::NetworkGlobals; @@ -106,7 +108,7 @@ use types::{ }; use validator::execution_payload_envelopes::get_validator_execution_payload_envelopes; use version::{ - ResponseIncludesVersion, V1, V2, add_consensus_version_header, add_ssz_content_type_header, + ResponseIncludesVersion, V1, V2, V4, add_consensus_version_header, add_ssz_content_type_header, execution_optimistic_finalized_beacon_response, inconsistent_fork_rejection, unsupported_version_rejection, }; @@ -384,6 +386,7 @@ pub async fn serve( let eth_v1 = single_version(any_version.clone(), V1); let eth_v2 = single_version(any_version.clone(), V2); + let eth_v4 = single_version(any_version.clone(), V4); // Create a `warp` filter that provides access to the network globals. let inner_network_globals = ctx.network_globals.clone(); @@ -819,6 +822,9 @@ pub async fn serve( */ let consensus_version_header_filter = warp::header::header::(CONSENSUS_VERSION_HEADER).boxed(); + // The winning builder's URL echoed by the VC on a Gloas block publish (beacon-APIs #630), so the + // node forwards the block to that builder. Optional: absent for self-build / p2p-won blocks. + let builder_url_header_filter = warp::header::optional::(BUILDER_URL_HEADER).boxed(); let optional_consensus_version_header_filter = warp::header::optional::(CONSENSUS_VERSION_HEADER).boxed(); @@ -855,6 +861,8 @@ pub async fn serve( &network_tx, BroadcastValidation::default(), duplicate_block_status_code, + // Legacy v1 publish: no builder-URL provenance (VC uses v2 for Gloas). + None, ) .await }) @@ -892,6 +900,8 @@ pub async fn serve( &network_tx, BroadcastValidation::default(), duplicate_block_status_code, + // Legacy v1 publish: no builder-URL provenance (VC uses v2 for Gloas). + None, ) .await }) @@ -909,13 +919,15 @@ pub async fn serve( .and(task_spawner_filter.clone()) .and(chain_filter.clone()) .and(network_tx_filter.clone()) + .and(builder_url_header_filter.clone()) .then( move |validation_level: api_types::BroadcastValidationQuery, value: serde_json::Value, consensus_version: ForkName, task_spawner: TaskSpawner, chain: Arc>, - network_tx: UnboundedSender>| { + network_tx: UnboundedSender>, + builder_url: Option| { task_spawner.spawn_async_with_rejection(Priority::P0, async move { let request = PublishBlockRequest::::context_deserialize( &value, @@ -932,6 +944,7 @@ pub async fn serve( &network_tx, validation_level.broadcast_validation, duplicate_block_status_code, + builder_url, ) .await }) @@ -949,13 +962,15 @@ pub async fn serve( .and(task_spawner_filter.clone()) .and(chain_filter.clone()) .and(network_tx_filter.clone()) + .and(builder_url_header_filter.clone()) .then( move |validation_level: api_types::BroadcastValidationQuery, block_bytes: Bytes, consensus_version: ForkName, task_spawner: TaskSpawner, chain: Arc>, - network_tx: UnboundedSender>| { + network_tx: UnboundedSender>, + builder_url: Option| { task_spawner.spawn_async_with_rejection(Priority::P0, async move { let block_contents = PublishBlockRequest::::from_ssz_bytes( &block_bytes, @@ -971,6 +986,7 @@ pub async fn serve( &network_tx, validation_level.broadcast_validation, duplicate_block_status_code, + builder_url, ) .await }) @@ -2570,6 +2586,14 @@ pub async fn serve( task_spawner_filter.clone(), ); + // POST v4/validator/blocks/{slot} + let post_validator_blocks_v4 = post_validator_blocks_v4( + eth_v4.clone(), + chain_filter.clone(), + not_while_syncing_filter.clone(), + task_spawner_filter.clone(), + ); + // GET validator/blinded_blocks/{slot} let get_validator_blinded_blocks = get_validator_blinded_blocks( eth_v1.clone(), @@ -2683,6 +2707,12 @@ pub async fn serve( chain_filter.clone(), task_spawner_filter.clone(), ); + // POST validator/builder_preferences + let post_validator_builder_preferences = post_validator_builder_preferences( + eth_v1.clone(), + chain_filter.clone(), + task_spawner_filter.clone(), + ); // POST validator/sync_committee_subscriptions let post_validator_sync_committee_subscriptions = post_validator_sync_committee_subscriptions( eth_v1.clone(), @@ -3496,6 +3526,8 @@ pub async fn serve( .uor(post_validator_sync_committee_subscriptions) .uor(post_validator_prepare_beacon_proposer) .uor(post_validator_register_validator) + .uor(post_validator_builder_preferences) + .uor(post_validator_blocks_v4) .uor(post_validator_liveness_epoch) .uor(post_lighthouse_liveness) .uor(post_lighthouse_database_reconstruct) diff --git a/beacon_node/http_api/src/produce_block.rs b/beacon_node/http_api/src/produce_block.rs index 63420fbe2d0..49315790da0 100644 --- a/beacon_node/http_api/src/produce_block.rs +++ b/beacon_node/http_api/src/produce_block.rs @@ -1,10 +1,10 @@ use crate::{ build_block_contents, version::{ - ResponseIncludesVersion, add_consensus_block_value_header, add_consensus_version_header, - add_execution_payload_blinded_header, add_execution_payload_included_header, - add_execution_payload_value_header, add_ssz_content_type_header, beacon_response, - inconsistent_fork_rejection, + ResponseIncludesVersion, add_builder_url_header, add_consensus_block_value_header, + add_consensus_version_header, add_execution_payload_blinded_header, + add_execution_payload_included_header, add_execution_payload_value_header, + add_ssz_content_type_header, beacon_response, inconsistent_fork_rejection, }, }; use beacon_chain::graffiti_calculator::GraffitiSettings; @@ -17,9 +17,10 @@ use eth2::{ beacon_response::ForkVersionedResponse, types::{BlockAndEnvelope, ProduceBlockV4Metadata}, }; +use sensitive_url::SensitiveUrl; use ssz::Encode; use std::sync::Arc; -use tracing::instrument; +use tracing::{debug, instrument}; use types::{execution::BlockProductionVersion, *}; use warp::{ http::response::Builder, @@ -58,13 +59,30 @@ pub async fn produce_block_v4( chain: Arc>, slot: Slot, query: api_types::ValidatorBlocksQuery, + builder_config: api_types::BuilderConfig, ) -> Result { + // `produceBlockV4` is the Gloas block-production endpoint. + let fork_name = chain.spec.fork_name_at_slot::(slot); + if !fork_name.gloas_enabled() { + return Err(warp_utils::reject::custom_bad_request( + "produceBlockV4 is only valid for Gloas and later".to_string(), + )); + } + let include_payload = query.include_payload.ok_or_else(|| { warp_utils::reject::custom_bad_request( "include_payload query parameter is required".to_string(), ) })?; + // The resolved builder config is threaded into block production, where it drives direct-builder + // bid requests and the gossip/direct bid policy (see `produce_block_on_state_gloas`). + debug!( + %slot, + builders = builder_config.builders.len(), + "Received produceBlockV4 request" + ); + let randao_reveal = query.randao_reveal.decompress().map_err(|e| { warp_utils::reject::custom_bad_request(format!( "randao reveal is not a valid BLS signature: {:?}", @@ -73,14 +91,9 @@ pub async fn produce_block_v4( })?; let randao_verification = get_randao_verification(&query, randao_reveal.is_infinity())?; - // The GET route carries only a boost factor; direct builders arrive with the `BuilderConfig` - // body once this route is converted to POST (later in this PR stack). Until then the winning - // bid's builder URL is unused (`Eth-Builder-Url` also lands with the POST conversion). - let builder_config = api_types::BuilderConfig { - builder_boost_factor: query.builder_boost_factor.unwrap_or(DEFAULT_BOOST_FACTOR), - ..api_types::BuilderConfig::empty() - }; + // Gloas takes its bid boost policy from `builder_config` (global for gossip, per-builder for + // direct), so the V3-style `builder_boost_factor` query param is not used on this path. let graffiti_settings = GraffitiSettings::new(query.graffiti, query.graffiti_policy); let ( @@ -89,7 +102,7 @@ pub async fn produce_block_v4( consensus_block_value, execution_payload_value, payload_contents, - _builder_url, + builder_url, ) = chain .produce_block_with_verification_gloas( randao_reveal, @@ -110,6 +123,7 @@ pub async fn produce_block_v4( consensus_block_value, execution_payload_value, payload_contents, + builder_url, accept_header, &chain.spec, ) @@ -164,9 +178,13 @@ pub fn build_response_v4( consensus_block_value: u64, execution_payload_value: Uint256, payload_contents: Option>, + builder_url: Option, accept_header: Option, spec: &ChainSpec, ) -> Result { + // Stringify the winning builder's URL only here, at the `Eth-Builder-Url` header boundary; it is + // kept as a redacted `SensitiveUrl` everywhere upstream. + let builder_url = builder_url.map(|url| url.expose_full().to_string()); let fork_name = block .to_ref() .fork_name(spec) @@ -180,14 +198,15 @@ pub fn build_response_v4( consensus_block_value: consensus_block_value_wei, execution_payload_value, execution_payload_included, - builder_url: None, + builder_url: builder_url.clone(), }; let add_v4_headers = |res: Response| { let res = add_consensus_version_header(res, fork_name); let res = add_consensus_block_value_header(res, consensus_block_value_wei); let res = add_execution_payload_value_header(res, execution_payload_value); - add_execution_payload_included_header(res, execution_payload_included) + let res = add_execution_payload_included_header(res, execution_payload_included); + add_builder_url_header(res, builder_url.as_deref()) }; // When the payload is included, bundle the block with the execution payload envelope, blobs and diff --git a/beacon_node/http_api/src/publish_blocks.rs b/beacon_node/http_api/src/publish_blocks.rs index a7336f2f6eb..5279a25b3be 100644 --- a/beacon_node/http_api/src/publish_blocks.rs +++ b/beacon_node/http_api/src/publish_blocks.rs @@ -19,6 +19,7 @@ use logging::crit; use network::NetworkMessage; use rand::prelude::SliceRandom; use reqwest::StatusCode; +use sensitive_url::SensitiveUrl; use slot_clock::SlotClock; use std::marker::PhantomData; use std::sync::Arc; @@ -73,6 +74,62 @@ impl ProvenancedBlock> } } +/// If a direct builder won this block's payload bid, forward the signed block to that builder via +/// `submitSignedBeaconBlock` so it reveals the execution payload envelope. +/// +/// The builder's URL is the `Eth-Builder-Url` request header the VC echoed on publish (beacon-APIs +/// #630), so this works even on a beacon node that did not produce the block. `None` (self-built or +/// p2p-won), no configured builders, or a malformed URL are all no-ops. +/// +/// Fire-and-forget: the submission runs in a detached task; a failure is logged at high severity +/// (the validator has already signed the commitment) but never blocks the publish response. Runs +/// only once per block since it hangs off the single p2p-publish point. +fn forward_signed_block_to_winning_builder( + chain: &Arc>, + block: Arc>, + builder_url: Option<&str>, +) { + // The VC echoes the winning builder's URL in the `Eth-Builder-Url` request header (beacon-APIs + // #630); absent for a self-built block or a p2p-won bid, in which case there's nothing to forward. + let Some(builder_url) = builder_url else { + return; + }; + let Some(builders) = chain.builders.as_ref() else { + return; + }; + let url = match SensitiveUrl::parse(builder_url) { + Ok(url) => url, + Err(e) => { + warn!(error = ?e, "Ignoring malformed Eth-Builder-Url header"); + return; + } + }; + + let builders = builders.clone(); + let slot = block.slot(); + let block_root = block.canonical_root(); + + chain.task_executor.spawn( + async move { + match builders.forward_signed_block(&url, &block).await { + Ok(()) => info!( + %slot, + %block_root, + "Forwarded signed block to winning builder" + ), + Err(e) => error!( + %slot, + %block_root, + builder_url = ?url, + error = ?e, + "Failed to forward signed block to winning builder" + ), + } + }, + "forward_signed_block_to_builder", + ); +} + /// Handles a request from the HTTP API for full blocks. #[allow(clippy::too_many_arguments)] #[instrument( @@ -88,6 +145,9 @@ pub async fn publish_block>( network_tx: &UnboundedSender>, validation_level: BroadcastValidation, duplicate_status_code: StatusCode, + // The `Eth-Builder-Url` request header (beacon-APIs #630): when a direct builder won the block's + // payload bid, its URL, so the block is forwarded there for envelope reveal. + builder_url: Option, ) -> Result { let seen_timestamp = chain.slot_clock.now_duration().unwrap_or_default(); let block_publishing_delay_for_testing = chain.config.block_publishing_delay; @@ -141,6 +201,14 @@ pub async fn publish_block>( BlockError::BeaconChainError(Box::new(BeaconChainError::UnableToPublish)) })?; + // If a direct builder won this block's payload bid, forward the signed block to it so it + // reveals the execution payload envelope. + forward_signed_block_to_winning_builder( + &publish_chain, + block.clone(), + builder_url.as_deref(), + ); + Ok(()) }; @@ -570,6 +638,8 @@ pub async fn publish_blinded_block( network_tx, validation_level, duplicate_status_code, + // Blinded (mev-boost) publish predates the Gloas builder-URL round-trip. + None, ) .await } else { diff --git a/beacon_node/http_api/src/validator/mod.rs b/beacon_node/http_api/src/validator/mod.rs index 7b904260b8e..ecb362985e4 100644 --- a/beacon_node/http_api/src/validator/mod.rs +++ b/beacon_node/http_api/src/validator/mod.rs @@ -6,7 +6,7 @@ use crate::utils::{ AnyVersionFilter, ChainFilter, EthV1Filter, NetworkTxFilter, NotWhileSyncingFilter, ResponseFilter, TaskSpawnerFilter, ValidatorSubscriptionTxFilter, publish_network_message, }; -use crate::version::{V1, V2, V3, V4, add_ssz_content_type_header, unsupported_version_rejection}; +use crate::version::{V1, V2, V3, add_ssz_content_type_header, unsupported_version_rejection}; use crate::{StateId, attester_duties, proposer_duties, ptc_duties, sync_committees}; use beacon_chain::attestation_verification::VerifiedAttestation; use beacon_chain::proposer_preferences_verification::ProposerPreferencesError; @@ -14,12 +14,13 @@ use beacon_chain::{AttestationError, BeaconChain, BeaconChainError, BeaconChainT use bls::PublicKeyBytes; use bytes::Bytes; use context_deserialize::ContextDeserialize; -use eth2::CONSENSUS_VERSION_HEADER; use eth2::types::{ - Accept, BeaconCommitteeSubscription, EndpointVersion, Failure, GenericResponse, - StandardLivenessResponseData, StateId as CoreStateId, ValidatorAggregateAttestationQuery, - ValidatorAttestationDataQuery, ValidatorBlocksQuery, ValidatorIndexData, ValidatorStatus, + Accept, BeaconCommitteeSubscription, BuilderConfig, BuilderPreferenceEntry, EndpointVersion, + Failure, GenericResponse, MAX_SUBMITTED_BUILDER_PREFERENCES, StandardLivenessResponseData, + StateId as CoreStateId, ValidatorAggregateAttestationQuery, ValidatorAttestationDataQuery, + ValidatorBlocksQuery, ValidatorIndexData, ValidatorStatus, }; +use eth2::{CONSENSUS_VERSION_HEADER, CONTENT_TYPE_HEADER, SSZ_CONTENT_TYPE_HEADER}; use lighthouse_network::PubsubMessage; use network::{NetworkMessage, ValidatorSubscriptionMessage}; use reqwest::StatusCode; @@ -483,8 +484,12 @@ pub fn get_validator_blocks( not_synced_filter?; - if endpoint_version == V4 { - produce_block_v4(accept_header, chain, slot, query).await + // Gloas block production is served via `POST v4/validator/blocks`. + let fork_name = chain.spec.fork_name_at_slot::(slot); + if fork_name.gloas_enabled() { + Err(warp_utils::reject::custom_bad_request( + "Gloas block production requires POST v4/validator/blocks".to_string(), + )) } else if endpoint_version == V3 { produce_block_v3(accept_header, chain, slot, query).await } else { @@ -496,6 +501,95 @@ pub fn get_validator_blocks( .boxed() } +/// Does the request's `Content-Type` header select SSZ? +/// +/// Tolerates media-type parameters (`application/octet-stream; ...`) and surrounding whitespace; +/// anything else (including an absent header) selects JSON. +fn is_ssz_content_type(content_type: Option<&str>) -> bool { + content_type + .and_then(|header| header.split(';').next()) + .is_some_and(|media_type| media_type.trim() == SSZ_CONTENT_TYPE_HEADER) +} + +// POST v4/validator/blocks/{slot} +// +// The Gloas block-production endpoint. Carries the validator's resolved `BuilderConfig` as the +// request body, accepted as either JSON or SSZ (selected by `Content-Type`; `application/octet-stream` +// => SSZ). The `Eth-Consensus-Version` request header is required (per beacon-APIs #630); the body +// is not fork-versioned, so like the builder-preferences endpoint the header is validated but only +// logged. +pub fn post_validator_blocks_v4( + eth_v4: EthV1Filter, + chain_filter: ChainFilter, + not_while_syncing_filter: NotWhileSyncingFilter, + task_spawner_filter: TaskSpawnerFilter, +) -> ResponseFilter { + eth_v4 + .and(warp::path("validator")) + .and(warp::path("blocks")) + .and(warp::path::param::().or_else(|_| async { + Err(warp_utils::reject::custom_bad_request( + "Invalid slot".to_string(), + )) + })) + .and(warp::path::end()) + .and(warp::header::optional::("accept")) + .and(warp::header::(CONSENSUS_VERSION_HEADER)) + .and(not_while_syncing_filter) + .and(warp::query::()) + .and( + warp::header::optional::(CONTENT_TYPE_HEADER) + .and(warp::body::bytes()) + .and_then(|content_type: Option, body: Bytes| async move { + let builder_config: BuilderConfig = + if is_ssz_content_type(content_type.as_deref()) { + BuilderConfig::from_ssz_bytes(&body).map_err(|e| { + warp_utils::reject::custom_bad_request(format!( + "invalid SSZ: {e:?}" + )) + })? + } else { + serde_json::from_slice(&body).map_err(|e| { + warp_utils::reject::custom_deserialize_error(format!("{e:?}")) + })? + }; + // A zero-length `url` or auth `data` makes the body itself invalid (beacon-APIs + // #630) — a 400, unlike per-entry bid failures, which are isolated. + for entry in builder_config.builders.iter() { + entry.validate().map_err(|e| { + warp_utils::reject::custom_bad_request(format!( + "invalid builder entry: {e}" + )) + })?; + } + Ok::<_, Rejection>(builder_config) + }), + ) + .and(task_spawner_filter) + .and(chain_filter) + .then( + |slot: Slot, + accept_header: Option, + consensus_version: ForkName, + not_synced_filter: Result<(), Rejection>, + query: ValidatorBlocksQuery, + builder_config: BuilderConfig, + task_spawner: TaskSpawner, + chain: Arc>| { + task_spawner.spawn_async_with_rejection(Priority::P0, async move { + debug!( + ?slot, + %consensus_version, + "Block production request from HTTP API (v4)" + ); + not_synced_filter?; + produce_block_v4(accept_header, chain, slot, query, builder_config).await + }) + }, + ) + .boxed() +} + // POST validator/liveness/{epoch} pub fn post_validator_liveness_epoch( eth_v1: EthV1Filter, @@ -770,6 +864,133 @@ pub fn post_validator_register_validator( .boxed() } +// POST validator/builder_preferences +// +// Accepts the `BuilderPreferenceEntry` list as either JSON or SSZ. A required +// `Eth-Consensus-Version` header carries the consensus version the preferences belong to (per +// beacon-APIs #630); it is not needed to decode the (currently single-fork) body, so it is only +// logged. +pub fn post_validator_builder_preferences( + eth_v1: EthV1Filter, + chain_filter: ChainFilter, + task_spawner_filter: TaskSpawnerFilter, +) -> ResponseFilter { + eth_v1 + .and(warp::path("validator")) + .and(warp::path("builder_preferences")) + .and(warp::path::end()) + .and(warp::header::(CONSENSUS_VERSION_HEADER)) + .and(task_spawner_filter.clone()) + .and(chain_filter.clone()) + .and( + warp::header::optional::(CONTENT_TYPE_HEADER) + .and(warp::body::bytes()) + .and_then(|content_type: Option, body: Bytes| async move { + let entries: Vec = + if is_ssz_content_type(content_type.as_deref()) { + Vec::from_ssz_bytes(&body).map_err(|e| { + warp_utils::reject::custom_bad_request(format!( + "invalid SSZ: {e:?}" + )) + })? + } else { + serde_json::from_slice(&body).map_err(|e| { + warp_utils::reject::custom_deserialize_error(format!("{e:?}")) + })? + }; + // The submission list is bounded (SSZ `List[BuilderPreferencesEntry, 4096]`, + // JSON `maxItems: 4096`, per beacon-APIs #630); a longer body is invalid. + if entries.len() > MAX_SUBMITTED_BUILDER_PREFERENCES { + return Err(warp_utils::reject::custom_bad_request(format!( + "too many builder preference entries: {} exceeds the limit of {}", + entries.len(), + MAX_SUBMITTED_BUILDER_PREFERENCES + ))); + } + // A zero-length `url` or auth `data` makes the body itself invalid (beacon-APIs + // #630) — a 400, unlike per-entry submission failures, which are isolated. + for entry in &entries { + entry.validate().map_err(|e| { + warp_utils::reject::custom_bad_request(format!( + "invalid builder preference entry: {e}" + )) + })?; + } + Ok::<_, Rejection>(entries) + }), + ) + .then( + |consensus_version: ForkName, + task_spawner: TaskSpawner, + chain: Arc>, + entries: Vec| async move { + let (tx, rx) = oneshot::channel(); + + let initial_result = task_spawner + .spawn_async_with_rejection_no_conversion(Priority::P0, async move { + // The builder service is only present when the Gloas fork is scheduled; a + // node without one can't submit preferences anywhere, which is the + // caller's misconfiguration (not a server fault), so reject with a 400. + let builders = chain + .builders + .as_ref() + .ok_or_else(|| { + warp_utils::reject::custom_bad_request( + "this beacon node has no builder service (the Gloas fork is \ + not scheduled on its network)" + .to_string(), + ) + })? + .clone(); + + debug!( + count = entries.len(), + %consensus_version, + "Received submit builder preferences request" + ); + + // Submitting to a builder can be slow (they frequently time out), so the + // fan-out runs in a detached task rather than holding a `BeaconProcessor` + // worker. The service submits each entry independently and best-effort, + // returning the failures by index (per beacon-APIs #630). + tokio::task::spawn(async move { + let response = match builders + .submit_builder_preferences(entries, consensus_version) + .await + { + Ok(()) => Ok(warp::reply::reply().into_response()), + Err(failures) => Err(warp_utils::reject::indexed_bad_request( + "error submitting builder preferences".to_string(), + failures + .into_iter() + .map(|f| Failure::new(f.index, f.error.to_string())) + .collect(), + )), + }; + let _ = tx.send(response); + }); + + Ok(warp::reply::reply().into_response()) + }) + .await; + + if initial_result.is_err() { + return convert_rejection(initial_result).await; + } + + convert_rejection(rx.await.unwrap_or_else(|_| { + Ok(warp::reply::with_status( + warp::reply::json(&"No response from channel"), + warp::http::StatusCode::INTERNAL_SERVER_ERROR, + ) + .into_response()) + })) + .await + }, + ) + .boxed() +} + // POST validator/prepare_beacon_proposer pub fn post_validator_prepare_beacon_proposer( eth_v1: EthV1Filter, diff --git a/beacon_node/http_api/src/version.rs b/beacon_node/http_api/src/version.rs index 6f441636b49..63914feb049 100644 --- a/beacon_node/http_api/src/version.rs +++ b/beacon_node/http_api/src/version.rs @@ -4,8 +4,8 @@ use eth2::beacon_response::{ ExecutionOptimisticFinalizedMetadata, ForkVersionedResponse, UnversionedResponse, }; use eth2::{ - CONSENSUS_BLOCK_VALUE_HEADER, CONSENSUS_VERSION_HEADER, CONTENT_TYPE_HEADER, - EXECUTION_PAYLOAD_BLINDED_HEADER, EXECUTION_PAYLOAD_INCLUDED_HEADER, + BUILDER_URL_HEADER, CONSENSUS_BLOCK_VALUE_HEADER, CONSENSUS_VERSION_HEADER, + CONTENT_TYPE_HEADER, EXECUTION_PAYLOAD_BLINDED_HEADER, EXECUTION_PAYLOAD_INCLUDED_HEADER, EXECUTION_PAYLOAD_VALUE_HEADER, SSZ_CONTENT_TYPE_HEADER, }; use serde::Serialize; @@ -116,6 +116,15 @@ pub fn add_execution_payload_value_header( .into_response() } +/// Add the `Eth-Builder-Url` header (the winning builder's URL) to a response, when present. +/// Absent for a self-built block or a block won by a p2p bid. +pub fn add_builder_url_header(reply: T, builder_url: Option<&str>) -> Response { + match builder_url { + Some(url) => reply::with_header(reply, BUILDER_URL_HEADER, url).into_response(), + None => reply.into_response(), + } +} + /// Add the `Eth-Consensus-Block-Value` header to a response. pub fn add_consensus_block_value_header( reply: T, diff --git a/beacon_node/http_api/tests/broadcast_validation_tests.rs b/beacon_node/http_api/tests/broadcast_validation_tests.rs index 4d2be52a0d5..5db04d7d136 100644 --- a/beacon_node/http_api/tests/broadcast_validation_tests.rs +++ b/beacon_node/http_api/tests/broadcast_validation_tests.rs @@ -433,6 +433,7 @@ pub async fn consensus_partial_pass_only_consensus() { &channel.0, validation_level, StatusCode::ACCEPTED, + None, ) .await; @@ -610,7 +611,7 @@ pub async fn equivocation_consensus_early_equivocation() { .post_beacon_blocks_v2_ssz( &PublishBlockRequest::new(block_a.clone(), blobs_a), validation_level, - None + None, ) .await .is_ok() @@ -763,6 +764,7 @@ pub async fn equivocation_consensus_late_equivocation() { &channel.0, validation_level, StatusCode::ACCEPTED, + None, ) .await; diff --git a/beacon_node/http_api/tests/gloas_reorg_tests.rs b/beacon_node/http_api/tests/gloas_reorg_tests.rs index f1a4ed02bb4..f31746ddf70 100644 --- a/beacon_node/http_api/tests/gloas_reorg_tests.rs +++ b/beacon_node/http_api/tests/gloas_reorg_tests.rs @@ -27,7 +27,7 @@ use std::collections::HashMap; use std::sync::Arc; use std::time::Duration; use types::{ - Address, BeaconBlockRef, EthSpec, ExecutionBlockHash, Hash256, MinimalEthSpec, + Address, BeaconBlockRef, EthSpec, ExecutionBlockHash, ForkName, Hash256, MinimalEthSpec, ProposerPreparationData, Slot, }; @@ -727,7 +727,15 @@ pub async fn proposer_boost_re_org_test( let (block_c, block_c_blobs) = { let (response, _) = tester .client - .get_validator_blocks_v4::(slot_c, &randao_reveal, None, false, None, None) + .post_validator_blocks_v4::( + slot_c, + &randao_reveal, + None, + false, + ð2::types::BuilderConfig::empty(), + None, + ForkName::Gloas, + ) .await .unwrap(); ( diff --git a/beacon_node/http_api/tests/interactive_tests.rs b/beacon_node/http_api/tests/interactive_tests.rs index 90c7f37b02c..df4fdaf0904 100644 --- a/beacon_node/http_api/tests/interactive_tests.rs +++ b/beacon_node/http_api/tests/interactive_tests.rs @@ -817,7 +817,15 @@ pub async fn fork_choice_before_proposal() { let block_d = if harness.spec.fork_name_at_slot::(slot_d).gloas_enabled() { tester .client - .get_validator_blocks_v4::(slot_d, &randao_reveal, None, false, None, None) + .post_validator_blocks_v4::( + slot_d, + &randao_reveal, + None, + false, + ð2::types::BuilderConfig::empty(), + None, + ForkName::Gloas, + ) .await .unwrap() .0 diff --git a/beacon_node/http_api/tests/tests.rs b/beacon_node/http_api/tests/tests.rs index 8f9dcac5a04..0fc3fd17cec 100644 --- a/beacon_node/http_api/tests/tests.rs +++ b/beacon_node/http_api/tests/tests.rs @@ -4813,7 +4813,15 @@ impl ApiTester { let (response, _metadata) = self .client - .get_validator_blocks_v4::(slot, &randao_reveal, None, false, None, None) + .post_validator_blocks_v4::( + slot, + &randao_reveal, + None, + false, + ð2::types::BuilderConfig::empty(), + None, + ForkName::Gloas, + ) .await .unwrap(); let block = response.into_block(); @@ -5005,6 +5013,248 @@ impl ApiTester { self } + pub async fn test_block_production_v4_missing_consensus_version_header_returns_400( + self, + ) -> Self { + if !self.chain.spec.is_gloas_scheduled() { + return self; + } + + let fork = self.chain.canonical_head.cached_head().head_fork(); + let genesis_validators_root = self.chain.genesis_validators_root; + let Some((slot, epoch, _fork_name)) = self.advance_to_gloas_slot() else { + return self; + }; + + let (_sk, randao_reveal) = self + .proposer_setup(slot, epoch, &fork, genesis_validators_root) + .await; + + let url = self + .client + .post_validator_blocks_v4_path( + slot, + &randao_reveal, + None, + SkipRandaoVerification::No, + false, + None, + ) + .await + .unwrap(); + + // A valid body, but no `Eth-Consensus-Version` header: the header is required + // (beacon-APIs #630), so the request must fail with a 400. + let response = reqwest::Client::new() + .post(url) + .json(ð2::types::BuilderConfig::empty()) + .send() + .await + .unwrap(); + assert_eq!(response.status(), StatusCode::BAD_REQUEST); + + self.chain.slot_clock.set_slot(slot.as_u64() + 1); + + self + } + + pub async fn test_block_production_v4_zero_length_entry_fields_return_400(self) -> Self { + if !self.chain.spec.is_gloas_scheduled() { + return self; + } + + let fork = self.chain.canonical_head.cached_head().head_fork(); + let genesis_validators_root = self.chain.genesis_validators_root; + let Some((slot, epoch, _fork_name)) = self.advance_to_gloas_slot() else { + return self; + }; + + let (_sk, randao_reveal) = self + .proposer_setup(slot, epoch, &fork, genesis_validators_root) + .await; + + let url = self + .client + .post_validator_blocks_v4_path( + slot, + &randao_reveal, + None, + SkipRandaoVerification::No, + false, + None, + ) + .await + .unwrap(); + + let valid_auth = eth2::types::SignedRequestAuth { + message: eth2::types::RequestAuth { + data: eth2::types::RequestAuthData::new(b"http://builder.example.com".to_vec()) + .unwrap(), + slot, + }, + signature: Signature::empty(), + }; + let entry = |url: &str, auth: eth2::types::SignedRequestAuth| eth2::types::BuilderEntry { + url: url.parse().unwrap(), + auth, + builder_pubkeys: <_>::default(), + max_execution_payment: 0, + min_bid: 0, + builder_boost_factor: 100, + }; + + // A zero-length `url` and a zero-length auth `data` each make the body invalid + // (beacon-APIs #630), so the request must fail with a 400. + let empty_url_entry = entry("", valid_auth.clone()); + let mut empty_data_auth = valid_auth; + empty_data_auth.message.data = eth2::types::RequestAuthData::default(); + let empty_data_entry = entry("http://builder.example.com", empty_data_auth); + + for bad_entry in [empty_url_entry, empty_data_entry] { + let config = serde_json::json!({ + "min_bid": "0", + "builder_boost_factor": "100", + "builders": [bad_entry], + }); + let response = reqwest::Client::new() + .post(url.clone()) + .header(eth2::CONSENSUS_VERSION_HEADER, "gloas") + .json(&config) + .send() + .await + .unwrap(); + assert_eq!(response.status(), StatusCode::BAD_REQUEST); + } + + self.chain.slot_clock.set_slot(slot.as_u64() + 1); + + self + } + + /// The `POST validator/builder_preferences` URL, for raw requests that bypass the eth2 client + /// (which always sets the required header). + fn builder_preferences_url(&self) -> reqwest::Url { + let mut url = self.client.server().expose_full().clone(); + url.path_segments_mut() + .unwrap() + .push("eth") + .push("v1") + .push("validator") + .push("builder_preferences"); + url + } + + /// A `BuilderPreferenceEntry` that passes the endpoint's body validation. + fn valid_builder_preference_entry() -> eth2::types::BuilderPreferenceEntry { + eth2::types::BuilderPreferenceEntry { + proposer_pubkey: PublicKeyBytes::empty(), + url: "http://builder.example.com".parse().unwrap(), + auth: eth2::types::SignedRequestAuth { + message: eth2::types::RequestAuth { + data: eth2::types::RequestAuthData::new(b"http://builder.example.com".to_vec()) + .unwrap(), + slot: Slot::new(0), + }, + signature: Signature::empty(), + }, + max_execution_payment: 0, + } + } + + pub async fn test_builder_preferences_missing_consensus_version_header_returns_400( + self, + ) -> Self { + if !self.chain.spec.is_gloas_scheduled() { + return self; + } + + // A valid body, but no `Eth-Consensus-Version` header: the header is required + // (beacon-APIs #630), so the request must fail with a 400. + let response = reqwest::Client::new() + .post(self.builder_preferences_url()) + .json(&vec![Self::valid_builder_preference_entry()]) + .send() + .await + .unwrap(); + assert_eq!(response.status(), StatusCode::BAD_REQUEST); + + self + } + + pub async fn test_builder_preferences_zero_length_entry_fields_return_400(self) -> Self { + if !self.chain.spec.is_gloas_scheduled() { + return self; + } + + // A zero-length `url` and a zero-length auth `data` each make the body invalid + // (beacon-APIs #630), so the request must fail with a 400. + let mut empty_url_entry = Self::valid_builder_preference_entry(); + empty_url_entry.url = "".parse().unwrap(); + let mut empty_data_entry = Self::valid_builder_preference_entry(); + empty_data_entry.auth.message.data = eth2::types::RequestAuthData::default(); + + for bad_entry in [empty_url_entry, empty_data_entry] { + let response = reqwest::Client::new() + .post(self.builder_preferences_url()) + .header(eth2::CONSENSUS_VERSION_HEADER, "gloas") + .json(&vec![bad_entry]) + .send() + .await + .unwrap(); + assert_eq!(response.status(), StatusCode::BAD_REQUEST); + } + + self + } + + pub async fn test_builder_preferences_oversize_list_returns_400(self) -> Self { + if !self.chain.spec.is_gloas_scheduled() { + return self; + } + + // The submission list is bounded at `MAX_SUBMITTED_BUILDER_PREFERENCES` entries + // (beacon-APIs #630); one more is an invalid body. + let entries = vec![ + Self::valid_builder_preference_entry(); + eth2::types::MAX_SUBMITTED_BUILDER_PREFERENCES + 1 + ]; + let response = reqwest::Client::new() + .post(self.builder_preferences_url()) + .header(eth2::CONSENSUS_VERSION_HEADER, "gloas") + .json(&entries) + .send() + .await + .unwrap(); + assert_eq!(response.status(), StatusCode::BAD_REQUEST); + + self + } + + pub async fn test_builder_preferences_without_builder_service_returns_400(self) -> Self { + if !self.chain.spec.is_gloas_scheduled() { + return self; + } + + // The test harness never wires a builder service into the chain, so a well-formed + // submission reaches the handler and must be rejected as a client-side misconfiguration + // (400 with a self-explanatory message), not a 500. + let response = reqwest::Client::new() + .post(self.builder_preferences_url()) + .header(eth2::CONSENSUS_VERSION_HEADER, "gloas") + .json(&vec![Self::valid_builder_preference_entry()]) + .send() + .await + .unwrap(); + assert_eq!(response.status(), StatusCode::BAD_REQUEST); + let body = response.text().await.unwrap(); + assert!( + body.contains("no builder service"), + "unexpected error body: {body}" + ); + + self + } + pub async fn test_envelope_post_when_syncing_returns_503(mut self) -> Self { if !self.chain.spec.is_gloas_scheduled() { return self; @@ -5178,7 +5428,15 @@ impl ApiTester { let (response, metadata) = self .client - .get_validator_blocks_v4::(slot, &randao_reveal, None, false, None, None) + .post_validator_blocks_v4::( + slot, + &randao_reveal, + None, + false, + ð2::types::BuilderConfig::empty(), + None, + ForkName::Gloas, + ) .await .unwrap(); let block = response.into_block(); @@ -5253,7 +5511,15 @@ impl ApiTester { let (response, metadata) = self .client - .get_validator_blocks_v4_ssz::(slot, &randao_reveal, None, false, None, None) + .post_validator_blocks_v4_ssz::( + slot, + &randao_reveal, + None, + false, + ð2::types::BuilderConfig::empty(), + None, + ForkName::Gloas, + ) .await .unwrap(); let block = response.into_block(); @@ -5325,12 +5591,28 @@ impl ApiTester { let (response, metadata) = if ssz { self.client - .get_validator_blocks_v4_ssz::(slot, &randao_reveal, None, true, None, None) + .post_validator_blocks_v4_ssz::( + slot, + &randao_reveal, + None, + true, + ð2::types::BuilderConfig::empty(), + None, + ForkName::Gloas, + ) .await .unwrap() } else { self.client - .get_validator_blocks_v4::(slot, &randao_reveal, None, true, None, None) + .post_validator_blocks_v4::( + slot, + &randao_reveal, + None, + true, + ð2::types::BuilderConfig::empty(), + None, + ForkName::Gloas, + ) .await .unwrap() }; @@ -5871,7 +6153,15 @@ impl ApiTester { // Produce and publish a block. let (response, _metadata) = self .client - .get_validator_blocks_v4::(slot, &randao_reveal, None, false, None, None) + .post_validator_blocks_v4::( + slot, + &randao_reveal, + None, + false, + ð2::types::BuilderConfig::empty(), + None, + ForkName::Gloas, + ) .await .unwrap(); let block = response.into_block(); @@ -5954,7 +6244,15 @@ impl ApiTester { // Produce and publish a block, but withhold its envelope. let (response, _metadata) = self .client - .get_validator_blocks_v4::(slot, &randao_reveal, None, false, None, None) + .post_validator_blocks_v4::( + slot, + &randao_reveal, + None, + false, + ð2::types::BuilderConfig::empty(), + None, + ForkName::Gloas, + ) .await .unwrap(); let block = response.into_block(); @@ -8923,7 +9221,15 @@ impl ApiTester { let (response, _metadata) = self .client - .get_validator_blocks_v4::(slot, &randao_reveal, None, false, None, None) + .post_validator_blocks_v4::( + slot, + &randao_reveal, + None, + false, + ð2::types::BuilderConfig::empty(), + None, + ForkName::Gloas, + ) .await .unwrap(); let block = response.into_block(); @@ -9229,7 +9535,6 @@ impl ApiTester { let epoch = self.chain.epoch().unwrap(); let (_, randao_reveal) = self.get_test_randao(slot, epoch).await; let graffiti = Some(Graffiti::from([0; GRAFFITI_BYTES_LEN])); - // When GraffitiPolicy is None let no_graffiti_policy_path = self .client @@ -10101,6 +10406,18 @@ async fn envelope_api() { .await .test_block_production_v4_missing_include_payload_returns_400() .await + .test_block_production_v4_missing_consensus_version_header_returns_400() + .await + .test_block_production_v4_zero_length_entry_fields_return_400() + .await + .test_builder_preferences_missing_consensus_version_header_returns_400() + .await + .test_builder_preferences_zero_length_entry_fields_return_400() + .await + .test_builder_preferences_oversize_list_returns_400() + .await + .test_builder_preferences_without_builder_service_returns_400() + .await .test_envelope_post_consensus_invalid_returns_400_no_broadcast() .await .test_envelope_post_gossip_partial_pass_returns_202()