Skip to main content

reth_engine_tree/tree/payload_processor/
bal_prewarm_pool.rs

1//! BAL read-set prewarming pool.
2
3use alloy_primitives::{Address, StorageKey};
4use reth_execution_cache::{CachedStateProvider, ExecutionCache, TxPoolPrewarmCacheSnapshot};
5use reth_provider::{EvmStateProvider, EvmStateProviderBox, ProviderResult};
6use std::{
7    sync::{
8        atomic::{AtomicUsize, Ordering},
9        Arc,
10    },
11    thread::JoinHandle,
12};
13use tokio::sync::oneshot;
14use tracing::trace;
15
16/// Builds a fresh EVM provider over the block's parent state. Type-erased so the pool is not
17/// generic over the provider factory; each worker builds its own per block.
18pub type BuildProviderFn = dyn Fn() -> ProviderResult<EvmStateProviderBox> + Send + Sync;
19
20/// A single warm request: an account and a batch of its storage slots, or a batch of storage slots
21/// on their own.
22enum PrewarmTarget {
23    Account(Address, Box<[StorageKey]>),
24    Storage(Address, Box<[StorageKey]>),
25}
26
27/// A message in a worker's queue. The per-block lifecycle is explicit and ordered (the queue is
28/// FIFO): one `BeginBlock`, then the worker's share of `Warm`s, then one `EndBlock`.
29enum PrewarmMsg {
30    /// Open a read txn for the new block: build a provider over the parent state and hold it.
31    BeginBlock {
32        build: Arc<BuildProviderFn>,
33        caches: ExecutionCache,
34        txpool_snapshot: Option<TxPoolPrewarmCacheSnapshot>,
35    },
36    /// Warm one target into the held provider's cache. Ignored if no provider is held.
37    Warm(PrewarmTarget),
38    /// Drop the held provider (and its read txn).
39    EndBlock(Arc<SendOnDrop>),
40}
41
42/// Long-lived pool of blocking threads that warm the BAL read-set into the shared execution cache.
43#[derive(Debug)]
44pub struct BalPrewarmPool {
45    /// One queue per worker. `BeginBlock`/`EndBlock` are broadcast to all; `Warm`s round-robin.
46    workers: Vec<crossbeam_channel::Sender<PrewarmMsg>>,
47    /// Round-robin cursor for distributing warm requests across workers.
48    next: AtomicUsize,
49    _handles: Vec<JoinHandle<()>>,
50}
51
52impl BalPrewarmPool {
53    /// Spawns `num_threads` long-lived blocking worker threads. Owned by the
54    /// [`PayloadProcessor`](super::PayloadProcessor); the threads exit when the pool is dropped.
55    pub fn new(num_threads: usize) -> Arc<Self> {
56        let mut workers = Vec::with_capacity(num_threads);
57        let mut handles = Vec::with_capacity(num_threads);
58        for i in 0..num_threads {
59            let (tx, rx) = crossbeam_channel::unbounded::<PrewarmMsg>();
60            workers.push(tx);
61            handles.push(
62                std::thread::Builder::new()
63                    .name(format!("bal-prewarm-{i:03}"))
64                    .spawn(move || prewarm_loop(rx))
65                    .expect("spawn bal-prewarm thread"),
66            );
67        }
68        trace!(target: "engine::tree::bal_prewarm_pool", num_threads, "BalPrewarmPool spawned");
69        Arc::new(Self { workers, next: AtomicUsize::new(0), _handles: handles })
70    }
71
72    /// Begins a block: hands every worker the provider builder and shared cache so each opens its
73    /// own read txn over the parent state. Pair with [`end_block`](Self::end_block).
74    pub fn begin_block(
75        &self,
76        build: Arc<BuildProviderFn>,
77        caches: ExecutionCache,
78        txpool_snapshot: Option<TxPoolPrewarmCacheSnapshot>,
79    ) {
80        for worker in &self.workers {
81            let _ = worker.send(PrewarmMsg::BeginBlock {
82                build: build.clone(),
83                caches: caches.clone(),
84                txpool_snapshot: txpool_snapshot.clone(),
85            });
86        }
87    }
88
89    /// Fire-and-forget: warm an account and its storage slots.
90    ///
91    /// The slots are dispatched in `WARM_BATCH_SIZE` chunks that are distributed independently,
92    /// so a single account with a large read-set does not serialize onto one worker;
93    /// [`end_block`](Self::end_block) waits for the slowest queue.
94    pub fn warm_account(&self, addr: Address, slots: impl IntoIterator<Item = StorageKey>) {
95        let mut slots = slots.into_iter();
96        let mut batch: Box<[StorageKey]> = slots.by_ref().take(WARM_BATCH_SIZE).collect();
97        self.send_warm(PrewarmTarget::Account(addr, batch));
98
99        loop {
100            batch = slots.by_ref().take(WARM_BATCH_SIZE).collect();
101            if batch.is_empty() {
102                break
103            }
104            self.send_warm(PrewarmTarget::Storage(addr, batch));
105        }
106    }
107
108    /// Ends the block: every worker drops its provider (and read txn) once it has drained the warm
109    /// requests queued ahead of this message.
110    ///
111    /// Blocks until all workers processed the end block message.
112    pub fn end_block(&self) {
113        let (tx, rx) = oneshot::channel();
114        let tx = Arc::new(SendOnDrop { sender: Some(tx) });
115
116        for worker in &self.workers {
117            let _ = worker.send(PrewarmMsg::EndBlock(tx.clone()));
118        }
119
120        drop(tx);
121        rx.blocking_recv().expect("BAL prewarm pool dropped without signaling completion");
122    }
123
124    fn send_warm(&self, target: PrewarmTarget) {
125        let i = self.next.fetch_add(1, Ordering::Relaxed) % self.workers.len();
126        let _ = self.workers[i].send(PrewarmMsg::Warm(target));
127    }
128}
129
130/// Number of warming threads.
131///
132/// The work performed on those threads boils down mostly to MDBX reads. An MDBX read consists of
133/// a tree traversal and major page faults causing I/O.
134///
135/// In order to utilize the parallelism of `NVMe` we have to give it enough work, or equally,
136/// maintain a high queue depth. Modern `NVMe` devices require in between 64-128 requests in-flight
137/// to achieve its peak performance. Ideally we don't grow past that but it's OK to do so, it just
138/// means that a request is going to wait in the `NVMe` queue rather than in memory.
139///
140/// MDBX piggy-backs on the OS page cache for its buffers. Oftentimes, the hit rate reaches 90-99%
141/// hit rate. At that point, the workload can be classified as CPU-bound. In that case, having
142/// a high number of threads is counterproductive due to the effects of context switching, core
143/// migration, contention, etc.
144///
145/// However, that overhead is considered negligible compared to the benefits of fully utilizing
146/// `NVMe` resources. For example, with request latency of 100µs, 100k IO requests the expected
147/// time to finish is 312.5ms at QD=32 and 156.26ms at QD=64.
148///
149/// This should explain why this particular value is picked.
150pub const DEFAULT_BAL_PREWARM_THREADS: usize = 128;
151
152/// Number of storage slots carried by one warm message.
153///
154/// Batching amortizes the send over many slots and hands the worker a run of slots that live
155/// close together in the storage table, while the cap keeps enough messages in flight to saturate
156/// the workers on blocks whose read-set is concentrated in a few accounts.
157const WARM_BATCH_SIZE: usize = 8;
158
159fn prewarm_loop(rx: crossbeam_channel::Receiver<PrewarmMsg>) {
160    // The provider (and its MDBX read txn) held for the current block, between `BeginBlock` and
161    // `EndBlock`. `None` while idle, so no read txn is pinned across the inter-block gap.
162    let mut provider: Option<CachedStateProvider<EvmStateProviderBox>> = None;
163
164    // Blocks when idle; the channel disconnects (and the loop ends) when the pool is dropped.
165    while let Ok(msg) = rx.recv() {
166        match msg {
167            PrewarmMsg::BeginBlock { build, caches, txpool_snapshot } => {
168                provider = match (build)() {
169                    Ok(inner) => Some(
170                        CachedStateProvider::new_prewarm(inner, caches)
171                            .with_txpool_snapshot(txpool_snapshot),
172                    ),
173                    Err(err) => {
174                        trace!(target: "engine::tree::bal_prewarm_pool", %err, "failed to build provider");
175                        None
176                    }
177                };
178            }
179            PrewarmMsg::Warm(target) => {
180                let Some(provider) = provider.as_ref() else { continue };
181                match target {
182                    PrewarmTarget::Account(addr, slots) => {
183                        let _ = provider.basic_account(&addr);
184                        for &slot in &slots {
185                            let _ = provider.storage(addr, slot);
186                        }
187                    }
188                    PrewarmTarget::Storage(addr, slots) => {
189                        for &slot in &slots {
190                            let _ = provider.storage(addr, slot);
191                        }
192                    }
193                }
194            }
195            PrewarmMsg::EndBlock(end_tx) => {
196                provider = None;
197                drop(end_tx);
198            }
199        }
200    }
201}
202
203struct SendOnDrop {
204    sender: Option<oneshot::Sender<()>>,
205}
206
207impl Drop for SendOnDrop {
208    fn drop(&mut self) {
209        if let Some(sender) = self.sender.take() {
210            let _ = sender.send(());
211        }
212    }
213}