1#![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
53pub type HeaderForPayload<P> = <<P as BuiltPayload>::Primitives as NodePrimitives>::BlockHeader;
55
56#[derive(Debug)]
58pub struct BasicPayloadJobGenerator<Client, Builder> {
59 client: Client,
61 executor: Runtime,
63 config: BasicPayloadJobGeneratorConfig,
65 payload_task_guard: PayloadTaskGuard,
67 builder: Builder,
71 pre_cached: Option<PrecachedState>,
73 pre_cached_parent_block_info: Option<PrecachedParentBlockInfo>,
75}
76
77impl<Client, Builder> BasicPayloadJobGenerator<Client, Builder> {
80 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 #[inline]
108 fn max_job_duration(&self, unix_timestamp: u64) -> Duration {
109 let duration_until_timestamp = duration_until(unix_timestamp);
110
111 let duration_until_timestamp = duration_until_timestamp.min(self.config.deadline * 3);
113
114 self.config.deadline + duration_until_timestamp
115 }
116
117 #[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 pub const fn tasks(&self) -> &Runtime {
126 &self.executor
127 }
128
129 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 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
148impl<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 self.client
172 .latest_header()
173 .map_err(PayloadBuilderError::from)?
174 .ok_or_else(|| PayloadBuilderError::MissingParentHeader(B256::ZERO))?
175 } else {
176 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 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 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 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 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#[derive(Debug, Clone)]
252pub struct PrecachedState {
253 pub block: B256,
255 pub cached: CachedReads,
257}
258
259#[derive(Debug, Clone, Copy)]
261struct PrecachedParentBlockInfo {
262 block: B256,
264 parent_block_info: PayloadParentBlockInfo,
266}
267
268#[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
280impl PayloadTaskGuard {
283 pub fn new(max_payload_tasks: usize) -> Self {
285 Self(Arc::new(Semaphore::new(max_payload_tasks)))
286 }
287
288 async fn acquire_owned(&self) -> tokio::sync::OwnedSemaphorePermit {
290 self.0.clone().acquire_owned().await.expect("payload task semaphore closed")
291 }
292}
293
294#[derive(Debug, Clone)]
296pub struct BasicPayloadJobGeneratorConfig {
297 interval: Duration,
299 deadline: Duration,
303 max_payload_tasks: usize,
305 pre_cache_state: bool,
307}
308
309impl BasicPayloadJobGeneratorConfig {
312 pub const fn interval(mut self, interval: Duration) -> Self {
314 self.interval = interval;
315 self
316 }
317
318 pub const fn deadline(mut self, deadline: Duration) -> Self {
320 self.deadline = deadline;
321 self
322 }
323
324 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 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 deadline: SLOT_DURATION,
351 max_payload_tasks: 3,
352 pre_cache_state: true,
353 }
354 }
355}
356
357#[derive(Debug)]
368pub struct BasicPayloadJob<Builder>
369where
370 Builder: PayloadBuilder,
371{
372 config: PayloadConfig<Builder::Attributes, HeaderForPayload<Builder::BuiltPayload>>,
374 executor: Runtime,
376 deadline: Pin<Box<Sleep>>,
378 interval: Interval,
380 best_payload: PayloadState<Builder::BuiltPayload>,
382 pending_block: Option<PendingPayload<Builder::BuiltPayload>>,
384 payload_task_guard: PayloadTaskGuard,
386 cached_reads: Option<CachedReads>,
391 execution_cache: Option<SavedCache>,
393 state_root_handle: Option<PayloadStateRootHandle>,
395 leases: Vec<PayloadBuilderLease>,
400 metrics: PayloadBuilderMetrics,
402 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 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 let permit = guard.acquire_owned().await;
436 executor.spawn_blocking_named_or_tokio(PAYLOAD_BUILDER_THREAD_NAME, move || {
437 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 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 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 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 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 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 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 self.metrics.inc_requested_empty_payload();
606 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 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#[derive(Debug, Clone)]
657pub enum PayloadState<P> {
658 Missing,
660 Best(P),
662 Frozen(P),
666}
667
668impl<P> PayloadState<P> {
669 pub const fn is_frozen(&self) -> bool {
671 matches!(self, Self::Frozen(_))
672 }
673
674 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#[derive(Debug)]
693pub struct ResolveBestPayload<Payload> {
694 pub best_payload: Option<Payload>,
696 pub maybe_better: Option<PendingPayload<Payload>>,
698 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 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#[derive(Debug)]
762pub struct PendingPayload<P> {
763 cancel: CancelOnDrop,
765 payload: oneshot::Receiver<Result<BuildOutcome<P>, PayloadBuilderError>>,
767}
768
769impl<P> PendingPayload<P> {
770 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#[derive(Clone, Debug)]
790pub struct PayloadConfig<Attributes, Header = alloy_consensus::Header> {
791 pub parent_header: Arc<SealedHeader<Header>>,
793 pub parent_block_info: Option<PayloadParentBlockInfo>,
795 pub attributes: Attributes,
797 pub payload_id: PayloadId,
799}
800
801#[derive(Clone, Copy, Debug, PartialEq, Eq)]
803pub struct PayloadParentBlockInfo {
804 pub transaction_count: usize,
806}
807
808impl<Attributes, Header> PayloadConfig<Attributes, Header>
809where
810 Attributes: PayloadAttributes,
811{
812 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 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 pub const fn payload_id(&self) -> PayloadId {
832 self.payload_id
833 }
834}
835
836#[derive(Debug)]
838pub enum BuildOutcome<Payload> {
839 Better {
841 payload: Payload,
843 cached_reads: CachedReads,
845 },
846 Aborted {
848 fees: U256,
850 cached_reads: CachedReads,
852 },
853 Cancelled,
855
856 Freeze(Payload),
858}
859
860impl<Payload> BuildOutcome<Payload> {
861 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 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 pub const fn is_better(&self) -> bool {
879 matches!(self, Self::Better { .. })
880 }
881
882 pub const fn is_frozen(&self) -> bool {
884 matches!(self, Self::Freeze { .. })
885 }
886
887 pub const fn is_aborted(&self) -> bool {
889 matches!(self, Self::Aborted { .. })
890 }
891
892 pub const fn is_cancelled(&self) -> bool {
894 matches!(self, Self::Cancelled)
895 }
896
897 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#[derive(Debug)]
915pub enum BuildOutcomeKind<Payload> {
916 Better {
918 payload: Payload,
920 },
921 Aborted {
923 fees: U256,
925 },
926 Cancelled,
928 Freeze(Payload),
930}
931
932impl<Payload> BuildOutcomeKind<Payload> {
933 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#[derive(Debug)]
950pub struct BuildArguments<Attributes, Payload: BuiltPayload> {
951 pub cached_reads: CachedReads,
953 pub execution_cache: Option<SavedCache>,
955 pub state_root_handle: Option<PayloadStateRootHandle>,
960 pub config: PayloadConfig<Attributes, HeaderTy<Payload::Primitives>>,
962 pub cancel: CancelOnDrop,
964 pub best_payload: Option<Payload>,
966}
967
968impl<Attributes, Payload: BuiltPayload> BuildArguments<Attributes, Payload> {
969 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
982pub trait PayloadBuilder: Send + Sync + Clone {
991 type Attributes: PayloadAttributes;
993 type BuiltPayload: BuiltPayload;
995
996 fn try_build(
1009 &self,
1010 args: BuildArguments<Self::Attributes, Self::BuiltPayload>,
1011 ) -> Result<BuildOutcome<Self::BuiltPayload>, PayloadBuilderError>;
1012
1013 fn on_missing_payload(
1017 &self,
1018 _args: BuildArguments<Self::Attributes, Self::BuiltPayload>,
1019 ) -> MissingPayloadBehaviour<Self::BuiltPayload> {
1020 MissingPayloadBehaviour::RaceEmptyPayload
1021 }
1022
1023 fn build_empty_payload(
1025 &self,
1026 config: PayloadConfig<Self::Attributes, HeaderForPayload<Self::BuiltPayload>>,
1027 ) -> Result<Self::BuiltPayload, PayloadBuilderError>;
1028}
1029
1030#[derive(Default)]
1034pub enum MissingPayloadBehaviour<Payload> {
1035 AwaitInProgress,
1037 #[default]
1039 RaceEmptyPayload,
1040 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#[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
1068fn 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}