diff --git a/beacon_node/beacon_chain/src/test_utils.rs b/beacon_node/beacon_chain/src/test_utils.rs index f1c0ef7caa8..77772b07f5c 100644 --- a/beacon_node/beacon_chain/src/test_utils.rs +++ b/beacon_node/beacon_chain/src/test_utils.rs @@ -43,7 +43,7 @@ use logging::create_test_tracing_subscriber; use merkle_proof::MerkleTree; use operation_pool::ReceivedPreCapella; use parking_lot::{Mutex, RwLockWriteGuard}; -use proto_array::{PayloadBlockHash, PayloadStatus}; +use proto_array::PayloadStatus; use rand::Rng; use rand::SeedableRng; use rand::rngs::StdRng; @@ -985,21 +985,6 @@ where self.chain.canonical_head.cached_head().head_block_root() } - /// The execution payload hash that `block_root` commits to, read from fork choice. - pub fn execution_block_hash(&self, block_root: Hash256) -> ExecutionBlockHash { - match self - .chain - .canonical_head - .fork_choice_read_lock() - .get_block(&block_root) - .expect("block should be in fork choice") - .block_hash() - { - PayloadBlockHash::Hash(block_hash) => block_hash, - PayloadBlockHash::PreMerge => panic!("block {block_root:?} has no payload"), - } - } - pub fn finalized_checkpoint(&self) -> Checkpoint { self.chain .canonical_head diff --git a/beacon_node/beacon_chain/tests/payload_invalidation.rs b/beacon_node/beacon_chain/tests/payload_invalidation.rs index 6a54036a3f8..39b3d24e1d0 100644 --- a/beacon_node/beacon_chain/tests/payload_invalidation.rs +++ b/beacon_node/beacon_chain/tests/payload_invalidation.rs @@ -343,8 +343,7 @@ impl InvalidPayloadRig { block_root } - async fn invalidate_manually(&self, block_root: Hash256) { - let head_hash = self.block_hash(block_root); + async fn invalidate_manually(&self, head_hash: ExecutionBlockHash) { self.harness .chain .process_invalid_execution_payload(&InvalidationOperation::InvalidateOne { head_hash }) @@ -1051,7 +1050,7 @@ async fn invalid_parent() { assert_eq!(block.parent_root(), parent_root); // Invalidate the parent block. - rig.invalidate_manually(parent_root).await; + rig.invalidate_manually(rig.block_hash(parent_root)).await; assert!(rig.execution_status(parent_root).is_invalid()); // Ensure the block built atop an invalid payload is invalid for gossip. @@ -1272,7 +1271,7 @@ impl InvalidHeadSetup { .set_current_slot(new_wall_clock_epoch.start_slot(slots_per_epoch)); // Invalidate the head block. - rig.invalidate_manually(invalid_head.head_block_root()) + rig.invalidate_manually(invalid_head.head_hash().unwrap()) .await; // Ensure the justified root is the head. This is the spec-correct choice of head when @@ -1407,7 +1406,7 @@ async fn weights_after_resetting_optimistic_status() { .map(|node| (node.root(), node.weight())) .collect::>(); - rig.invalidate_manually(roots[1]).await; + rig.invalidate_manually(rig.block_hash(roots[1])).await; rig.harness .chain diff --git a/beacon_node/network/src/network_beacon_processor/sync_methods.rs b/beacon_node/network/src/network_beacon_processor/sync_methods.rs index c9d40e7b8b7..f5549548c39 100644 --- a/beacon_node/network/src/network_beacon_processor/sync_methods.rs +++ b/beacon_node/network/src/network_beacon_processor/sync_methods.rs @@ -271,13 +271,14 @@ impl NetworkBeaconProcessor { match &result { Ok(availability) => match availability { - AvailabilityProcessingStatus::Imported(_, hash) => { + AvailabilityProcessingStatus::Imported(slot, hash) => { debug!( result = "imported block and custody columns", block_hash = %hash, "Block components retrieved" ); self.chain.recompute_head_at_current_slot().await; + self.notify_import_after_column(*slot, *hash, EnvelopeSource::Rpc); } AvailabilityProcessingStatus::MissingComponents(_, _) => { debug!( diff --git a/beacon_node/network/src/network_beacon_processor/tests.rs b/beacon_node/network/src/network_beacon_processor/tests.rs index 9052a5ff40b..71e016d6cc5 100644 --- a/beacon_node/network/src/network_beacon_processor/tests.rs +++ b/beacon_node/network/src/network_beacon_processor/tests.rs @@ -11,7 +11,7 @@ use crate::{ use beacon_chain::block_verification_types::LookupBlock; use beacon_chain::custody_context::NodeCustodyType; use beacon_chain::data_column_verification::GossipVerifiedDataColumn; -use beacon_chain::kzg_utils::blobs_to_data_column_sidecars; +use beacon_chain::kzg_utils::{blobs_to_data_column_sidecars, blobs_to_data_column_sidecars_gloas}; use beacon_chain::observed_data_sidecars::DoNotObserve; use beacon_chain::test_utils::{ AttestationStrategy, BeaconChainHarness, BlockStrategy, EphemeralHarnessType, get_kzg, @@ -517,6 +517,25 @@ impl TestRig { } } + fn processor_with_reprocess_receiver( + &self, + ) -> (Arc>, mpsc::Receiver>) { + let (beacon_processor_tx, beacon_processor_rx) = mpsc::channel(8); + let (network_tx, _network_rx) = mpsc::unbounded_channel(); + let (sync_tx, _sync_rx) = mpsc::unbounded_channel(); + let processor = Arc::new(NetworkBeaconProcessor { + beacon_processor_send: BeaconProcessorSend(beacon_processor_tx), + duplicate_cache: DuplicateCache::default(), + chain: self.chain.clone(), + network_tx, + sync_tx, + network_globals: self.network_beacon_processor.network_globals.clone(), + invalid_block_storage: InvalidBlockStorage::Disabled, + executor: self.network_beacon_processor.executor.clone(), + }); + (processor, beacon_processor_rx) + } + pub fn enqueue_blobs_by_range_request(&self, start_slot: u64, count: u64) { self.network_beacon_processor .send_blobs_by_range_request( @@ -1769,6 +1788,189 @@ async fn requeue_early_gossip_payload_envelope() { ); } +/// An RPC envelope can finish execution verification before its custody columns arrive. When the +/// columns make the envelope available, the processor must notify the reprocessing queue. +#[tokio::test] +async fn rpc_columns_notify_after_deferred_envelope_import() { + use beacon_chain::{AvailabilityProcessingStatus, NotifyExecutionLayer}; + use types::BlockImportSource; + + if test_spec::().gloas_fork_epoch.is_none() { + return; + } + + let rig = TestRig::new(SMALL_CHAIN).await; + let block_root = rig.next_block.canonical_root(); + + let block_result = rig + .chain + .process_block( + block_root, + LookupBlock::new(rig.next_block.clone()), + NotifyExecutionLayer::Yes, + BlockImportSource::Lookup, + || Ok(()), + ) + .await; + assert_matches!(block_result, Ok(AvailabilityProcessingStatus::Imported(..))); + + let (processor, mut beacon_processor_rx) = rig.processor_with_reprocess_receiver(); + processor + .clone() + .process_lookup_envelope( + block_root, + rig.next_block_envelope + .clone() + .expect("the next block should have an envelope post-Gloas"), + BlockProcessType::SinglePayloadEnvelope(1), + ) + .await; + assert!( + rig.chain + .get_payload_envelope(&block_root) + .unwrap() + .is_none(), + "envelope should await custody columns" + ); + assert!( + rig.chain + .pending_payload_cache + .get_executed_payload_envelope(&block_root) + .is_some(), + "executed envelope should remain pending on custody columns" + ); + assert_matches!( + beacon_processor_rx.try_recv(), + Err(mpsc::error::TryRecvError::Empty) + ); + + let bid = rig + .next_block + .message() + .body() + .signed_execution_payload_bid() + .expect("Gloas payload bid"); + let blobs = { + let generator = rig._harness.execution_block_generator(); + generator + .blobs_bundles + .values() + .find(|bundle| { + bundle + .commitments + .iter() + .eq(bid.message.blob_kzg_commitments.iter()) + }) + .expect("blobs for next block") + .blobs + .clone() + }; + let sampling_indices = rig + .chain + .custody_context + .sampling_columns_for_epoch(rig.next_block.epoch()); + let custody_columns = blobs_to_data_column_sidecars_gloas( + &blobs.iter().collect::>(), + block_root, + rig.next_block.slot(), + &rig.chain.kzg, + &rig.chain.spec, + ) + .expect("build Gloas columns") + .into_iter() + .filter(|column| sampling_indices.contains(column.index())) + .collect(); + processor + .clone() + .process_rpc_custody_columns( + block_root, + custody_columns, + BlockProcessType::SingleCustodyColumn(1), + ) + .await; + + assert!( + rig.chain + .get_payload_envelope(&block_root) + .unwrap() + .is_some(), + "columns should complete envelope import" + ); + assert_matches!( + beacon_processor_rx.try_recv(), + Ok(WorkEvent { + work: Work::Reprocess(ReprocessQueueMessage::PayloadEnvelopeImported { + block_root: notified_root + }), + .. + }) if notified_root == block_root + ); + assert_matches!( + beacon_processor_rx.try_recv(), + Err(mpsc::error::TryRecvError::Empty) + ); +} + +/// Before Gloas, RPC custody columns complete a block import and should release work waiting for +/// that block rather than work waiting for a payload envelope. +#[tokio::test] +async fn rpc_columns_notify_after_deferred_block_import() { + use beacon_chain::{AvailabilityProcessingStatus, NotifyExecutionLayer}; + use types::BlockImportSource; + + let spec = test_spec::(); + if spec.fulu_fork_epoch.is_none() || spec.gloas_fork_epoch.is_some() { + return; + } + + let rig = TestRig::new(SMALL_CHAIN).await; + let block_root = rig.next_block.canonical_root(); + let custody_columns = rig + .next_data_columns + .clone() + .expect("the next block should have data columns pre-Gloas"); + + let block_result = rig + .chain + .process_block( + block_root, + LookupBlock::new(rig.next_block.clone()), + NotifyExecutionLayer::Yes, + BlockImportSource::Lookup, + || Ok(()), + ) + .await; + assert_matches!( + block_result, + Ok(AvailabilityProcessingStatus::MissingComponents(_, pending_root)) + if pending_root == block_root + ); + + let (processor, mut beacon_processor_rx) = rig.processor_with_reprocess_receiver(); + processor + .clone() + .process_rpc_custody_columns( + block_root, + custody_columns, + BlockProcessType::SingleCustodyColumn(1), + ) + .await; + + assert_matches!( + beacon_processor_rx.try_recv(), + Ok(WorkEvent { + work: Work::Reprocess(ReprocessQueueMessage::BlockImported { + block_root: notified_root + }), + .. + }) if notified_root == block_root + ); + assert_matches!( + beacon_processor_rx.try_recv(), + Err(mpsc::error::TryRecvError::Empty) + ); +} + /// Ensure that attestations that reference an unknown block get properly re-queued and re-processed /// when the block is not seen. #[tokio::test] diff --git a/consensus/proto_array/src/fork_choice_test_definition.rs b/consensus/proto_array/src/fork_choice_test_definition.rs index df3bd8c3c5c..c4013ca0d26 100644 --- a/consensus/proto_array/src/fork_choice_test_definition.rs +++ b/consensus/proto_array/src/fork_choice_test_definition.rs @@ -84,8 +84,8 @@ pub enum Operation { expected_len: usize, }, InvalidatePayload { - head_block_root: Hash256, - latest_valid_ancestor_root: Option, + head_hash: ExecutionBlockHash, + latest_valid_ancestor: Option, }, AssertWeight { block_root: Hash256, @@ -424,13 +424,10 @@ impl ForkChoiceTestDefinition { ); } Operation::InvalidatePayload { - head_block_root, - latest_valid_ancestor_root, + head_hash, + latest_valid_ancestor, } => { - // Operations name payloads. Test blocks commit to `from_root(root)`, as - // `get_hash` spells it. - let head_hash = ExecutionBlockHash::from_root(head_block_root); - let op = if let Some(latest_valid_ancestor) = latest_valid_ancestor_root { + let op = if let Some(latest_valid_ancestor) = latest_valid_ancestor { InvalidationOperation::InvalidateMany { head_hash, always_invalidate_head: true, diff --git a/consensus/proto_array/src/fork_choice_test_definition/execution_status.rs b/consensus/proto_array/src/fork_choice_test_definition/execution_status.rs index 9809f853952..029dda273fa 100644 --- a/consensus/proto_array/src/fork_choice_test_definition/execution_status.rs +++ b/consensus/proto_array/src/fork_choice_test_definition/execution_status.rs @@ -295,8 +295,8 @@ pub fn get_execution_status_test_definition_01() -> ForkChoiceTestDefinition { // | // 3 <- INVALID Operation::InvalidatePayload { - head_block_root: get_root(3), - latest_valid_ancestor_root: Some(get_hash(1)), + head_hash: get_hash(3), + latest_valid_ancestor: Some(get_hash(1)), }, // Ensure that the head is still 2. // @@ -714,8 +714,8 @@ pub fn get_execution_status_test_definition_02() -> ForkChoiceTestDefinition { // | // 3 <- INVALID Operation::InvalidatePayload { - head_block_root: get_root(3), - latest_valid_ancestor_root: Some(get_hash(1)), + head_hash: get_hash(3), + latest_valid_ancestor: Some(get_hash(1)), }, // Ensure that the head is now 2. // @@ -1023,8 +1023,8 @@ pub fn get_execution_status_test_definition_03() -> ForkChoiceTestDefinition { // | // 3 <- INVALID Operation::InvalidatePayload { - head_block_root: get_root(3), - latest_valid_ancestor_root: Some(get_hash(1)), + head_hash: get_hash(3), + latest_valid_ancestor: Some(get_hash(1)), }, // Ensure that the head is now 1, maintaining the proposer boost on the invalid block. // diff --git a/consensus/proto_array/src/fork_choice_test_definition/gloas_payload.rs b/consensus/proto_array/src/fork_choice_test_definition/gloas_payload.rs index 8f65419596f..3d9596247ef 100644 --- a/consensus/proto_array/src/fork_choice_test_definition/gloas_payload.rs +++ b/consensus/proto_array/src/fork_choice_test_definition/gloas_payload.rs @@ -1307,8 +1307,8 @@ mod tests { // Invalidate block 1 (V17). filter_block_tree excludes the entire branch. ops.push(Operation::InvalidatePayload { - head_block_root: get_root(1), - latest_valid_ancestor_root: Some(get_hash(0)), + head_hash: get_hash(1), + latest_valid_ancestor: Some(get_hash(0)), }); // Head falls back to genesis — the invalid branch is no longer selectable. diff --git a/testing/ef_tests/src/cases/gossip_validation.rs b/testing/ef_tests/src/cases/gossip_validation.rs index 40728a5a2c0..cb9ad26ae14 100644 --- a/testing/ef_tests/src/cases/gossip_validation.rs +++ b/testing/ef_tests/src/cases/gossip_validation.rs @@ -651,7 +651,15 @@ impl GossipTester { if payload_status == Some(PayloadStatus::Invalidated) { // The block has been imported optimistically. Mark its payload invalid in fork // choice so descendants observe an invalid execution parent. - let head_hash = self.harness.execution_block_hash(block_root); + let head_hash = block + .message() + .execution_payload() + .map_err(|e| { + Error::InternalError(format!( + "setup block {block_root:?} has no execution payload: {e:?}" + )) + })? + .block_hash(); self.block_on_dangerous(self.harness.chain.process_invalid_execution_payload( &InvalidationOperation::InvalidateOne { head_hash }, ))?