Skip to main content

reth_basic_payload_builder/
lib.rs

1//! A basic payload generator for reth.
2
3#![doc(
4    html_logo_url = "https://raw.githubusercontent.com/paradigmxyz/reth/main/assets/reth-docs.png",
5    html_favicon_url = "https://avatars0.githubusercontent.com/u/97369466?s=256",
6    issue_tracker_base_url = "https://github.com/paradigmxyz/reth/issues/"
7)]
8#![cfg_attr(not(test), warn(unused_crate_dependencies))]
9#![cfg_attr(docsrs, feature(doc_cfg))]
10
11use crate::metrics::PayloadBuilderMetrics;
12use alloy_eips::merge::SLOT_DURATION;
13use alloy_primitives::{B256, U256};
14use futures_core::ready;
15use futures_util::FutureExt;
16use reth_chain_state::CanonStateNotification;
17use reth_execution_cache::SavedCache;
18use reth_payload_builder::{
19    BuildNewPayload, KeepPayloadJobAlive, PayloadBuilderLease, PayloadId, PayloadJob,
20    PayloadJobGenerator,
21};
22use reth_payload_builder_primitives::PayloadBuilderError;
23use reth_payload_primitives::{BuiltPayload, PayloadAttributes, PayloadKind};
24use reth_primitives_traits::{AlloyBlockHeader, HeaderTy, NodePrimitives, SealedHeader};
25use reth_revm::cached::CachedReads;
26use reth_storage_api::{BlockReaderIdExt, StateProviderFactory};
27use reth_tasks::{CancelOnDrop, Runtime};
28use reth_trie_parallel::state_root_task::PayloadStateRootHandle;
29use std::{
30    fmt,
31    future::Future,
32    ops::Deref,
33    pin::Pin,
34    sync::Arc,
35    task::{Context, Poll},
36    time::{Duration, SystemTime, UNIX_EPOCH},
37};
38use tokio::{
39    sync::{oneshot, Semaphore},
40    time::{Interval, Sleep},
41};
42use tracing::{debug, debug_span, trace, warn, Span};
43
44mod better_payload_emitter;
45mod metrics;
46mod stack;
47
48pub use better_payload_emitter::BetterPayloadEmitter;
49pub use stack::PayloadBuilderStack;
50
51const PAYLOAD_BUILDER_THREAD_NAME: &str = "payload-builder";
52
53/// Helper to access [`NodePrimitives::BlockHeader`] from [`PayloadBuilder::BuiltPayload`].
54pub type HeaderForPayload<P> = <<P as BuiltPayload>::Primitives as NodePrimitives>::BlockHeader;
55
56/// The [`PayloadJobGenerator`] that creates [`BasicPayloadJob`]s.
57#[derive(Debug)]
58pub struct BasicPayloadJobGenerator<Client, Builder> {
59    /// The client that can interact with the chain.
60    client: Client,
61    /// The task executor to spawn payload building tasks on.
62    executor: Runtime,
63    /// The configuration for the job generator.
64    config: BasicPayloadJobGeneratorConfig,
65    /// Restricts how many generator tasks can be executed at once.
66    payload_task_guard: PayloadTaskGuard,
67    /// The type responsible for building payloads.
68    ///
69    /// See [`PayloadBuilder`]
70    builder: Builder,
71    /// Stored `cached_reads` for new payload jobs.
72    pre_cached: Option<PrecachedState>,
73    /// Stored parent block information for new payload jobs.
74    pre_cached_parent_block_info: Option<PrecachedParentBlockInfo>,
75}
76
77// === impl BasicPayloadJobGenerator ===
78
79impl<Client, Builder> BasicPayloadJobGenerator<Client, Builder> {
80    /// Creates a new [`BasicPayloadJobGenerator`] with the given config and custom
81    /// [`PayloadBuilder`]
82    pub fn with_builder(
83        client: Client,
84        executor: Runtime,
85        config: BasicPayloadJobGeneratorConfig,
86        builder: Builder,
87    ) -> Self {
88        Self {
89            client,
90            executor,
91            payload_task_guard: PayloadTaskGuard::new(config.max_payload_tasks),
92            config,
93            builder,
94            pre_cached: None,
95            pre_cached_parent_block_info: None,
96        }
97    }
98
99    /// Returns the maximum duration a job should be allowed to run.
100    ///
101    /// This adheres to the following specification:
102    /// > Client software SHOULD stop the updating process when either a call to engine_getPayload
103    /// > with the build process's payloadId is made or SECONDS_PER_SLOT (12s in the Mainnet
104    /// > configuration) have passed since the point in time identified by the timestamp parameter.
105    ///
106    /// See also <https://github.com/ethereum/execution-apis/blob/431cf72fd3403d946ca3e3afc36b973fc87e0e89/src/engine/paris.md?plain=1#L137>
107    #[inline]
108    fn max_job_duration(&self, unix_timestamp: u64) -> Duration {
109        let duration_until_timestamp = duration_until(unix_timestamp);
110
111        // safety in case clocks are bad
112        let duration_until_timestamp = duration_until_timestamp.min(self.config.deadline * 3);
113
114        self.config.deadline + duration_until_timestamp
115    }
116
117    /// Returns the [Instant](tokio::time::Instant) at which the job should be terminated because it
118    /// is considered timed out.
119    #[inline]
120    fn job_deadline(&self, unix_timestamp: u64) -> tokio::time::Instant {
121        tokio::time::Instant::now() + self.max_job_duration(unix_timestamp)
122    }
123
124    /// Returns a reference to the tasks type
125    pub const fn tasks(&self) -> &Runtime {
126        &self.executor
127    }
128
129    /// Returns the pre-cached reads for the given parent header if it matches the cached state's
130    /// block.
131    fn maybe_pre_cached(&self, parent: B256) -> Option<CachedReads> {
132        if !self.config.pre_cache_state {
133            return None
134        }
135
136        self.pre_cached.as_ref().filter(|pc| pc.block == parent).map(|pc| pc.cached.clone())
137    }
138
139    /// Returns the cached parent block information if it matches the requested parent.
140    fn maybe_parent_block_info(&self, parent: B256) -> Option<PayloadParentBlockInfo> {
141        self.pre_cached_parent_block_info
142            .as_ref()
143            .filter(|info| info.block == parent)
144            .map(|info| info.parent_block_info)
145    }
146}
147
148// === impl BasicPayloadJobGenerator ===
149
150impl<Client, Builder> PayloadJobGenerator for BasicPayloadJobGenerator<Client, Builder>
151where
152    Client: StateProviderFactory
153        + BlockReaderIdExt<Header = HeaderForPayload<Builder::BuiltPayload>>
154        + Clone
155        + Unpin
156        + 'static,
157    Builder: PayloadBuilder + Unpin + 'static,
158    Builder::Attributes: Unpin + Clone,
159    Builder::BuiltPayload: Unpin + Clone,
160{
161    type Job = BasicPayloadJob<Builder>;
162
163    fn new_payload_job(
164        &self,
165        input: BuildNewPayload<Builder::Attributes>,
166        id: PayloadId,
167    ) -> Result<Self::Job, PayloadBuilderError> {
168        let BuildNewPayload { attributes, parent_hash, mut resources } = input;
169        let parent_header = if parent_hash.is_zero() {
170            // Use latest header for genesis block case
171            self.client
172                .latest_header()
173                .map_err(PayloadBuilderError::from)?
174                .ok_or_else(|| PayloadBuilderError::MissingParentHeader(B256::ZERO))?
175        } else {
176            // Fetch specific header by hash
177            self.client
178                .sealed_header_by_hash(parent_hash)
179                .map_err(PayloadBuilderError::from)?
180                .ok_or_else(|| PayloadBuilderError::MissingParentHeader(parent_hash))?
181        };
182
183        let parent_hash = parent_header.hash();
184        let cached_reads = self.maybe_pre_cached(parent_hash);
185        let parent_block_info = self.maybe_parent_block_info(parent_hash);
186
187        let config = PayloadConfig::new(Arc::new(parent_header), attributes, id)
188            .with_parent_block_info(parent_block_info);
189
190        let until = self.job_deadline(config.attributes.timestamp());
191        let deadline = Box::pin(tokio::time::sleep_until(until));
192
193        let mut job = BasicPayloadJob {
194            config,
195            executor: self.executor.clone(),
196            deadline,
197            // ticks immediately
198            interval: tokio::time::interval(self.config.interval),
199            best_payload: PayloadState::Missing,
200            pending_block: None,
201            cached_reads,
202            execution_cache: resources.take_execution_cache(),
203            state_root_handle: resources.take_state_root_handle(),
204            leases: resources.take_leases(),
205            payload_task_guard: self.payload_task_guard.clone(),
206            metrics: Default::default(),
207            builder: self.builder.clone(),
208        };
209
210        // start the first job right away
211        job.spawn_build_job();
212
213        Ok(job)
214    }
215
216    fn on_new_state<N: NodePrimitives>(&mut self, new_state: CanonStateNotification<N>) {
217        if !self.config.pre_cache_state {
218            self.pre_cached = None;
219            return
220        }
221
222        let mut cached = CachedReads::default();
223
224        // extract the state from the notification and put it into the cache
225        let committed = new_state.committed();
226        let new_execution_outcome = committed.execution_outcome();
227        for (addr, acc) in new_execution_outcome.bundle_accounts_iter() {
228            if let Some(info) = acc.info.clone() {
229                // we want pre cache existing accounts and their storage
230                // this only includes changed accounts and storage but is better than nothing
231                let storage =
232                    acc.storage.iter().map(|(key, slot)| (*key, slot.present_value)).collect();
233                cached.insert_account(addr, info, storage);
234            }
235        }
236
237        let tip = committed.tip();
238        let block = tip.hash();
239        let parent_block_info =
240            PayloadParentBlockInfo { transaction_count: tip.transaction_count() };
241
242        self.pre_cached = Some(PrecachedState { block, cached });
243        self.pre_cached_parent_block_info =
244            Some(PrecachedParentBlockInfo { block, parent_block_info });
245    }
246}
247
248/// Pre-filled [`CachedReads`] for a specific block.
249///
250/// This is extracted from the [`CanonStateNotification`] for the tip block.
251#[derive(Debug, Clone)]
252pub struct PrecachedState {
253    /// The block for which the state is pre-cached.
254    pub block: B256,
255    /// Cached state for the block.
256    pub cached: CachedReads,
257}
258
259/// Pre-filled parent block information for a specific block.
260#[derive(Debug, Clone, Copy)]
261struct PrecachedParentBlockInfo {
262    /// The block for which the parent block information is cached.
263    block: B256,
264    /// Cached parent block information.
265    parent_block_info: PayloadParentBlockInfo,
266}
267
268/// Restricts how many generator tasks can be executed at once.
269#[derive(Debug, Clone)]
270pub struct PayloadTaskGuard(Arc<Semaphore>);
271
272impl Deref for PayloadTaskGuard {
273    type Target = Semaphore;
274
275    fn deref(&self) -> &Self::Target {
276        &self.0
277    }
278}
279
280// === impl PayloadTaskGuard ===
281
282impl PayloadTaskGuard {
283    /// Constructs `Self` with a maximum task count of `max_payload_tasks`.
284    pub fn new(max_payload_tasks: usize) -> Self {
285        Self(Arc::new(Semaphore::new(max_payload_tasks)))
286    }
287
288    /// Acquires an owned permit for a payload build task.
289    async fn acquire_owned(&self) -> tokio::sync::OwnedSemaphorePermit {
290        self.0.clone().acquire_owned().await.expect("payload task semaphore closed")
291    }
292}
293
294/// Settings for the [`BasicPayloadJobGenerator`].
295#[derive(Debug, Clone)]
296pub struct BasicPayloadJobGeneratorConfig {
297    /// The interval at which the job should build a new payload after the last.
298    interval: Duration,
299    /// The deadline for when the payload builder job should resolve.
300    ///
301    /// By default this is [`SLOT_DURATION`]: 12s
302    deadline: Duration,
303    /// Maximum number of tasks to spawn for building a payload.
304    max_payload_tasks: usize,
305    /// Whether to pre-cache changed state from canonical state notifications.
306    pre_cache_state: bool,
307}
308
309// === impl BasicPayloadJobGeneratorConfig ===
310
311impl BasicPayloadJobGeneratorConfig {
312    /// Sets the interval at which the job should build a new payload after the last.
313    pub const fn interval(mut self, interval: Duration) -> Self {
314        self.interval = interval;
315        self
316    }
317
318    /// Sets the deadline when this job should resolve.
319    pub const fn deadline(mut self, deadline: Duration) -> Self {
320        self.deadline = deadline;
321        self
322    }
323
324    /// Sets the maximum number of tasks to spawn for building a payload(s).
325    ///
326    /// # Panics
327    ///
328    /// If `max_payload_tasks` is 0.
329    pub fn max_payload_tasks(mut self, max_payload_tasks: usize) -> Self {
330        assert!(max_payload_tasks > 0, "max_payload_tasks must be greater than 0");
331        self.max_payload_tasks = max_payload_tasks;
332        self
333    }
334
335    /// Sets whether to pre-cache changed state from canonical state notifications.
336    ///
337    /// This keeps the parent block's state changes in memory so payload jobs building on top of it
338    /// can reuse those reads.
339    pub const fn pre_cache_state(mut self, pre_cache_state: bool) -> Self {
340        self.pre_cache_state = pre_cache_state;
341        self
342    }
343}
344
345impl Default for BasicPayloadJobGeneratorConfig {
346    fn default() -> Self {
347        Self {
348            interval: Duration::from_secs(1),
349            // 12s slot time
350            deadline: SLOT_DURATION,
351            max_payload_tasks: 3,
352            pre_cache_state: true,
353        }
354    }
355}
356
357/// A basic payload job that continuously builds a payload with the best transactions from the pool.
358///
359/// This type is a [`PayloadJob`] and [`Future`] that terminates when the deadline is reached or
360/// when the job is resolved: [`PayloadJob::resolve`].
361///
362/// This basic job implementation will trigger new payload build task continuously until the job is
363/// resolved or the deadline is reached, or until the built payload is marked as frozen:
364/// [`BuildOutcome::Freeze`]. Once a frozen payload is returned, no additional payloads will be
365/// built and this future will wait to be resolved: [`PayloadJob::resolve`] or terminated if the
366/// deadline is reached.
367#[derive(Debug)]
368pub struct BasicPayloadJob<Builder>
369where
370    Builder: PayloadBuilder,
371{
372    /// The configuration for how the payload will be created.
373    config: PayloadConfig<Builder::Attributes, HeaderForPayload<Builder::BuiltPayload>>,
374    /// How to spawn building tasks
375    executor: Runtime,
376    /// The deadline when this job should resolve.
377    deadline: Pin<Box<Sleep>>,
378    /// The interval at which the job should build a new payload after the last.
379    interval: Interval,
380    /// The best payload so far and its state.
381    best_payload: PayloadState<Builder::BuiltPayload>,
382    /// Receiver for the block that is currently being built.
383    pending_block: Option<PendingPayload<Builder::BuiltPayload>>,
384    /// Restricts how many generator tasks can be executed at once.
385    payload_task_guard: PayloadTaskGuard,
386    /// Caches all disk reads for the state the new payloads builds on
387    ///
388    /// This is used to avoid reading the same state over and over again when new attempts are
389    /// triggered, because during the building process we'll repeatedly execute the transactions.
390    cached_reads: Option<CachedReads>,
391    /// Optional execution cache shared with the engine.
392    execution_cache: Option<SavedCache>,
393    /// Optional state-root task handle, shared with the engine.
394    state_root_handle: Option<PayloadStateRootHandle>,
395    /// Lifecycle leases shared with the payload-builder service.
396    ///
397    /// Every detached build task clones these so that the loaned resources remain available until
398    /// `try_build` completes, even if the payload job is resolved first.
399    leases: Vec<PayloadBuilderLease>,
400    /// metrics for this type
401    metrics: PayloadBuilderMetrics,
402    /// The type responsible for building payloads.
403    ///
404    /// See [`PayloadBuilder`]
405    builder: Builder,
406}
407
408impl<Builder> BasicPayloadJob<Builder>
409where
410    Builder: PayloadBuilder + Unpin + 'static,
411    Builder::Attributes: Unpin + Clone,
412    Builder::BuiltPayload: Unpin + Clone,
413{
414    /// Spawns a new payload build task.
415    fn spawn_build_job(&mut self) {
416        trace!(target: "payload_builder", id = %self.config.payload_id(), "spawn new payload build task");
417        let (tx, rx) = oneshot::channel();
418        let cancel = CancelOnDrop::default();
419        let pending_cancel = cancel.clone();
420        let guard = self.payload_task_guard.clone();
421        let payload_config = self.config.clone();
422        let best_payload = self.best_payload.payload().cloned();
423        self.metrics.inc_initiated_payload_builds();
424        let cached_reads = self.cached_reads.take().unwrap_or_default();
425        let execution_cache = self.execution_cache.clone();
426        let mut state_root_handle = self.state_root_handle.take();
427        let on_payload_built =
428            state_root_handle.as_mut().and_then(PayloadStateRootHandle::take_on_payload_built);
429        let leases = self.leases.clone();
430        let builder = self.builder.clone();
431        let executor = self.executor.clone();
432        let span = Span::current();
433        self.executor.spawn_task(async move {
434            // acquire the permit for executing the task
435            let permit = guard.acquire_owned().await;
436            executor.spawn_blocking_named_or_tokio(PAYLOAD_BUILDER_THREAD_NAME, move || {
437                // Restore the parent even when the worker span is filtered out.
438                let _parent = span.enter();
439                let _span = debug_span!(target: "payload_builder", "build_payload").entered();
440                let _permit = permit;
441                let args = BuildArguments {
442                    cached_reads,
443                    execution_cache,
444                    state_root_handle,
445                    config: payload_config,
446                    cancel,
447                    best_payload,
448                };
449                let result = builder.try_build(args);
450                if let Some(on_payload_built) = on_payload_built &&
451                    let Ok(outcome) = &result &&
452                    let Some(payload) = outcome.payload()
453                {
454                    on_payload_built(payload.block().hash(), payload.block().state_root());
455                }
456                drop(leases);
457                let _ = tx.send(result);
458            });
459        });
460
461        self.pending_block = Some(PendingPayload { cancel: pending_cancel, payload: rx });
462    }
463}
464
465impl<Builder> Future for BasicPayloadJob<Builder>
466where
467    Builder: PayloadBuilder + Unpin + 'static,
468    Builder::Attributes: Unpin + Clone,
469    Builder::BuiltPayload: Unpin + Clone,
470{
471    type Output = Result<(), PayloadBuilderError>;
472
473    fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
474        let this = self.get_mut();
475
476        // check if the deadline is reached
477        if this.deadline.as_mut().poll(cx).is_ready() {
478            trace!(target: "payload_builder", "payload building deadline reached");
479            return Poll::Ready(Ok(()))
480        }
481
482        loop {
483            // Wait for any pending build to complete before polling the next tick.
484            //
485            // This avoids consuming interval ticks while a build is still in-flight,
486            // which would delay the follow-up build by a full interval even though
487            // the current attempt has already finished.
488            if let Some(mut fut) = this.pending_block.take() {
489                match fut.poll_unpin(cx) {
490                    Poll::Ready(Ok(outcome)) => match outcome {
491                        BuildOutcome::Better { payload, cached_reads } => {
492                            this.cached_reads = Some(cached_reads);
493                            debug!(target: "payload_builder", value = %payload.fees(), "built better payload");
494                            this.best_payload = PayloadState::Best(payload);
495                        }
496                        BuildOutcome::Freeze(payload) => {
497                            debug!(target: "payload_builder", "payload frozen, no further building will occur");
498                            this.best_payload = PayloadState::Frozen(payload);
499                        }
500                        BuildOutcome::Aborted { fees, cached_reads } => {
501                            this.cached_reads = Some(cached_reads);
502                            trace!(target: "payload_builder", worse_fees = %fees, "skipped payload build of worse block");
503                        }
504                        BuildOutcome::Cancelled => {
505                            unreachable!("the cancel signal never fired")
506                        }
507                    },
508                    Poll::Ready(Err(error)) => {
509                        // job failed, but we simply try again next interval
510                        debug!(target: "payload_builder", %error, "payload build attempt failed");
511                        this.metrics.inc_failed_payload_builds();
512                    }
513                    Poll::Pending => {
514                        this.pending_block = Some(fut);
515                        return Poll::Pending
516                    }
517                }
518            }
519
520            if this.best_payload.is_frozen() {
521                return Poll::Pending
522            }
523
524            // Wait for the next build interval tick.
525            //
526            // The loop is needed because `poll_tick` does not register a waker
527            // when it returns `Ready`, so we must loop back after spawning a job
528            // to reach a point that *does* register one (the pending block poll above).
529            ready!(this.interval.poll_tick(cx));
530            this.spawn_build_job()
531        }
532    }
533}
534
535impl<Builder> PayloadJob for BasicPayloadJob<Builder>
536where
537    Builder: PayloadBuilder + Unpin + 'static,
538    Builder::Attributes: Unpin + Clone,
539    Builder::BuiltPayload: Unpin + Clone,
540{
541    type PayloadAttributes = Builder::Attributes;
542    type ResolvePayloadFuture = ResolveBestPayload<Self::BuiltPayload>;
543    type BuiltPayload = Builder::BuiltPayload;
544
545    fn best_payload(&self) -> Result<Self::BuiltPayload, PayloadBuilderError> {
546        if let Some(payload) = self.best_payload.payload() {
547            Ok(payload.clone())
548        } else {
549            // No payload has been built yet, but we need to return something that the CL then
550            // can deliver, so we need to return an empty payload.
551            //
552            // Note: it is assumed that this is unlikely to happen, as the payload job is
553            // started right away and the first full block should have been
554            // built by the time CL is requesting the payload.
555            self.metrics.inc_requested_empty_payload();
556            self.builder.build_empty_payload(self.config.clone())
557        }
558    }
559
560    fn payload_attributes(&self) -> Result<Self::PayloadAttributes, PayloadBuilderError> {
561        Ok(self.config.attributes.clone())
562    }
563
564    fn payload_timestamp(&self) -> Result<u64, PayloadBuilderError> {
565        Ok(self.config.attributes.timestamp())
566    }
567
568    fn resolve_kind(
569        &mut self,
570        kind: PayloadKind,
571    ) -> (Self::ResolvePayloadFuture, KeepPayloadJobAlive) {
572        let best_payload = self.best_payload.payload().cloned();
573        if best_payload.is_none() && self.pending_block.is_none() {
574            // ensure we have a job scheduled if we don't have a best payload yet and none is active
575            self.spawn_build_job();
576        }
577
578        let maybe_better = self.pending_block.take();
579        let mut empty_payload = None;
580
581        if best_payload.is_none() {
582            if let Some(pending) = maybe_better.as_ref() {
583                pending.cancel.request_finalization();
584            }
585
586            debug!(target: "payload_builder", id=%self.config.payload_id(), "no best payload yet to resolve, building empty payload");
587
588            let args = BuildArguments {
589                cached_reads: self.cached_reads.take().unwrap_or_default(),
590                execution_cache: self.execution_cache.clone(),
591                state_root_handle: None,
592                config: self.config.clone(),
593                cancel: CancelOnDrop::default(),
594                best_payload: None,
595            };
596
597            match self.builder.on_missing_payload(args) {
598                MissingPayloadBehaviour::AwaitInProgress => {
599                    debug!(target: "payload_builder", id=%self.config.payload_id(), "awaiting in progress payload build job");
600                }
601                MissingPayloadBehaviour::RaceEmptyPayload => {
602                    debug!(target: "payload_builder", id=%self.config.payload_id(), "racing empty payload");
603
604                    // if no payload has been built yet
605                    self.metrics.inc_requested_empty_payload();
606                    // no payload built yet, so we need to return an empty payload
607                    let (tx, rx) = oneshot::channel();
608                    let config = self.config.clone();
609                    let builder = self.builder.clone();
610                    let span = Span::current();
611                    self.executor.spawn_blocking_named_or_tokio(
612                        PAYLOAD_BUILDER_THREAD_NAME,
613                        move || {
614                            let _parent = span.enter();
615                            let _span =
616                                debug_span!(target: "payload_builder", "build_empty_payload")
617                                    .entered();
618                            let res = builder.build_empty_payload(config);
619                            let _ = tx.send(res);
620                        },
621                    );
622
623                    empty_payload = Some(rx);
624                }
625                MissingPayloadBehaviour::RacePayload(job) => {
626                    debug!(target: "payload_builder", id=%self.config.payload_id(), "racing fallback payload");
627                    // race the in progress job with this job
628                    let (tx, rx) = oneshot::channel();
629                    let span = Span::current();
630                    self.executor.spawn_blocking_named_or_tokio(
631                        PAYLOAD_BUILDER_THREAD_NAME,
632                        move || {
633                            let _parent = span.enter();
634                            let _span =
635                                debug_span!(target: "payload_builder", "build_fallback_payload")
636                                    .entered();
637                            let _ = tx.send(job());
638                        },
639                    );
640                    empty_payload = Some(rx);
641                }
642            };
643        }
644
645        let fut = ResolveBestPayload {
646            best_payload,
647            maybe_better,
648            empty_payload: empty_payload.filter(|_| kind != PayloadKind::WaitForPending),
649        };
650
651        (fut, KeepPayloadJobAlive::No)
652    }
653}
654
655/// Represents the current state of a payload being built.
656#[derive(Debug, Clone)]
657pub enum PayloadState<P> {
658    /// No payload has been built yet.
659    Missing,
660    /// The best payload built so far, which may still be improved upon.
661    Best(P),
662    /// The payload is frozen and no further building should occur.
663    ///
664    /// Contains the final payload `P` that should be used.
665    Frozen(P),
666}
667
668impl<P> PayloadState<P> {
669    /// Checks if the payload is frozen.
670    pub const fn is_frozen(&self) -> bool {
671        matches!(self, Self::Frozen(_))
672    }
673
674    /// Returns the payload if it exists (either Best or Frozen).
675    pub const fn payload(&self) -> Option<&P> {
676        match self {
677            Self::Missing => None,
678            Self::Best(p) | Self::Frozen(p) => Some(p),
679        }
680    }
681}
682
683/// The future that returns the best payload to be served to the consensus layer.
684///
685/// This returns the payload that's supposed to be sent to the CL.
686///
687/// If payload has been built so far, it will return that, but it will check if there's a better
688/// payload available from an in progress build job. If so it will return that.
689///
690/// If no payload has been built so far, it will either return an empty payload or the result of the
691/// in progress build job, whatever finishes first.
692#[derive(Debug)]
693pub struct ResolveBestPayload<Payload> {
694    /// Best payload so far.
695    pub best_payload: Option<Payload>,
696    /// Regular payload job that's currently running that might produce a better payload.
697    pub maybe_better: Option<PendingPayload<Payload>>,
698    /// The empty payload building job in progress, if any.
699    pub empty_payload: Option<oneshot::Receiver<Result<Payload, PayloadBuilderError>>>,
700}
701
702impl<Payload> ResolveBestPayload<Payload> {
703    const fn is_empty(&self) -> bool {
704        self.best_payload.is_none() && self.maybe_better.is_none() && self.empty_payload.is_none()
705    }
706}
707
708impl<Payload> Future for ResolveBestPayload<Payload>
709where
710    Payload: Unpin,
711{
712    type Output = Result<Payload, PayloadBuilderError>;
713
714    fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
715        let this = self.get_mut();
716
717        // check if there is a better payload before returning the best payload
718        if let Some(fut) = Pin::new(&mut this.maybe_better).as_pin_mut() &&
719            let Poll::Ready(res) = fut.poll(cx)
720        {
721            this.maybe_better = None;
722            if let Ok(Some(payload)) = res.map(|out| out.into_payload()).inspect_err(
723                |err| warn!(target: "payload_builder", %err, "failed to resolve pending payload"),
724            ) {
725                debug!(target: "payload_builder", "resolving better payload");
726                return Poll::Ready(Ok(payload))
727            }
728        }
729
730        if let Some(best) = this.best_payload.take() {
731            debug!(target: "payload_builder", "resolving best payload");
732            return Poll::Ready(Ok(best))
733        }
734
735        if let Some(fut) = Pin::new(&mut this.empty_payload).as_pin_mut() &&
736            let Poll::Ready(res) = fut.poll(cx)
737        {
738            this.empty_payload = None;
739            return match res {
740                Ok(res) => {
741                    if let Err(err) = &res {
742                        warn!(target: "payload_builder", %err, "failed to resolve empty payload");
743                    } else {
744                        debug!(target: "payload_builder", "resolving empty payload");
745                    }
746                    Poll::Ready(res)
747                }
748                Err(err) => Poll::Ready(Err(err.into())),
749            }
750        }
751
752        if this.is_empty() {
753            return Poll::Ready(Err(PayloadBuilderError::MissingPayload))
754        }
755
756        Poll::Pending
757    }
758}
759
760/// A future that resolves to the result of the block building job.
761#[derive(Debug)]
762pub struct PendingPayload<P> {
763    /// Cancels the job on drop and carries cooperative control signals.
764    cancel: CancelOnDrop,
765    /// The channel to send the result to.
766    payload: oneshot::Receiver<Result<BuildOutcome<P>, PayloadBuilderError>>,
767}
768
769impl<P> PendingPayload<P> {
770    /// Constructs a `PendingPayload` future.
771    pub const fn new(
772        cancel: CancelOnDrop,
773        payload: oneshot::Receiver<Result<BuildOutcome<P>, PayloadBuilderError>>,
774    ) -> Self {
775        Self { cancel, payload }
776    }
777}
778
779impl<P> Future for PendingPayload<P> {
780    type Output = Result<BuildOutcome<P>, PayloadBuilderError>;
781
782    fn poll(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
783        let res = ready!(self.payload.poll_unpin(cx));
784        Poll::Ready(res.map_err(Into::into).and_then(|res| res))
785    }
786}
787
788/// Static config for how to build a payload.
789#[derive(Clone, Debug)]
790pub struct PayloadConfig<Attributes, Header = alloy_consensus::Header> {
791    /// The parent header.
792    pub parent_header: Arc<SealedHeader<Header>>,
793    /// Additional parent block information, if available.
794    pub parent_block_info: Option<PayloadParentBlockInfo>,
795    /// Requested attributes for the payload.
796    pub attributes: Attributes,
797    /// The payload id.
798    pub payload_id: PayloadId,
799}
800
801/// Additional information about the parent block.
802#[derive(Clone, Copy, Debug, PartialEq, Eq)]
803pub struct PayloadParentBlockInfo {
804    /// Number of transactions in the parent block.
805    pub transaction_count: usize,
806}
807
808impl<Attributes, Header> PayloadConfig<Attributes, Header>
809where
810    Attributes: PayloadAttributes,
811{
812    /// Create new payload config.
813    pub const fn new(
814        parent_header: Arc<SealedHeader<Header>>,
815        attributes: Attributes,
816        payload_id: PayloadId,
817    ) -> Self {
818        Self { parent_header, parent_block_info: None, attributes, payload_id }
819    }
820
821    /// Attaches cached parent block information.
822    pub const fn with_parent_block_info(
823        mut self,
824        parent_block_info: Option<PayloadParentBlockInfo>,
825    ) -> Self {
826        self.parent_block_info = parent_block_info;
827        self
828    }
829
830    /// Returns the payload id.
831    pub const fn payload_id(&self) -> PayloadId {
832        self.payload_id
833    }
834}
835
836/// The possible outcomes of a payload building attempt.
837#[derive(Debug)]
838pub enum BuildOutcome<Payload> {
839    /// Successfully built a better block.
840    Better {
841        /// The new payload that was built.
842        payload: Payload,
843        /// The cached reads that were used to build the payload.
844        cached_reads: CachedReads,
845    },
846    /// Aborted payload building because resulted in worse block wrt. fees.
847    Aborted {
848        /// The total fees associated with the attempted payload.
849        fees: U256,
850        /// The cached reads that were used to build the payload.
851        cached_reads: CachedReads,
852    },
853    /// Build job was cancelled
854    Cancelled,
855
856    /// The payload is final and no further building should occur
857    Freeze(Payload),
858}
859
860impl<Payload> BuildOutcome<Payload> {
861    /// Consumes the type and returns the payload if the outcome is `Better` or `Freeze`.
862    pub fn into_payload(self) -> Option<Payload> {
863        match self {
864            Self::Better { payload, .. } | Self::Freeze(payload) => Some(payload),
865            _ => None,
866        }
867    }
868
869    /// Consumes the type and returns the payload if the outcome is `Better` or `Freeze`.
870    pub const fn payload(&self) -> Option<&Payload> {
871        match self {
872            Self::Better { payload, .. } | Self::Freeze(payload) => Some(payload),
873            _ => None,
874        }
875    }
876
877    /// Returns true if the outcome is `Better`.
878    pub const fn is_better(&self) -> bool {
879        matches!(self, Self::Better { .. })
880    }
881
882    /// Returns true if the outcome is `Freeze`.
883    pub const fn is_frozen(&self) -> bool {
884        matches!(self, Self::Freeze { .. })
885    }
886
887    /// Returns true if the outcome is `Aborted`.
888    pub const fn is_aborted(&self) -> bool {
889        matches!(self, Self::Aborted { .. })
890    }
891
892    /// Returns true if the outcome is `Cancelled`.
893    pub const fn is_cancelled(&self) -> bool {
894        matches!(self, Self::Cancelled)
895    }
896
897    /// Applies a fn on the current payload.
898    pub fn map_payload<F, P>(self, f: F) -> BuildOutcome<P>
899    where
900        F: FnOnce(Payload) -> P,
901    {
902        match self {
903            Self::Better { payload, cached_reads } => {
904                BuildOutcome::Better { payload: f(payload), cached_reads }
905            }
906            Self::Aborted { fees, cached_reads } => BuildOutcome::Aborted { fees, cached_reads },
907            Self::Cancelled => BuildOutcome::Cancelled,
908            Self::Freeze(payload) => BuildOutcome::Freeze(f(payload)),
909        }
910    }
911}
912
913/// The possible outcomes of a payload building attempt without reused [`CachedReads`]
914#[derive(Debug)]
915pub enum BuildOutcomeKind<Payload> {
916    /// Successfully built a better block.
917    Better {
918        /// The new payload that was built.
919        payload: Payload,
920    },
921    /// Aborted payload building because resulted in worse block wrt. fees.
922    Aborted {
923        /// The total fees associated with the attempted payload.
924        fees: U256,
925    },
926    /// Build job was cancelled
927    Cancelled,
928    /// The payload is final and no further building should occur
929    Freeze(Payload),
930}
931
932impl<Payload> BuildOutcomeKind<Payload> {
933    /// Attaches the [`CachedReads`] to the outcome.
934    pub fn with_cached_reads(self, cached_reads: CachedReads) -> BuildOutcome<Payload> {
935        match self {
936            Self::Better { payload } => BuildOutcome::Better { payload, cached_reads },
937            Self::Aborted { fees } => BuildOutcome::Aborted { fees, cached_reads },
938            Self::Cancelled => BuildOutcome::Cancelled,
939            Self::Freeze(payload) => BuildOutcome::Freeze(payload),
940        }
941    }
942}
943
944/// A collection of arguments used for building payloads.
945///
946/// This struct encapsulates the essential components and configuration required for the payload
947/// building process. It holds references to the Ethereum client, transaction pool, cached reads,
948/// payload configuration, cancellation status, and the best payload achieved so far.
949#[derive(Debug)]
950pub struct BuildArguments<Attributes, Payload: BuiltPayload> {
951    /// Previously cached disk reads
952    pub cached_reads: CachedReads,
953    /// Optional execution cache shared with the engine.
954    pub execution_cache: Option<SavedCache>,
955    /// Optional state-root task handle, shared with the engine.
956    ///
957    /// A successful build returns its retained trie through the handle's completion callback,
958    /// associated with the built block's hash and state root.
959    pub state_root_handle: Option<PayloadStateRootHandle>,
960    /// How to configure the payload.
961    pub config: PayloadConfig<Attributes, HeaderTy<Payload::Primitives>>,
962    /// A marker that can be used to cancel the job.
963    pub cancel: CancelOnDrop,
964    /// The best payload achieved so far.
965    pub best_payload: Option<Payload>,
966}
967
968impl<Attributes, Payload: BuiltPayload> BuildArguments<Attributes, Payload> {
969    /// Create new build arguments.
970    pub const fn new(
971        cached_reads: CachedReads,
972        execution_cache: Option<SavedCache>,
973        state_root_handle: Option<PayloadStateRootHandle>,
974        config: PayloadConfig<Attributes, HeaderTy<Payload::Primitives>>,
975        cancel: CancelOnDrop,
976        best_payload: Option<Payload>,
977    ) -> Self {
978        Self { cached_reads, execution_cache, state_root_handle, config, cancel, best_payload }
979    }
980}
981
982/// A trait for building payloads that encapsulate Ethereum transactions.
983///
984/// This trait provides the `try_build` method to construct a transaction payload
985/// using `BuildArguments`. It returns a `Result` indicating success or a
986/// `PayloadBuilderError` if building fails.
987///
988/// Generic parameters `Pool` and `Client` represent the transaction pool and
989/// Ethereum client types.
990pub trait PayloadBuilder: Send + Sync + Clone {
991    /// The payload attributes type to accept for building.
992    type Attributes: PayloadAttributes;
993    /// The type of the built payload.
994    type BuiltPayload: BuiltPayload;
995
996    /// Tries to build a transaction payload using provided arguments.
997    ///
998    /// Constructs a transaction payload based on the given arguments,
999    /// returning a `Result` indicating success or an error if building fails.
1000    ///
1001    /// # Arguments
1002    ///
1003    /// - `args`: Build arguments containing necessary components.
1004    ///
1005    /// # Returns
1006    ///
1007    /// A `Result` indicating the build outcome or an error.
1008    fn try_build(
1009        &self,
1010        args: BuildArguments<Self::Attributes, Self::BuiltPayload>,
1011    ) -> Result<BuildOutcome<Self::BuiltPayload>, PayloadBuilderError>;
1012
1013    /// Invoked when the payload job is being resolved and there is no payload yet.
1014    ///
1015    /// This can happen if the CL requests a payload before the first payload has been built.
1016    fn on_missing_payload(
1017        &self,
1018        _args: BuildArguments<Self::Attributes, Self::BuiltPayload>,
1019    ) -> MissingPayloadBehaviour<Self::BuiltPayload> {
1020        MissingPayloadBehaviour::RaceEmptyPayload
1021    }
1022
1023    /// Builds an empty payload without any transaction.
1024    fn build_empty_payload(
1025        &self,
1026        config: PayloadConfig<Self::Attributes, HeaderForPayload<Self::BuiltPayload>>,
1027    ) -> Result<Self::BuiltPayload, PayloadBuilderError>;
1028}
1029
1030/// Tells the payload builder how to react to payload request if there's no payload available yet.
1031///
1032/// This situation can occur if the CL requests a payload before the first payload has been built.
1033#[derive(Default)]
1034pub enum MissingPayloadBehaviour<Payload> {
1035    /// Await the regular scheduled payload process.
1036    AwaitInProgress,
1037    /// Race the in progress payload process with an empty payload.
1038    #[default]
1039    RaceEmptyPayload,
1040    /// Race the in progress payload process with this job.
1041    RacePayload(Box<dyn FnOnce() -> Result<Payload, PayloadBuilderError> + Send>),
1042}
1043
1044impl<Payload> fmt::Debug for MissingPayloadBehaviour<Payload> {
1045    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1046        match self {
1047            Self::AwaitInProgress => write!(f, "AwaitInProgress"),
1048            Self::RaceEmptyPayload => {
1049                write!(f, "RaceEmptyPayload")
1050            }
1051            Self::RacePayload(_) => write!(f, "RacePayload"),
1052        }
1053    }
1054}
1055
1056/// Checks if the new payload is better than the current best.
1057///
1058/// This compares the total fees of the blocks, higher is better.
1059#[inline(always)]
1060pub fn is_better_payload<T: BuiltPayload>(best_payload: Option<&T>, new_fees: U256) -> bool {
1061    if let Some(best_payload) = best_payload {
1062        new_fees > best_payload.fees()
1063    } else {
1064        true
1065    }
1066}
1067
1068/// Returns the duration until the given unix timestamp in seconds.
1069///
1070/// Returns `Duration::ZERO` if the given timestamp is in the past.
1071fn duration_until(unix_timestamp_secs: u64) -> Duration {
1072    let unix_now = SystemTime::now().duration_since(UNIX_EPOCH).unwrap_or_default();
1073    let timestamp = Duration::from_secs(unix_timestamp_secs);
1074    timestamp.saturating_sub(unix_now)
1075}