1use crate::{
30 eth::{
31 error::PoolError,
32 pool::transactions::{
33 PendingPoolTransaction, PendingTransactions, PoolTransaction, ReadyTransactions,
34 TransactionsIterator, TxMarker,
35 },
36 },
37 mem::storage::MinedBlockOutcome,
38};
39use alloy_consensus::Transaction;
40use alloy_primitives::{Address, TxHash};
41use alloy_rpc_types::txpool::TxpoolStatus;
42use anvil_core::eth::transaction::PendingTransaction;
43use futures::channel::mpsc::{Receiver, Sender, channel};
44use parking_lot::{Mutex, RwLock};
45use std::{collections::VecDeque, fmt, sync::Arc};
46
47#[cfg(feature = "base")]
48use alloy_consensus::Typed2718;
49
50pub mod transactions;
51
52pub struct Pool<T> {
54 inner: RwLock<PoolInner<T>>,
56 transaction_listener: Mutex<Vec<Sender<TxHash>>>,
58}
59
60#[derive(Debug)]
62pub(crate) struct PoolSnapshot<T> {
63 inner: PoolInner<T>,
64}
65
66impl<T> Default for Pool<T> {
67 fn default() -> Self {
68 Self { inner: RwLock::new(PoolInner::default()), transaction_listener: Default::default() }
69 }
70}
71
72impl<T> Pool<T> {
75 pub(crate) fn snapshot(&self) -> PoolSnapshot<T> {
77 let pool = self.inner.read();
78 PoolSnapshot {
79 inner: PoolInner {
80 ready_transactions: pool.ready_transactions.snapshot(),
81 pending_transactions: pool.pending_transactions.snapshot(),
82 },
83 }
84 }
85
86 pub(crate) fn restore(&self, snapshot: PoolSnapshot<T>) {
88 let ready = {
89 let mut pool = self.inner.write();
90 *pool = snapshot.inner;
91 pool.ready_transactions().map(|tx| tx.hash()).collect::<Vec<_>>()
92 };
93 for hash in ready {
94 self.notify_listener(hash);
95 }
96 }
97
98 pub fn ready_transactions(&self) -> TransactionsIterator<T> {
100 self.inner.read().ready_transactions()
101 }
102
103 pub fn pending_transactions(&self) -> Vec<Arc<PoolTransaction<T>>> {
105 self.inner.read().pending_transactions.transactions().collect()
106 }
107
108 #[cfg(feature = "base")]
110 pub fn all_transactions(&self) -> Vec<Arc<PoolTransaction<T>>> {
111 let pool = self.inner.read();
112 pool.pending_transactions.transactions().chain(pool.ready_transactions()).collect()
113 }
114
115 pub fn txpool_status(&self) -> TxpoolStatus {
117 let pending: u64 = self.inner.read().ready_transactions.len().try_into().unwrap_or(0);
119 let queued: u64 = self.inner.read().pending_transactions.len().try_into().unwrap_or(0);
120 TxpoolStatus { pending, queued }
121 }
122
123 pub fn add_ready_listener(&self) -> Receiver<TxHash> {
126 const TX_LISTENER_BUFFER_SIZE: usize = 2048;
127 let (tx, rx) = channel(TX_LISTENER_BUFFER_SIZE);
128 self.transaction_listener.lock().push(tx);
129 rx
130 }
131
132 pub fn contains(&self, tx_hash: &TxHash) -> bool {
134 self.inner.read().contains(tx_hash)
135 }
136
137 pub fn contains_sender_nonce(&self, sender: Address, nonce: u64) -> bool
139 where
140 T: Transaction,
141 {
142 self.inner
143 .read()
144 .transactions_by_sender(sender)
145 .any(|tx| tx.pending_transaction.nonce() == nonce)
146 }
147
148 #[cfg(feature = "base")]
150 pub fn transaction_with_markers(
151 &self,
152 sender: Address,
153 markers: &[TxMarker],
154 ) -> Option<Arc<PoolTransaction<T>>> {
155 self.inner.read().transactions_by_sender(sender).find(|tx| tx.provides == markers)
156 }
157
158 pub fn clear(&self) {
160 let mut pool = self.inner.write();
161 pool.clear();
162 }
163
164 pub fn remove_invalid(&self, tx_hashes: Vec<TxHash>) -> Vec<Arc<PoolTransaction<T>>> {
166 self.inner.write().remove_invalid(tx_hashes)
167 }
168
169 pub fn remove_transactions_by_address(&self, sender: Address) -> Vec<Arc<PoolTransaction<T>>> {
171 self.inner.write().remove_transactions_by_address(sender)
172 }
173
174 pub fn drop_transaction(&self, tx: TxHash) -> Option<Arc<PoolTransaction<T>>> {
180 trace!(target: "txpool", "Dropping transaction: [{:?}]", tx);
181 let removed = {
182 let mut pool = self.inner.write();
183 let mut removed = pool.ready_transactions.remove_with_markers(vec![tx], None);
184 let invalidated = removed.iter().flat_map(|tx| tx.provides.iter().cloned());
185 removed.extend(pool.pending_transactions.remove_with_dependents(vec![tx], invalidated));
186 removed
187 };
188 trace!(target: "txpool", "Dropped transactions: {:?}", removed.iter().map(|tx| tx.hash()).collect::<Vec<_>>());
189
190 if removed.is_empty() {
191 None
192 } else {
193 removed.into_iter().find(|t| *t.pending_transaction.hash() == tx)
194 }
195 }
196
197 fn notify_ready(&self, tx: &AddedTransaction<T>) {
199 if let AddedTransaction::Ready(ready) = tx {
200 self.notify_listener(ready.hash);
201 for promoted in ready.promoted.iter().copied() {
202 self.notify_listener(promoted);
203 }
204 }
205 }
206
207 fn notify_listener(&self, hash: TxHash) {
209 let mut listener = self.transaction_listener.lock();
210 for n in (0..listener.len()).rev() {
212 let mut listener_tx = listener.swap_remove(n);
213 let retain = match listener_tx.try_send(hash) {
214 Ok(()) => true,
215 Err(e) => {
216 if e.is_full() {
217 warn!(
218 target: "txpool",
219 "[{:?}] Failed to send tx notification because channel is full",
220 hash,
221 );
222 true
223 } else {
224 false
225 }
226 }
227 };
228 if retain {
229 listener.push(listener_tx)
230 }
231 }
232 }
233}
234
235impl<T: Clone> Pool<T> {
236 pub fn get_transaction(&self, hash: TxHash) -> Option<PendingTransaction<T>> {
238 self.inner.read().get_transaction(hash)
239 }
240}
241
242impl<T: Transaction> Pool<T> {
243 pub fn on_mined_block(self: &Arc<Self>, outcome: MinedBlockOutcome<T>) -> bool {
250 let MinedBlockOutcome { block_number, included, stale, invalid, not_yet_valid } = outcome;
251 let made_progress = !included.is_empty() || !stale.is_empty() || !invalid.is_empty();
255 let retry_ready = made_progress && not_yet_valid.is_empty();
256
257 self.remove_invalid(invalid.into_iter().map(|tx| tx.hash()).collect());
259
260 let res = self.prune_markers(
262 block_number,
263 included.into_iter().chain(stale).flat_map(|tx| tx.provides.clone()),
264 );
265 trace!(target: "txpool", "pruned transaction markers {:?}", res);
266
267 if !not_yet_valid.is_empty() {
270 let tx_hashes: Vec<_> = not_yet_valid.iter().map(|tx| tx.hash()).collect();
271 let pool = Arc::clone(self);
272 tokio::spawn(async move {
273 tokio::time::sleep(std::time::Duration::from_secs(1)).await;
274 for hash in tx_hashes {
275 trace!(target: "txpool", "re-notifying for not-yet-valid tx: {:?}", hash);
276 pool.notify_listener(hash);
277 }
278 });
279 }
280
281 retry_ready && !self.inner.read().ready_transactions.is_empty()
282 }
283
284 pub fn prune_markers(
289 &self,
290 block_number: u64,
291 markers: impl IntoIterator<Item = TxMarker>,
292 ) -> PruneResult<T> {
293 debug!(target: "txpool", ?block_number, "pruning transactions");
294 let res = self.inner.write().prune_markers(markers);
295 for tx in &res.promoted {
296 self.notify_ready(tx);
297 }
298 res
299 }
300
301 pub fn add_transaction(
303 &self,
304 tx: PoolTransaction<T>,
305 ) -> Result<AddedTransaction<T>, PoolError> {
306 let added = self.inner.write().add_transaction(tx)?;
307 self.notify_ready(&added);
308 Ok(added)
309 }
310}
311
312#[cfg(feature = "base")]
313impl<T: Typed2718> Pool<T> {
314 pub fn clear_transaction_type(&self, tx_type: u8) -> Vec<Arc<PoolTransaction<T>>> {
316 let hashes = {
317 let pool = self.inner.read();
318 pool.pending_transactions
319 .transactions()
320 .chain(pool.ready_transactions())
321 .filter_map(|tx| {
322 (tx.pending_transaction.transaction.ty() == tx_type).then_some(tx.hash())
323 })
324 .collect()
325 };
326 self.remove_invalid(hashes)
327 }
328}
329
330#[derive(Debug)]
334struct PoolInner<T> {
335 ready_transactions: ReadyTransactions<T>,
336 pending_transactions: PendingTransactions<T>,
337}
338
339impl<T> Default for PoolInner<T> {
340 fn default() -> Self {
341 Self { ready_transactions: Default::default(), pending_transactions: Default::default() }
342 }
343}
344
345impl<T> PoolInner<T> {
348 fn ready_transactions(&self) -> TransactionsIterator<T> {
350 self.ready_transactions.get_transactions()
351 }
352
353 fn clear(&mut self) {
355 self.ready_transactions.clear();
356 self.pending_transactions.clear();
357 }
358
359 pub fn transactions_by_sender(
361 &self,
362 sender: Address,
363 ) -> impl Iterator<Item = Arc<PoolTransaction<T>>> + '_ {
364 let pending_txs = self
365 .pending_transactions
366 .transactions()
367 .filter(move |tx| tx.pending_transaction.sender().eq(&sender));
368
369 let ready_txs = self
370 .ready_transactions
371 .get_transactions()
372 .filter(move |tx| tx.pending_transaction.sender().eq(&sender));
373
374 pending_txs.chain(ready_txs)
375 }
376
377 fn contains(&self, tx_hash: &TxHash) -> bool {
379 self.pending_transactions.contains(tx_hash) || self.ready_transactions.contains(tx_hash)
380 }
381
382 fn remove_invalid(&mut self, tx_hashes: Vec<TxHash>) -> Vec<Arc<PoolTransaction<T>>> {
384 if tx_hashes.is_empty() {
386 return vec![];
387 }
388 trace!(target: "txpool", "Removing invalid transactions: {:?}", tx_hashes);
389
390 let mut removed = self.ready_transactions.remove_with_markers(tx_hashes.clone(), None);
391 removed.extend(self.pending_transactions.remove(tx_hashes));
392
393 trace!(target: "txpool", "Removed invalid transactions: {:?}", removed.iter().map(|tx| tx.hash()).collect::<Vec<_>>());
394
395 removed
396 }
397
398 fn remove_transactions_by_address(&mut self, sender: Address) -> Vec<Arc<PoolTransaction<T>>> {
400 let tx_hashes =
401 self.transactions_by_sender(sender).map(move |tx| tx.hash()).collect::<Vec<TxHash>>();
402
403 if tx_hashes.is_empty() {
404 return vec![];
405 }
406
407 trace!(target: "txpool", "Removing transactions: {:?}", tx_hashes);
408
409 let mut removed = self.ready_transactions.remove_with_markers(tx_hashes.clone(), None);
410 removed.extend(self.pending_transactions.remove(tx_hashes));
411
412 trace!(target: "txpool", "Removed transactions: {:?}", removed.iter().map(|tx| tx.hash()).collect::<Vec<_>>());
413
414 removed
415 }
416}
417
418impl<T: Clone> PoolInner<T> {
419 fn get_transaction(&self, hash: TxHash) -> Option<PendingTransaction<T>> {
423 if let Some(pending) = self.pending_transactions.get(&hash) {
424 return Some(pending.transaction.pending_transaction.clone());
425 }
426 Some(
427 self.ready_transactions.get(&hash)?.transaction.transaction.pending_transaction.clone(),
428 )
429 }
430}
431
432impl<T: Transaction> PoolInner<T> {
433 fn add_transaction(
434 &mut self,
435 tx: PoolTransaction<T>,
436 ) -> Result<AddedTransaction<T>, PoolError> {
437 if self.contains(&tx.hash()) {
438 debug!(target: "txpool", "[{:?}] Already imported", tx.hash());
439 return Err(PoolError::AlreadyImported(tx.hash()));
440 }
441
442 let tx = PendingPoolTransaction::new(tx, self.ready_transactions.provided_markers());
443 trace!(target: "txpool", "[{:?}] ready={}", tx.transaction.hash(), tx.is_ready());
444
445 if !tx.is_ready() {
447 let hash = tx.transaction.hash();
448 self.pending_transactions.add_transaction(tx)?;
449 return Ok(AddedTransaction::Pending { hash });
450 }
451 self.add_ready_transaction(tx)
452 }
453
454 fn add_ready_transaction(
456 &mut self,
457 tx: PendingPoolTransaction<T>,
458 ) -> Result<AddedTransaction<T>, PoolError> {
459 let hash = tx.transaction.hash();
460 trace!(target: "txpool", "adding ready transaction [{:?}]", hash);
461 let mut ready = ReadyTransaction::new(hash);
462
463 let mut tx_queue = VecDeque::from([tx]);
464 let mut is_new_tx = true;
466
467 while let Some(current_tx) = tx_queue.pop_front() {
469 tx_queue.extend(
471 self.pending_transactions.mark_and_unlock(¤t_tx.transaction.provides),
472 );
473
474 let current_hash = current_tx.transaction.hash();
475 match self.ready_transactions.add_transaction(current_tx) {
477 Ok(replaced_transactions) => {
478 if !is_new_tx {
479 ready.promoted.push(current_hash);
480 }
481 ready.removed.extend(replaced_transactions);
483 }
484 Err(err) => {
485 if is_new_tx {
487 debug!(target: "txpool", "[{:?}] Failed to add tx: {:?}", current_hash,
488 err);
489 return Err(err);
490 }
491 ready.discarded.push(current_hash);
492 }
493 }
494 is_new_tx = false;
495 }
496
497 if ready.removed.iter().any(|tx| *tx.hash() == hash) {
501 self.ready_transactions.clear_transactions(&ready.promoted);
502 return Err(PoolError::CyclicTransaction);
503 }
504
505 Ok(AddedTransaction::Ready(ready))
506 }
507
508 pub fn prune_markers(&mut self, markers: impl IntoIterator<Item = TxMarker>) -> PruneResult<T> {
513 let mut imports = vec![];
514 let mut pruned = vec![];
515
516 for marker in markers {
517 imports.extend(self.pending_transactions.mark_and_unlock(Some(&marker)));
519 pruned.extend(self.ready_transactions.prune_tags(marker.clone()));
521 }
522
523 let mut promoted = vec![];
524 let mut failed = vec![];
525 for tx in imports {
526 let hash = tx.transaction.hash();
527 match self.add_ready_transaction(tx) {
528 Ok(res) => promoted.push(res),
529 Err(e) => {
530 warn!(target: "txpool", "Failed to promote tx [{:?}] : {:?}", hash, e);
531 failed.push(hash)
532 }
533 }
534 }
535
536 PruneResult { pruned, failed, promoted }
537 }
538}
539
540pub struct PruneResult<T> {
542 pub promoted: Vec<AddedTransaction<T>>,
544 pub failed: Vec<TxHash>,
546 pub pruned: Vec<Arc<PoolTransaction<T>>>,
548}
549
550impl<T> fmt::Debug for PruneResult<T> {
551 fn fmt(&self, fmt: &mut fmt::Formatter<'_>) -> fmt::Result {
552 write!(fmt, "PruneResult {{ ")?;
553 write!(
554 fmt,
555 "promoted: {:?}, ",
556 self.promoted.iter().map(|tx| *tx.hash()).collect::<Vec<_>>()
557 )?;
558 write!(fmt, "failed: {:?}, ", self.failed)?;
559 write!(
560 fmt,
561 "pruned: {:?}, ",
562 self.pruned.iter().map(|tx| *tx.pending_transaction.hash()).collect::<Vec<_>>()
563 )?;
564 write!(fmt, "}}")?;
565 Ok(())
566 }
567}
568
569#[derive(Clone, Debug)]
570pub struct ReadyTransaction<T> {
571 hash: TxHash,
573 promoted: Vec<TxHash>,
575 discarded: Vec<TxHash>,
577 removed: Vec<Arc<PoolTransaction<T>>>,
579}
580
581impl<T> ReadyTransaction<T> {
582 pub fn new(hash: TxHash) -> Self {
583 Self {
584 hash,
585 promoted: Default::default(),
586 discarded: Default::default(),
587 removed: Default::default(),
588 }
589 }
590}
591
592#[derive(Clone, Debug)]
593pub enum AddedTransaction<T> {
594 Ready(ReadyTransaction<T>),
596 Pending {
598 hash: TxHash,
600 },
601}
602
603impl<T> AddedTransaction<T> {
604 pub const fn hash(&self) -> &TxHash {
605 match self {
606 Self::Ready(tx) => &tx.hash,
607 Self::Pending { hash } => hash,
608 }
609 }
610}