Skip to main content

reth_engine_primitives/
message.rs

1use crate::{
2    error::BeaconForkChoiceUpdateError, BeaconOnNewPayloadError, ExecutionPayload, ForkchoiceStatus,
3};
4use alloy_eips::eip4895::Withdrawal;
5use alloy_primitives::{Bytes, B256};
6use alloy_rpc_types_engine::{
7    ExecutionData, ForkChoiceUpdateResult, ForkchoiceState, ForkchoiceUpdateError,
8    ForkchoiceUpdated, PayloadId, PayloadStatus, PayloadStatusEnum,
9};
10use core::{
11    fmt::{self, Display},
12    future::Future,
13    pin::Pin,
14    task::{ready, Context, Poll},
15};
16use futures::{future::Either, FutureExt, TryFutureExt};
17use reth_errors::RethResult;
18use reth_payload_builder_primitives::PayloadBuilderError;
19use reth_payload_primitives::PayloadTypes;
20use std::time::{Duration, Instant};
21use tokio::sync::{mpsc::UnboundedSender, oneshot};
22use tracing::Span;
23
24/// Type alias for backwards compat
25#[deprecated(note = "Use ConsensusEngineHandle instead")]
26pub type BeaconConsensusEngineHandle<Payload> = ConsensusEngineHandle<Payload>;
27
28/// Represents the outcome of forkchoice update.
29///
30/// This is a future that resolves to [`ForkChoiceUpdateResult`]
31#[must_use = "futures do nothing unless you `.await` or poll them"]
32#[derive(Debug)]
33pub struct OnForkChoiceUpdated {
34    /// Represents the status of the forkchoice update.
35    ///
36    /// Note: This is separate from the response `fut`, because we still can return an error
37    /// depending on the payload attributes, even if the forkchoice update itself is valid.
38    forkchoice_status: ForkchoiceStatus,
39    /// Returns the result of the forkchoice update.
40    fut: Either<futures::future::Ready<ForkChoiceUpdateResult>, PendingPayloadId>,
41}
42
43// === impl OnForkChoiceUpdated ===
44
45impl OnForkChoiceUpdated {
46    /// Returns the determined status of the received `ForkchoiceState`.
47    pub const fn forkchoice_status(&self) -> ForkchoiceStatus {
48        self.forkchoice_status
49    }
50
51    /// Creates a new instance of `OnForkChoiceUpdated` for the `SYNCING` state
52    pub fn syncing() -> Self {
53        let status = PayloadStatus::from_status(PayloadStatusEnum::Syncing);
54        Self {
55            forkchoice_status: ForkchoiceStatus::from_payload_status(&status.status),
56            fut: Either::Left(futures::future::ready(Ok(ForkchoiceUpdated::new(status)))),
57        }
58    }
59
60    /// Creates a new instance of `OnForkChoiceUpdated` if the forkchoice update succeeded and no
61    /// payload attributes were provided.
62    pub fn valid(status: PayloadStatus) -> Self {
63        Self {
64            forkchoice_status: ForkchoiceStatus::from_payload_status(&status.status),
65            fut: Either::Left(futures::future::ready(Ok(ForkchoiceUpdated::new(status)))),
66        }
67    }
68
69    /// Creates a new instance of `OnForkChoiceUpdated` with the given payload status, if the
70    /// forkchoice update failed due to an invalid payload.
71    pub fn with_invalid(status: PayloadStatus) -> Self {
72        Self {
73            forkchoice_status: ForkchoiceStatus::from_payload_status(&status.status),
74            fut: Either::Left(futures::future::ready(Ok(ForkchoiceUpdated::new(status)))),
75        }
76    }
77
78    /// Creates a new instance of `OnForkChoiceUpdated` if the forkchoice update failed because the
79    /// given state is considered invalid
80    pub fn invalid_state() -> Self {
81        Self {
82            forkchoice_status: ForkchoiceStatus::Invalid,
83            fut: Either::Left(futures::future::ready(Err(ForkchoiceUpdateError::InvalidState))),
84        }
85    }
86
87    /// Creates a new instance of `OnForkChoiceUpdated` if the forkchoice update failed because the
88    /// requested reorg to the head block exceeds the supported reorg depth.
89    pub fn too_deep_reorg() -> Self {
90        Self {
91            forkchoice_status: ForkchoiceStatus::Invalid,
92            fut: Either::Left(futures::future::ready(Err(ForkchoiceUpdateError::TooDeepReorg))),
93        }
94    }
95
96    /// Creates a new instance of `OnForkChoiceUpdated` if the forkchoice update was successful but
97    /// payload attributes were invalid.
98    pub fn invalid_payload_attributes() -> Self {
99        Self {
100            // This is valid because this is only reachable if the state and payload is valid
101            forkchoice_status: ForkchoiceStatus::Valid,
102            fut: Either::Left(futures::future::ready(Err(
103                ForkchoiceUpdateError::UpdatedInvalidPayloadAttributes,
104            ))),
105        }
106    }
107
108    /// If the forkchoice update was successful and no payload attributes were provided, this method
109    pub const fn updated_with_pending_payload_id(
110        payload_status: PayloadStatus,
111        pending_payload_id: oneshot::Receiver<Result<PayloadId, PayloadBuilderError>>,
112    ) -> Self {
113        Self {
114            forkchoice_status: ForkchoiceStatus::from_payload_status(&payload_status.status),
115            fut: Either::Right(PendingPayloadId {
116                payload_status: Some(payload_status),
117                pending_payload_id,
118            }),
119        }
120    }
121}
122
123impl Future for OnForkChoiceUpdated {
124    type Output = ForkChoiceUpdateResult;
125
126    fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
127        self.get_mut().fut.poll_unpin(cx)
128    }
129}
130
131/// A future that returns the payload id of a yet to be initiated payload job after a successful
132/// forkchoice update
133#[derive(Debug)]
134struct PendingPayloadId {
135    payload_status: Option<PayloadStatus>,
136    pending_payload_id: oneshot::Receiver<Result<PayloadId, PayloadBuilderError>>,
137}
138
139impl Future for PendingPayloadId {
140    type Output = ForkChoiceUpdateResult;
141
142    fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
143        let this = self.get_mut();
144        let res = ready!(this.pending_payload_id.poll_unpin(cx));
145        match res {
146            Ok(Ok(payload_id)) => Poll::Ready(Ok(ForkchoiceUpdated {
147                payload_status: this.payload_status.take().expect("Polled after completion"),
148                payload_id: Some(payload_id),
149            })),
150            Err(_) | Ok(Err(_)) => {
151                // failed to initiate a payload build job
152                Poll::Ready(Err(ForkchoiceUpdateError::UpdatedInvalidPayloadAttributes))
153            }
154        }
155    }
156}
157
158/// Timing breakdown for `reth_newPayload` responses.
159#[derive(Debug, Clone, Copy)]
160pub struct NewPayloadTimings {
161    /// Server-side execution latency.
162    pub latency: Duration,
163    /// Time spent waiting on persistence, including both time this message spent queued
164    /// due to persistence backpressure and, when `wait_for_persistence` was requested,
165    /// the explicit wait for in-flight persistence to complete.
166    pub persistence_wait: Duration,
167    /// Time spent waiting for the execution cache lock.
168    ///
169    /// `None` when wasn't asked to wait for execution cache.
170    pub execution_cache_wait: Option<Duration>,
171    /// Time spent waiting for the sparse trie cache lock.
172    ///
173    /// `None` when wasn't asked to wait for sparse trie cache.
174    pub sparse_trie_wait: Option<Duration>,
175}
176
177/// Additional data for big block payloads that merge multiple real blocks.
178///
179/// This is used by the `reth_newPayload` endpoint to pass environment switches
180/// and prior block hashes needed for correct multi-segment execution.
181#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
182pub struct BigBlockData<ExecutionData> {
183    /// Environment switches at block boundaries.
184    /// Each entry is `(cumulative_tx_count, execution_data_of_next_block)`.
185    ///
186    /// The first entry at index 0 represents the **original unmutated** base block's
187    /// `ExecutionData`, which must be used to derive the initial EVM environment.
188    pub env_switches: Vec<ExecutionData>,
189    /// Block number → real block hash for blocks covered by previous big blocks in a sequence.
190    /// When replaying chained big blocks, the BLOCKHASH opcode needs real hashes for blocks
191    /// that were merged into earlier big blocks (and thus not individually persisted).
192    pub prior_block_hashes: Vec<(u64, alloy_primitives::B256)>,
193    /// Block number for this big block.
194    pub block_number: u64,
195    /// Merged block access list for this big block.
196    #[serde(default, skip_serializing_if = "Option::is_none")]
197    pub merged_block_access_list: Option<Bytes>,
198}
199
200impl ExecutionPayload for BigBlockData<ExecutionData> {
201    fn parent_hash(&self) -> B256 {
202        self.env_switches[0].parent_hash()
203    }
204
205    fn block_hash(&self) -> B256 {
206        self.env_switches.last().unwrap().block_hash()
207    }
208
209    fn block_number(&self) -> u64 {
210        self.block_number
211    }
212
213    fn withdrawals(&self) -> Option<&Vec<Withdrawal>> {
214        self.env_switches[0].withdrawals()
215    }
216
217    fn block_access_list(&self) -> Option<&Bytes> {
218        self.merged_block_access_list.as_ref()
219    }
220
221    fn parent_beacon_block_root(&self) -> Option<B256> {
222        self.env_switches[0].parent_beacon_block_root()
223    }
224
225    fn timestamp(&self) -> u64 {
226        self.env_switches[0].timestamp()
227    }
228
229    fn gas_used(&self) -> u64 {
230        self.env_switches.iter().map(|data| data.gas_used()).sum()
231    }
232
233    fn gas_limit(&self) -> u64 {
234        self.env_switches.iter().map(|data| data.gas_limit()).sum()
235    }
236
237    fn transaction_count(&self) -> usize {
238        self.env_switches.iter().map(|data| data.transaction_count()).sum()
239    }
240
241    fn slot_number(&self) -> Option<u64> {
242        self.env_switches[0].payload.slot_number()
243    }
244}
245
246/// A message for the beacon engine from other components of the node (engine RPC API invoked by the
247/// consensus layer).
248///
249/// The `cause` span carries in-process tracing context across the engine queue. Callers
250/// constructing messages directly should capture [`Span::current`], or use [`Span::none`] when no
251/// context exists.
252#[derive(Debug)]
253pub enum BeaconEngineMessage<Payload: PayloadTypes> {
254    /// Message with new payload.
255    NewPayload {
256        /// The caller span that caused this message.
257        cause: Span,
258        /// The execution payload received by Engine API.
259        payload: Payload::ExecutionData,
260        /// The sender for returning payload status result.
261        tx: oneshot::Sender<Result<PayloadStatus, BeaconOnNewPayloadError>>,
262    },
263    /// Message with new payload used by `reth_newPayload` endpoint.
264    ///
265    /// Supports independent control over waiting for persistence and cache locks before
266    /// processing, providing unbiased timing measurements when enabled.
267    ///
268    /// Returns detailed timing breakdown alongside the payload status.
269    RethNewPayload {
270        /// The caller span that caused this message.
271        cause: Span,
272        /// The execution payload received by Engine API.
273        payload: Payload::ExecutionData,
274        /// Whether to wait for in-flight persistence to complete before processing.
275        wait_for_persistence: bool,
276        /// Whether to wait for execution cache and sparse trie locks before processing.
277        wait_for_caches: bool,
278        /// The sender for returning payload status result and timing breakdown.
279        tx: oneshot::Sender<Result<(PayloadStatus, NewPayloadTimings), BeaconOnNewPayloadError>>,
280        /// When this message was enqueued, used to measure backpressure wait time.
281        enqueued_at: Instant,
282    },
283    /// Message with updated forkchoice state.
284    ForkchoiceUpdated {
285        /// The caller span that caused this message.
286        cause: Span,
287        /// The updated forkchoice state.
288        state: ForkchoiceState,
289        /// The payload attributes for block building.
290        payload_attrs: Option<Payload::PayloadAttributes>,
291        /// The sender for returning forkchoice updated result.
292        tx: oneshot::Sender<RethResult<OnForkChoiceUpdated>>,
293    },
294}
295
296impl<Payload: PayloadTypes> Display for BeaconEngineMessage<Payload> {
297    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
298        match self {
299            Self::NewPayload { payload, .. } => {
300                write!(
301                    f,
302                    "NewPayload(parent: {}, number: {}, hash: {})",
303                    payload.parent_hash(),
304                    payload.block_number(),
305                    payload.block_hash()
306                )
307            }
308            Self::RethNewPayload { payload, .. } => {
309                write!(
310                    f,
311                    "RethNewPayload(parent: {}, number: {}, hash: {})",
312                    payload.parent_hash(),
313                    payload.block_number(),
314                    payload.block_hash()
315                )
316            }
317            Self::ForkchoiceUpdated { state, payload_attrs, .. } => {
318                // we don't want to print the entire payload attributes, because for OP this
319                // includes all txs
320                write!(
321                    f,
322                    "ForkchoiceUpdated {{ state: {state:?}, has_payload_attributes: {} }}",
323                    payload_attrs.is_some()
324                )
325            }
326        }
327    }
328}
329
330/// A cloneable sender type that can be used to send engine API messages.
331///
332/// This type mirrors consensus related functions of the engine API.
333#[derive(Debug, Clone)]
334pub struct ConsensusEngineHandle<Payload>
335where
336    Payload: PayloadTypes,
337{
338    to_engine: UnboundedSender<BeaconEngineMessage<Payload>>,
339}
340
341impl<Payload> ConsensusEngineHandle<Payload>
342where
343    Payload: PayloadTypes,
344{
345    /// Creates a new beacon consensus engine handle.
346    pub const fn new(to_engine: UnboundedSender<BeaconEngineMessage<Payload>>) -> Self {
347        Self { to_engine }
348    }
349
350    /// Sends a new payload message to the beacon consensus engine and waits for a response.
351    ///
352    /// See also <https://github.com/ethereum/execution-apis/blob/3d627c95a4d3510a8187dd02e0250ecb4331d27e/src/engine/shanghai.md#engine_newpayloadv2>
353    pub async fn new_payload(
354        &self,
355        payload: Payload::ExecutionData,
356    ) -> Result<PayloadStatus, BeaconOnNewPayloadError> {
357        let (tx, rx) = oneshot::channel();
358        let _ = self.to_engine.send(BeaconEngineMessage::NewPayload {
359            cause: Span::current(),
360            payload,
361            tx,
362        });
363        rx.await.map_err(|_| BeaconOnNewPayloadError::EngineUnavailable)?
364    }
365
366    /// Sends a new payload message used by `reth_newPayload` endpoint.
367    ///
368    /// `wait_for_persistence`: waits for in-flight persistence to complete.
369    /// `wait_for_caches`: waits for execution cache and sparse trie locks, excluding destruction
370    /// of removed execution-cache allocations after unlocking.
371    ///
372    /// Returns detailed timing breakdown alongside the payload status.
373    pub async fn reth_new_payload(
374        &self,
375        payload: Payload::ExecutionData,
376        wait_for_persistence: bool,
377        wait_for_caches: bool,
378    ) -> Result<(PayloadStatus, NewPayloadTimings), BeaconOnNewPayloadError> {
379        let (tx, rx) = oneshot::channel();
380        let _ = self.to_engine.send(BeaconEngineMessage::RethNewPayload {
381            cause: Span::current(),
382            payload,
383            wait_for_persistence,
384            wait_for_caches,
385            tx,
386            enqueued_at: Instant::now(),
387        });
388        rx.await.map_err(|_| BeaconOnNewPayloadError::EngineUnavailable)?
389    }
390
391    /// Sends a forkchoice update message to the beacon consensus engine and waits for a response.
392    ///
393    /// See also <https://github.com/ethereum/execution-apis/blob/3d627c95a4d3510a8187dd02e0250ecb4331d27e/src/engine/shanghai.md#engine_forkchoiceupdatedv2>
394    pub async fn fork_choice_updated(
395        &self,
396        state: ForkchoiceState,
397        payload_attrs: Option<Payload::PayloadAttributes>,
398    ) -> Result<ForkchoiceUpdated, BeaconForkChoiceUpdateError> {
399        Ok(self
400            .send_fork_choice_updated(state, payload_attrs)
401            .map_err(|_| BeaconForkChoiceUpdateError::EngineUnavailable)
402            .await?
403            .map_err(BeaconForkChoiceUpdateError::internal)?
404            .await?)
405    }
406
407    /// Sends a forkchoice update message to the beacon consensus engine and returns the receiver to
408    /// wait for a response.
409    fn send_fork_choice_updated(
410        &self,
411        state: ForkchoiceState,
412        payload_attrs: Option<Payload::PayloadAttributes>,
413    ) -> oneshot::Receiver<RethResult<OnForkChoiceUpdated>> {
414        let (tx, rx) = oneshot::channel();
415        let _ = self.to_engine.send(BeaconEngineMessage::ForkchoiceUpdated {
416            cause: Span::current(),
417            state,
418            payload_attrs,
419            tx,
420        });
421        rx
422    }
423}