1use crate::{
7 metrics::PayloadBuilderServiceMetrics, traits::PayloadJobGenerator, KeepPayloadJobAlive,
8 PayloadJob,
9};
10use alloy_consensus::BlockHeader;
11use alloy_primitives::{BlockTimestamp, B256};
12use alloy_rpc_types::engine::PayloadId;
13use futures_util::{future::FutureExt, Stream, StreamExt};
14use reth_chain_state::CanonStateNotification;
15use reth_execution_cache::SavedCache;
16use reth_payload_builder_primitives::{Events, PayloadBuilderError, PayloadEvents};
17use reth_payload_primitives::{BuiltPayload, PayloadAttributes, PayloadKind, PayloadTypes};
18use reth_primitives_traits::{FastInstant as Instant, NodePrimitives};
19use reth_trie_parallel::state_root_task::PayloadStateRootHandle;
20use std::{
21 future::Future,
22 pin::Pin,
23 sync::Arc,
24 task::{Context, Poll},
25};
26use tokio::sync::{
27 broadcast, mpsc,
28 oneshot::{self, Receiver},
29 watch,
30};
31use tokio_stream::wrappers::UnboundedReceiverStream;
32use tracing::{debug, debug_span, info, trace, warn, Span};
33
34type PayloadFuture<P> = Pin<Box<dyn Future<Output = Result<P, PayloadBuilderError>> + Send>>;
35type ResolvePayloadResult<P, Job> = (Option<PayloadFuture<P>>, Option<PayloadJobEntry<Job>>);
36
37#[derive(Debug)]
42pub struct PayloadStore<T: PayloadTypes> {
43 inner: Arc<PayloadBuilderHandle<T>>,
44}
45
46impl<T> PayloadStore<T>
47where
48 T: PayloadTypes,
49{
50 pub fn resolve_kind(
55 &self,
56 id: PayloadId,
57 kind: PayloadKind,
58 ) -> impl Future<Output = Option<Result<T::BuiltPayload, PayloadBuilderError>>> {
59 self.inner.resolve_kind(id, kind)
60 }
61
62 pub async fn resolve(
64 &self,
65 id: PayloadId,
66 ) -> Option<Result<T::BuiltPayload, PayloadBuilderError>> {
67 self.resolve_kind(id, PayloadKind::Earliest).await
68 }
69
70 pub async fn best_payload(
74 &self,
75 id: PayloadId,
76 ) -> Option<Result<T::BuiltPayload, PayloadBuilderError>> {
77 self.inner.best_payload(id).await
78 }
79
80 pub async fn payload_timestamp(
84 &self,
85 id: PayloadId,
86 ) -> Option<Result<u64, PayloadBuilderError>> {
87 self.inner.payload_timestamp(id).await
88 }
89
90 pub fn new(inner: PayloadBuilderHandle<T>) -> Self {
92 Self { inner: Arc::new(inner) }
93 }
94}
95
96impl<T> From<PayloadBuilderHandle<T>> for PayloadStore<T>
97where
98 T: PayloadTypes,
99{
100 fn from(inner: PayloadBuilderHandle<T>) -> Self {
101 Self::new(inner)
102 }
103}
104
105#[derive(Debug)]
109pub struct PayloadBuilderHandle<T: PayloadTypes> {
110 to_service: mpsc::UnboundedSender<PayloadServiceCommand<T>>,
112}
113
114impl<T: PayloadTypes> PayloadBuilderHandle<T> {
115 pub const fn new(to_service: mpsc::UnboundedSender<PayloadServiceCommand<T>>) -> Self {
120 Self { to_service }
121 }
122
123 pub fn send_new_payload(
127 &self,
128 input: BuildNewPayload<T::PayloadAttributes>,
129 ) -> Receiver<Result<PayloadId, PayloadBuilderError>> {
130 let (tx, rx) = oneshot::channel();
131 let span = debug_span!(parent: Span::current(), "payload_job");
132 let _ =
133 self.to_service.send(PayloadServiceCommand::BuildNewPayload(input.into(), span, tx));
134 rx
135 }
136
137 pub async fn best_payload(
140 &self,
141 id: PayloadId,
142 ) -> Option<Result<T::BuiltPayload, PayloadBuilderError>> {
143 let (tx, rx) = oneshot::channel();
144 self.to_service.send(PayloadServiceCommand::BestPayload(id, tx)).ok()?;
145 rx.await.ok()?
146 }
147
148 pub fn resolve_kind(
156 &self,
157 id: PayloadId,
158 kind: PayloadKind,
159 ) -> impl Future<Output = Option<Result<T::BuiltPayload, PayloadBuilderError>>> {
160 let (tx, rx) = oneshot::channel();
161 let sent = self.to_service.send(PayloadServiceCommand::Resolve(id, kind, tx)).is_ok();
162 async move {
163 if !sent {
164 return None
165 }
166
167 match rx.await.transpose()? {
168 Ok(fut) => Some(fut.await),
169 Err(e) => Some(Err(e.into())),
170 }
171 }
172 }
173
174 pub async fn subscribe(&self) -> Result<PayloadEvents<T>, PayloadBuilderError> {
177 let (tx, rx) = oneshot::channel();
178 let _ = self.to_service.send(PayloadServiceCommand::Subscribe(tx));
179 Ok(PayloadEvents { receiver: rx.await? })
180 }
181
182 pub async fn payload_timestamp(
186 &self,
187 id: PayloadId,
188 ) -> Option<Result<u64, PayloadBuilderError>> {
189 let (tx, rx) = oneshot::channel();
190 self.to_service.send(PayloadServiceCommand::PayloadTimestamp(id, tx)).ok()?;
191 rx.await.ok()?
192 }
193}
194
195impl<T> Clone for PayloadBuilderHandle<T>
196where
197 T: PayloadTypes,
198{
199 fn clone(&self) -> Self {
200 Self { to_service: self.to_service.clone() }
201 }
202}
203
204#[derive(Debug)]
213#[must_use = "futures do nothing unless you `.await` or poll them"]
214pub struct PayloadBuilderService<Gen, St, T>
215where
216 T: PayloadTypes,
217 Gen: PayloadJobGenerator,
218 Gen::Job: PayloadJob<PayloadAttributes = T::PayloadAttributes>,
219{
220 generator: Gen,
222 payload_jobs: Vec<PayloadJobEntry<Gen::Job>>,
226 service_tx: mpsc::UnboundedSender<PayloadServiceCommand<T>>,
228 command_rx: UnboundedReceiverStream<PayloadServiceCommand<T>>,
230 metrics: PayloadBuilderServiceMetrics,
232 chain_events: St,
234 payload_events: broadcast::Sender<Events<T>>,
236 cached_payload_rx: watch::Receiver<Option<(PayloadId, BlockTimestamp, T::BuiltPayload)>>,
239 cached_payload_tx: watch::Sender<Option<(PayloadId, BlockTimestamp, T::BuiltPayload)>>,
241}
242
243const PAYLOAD_EVENTS_BUFFER_SIZE: usize = 20;
244
245impl<Gen, St, T> PayloadBuilderService<Gen, St, T>
248where
249 T: PayloadTypes,
250 Gen: PayloadJobGenerator,
251 Gen::Job: PayloadJob<PayloadAttributes = T::PayloadAttributes>,
252 <Gen::Job as PayloadJob>::BuiltPayload: Into<T::BuiltPayload>,
253{
254 pub fn new(generator: Gen, chain_events: St) -> (Self, PayloadBuilderHandle<T>) {
261 let (service_tx, command_rx) = mpsc::unbounded_channel();
262 let (payload_events, _) = broadcast::channel(PAYLOAD_EVENTS_BUFFER_SIZE);
263
264 let (cached_payload_tx, cached_payload_rx) = watch::channel(None);
265
266 let service = Self {
267 generator,
268 payload_jobs: Vec::new(),
269 service_tx,
270 command_rx: UnboundedReceiverStream::new(command_rx),
271 metrics: Default::default(),
272 chain_events,
273 payload_events,
274 cached_payload_rx,
275 cached_payload_tx,
276 };
277
278 let handle = service.handle();
279 (service, handle)
280 }
281
282 pub fn handle(&self) -> PayloadBuilderHandle<T> {
284 PayloadBuilderHandle::new(self.service_tx.clone())
285 }
286
287 pub fn payload_events_handle(&self) -> broadcast::Sender<Events<T>> {
290 self.payload_events.clone()
291 }
292
293 fn contains_payload(&self, id: PayloadId) -> bool {
295 self.payload_jobs.iter().any(|entry| entry.id == id)
296 }
297
298 fn best_payload(&self, id: PayloadId) -> Option<Result<T::BuiltPayload, PayloadBuilderError>> {
300 let res = self
301 .payload_jobs
302 .iter()
303 .find(|entry| entry.id == id)
304 .map(|entry| entry.job.best_payload().map(|payload| payload.into()));
305 if let Some(Ok(ref best)) = res {
306 self.metrics.set_best_revenue(best.block().number(), f64::from(best.fees()));
307 }
308
309 res
310 }
311
312 fn resolve(
317 &mut self,
318 id: PayloadId,
319 kind: PayloadKind,
320 ) -> ResolvePayloadResult<T::BuiltPayload, Gen::Job> {
321 let start = Instant::now();
322 debug!(target: "payload_builder", %id, "resolving payload job");
323
324 if let Some((cached, _, payload)) = &*self.cached_payload_rx.borrow() &&
325 *cached == id
326 {
327 self.metrics.resolve_duration_seconds.record(start.elapsed());
328 return (Some(Box::pin(core::future::ready(Ok(payload.clone())))), None);
329 }
330
331 let Some(job) = self.payload_jobs.iter().position(|entry| entry.id == id) else {
332 return (None, None)
333 };
334 let (fut, keep_alive) = self.payload_jobs[job].job.resolve_kind(kind);
335 let payload_timestamp = self.payload_jobs[job].job.payload_timestamp();
336
337 let mut resolved_job =
338 (keep_alive == KeepPayloadJobAlive::No).then(|| self.payload_jobs.swap_remove(job));
339 let leases = resolved_job
340 .as_mut()
341 .map(|entry| std::mem::take(&mut entry.leases))
342 .unwrap_or_default();
343
344 let resolved_metrics = self.metrics.clone();
347 let payload_events = self.payload_events.clone();
348 let cached_payload_tx = self.cached_payload_tx.clone();
349
350 let fut = async move {
351 let _leases = leases;
352 let res = fut.await;
353 resolved_metrics.resolve_duration_seconds.record(start.elapsed());
354 if let Ok(payload) = &res {
355 if payload_events.receiver_count() > 0 {
356 payload_events.send(Events::BuiltPayload(payload.clone().into())).ok();
357 }
358
359 if let Ok(timestamp) = payload_timestamp {
360 let _ = cached_payload_tx.send(Some((id, timestamp, payload.clone().into())));
361 }
362
363 resolved_metrics
364 .set_resolved_revenue(payload.block().number(), f64::from(payload.fees()));
365 }
366 res.map(|p| p.into())
367 };
368
369 (Some(Box::pin(fut)), resolved_job)
370 }
371
372 fn payload_timestamp(&self, id: PayloadId) -> Option<Result<u64, PayloadBuilderError>> {
374 if let Some((cached_id, timestamp, _)) = *self.cached_payload_rx.borrow() &&
375 cached_id == id
376 {
377 return Some(Ok(timestamp));
378 }
379
380 let timestamp = self
381 .payload_jobs
382 .iter()
383 .find(|entry| entry.id == id)
384 .map(|entry| entry.job.payload_timestamp());
385
386 if timestamp.is_none() {
387 trace!(target: "payload_builder", %id, "no matching payload job found to get timestamp for");
388 }
389
390 timestamp
391 }
392}
393
394impl<Gen, St, T, N> Future for PayloadBuilderService<Gen, St, T>
395where
396 T: PayloadTypes,
397 N: NodePrimitives,
398 Gen: PayloadJobGenerator + Unpin + 'static,
399 <Gen as PayloadJobGenerator>::Job: Unpin + 'static,
400 St: Stream<Item = CanonStateNotification<N>> + Send + Unpin + 'static,
401 Gen::Job: PayloadJob<PayloadAttributes = T::PayloadAttributes>,
402 <Gen::Job as PayloadJob>::BuiltPayload: Into<T::BuiltPayload>,
403{
404 type Output = ();
405
406 fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
407 let this = self.get_mut();
408 loop {
409 while let Poll::Ready(Some(new_head)) = this.chain_events.poll_next_unpin(cx) {
411 this.generator.on_new_state(new_head);
412 }
413
414 for idx in (0..this.payload_jobs.len()).rev() {
418 let PayloadJobEntry { mut job, id, span, leases } =
419 this.payload_jobs.swap_remove(idx);
420
421 let poll_result = {
422 let _entered = span.enter();
423 job.poll_unpin(cx)
424 };
425
426 match poll_result {
427 Poll::Ready(Ok(_)) => {
428 this.metrics.set_active_jobs(this.payload_jobs.len());
429 trace!(target: "payload_builder", %id, "payload job finished");
430 }
431 Poll::Ready(Err(err)) => {
432 warn!(target: "payload_builder",%err, ?id, "Payload builder job failed; resolving payload");
433 this.metrics.inc_failed_jobs();
434 this.metrics.set_active_jobs(this.payload_jobs.len());
435 }
436 Poll::Pending => {
437 this.payload_jobs.push(PayloadJobEntry { job, id, span, leases });
438 }
439 }
440 }
441
442 let mut new_job = false;
444
445 while let Poll::Ready(Some(cmd)) = this.command_rx.poll_next_unpin(cx) {
447 match cmd {
448 PayloadServiceCommand::BuildNewPayload(input, job_span, tx) => {
449 let id = input.payload_id();
450 let mut res = Ok(id);
451 let parent = input.parent_hash;
452
453 if this.contains_payload(id) {
454 debug!(target: "payload_builder", %id, %parent, "Payload job already in progress, ignoring.");
455 } else {
456 let start = Instant::now();
457 let attributes = input.attributes.clone();
458 let leases = input.resources.clone_leases();
462 let job_result = {
463 let _entered = job_span.enter();
464 this.generator.new_payload_job(*input, id)
465 };
466
467 match job_result {
468 Ok(job) => {
469 this.metrics.new_job_duration_seconds.record(start.elapsed());
470 info!(target: "payload_builder", %id, %parent, "New payload job created");
471 this.metrics.inc_initiated_jobs();
472 new_job = true;
473 this.payload_jobs.push(PayloadJobEntry {
474 job,
475 id,
476 span: job_span,
477 leases,
478 });
479 this.payload_events.send(Events::Attributes(attributes)).ok();
480
481 if this
485 .cached_payload_rx
486 .borrow()
487 .as_ref()
488 .is_some_and(|(cached_id, _, _)| *cached_id == id)
489 {
490 trace!(target: "payload_builder", %id, "clearing stale cached payload for reused payload id");
491 let _ = this.cached_payload_tx.send(None);
492 }
493 }
494 Err(err) => {
495 this.metrics.new_job_duration_seconds.record(start.elapsed());
496 this.metrics.inc_failed_jobs();
497 warn!(target: "payload_builder", %err, %id, "Failed to create payload builder job");
498 res = Err(err);
499 }
500 }
501 }
502
503 let _ = tx.send(res);
504 }
505 PayloadServiceCommand::BestPayload(id, tx) => {
506 let _ = tx.send(this.best_payload(id));
507 }
508 PayloadServiceCommand::PayloadTimestamp(id, tx) => {
509 let timestamp = this.payload_timestamp(id);
510 let _ = tx.send(timestamp);
511 }
512 PayloadServiceCommand::Resolve(id, strategy, tx) => {
513 let (payload_fut, resolved_job) = this.resolve(id, strategy);
514 let _ = tx.send(payload_fut);
515
516 if let Some(entry) = resolved_job {
517 debug!(target: "payload_builder", id = %entry.id, "terminated resolved job");
518 }
519 }
520 PayloadServiceCommand::Subscribe(tx) => {
521 let new_rx = this.payload_events.subscribe();
522 let _ = tx.send(new_rx);
523 }
524 }
525 }
526
527 if !new_job {
528 return Poll::Pending
529 }
530 }
531 }
532}
533
534#[derive(derive_more::Debug)]
536pub enum PayloadServiceCommand<T: PayloadTypes> {
537 BuildNewPayload(
542 Box<BuildNewPayload<T::PayloadAttributes>>,
543 Span,
544 oneshot::Sender<Result<PayloadId, PayloadBuilderError>>,
545 ),
546 BestPayload(PayloadId, oneshot::Sender<Option<Result<T::BuiltPayload, PayloadBuilderError>>>),
548 PayloadTimestamp(PayloadId, oneshot::Sender<Option<Result<u64, PayloadBuilderError>>>),
550 Resolve(
552 PayloadId,
553 PayloadKind,
554 #[debug(skip)] oneshot::Sender<Option<PayloadFuture<T::BuiltPayload>>>,
555 ),
556 Subscribe(oneshot::Sender<broadcast::Receiver<Events<T>>>),
558}
559
560#[derive(Debug)]
562pub struct BuildNewPayload<T> {
563 pub attributes: T,
565 pub parent_hash: B256,
567 pub resources: PayloadBuilderResources,
569}
570
571impl<T: PayloadAttributes> BuildNewPayload<T> {
572 pub fn payload_id(&self) -> PayloadId {
574 self.attributes.payload_id(&self.parent_hash)
575 }
576}
577
578#[derive(Debug, Default)]
580pub struct PayloadBuilderResources {
581 execution_cache: Option<SavedCache>,
585 state_root_handle: Option<PayloadStateRootHandle>,
587 leases: Vec<PayloadBuilderLease>,
589}
590
591impl PayloadBuilderResources {
592 pub const fn new(
594 execution_cache: Option<SavedCache>,
595 state_root_handle: Option<PayloadStateRootHandle>,
596 ) -> Self {
597 Self { execution_cache, state_root_handle, leases: Vec::new() }
598 }
599
600 pub fn with_lease(mut self, lease: PayloadBuilderLease) -> Self {
602 self.leases.push(lease);
603 self
604 }
605
606 pub const fn execution_cache(&self) -> Option<&SavedCache> {
608 self.execution_cache.as_ref()
609 }
610
611 pub const fn take_execution_cache(&mut self) -> Option<SavedCache> {
613 self.execution_cache.take()
614 }
615
616 pub const fn state_root_handle(&self) -> Option<&PayloadStateRootHandle> {
618 self.state_root_handle.as_ref()
619 }
620
621 pub const fn take_state_root_handle(&mut self) -> Option<PayloadStateRootHandle> {
623 self.state_root_handle.take()
624 }
625
626 pub fn take_leases(&mut self) -> Vec<PayloadBuilderLease> {
628 std::mem::take(&mut self.leases)
629 }
630
631 fn clone_leases(&self) -> Vec<PayloadBuilderLease> {
633 self.leases.clone()
634 }
635}
636
637#[derive(Clone)]
639pub struct PayloadBuilderLease {
640 _lease: Arc<dyn Send + Sync>,
641}
642
643impl PayloadBuilderLease {
644 pub fn new(lease: impl Send + Sync + 'static) -> Self {
646 Self { _lease: Arc::new(lease) }
647 }
648}
649
650impl std::fmt::Debug for PayloadBuilderLease {
651 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
652 f.debug_struct("PayloadBuilderLease").finish_non_exhaustive()
653 }
654}
655
656#[derive(Debug)]
658struct PayloadJobEntry<Job> {
659 job: Job,
660 id: PayloadId,
661 span: Span,
662 leases: Vec<PayloadBuilderLease>,
663}
664
665#[cfg(test)]
666mod tests {
667 use super::*;
668 use crate::test_utils::test_payload_service;
669 use alloy_primitives::Address;
670 use reth_ethereum_engine_primitives::{EthEngineTypes, EthPayloadAttributes};
671 use std::sync::atomic::{AtomicBool, Ordering};
672
673 struct DropProbe(Arc<AtomicBool>);
674
675 impl Drop for DropProbe {
676 fn drop(&mut self) {
677 self.0.store(true, Ordering::Release);
678 }
679 }
680
681 #[test]
682 fn payload_builder_lease_is_held_until_resolve_finishes() {
683 tokio::runtime::Builder::new_current_thread().build().unwrap().block_on(async {
684 let (service, handle) = test_payload_service::<EthEngineTypes>();
685 let service = tokio::spawn(service);
686 let dropped = Arc::new(AtomicBool::new(false));
687 let lease = PayloadBuilderLease::new(DropProbe(Arc::clone(&dropped)));
688 let input = BuildNewPayload {
689 attributes: EthPayloadAttributes {
690 timestamp: 1,
691 prev_randao: B256::ZERO,
692 suggested_fee_recipient: Address::ZERO,
693 withdrawals: None,
694 parent_beacon_block_root: None,
695 slot_number: None,
696 ..Default::default()
697 },
698 parent_hash: B256::ZERO,
699 resources: PayloadBuilderResources::default().with_lease(lease),
700 };
701
702 let id = handle.send_new_payload(input).await.unwrap().unwrap();
703 assert!(!dropped.load(Ordering::Acquire));
704
705 handle.resolve_kind(id, PayloadKind::Earliest).await.unwrap().unwrap();
706 assert!(dropped.load(Ordering::Acquire));
707 service.abort();
708 });
709 }
710}