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
116 changes: 115 additions & 1 deletion fork_choice_control/src/events.rs
Original file line number Diff line number Diff line change
Expand Up @@ -25,7 +25,10 @@ use types::{
primitives::{BlobIndex, KzgCommitment, VersionedHash},
},
fulu::primitives::ColumnIndex,
gloas::containers::{PayloadAttestationMessage, SignedExecutionPayloadBid},
gloas::{
containers::{PayloadAttestationMessage, SignedExecutionPayloadBid},
primitives::BuilderIndex,
},
nonstandard::Phase,
phase0::{
containers::{Checkpoint, ProposerSlashing, SignedVoluntaryExit},
Expand All @@ -52,8 +55,10 @@ pub enum Topic {
ChainReorg,
ContributionAndProof,
DataColumnSidecar,
ExecutionPayload,
ExecutionPayloadBid,
ExecutionPayloadAvailable,
ExecutionPayloadGossip,
FinalizedCheckpoint,
Head,
PayloadAttestation,
Expand All @@ -73,8 +78,10 @@ pub enum Event<P: Preset> {
ChainReorg(ChainReorgEvent),
ContributionAndProof(Box<SignedContributionAndProof<P>>),
DataColumnSidecar(DataColumnSidecarEvent<P>),
ExecutionPayload(ExecutionPayloadEvent),
ExecutionPayloadAvailable(ExecutionPayloadAvailableEvent),
ExecutionPayloadBid(ExecutionPayloadBidEvent<P>),
ExecutionPayloadGossip(ExecutionPayloadGossipEvent),
FinalizedCheckpoint(FinalizedCheckpointEvent),
Head(HeadEvent),
PayloadAttestation(PayloadAttestationEvent),
Expand All @@ -96,8 +103,10 @@ impl<P: Preset> Event<P> {
Self::ChainReorg(_) => Topic::ChainReorg,
Self::ContributionAndProof(_) => Topic::ContributionAndProof,
Self::DataColumnSidecar(_) => Topic::DataColumnSidecar,
Self::ExecutionPayload(_) => Topic::ExecutionPayload,
Self::ExecutionPayloadAvailable(_) => Topic::ExecutionPayloadAvailable,
Self::ExecutionPayloadBid(_) => Topic::ExecutionPayloadBid,
Self::ExecutionPayloadGossip(_) => Topic::ExecutionPayloadGossip,
Self::FinalizedCheckpoint(_) => Topic::FinalizedCheckpoint,
Self::Head(_) => Topic::Head,
Self::PayloadAttestation(_) => Topic::PayloadAttestation,
Expand All @@ -120,8 +129,10 @@ pub struct EventChannels<P: Preset> {
pub chain_reorgs: Sender<Event<P>>,
pub contribution_and_proofs: Sender<Event<P>>,
pub data_column_sidecars: Sender<Event<P>>,
pub execution_payloads: Sender<Event<P>>,
pub execution_payload_available: Sender<Event<P>>,
pub execution_payload_bids: Sender<Event<P>>,
pub execution_payloads_gossip: Sender<Event<P>>,
pub finalized_checkpoints: Sender<Event<P>>,
pub heads: Sender<Event<P>>,
pub payload_attestations: Sender<Event<P>>,
Expand Down Expand Up @@ -151,8 +162,10 @@ impl<P: Preset> EventChannels<P> {
chain_reorgs: broadcast::channel(max_events).0,
contribution_and_proofs: broadcast::channel(max_events).0,
data_column_sidecars: broadcast::channel(max_events).0,
execution_payloads: broadcast::channel(max_events).0,
execution_payload_available: broadcast::channel(max_events).0,
execution_payload_bids: broadcast::channel(max_events).0,
execution_payloads_gossip: broadcast::channel(max_events).0,
finalized_checkpoints: broadcast::channel(max_events).0,
heads: broadcast::channel(max_events).0,
payload_attestations: broadcast::channel(max_events).0,
Expand All @@ -175,8 +188,10 @@ impl<P: Preset> EventChannels<P> {
Topic::ChainReorg => &self.chain_reorgs,
Topic::ContributionAndProof => &self.contribution_and_proofs,
Topic::DataColumnSidecar => &self.data_column_sidecars,
Topic::ExecutionPayload => &self.execution_payloads,
Topic::ExecutionPayloadAvailable => &self.execution_payload_available,
Topic::ExecutionPayloadBid => &self.execution_payload_bids,
Topic::ExecutionPayloadGossip => &self.execution_payloads_gossip,
Topic::FinalizedCheckpoint => &self.finalized_checkpoints,
Topic::Head => &self.heads,
Topic::PayloadAttestation => &self.payload_attestations,
Expand Down Expand Up @@ -271,6 +286,42 @@ impl<P: Preset> EventChannels<P> {
}
}

pub fn send_execution_payload_event(
&self,
slot: Slot,
builder_index: BuilderIndex,
block_hash: ExecutionBlockHash,
block_root: H256,
execution_optimistic: bool,
) {
if let Err(error) = self.send_execution_payload_event_internal(
slot,
builder_index,
block_hash,
block_root,
execution_optimistic,
) {
warn_with_peers!("unable to send execution payload event: {error}");
}
}

pub fn send_execution_payload_gossip_event(
&self,
slot: Slot,
builder_index: BuilderIndex,
block_hash: ExecutionBlockHash,
block_root: H256,
) {
if let Err(error) = self.send_execution_payload_gossip_event_internal(
slot,
builder_index,
block_hash,
block_root,
) {
warn_with_peers!("unable to send execution payload gossip event: {error}");
}
}

pub fn send_execution_payload_available_event(&self, slot: Slot, block_root: H256) {
if let Err(error) = self.send_execution_payload_available_event_internal(slot, block_root) {
warn_with_peers!("unable to send execution payload available event: {error}");
Expand Down Expand Up @@ -505,6 +556,48 @@ impl<P: Preset> EventChannels<P> {
Ok(())
}

fn send_execution_payload_event_internal(
&self,
slot: Slot,
builder_index: BuilderIndex,
block_hash: ExecutionBlockHash,
block_root: H256,
execution_optimistic: bool,
) -> Result<()> {
if self.execution_payloads.receiver_count() > 0 {
let event = Event::ExecutionPayload(ExecutionPayloadEvent {
slot,
builder_index,
block_hash,
block_root,
execution_optimistic,
});
self.execution_payloads.send(event)?;
}

Ok(())
}

fn send_execution_payload_gossip_event_internal(
&self,
slot: Slot,
builder_index: BuilderIndex,
block_hash: ExecutionBlockHash,
block_root: H256,
) -> Result<()> {
if self.execution_payloads_gossip.receiver_count() > 0 {
let event = Event::ExecutionPayloadGossip(ExecutionPayloadGossipEvent {
slot,
builder_index,
block_hash,
block_root,
});
self.execution_payloads_gossip.send(event)?;
}

Ok(())
}

fn send_execution_payload_available_event_internal(
&self,
slot: Slot,
Expand Down Expand Up @@ -713,6 +806,27 @@ impl<P: Preset> DataColumnSidecarEvent<P> {
}
}

#[derive(Clone, Copy, Debug, Serialize)]
pub struct ExecutionPayloadEvent {
#[serde(with = "serde_utils::string_or_native")]
pub slot: Slot,
#[serde(with = "serde_utils::string_or_native")]
pub builder_index: BuilderIndex,
pub block_hash: ExecutionBlockHash,
pub block_root: H256,
pub execution_optimistic: bool,
}

#[derive(Clone, Copy, Debug, Serialize)]
pub struct ExecutionPayloadGossipEvent {
#[serde(with = "serde_utils::string_or_native")]
pub slot: Slot,
#[serde(with = "serde_utils::string_or_native")]
pub builder_index: BuilderIndex,
pub block_hash: ExecutionBlockHash,
pub block_root: H256,
}

#[derive(Clone, Copy, Debug, Serialize)]
pub struct ExecutionPayloadAvailableEvent {
#[serde(with = "serde_utils::string_or_native")]
Expand Down
27 changes: 26 additions & 1 deletion fork_choice_control/src/mutator.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2005,16 +2005,29 @@ where

let beacon_block_root = envelope.block_root();
let slot = envelope.slot();
let builder_index = envelope.builder_index();
let block_hash = envelope.message.payload.block_hash;
let should_send_gossip_event = origin.should_send_gossip_event();
let should_generate_event = origin.should_generate_event();

debug_with_peers!(
"execution payload envelope accepted (beacon_block_root: {beacon_block_root:?}, slot: {slot})"
);

if origin.should_generate_event() {
if should_generate_event {
self.event_channels
.send_execution_payload_available_event(slot, beacon_block_root);
}

if should_send_gossip_event {
self.event_channels.send_execution_payload_gossip_event(
slot,
builder_index,
block_hash,
beacon_block_root,
);
}

let (gossip_id, sender) = origin.split();

if let Some(gossip_id) = gossip_id {
Expand All @@ -2027,6 +2040,18 @@ where
);

self.accept_execution_payload_envelope(&wait_group, envelope);

if should_generate_event {
if let Some(chain_link) = self.store.chain_link(beacon_block_root) {
self.event_channels.send_execution_payload_event(
slot,
builder_index,
block_hash,
beacon_block_root,
chain_link.is_optimistic(),
);
}
}
}
Ok(ExecutionPayloadEnvelopeAction::Ignore(publishable)) => {
if let Some(metrics) = self.metrics.as_ref() {
Expand Down
5 changes: 5 additions & 0 deletions fork_choice_store/src/misc.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1129,6 +1129,11 @@ impl ExecutionPayloadEnvelopeOrigin {
matches!(self, Self::Gossip(_) | Self::Api(_) | Self::Own)
}

#[must_use]
pub const fn should_send_gossip_event(&self) -> bool {
matches!(self, Self::Gossip(_) | Self::Api(_))
}

#[must_use]
pub const fn verify_signatures(&self) -> bool {
match self {
Expand Down
2 changes: 2 additions & 0 deletions http_api/src/standard.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2586,8 +2586,10 @@ pub async fn beacon_events<P: Preset>(
Event::ChainReorg(data) => ssevent.json_data(data),
Event::ContributionAndProof(data) => ssevent.json_data(data),
Event::DataColumnSidecar(data) => ssevent.json_data(data),
Event::ExecutionPayload(data) => ssevent.json_data(data),
Event::ExecutionPayloadAvailable(data) => ssevent.json_data(data),
Event::ExecutionPayloadBid(data) => ssevent.json_data(data),
Event::ExecutionPayloadGossip(data) => ssevent.json_data(data),
Event::FinalizedCheckpoint(data) => ssevent.json_data(data),
Event::Head(data) => ssevent.json_data(data),
Event::PayloadAttestation(data) => ssevent.json_data(data),
Expand Down
Loading