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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
17 changes: 1 addition & 16 deletions beacon_node/beacon_chain/src/test_utils.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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
Expand Down
9 changes: 4 additions & 5 deletions beacon_node/beacon_chain/tests/payload_invalidation.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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 })
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -1407,7 +1406,7 @@ async fn weights_after_resetting_optimistic_status() {
.map(|node| (node.root(), node.weight()))
.collect::<HashMap<_, _>>();

rig.invalidate_manually(roots[1]).await;
rig.invalidate_manually(rig.block_hash(roots[1])).await;

rig.harness
.chain
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -271,13 +271,14 @@ impl<T: BeaconChainTypes> NetworkBeaconProcessor<T> {

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!(
Expand Down
204 changes: 203 additions & 1 deletion beacon_node/network/src/network_beacon_processor/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -517,6 +517,25 @@ impl TestRig {
}
}

fn processor_with_reprocess_receiver(
&self,
) -> (Arc<NetworkBeaconProcessor<T>>, mpsc::Receiver<WorkEvent<E>>) {
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(
Expand Down Expand Up @@ -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::<E>().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::<Vec<_>>(),
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::<E>();
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]
Expand Down
13 changes: 5 additions & 8 deletions consensus/proto_array/src/fork_choice_test_definition.rs
Original file line number Diff line number Diff line change
Expand Up @@ -84,8 +84,8 @@ pub enum Operation {
expected_len: usize,
},
InvalidatePayload {
head_block_root: Hash256,
latest_valid_ancestor_root: Option<ExecutionBlockHash>,
head_hash: ExecutionBlockHash,
latest_valid_ancestor: Option<ExecutionBlockHash>,
},
AssertWeight {
block_root: Hash256,
Expand Down Expand Up @@ -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,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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.
//
Expand Down Expand Up @@ -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.
//
Expand Down Expand Up @@ -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.
//
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
10 changes: 9 additions & 1 deletion testing/ef_tests/src/cases/gossip_validation.rs
Original file line number Diff line number Diff line change
Expand Up @@ -651,7 +651,15 @@ impl<E: EthSpec> GossipTester<E> {
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 },
))?
Expand Down
Loading