1use crate::eth::pool::{Pool, transactions::PoolTransaction};
4use alloy_primitives::TxHash;
5use futures::{
6 channel::mpsc::Receiver,
7 stream::{Fuse, StreamExt},
8 task::AtomicWaker,
9};
10use parking_lot::{RawRwLock, RwLock, lock_api::RwLockWriteGuard};
11use std::{
12 fmt,
13 marker::PhantomData,
14 pin::Pin,
15 sync::{
16 Arc,
17 atomic::{AtomicU64, Ordering},
18 },
19 task::{Context, Poll},
20 time::Duration,
21};
22use tokio::time::{Interval, MissedTickBehavior, Sleep};
23
24pub(crate) const INSTANT_COALESCE_WINDOW: Duration = Duration::from_millis(5);
27
28pub struct Miner<T> {
29 mode: Arc<RwLock<MiningMode>>,
31 generation: Arc<AtomicU64>,
33 inner: Arc<MinerInner>,
37 transaction: PhantomData<fn() -> T>,
39}
40
41impl<T> Clone for Miner<T> {
42 fn clone(&self) -> Self {
43 Self {
44 mode: self.mode.clone(),
45 generation: self.generation.clone(),
46 inner: self.inner.clone(),
47 transaction: PhantomData,
48 }
49 }
50}
51
52impl<T> fmt::Debug for Miner<T> {
53 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
54 f.debug_struct("Miner").field("mode", &self.mode).finish_non_exhaustive()
55 }
56}
57
58impl<T> Miner<T> {
59 pub fn new(mode: MiningMode) -> Self {
61 Self {
62 mode: Arc::new(RwLock::new(mode)),
63 generation: Default::default(),
64 inner: Default::default(),
65 transaction: PhantomData,
66 }
67 }
68
69 pub fn mode_write(&self) -> RwLockWriteGuard<'_, RawRwLock, MiningMode> {
71 self.mode.write()
72 }
73
74 pub fn is_auto_mine(&self) -> bool {
76 let mode = self.mode.read();
77 matches!(*mode, MiningMode::Auto(_))
78 }
79
80 pub fn get_interval(&self) -> Option<u64> {
81 let mode = self.mode.read();
82 if let MiningMode::FixedBlockTime(ref mm) = *mode {
83 return Some(mm.interval.period().as_secs());
84 }
85 None
86 }
87
88 pub(crate) fn block_interval(&self) -> Option<Duration> {
90 let mode = self.mode.read();
91 match &*mode {
92 MiningMode::FixedBlockTime(miner) | MiningMode::Mixed(_, miner) => {
93 Some(miner.interval.period())
94 }
95 MiningMode::None | MiningMode::Auto(_) => None,
96 }
97 }
98
99 pub fn set_mining_mode(&self, mode: MiningMode) {
101 let new_mode = format!("{mode:?}");
102 let mut current = self.mode_write();
103 let mode = std::mem::replace(&mut *current, mode);
104 self.generation.fetch_add(1, Ordering::Relaxed);
105 drop(current);
106 trace!(target: "miner", "updated mining mode from {:?} to {}", mode, new_mode);
107 self.inner.wake();
108 }
109
110 pub(crate) fn handle_failed_candidate(&self, generation: u64) {
112 let mut mode = self.mode.write();
113 if self.generation.load(Ordering::Relaxed) != generation {
114 if let MiningMode::Auto(miner) | MiningMode::Mixed(miner, _) = &mut *mode {
115 miner.has_pending_txs = Some(true);
116 miner.coalesce = None;
117 }
118 return;
119 }
120 if let MiningMode::Auto(miner) | MiningMode::Mixed(miner, _) = &mut *mode {
123 miner.coalesce = None;
124 }
125 match &mut *mode {
126 MiningMode::FixedBlockTime(miner) | MiningMode::Mixed(_, miner) => {
127 let period = miner.interval.period();
128 *miner = FixedBlockTimeMiner::new(period);
129 }
130 MiningMode::None | MiningMode::Auto(_) => {}
131 }
132 }
133
134 pub(crate) fn poll(
139 &mut self,
140 pool: &Arc<Pool<T>>,
141 cx: &mut Context<'_>,
142 ) -> Poll<MiningWork<T>> {
143 self.inner.register(cx);
144 let mut mode = self.mode.write();
145 let generation = self.generation.load(Ordering::Relaxed);
146 mode.poll(pool, cx).map(|transactions| MiningWork { transactions, generation })
147 }
148
149 pub(crate) fn retry_ready_transactions(&self) {
153 if let MiningMode::Auto(miner) | MiningMode::Mixed(miner, _) = &mut *self.mode.write() {
154 miner.has_pending_txs = Some(true);
155 self.inner.wake();
156 }
157 }
158}
159
160pub(crate) struct MiningWork<T> {
162 pub(crate) transactions: Vec<Arc<PoolTransaction<T>>>,
163 pub(crate) generation: u64,
164}
165
166#[derive(Debug)]
168pub struct MinerInner {
169 waker: AtomicWaker,
170}
171
172impl MinerInner {
173 fn wake(&self) {
175 self.waker.wake();
176 }
177
178 fn register(&self, cx: &Context<'_>) {
179 self.waker.register(cx.waker());
180 }
181}
182
183impl Default for MinerInner {
184 fn default() -> Self {
185 Self { waker: AtomicWaker::new() }
186 }
187}
188
189#[derive(Debug)]
191pub enum MiningMode {
192 None,
194 Auto(ReadyTransactionMiner),
199 FixedBlockTime(FixedBlockTimeMiner),
201
202 Mixed(ReadyTransactionMiner, FixedBlockTimeMiner),
204}
205
206impl MiningMode {
207 pub fn instant(max_transactions: usize, listener: Receiver<TxHash>) -> Self {
208 Self::Auto(ReadyTransactionMiner {
209 max_transactions,
210 has_pending_txs: None,
211 rx: listener.fuse(),
212 coalesce: None,
213 coalesce_window: INSTANT_COALESCE_WINDOW,
214 })
215 }
216
217 pub fn interval(duration: Duration) -> Self {
218 Self::FixedBlockTime(FixedBlockTimeMiner::new(duration))
219 }
220
221 pub fn mixed(max_transactions: usize, listener: Receiver<TxHash>, duration: Duration) -> Self {
222 Self::Mixed(
223 ReadyTransactionMiner {
224 max_transactions,
225 has_pending_txs: None,
226 rx: listener.fuse(),
227 coalesce: None,
228 coalesce_window: INSTANT_COALESCE_WINDOW,
229 },
230 FixedBlockTimeMiner::new(duration),
231 )
232 }
233
234 #[must_use]
237 pub fn with_coalescing_window(mut self, window: Duration) -> Self {
238 if let Self::Auto(miner) | Self::Mixed(miner, _) = &mut self {
239 miner.coalesce_window = window;
240 miner.coalesce = None;
241 }
242 self
243 }
244
245 pub fn poll<T>(
247 &mut self,
248 pool: &Arc<Pool<T>>,
249 cx: &mut Context<'_>,
250 ) -> Poll<Vec<Arc<PoolTransaction<T>>>> {
251 match self {
252 Self::None => Poll::Pending,
253 Self::Auto(miner) => miner.poll(pool, cx),
254 Self::FixedBlockTime(miner) => miner.poll(pool, cx),
255 Self::Mixed(auto, fixed) => {
256 let auto_txs = auto.poll(pool, cx);
257 let fixed_txs = fixed.poll(pool, cx);
258
259 match (auto_txs, fixed_txs) {
260 (Poll::Ready(mut auto_txs), Poll::Ready(fixed_txs)) => {
262 for tx in fixed_txs {
263 if auto_txs.iter().any(|auto_tx| auto_tx.hash() == tx.hash()) {
265 continue;
266 }
267 auto_txs.push(tx);
268 }
269 Poll::Ready(auto_txs)
270 }
271 (Poll::Ready(auto_txs), Poll::Pending) => Poll::Ready(auto_txs),
273 (Poll::Pending, fixed_txs) => fixed_txs,
276 }
277 }
278 }
279 }
280}
281
282#[derive(Debug)]
287pub struct FixedBlockTimeMiner {
288 interval: Interval,
290}
291
292impl FixedBlockTimeMiner {
293 pub fn new(duration: Duration) -> Self {
295 let start = tokio::time::Instant::now() + duration;
296 let mut interval = tokio::time::interval_at(start, duration);
297 interval.set_missed_tick_behavior(MissedTickBehavior::Delay);
300 Self { interval }
301 }
302
303 fn poll<T>(
304 &mut self,
305 pool: &Arc<Pool<T>>,
306 cx: &mut Context<'_>,
307 ) -> Poll<Vec<Arc<PoolTransaction<T>>>> {
308 if self.interval.poll_tick(cx).is_ready() {
309 return Poll::Ready(pool.ready_transactions().collect());
311 }
312 Poll::Pending
313 }
314}
315
316impl Default for FixedBlockTimeMiner {
317 fn default() -> Self {
318 Self::new(Duration::from_secs(6))
319 }
320}
321
322pub struct ReadyTransactionMiner {
324 max_transactions: usize,
326 has_pending_txs: Option<bool>,
328 rx: Fuse<Receiver<TxHash>>,
330 coalesce_window: Duration,
332 coalesce: Option<Pin<Box<Sleep>>>,
334}
335
336impl ReadyTransactionMiner {
337 fn poll<T>(
338 &mut self,
339 pool: &Arc<Pool<T>>,
340 cx: &mut Context<'_>,
341 ) -> Poll<Vec<Arc<PoolTransaction<T>>>> {
342 let mut saw_new_ready = false;
344 while let Poll::Ready(Some(_hash)) = self.rx.poll_next_unpin(cx) {
345 saw_new_ready = true;
346 }
347
348 if saw_new_ready {
351 self.has_pending_txs = Some(true);
352 if self.coalesce.is_none() && !self.coalesce_window.is_zero() {
353 self.coalesce = Some(Box::pin(tokio::time::sleep(self.coalesce_window)));
354 }
355 }
356
357 if self.has_pending_txs == Some(false) {
358 return Poll::Pending;
359 }
360
361 if let Some(sleep) = self.coalesce.as_mut()
362 && sleep.as_mut().poll(cx).is_pending()
363 {
364 return Poll::Pending;
365 }
366 self.coalesce = None;
367
368 let transactions =
369 pool.ready_transactions().take(self.max_transactions).collect::<Vec<_>>();
370
371 self.has_pending_txs = Some(false);
374
375 if transactions.is_empty() {
376 return Poll::Pending;
377 }
378
379 Poll::Ready(transactions)
380 }
381}
382
383impl fmt::Debug for ReadyTransactionMiner {
384 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
385 f.debug_struct("ReadyTransactionMiner")
386 .field("max_transactions", &self.max_transactions)
387 .finish_non_exhaustive()
388 }
389}
390
391#[cfg(test)]
392mod tests {
393 use super::*;
394 use futures::{channel::mpsc, future::poll_fn};
395
396 #[test]
397 fn stale_failure_resumes_replacement_autominer() {
398 let (_tx, rx) = mpsc::channel(1);
399 let miner = Miner::<()>::new(MiningMode::None);
400 miner.set_mining_mode(MiningMode::instant(1, rx));
401
402 miner.handle_failed_candidate(0);
403
404 let mode = miner.mode.read();
405 let MiningMode::Auto(auto) = &*mode else { panic!("expected auto mining") };
406 assert_eq!(auto.has_pending_txs, Some(true));
407 }
408
409 #[test]
410 fn failure_keeps_retry_requested_after_candidate_selection() {
411 let (_tx, rx) = mpsc::channel(1);
412 let miner = Miner::<()>::new(MiningMode::instant(1, rx));
413 miner.retry_ready_transactions();
414
415 miner.handle_failed_candidate(0);
416
417 let mode = miner.mode.read();
418 let MiningMode::Auto(auto) = &*mode else { panic!("expected auto mining") };
419 assert_eq!(auto.has_pending_txs, Some(true));
420 }
421
422 #[tokio::test]
423 async fn failed_fixed_candidate_rearms_interval() {
424 let mut miner = Miner::<()>::new(MiningMode::interval(Duration::from_millis(10)));
425 let pool = Arc::new(Pool::default());
426 tokio::time::timeout(Duration::from_secs(1), poll_fn(|cx| miner.poll(&pool, cx)))
427 .await
428 .unwrap();
429
430 miner.handle_failed_candidate(0);
431
432 tokio::time::timeout(Duration::from_secs(1), poll_fn(|cx| miner.poll(&pool, cx)))
433 .await
434 .unwrap();
435 }
436
437 #[tokio::test]
438 async fn coalescing_window_controls_fresh_notifications() {
439 for mixed in [false, true] {
440 for window in [Duration::ZERO, Duration::from_secs(60)] {
441 let (mut tx, rx) = mpsc::channel(1);
442 let mode = if mixed {
443 MiningMode::mixed(1, rx, Duration::from_secs(120))
444 } else {
445 MiningMode::instant(1, rx)
446 };
447 let mut mode = mode.with_coalescing_window(window);
448 let pool = Arc::new(Pool::<()>::default());
449 tx.try_send(TxHash::ZERO).unwrap();
450 let mut cx = Context::from_waker(futures::task::noop_waker_ref());
451 assert!(mode.poll(&pool, &mut cx).is_pending());
452 let auto = match &mode {
453 MiningMode::Auto(auto) | MiningMode::Mixed(auto, _) => auto,
454 _ => unreachable!(),
455 };
456 if window.is_zero() {
457 assert!(auto.coalesce.is_none());
458 assert_eq!(auto.has_pending_txs, Some(false));
459 } else {
460 assert!(auto.coalesce.is_some());
461 assert_eq!(auto.has_pending_txs, Some(true));
462 }
463 }
464 }
465 }
466}