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#[deprecated(note = "Use ConsensusEngineHandle instead")]
26pub type BeaconConsensusEngineHandle<Payload> = ConsensusEngineHandle<Payload>;
27
28#[must_use = "futures do nothing unless you `.await` or poll them"]
32#[derive(Debug)]
33pub struct OnForkChoiceUpdated {
34 forkchoice_status: ForkchoiceStatus,
39 fut: Either<futures::future::Ready<ForkChoiceUpdateResult>, PendingPayloadId>,
41}
42
43impl OnForkChoiceUpdated {
46 pub const fn forkchoice_status(&self) -> ForkchoiceStatus {
48 self.forkchoice_status
49 }
50
51 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 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 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 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 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 pub fn invalid_payload_attributes() -> Self {
99 Self {
100 forkchoice_status: ForkchoiceStatus::Valid,
102 fut: Either::Left(futures::future::ready(Err(
103 ForkchoiceUpdateError::UpdatedInvalidPayloadAttributes,
104 ))),
105 }
106 }
107
108 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#[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 Poll::Ready(Err(ForkchoiceUpdateError::UpdatedInvalidPayloadAttributes))
153 }
154 }
155 }
156}
157
158#[derive(Debug, Clone, Copy)]
160pub struct NewPayloadTimings {
161 pub latency: Duration,
163 pub persistence_wait: Duration,
167 pub execution_cache_wait: Option<Duration>,
171 pub sparse_trie_wait: Option<Duration>,
175}
176
177#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
182pub struct BigBlockData<ExecutionData> {
183 pub env_switches: Vec<ExecutionData>,
189 pub prior_block_hashes: Vec<(u64, alloy_primitives::B256)>,
193 pub block_number: u64,
195 #[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#[derive(Debug)]
253pub enum BeaconEngineMessage<Payload: PayloadTypes> {
254 NewPayload {
256 cause: Span,
258 payload: Payload::ExecutionData,
260 tx: oneshot::Sender<Result<PayloadStatus, BeaconOnNewPayloadError>>,
262 },
263 RethNewPayload {
270 cause: Span,
272 payload: Payload::ExecutionData,
274 wait_for_persistence: bool,
276 wait_for_caches: bool,
278 tx: oneshot::Sender<Result<(PayloadStatus, NewPayloadTimings), BeaconOnNewPayloadError>>,
280 enqueued_at: Instant,
282 },
283 ForkchoiceUpdated {
285 cause: Span,
287 state: ForkchoiceState,
289 payload_attrs: Option<Payload::PayloadAttributes>,
291 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 write!(
321 f,
322 "ForkchoiceUpdated {{ state: {state:?}, has_payload_attributes: {} }}",
323 payload_attrs.is_some()
324 )
325 }
326 }
327 }
328}
329
330#[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 pub const fn new(to_engine: UnboundedSender<BeaconEngineMessage<Payload>>) -> Self {
347 Self { to_engine }
348 }
349
350 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 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 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 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}