Skip to main content

anvil/eth/pool/
mod.rs

1//! # Transaction Pool implementation
2//!
3//! The transaction pool is responsible for managing a set of transactions that can be included in
4//! upcoming blocks.
5//!
6//! The main task of the pool is to prepare an ordered list of transactions that are ready to be
7//! included in a new block.
8//!
9//! Each imported block can affect the validity of transactions already in the pool.
10//! The miner expects the most up-to-date transactions when attempting to create a new block.
11//! After being included in a block, a transaction should be removed from the pool, this process is
12//! called _pruning_ and due to separation of concerns is triggered externally.
13//! The pool essentially performs following services:
14//!   * import transactions
15//!   * order transactions
16//!   * provide ordered set of transactions that are ready for inclusion
17//!   * prune transactions
18//!
19//! Each transaction in the pool contains markers that it _provides_ or _requires_. This property is
20//! used to determine whether it can be included in a block (transaction is ready) or whether it
21//! still _requires_ other transactions to be mined first (transaction is pending).
22//! A transaction is associated with the nonce of the account it's sent from. A unique identifying
23//! marker for a transaction is therefore the pair `(nonce + account)`. An incoming transaction with
24//! a `nonce > nonce on chain` will _require_ `(nonce -1, account)` first, before it is ready to be
25//! included in a block.
26//!
27//! This implementation is adapted from <https://github.com/paritytech/substrate/tree/master/client/transaction-pool>
28
29use 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
52/// Transaction pool that performs validation.
53pub struct Pool<T> {
54    /// processes all pending transactions
55    inner: RwLock<PoolInner<T>>,
56    /// listeners for new ready transactions
57    transaction_listener: Mutex<Vec<Sender<TxHash>>>,
58}
59
60/// An independent snapshot of a transaction pool.
61#[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
72// == impl Pool ==
73
74impl<T> Pool<T> {
75    /// Returns an independent snapshot of the pool.
76    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    /// Restores the pool to a previous snapshot.
87    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    /// Returns an iterator that yields all transactions that are currently ready
99    pub fn ready_transactions(&self) -> TransactionsIterator<T> {
100        self.inner.read().ready_transactions()
101    }
102
103    /// Returns all transactions that are not ready to be included in a block yet
104    pub fn pending_transactions(&self) -> Vec<Arc<PoolTransaction<T>>> {
105        self.inner.read().pending_transactions.transactions().collect()
106    }
107
108    /// Returns every ready and queued transaction.
109    #[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    /// Returns the number of tx that are ready and queued for further execution
116    pub fn txpool_status(&self) -> TxpoolStatus {
117        // Note: naming differs here compared to geth's `TxpoolStatus`
118        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    /// Adds a new transaction listener to the pool that gets notified about every new ready
124    /// transaction
125    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    /// Returns true if this pool already contains the transaction
133    pub fn contains(&self, tx_hash: &TxHash) -> bool {
134        self.inner.read().contains(tx_hash)
135    }
136
137    /// Returns true if this pool contains a transaction from `sender` with `nonce`.
138    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    /// Returns a transaction from `sender` that provides exactly `markers`.
149    #[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    /// Removes all transactions from the pool
159    pub fn clear(&self) {
160        let mut pool = self.inner.write();
161        pool.clear();
162    }
163
164    /// Remove the given transactions from the pool
165    pub fn remove_invalid(&self, tx_hashes: Vec<TxHash>) -> Vec<Arc<PoolTransaction<T>>> {
166        self.inner.write().remove_invalid(tx_hashes)
167    }
168
169    /// Remove transactions by sender
170    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    /// Removes a single transaction from the pool
175    ///
176    /// This is similar to `[Pool::remove_invalid()]` but for a single transaction.
177    ///
178    /// **Note**: this will also drop any transaction that depend on the `tx`
179    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    /// Notifies listeners if the transaction was added to the ready queue.
198    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    /// notifies all listeners about the transaction
208    fn notify_listener(&self, hash: TxHash) {
209        let mut listener = self.transaction_listener.lock();
210        // this is basically a retain but with mut reference
211        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    /// Returns the _pending_ transaction for that `hash` if it exists in the mempool
237    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    /// Invoked when a set of transactions ([Self::ready_transactions()]) was executed.
244    ///
245    /// This will remove the transactions from the pool.
246    ///
247    /// Returns `true` if ready transactions left behind by the block can be included by mining
248    /// again right away, e.g. because the block hit `max_transactions` or ran out of gas.
249    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        // Requiring txs to leave the pool keeps this retry from mining empty blocks for txs that
252        // can never be included. Not-yet-valid txs and their dependents are retried by the delayed
253        // re-notify.
254        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        // remove invalid transactions from the pool
258        self.remove_invalid(invalid.into_iter().map(|tx| tx.hash()).collect());
259
260        // Prune mined and stale markers; both are satisfied by the resulting state.
261        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        // Re-notify the miner about not-yet-valid transactions so they'll be retried.
268        // Delay by 1 second to let time advance before the next mining attempt.
269        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    /// Removes ready transactions for the given iterator of identifying markers.
285    ///
286    /// For each marker we can remove transactions in the pool that either provide the marker
287    /// directly or are a dependency of the transaction associated with that marker.
288    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    /// Adds a new transaction to the pool
302    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    /// Removes every transaction with the given EIP-2718 type.
315    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/// A Transaction Pool
331///
332/// Contains all transactions that are ready to be executed
333#[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
345// == impl PoolInner ==
346
347impl<T> PoolInner<T> {
348    /// Returns an iterator over transactions that are ready.
349    fn ready_transactions(&self) -> TransactionsIterator<T> {
350        self.ready_transactions.get_transactions()
351    }
352
353    /// Clears
354    fn clear(&mut self) {
355        self.ready_transactions.clear();
356        self.pending_transactions.clear();
357    }
358
359    /// Returns an iterator over all transactions in the pool filtered by the sender
360    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    /// Returns true if this pool already contains the transaction
378    fn contains(&self, tx_hash: &TxHash) -> bool {
379        self.pending_transactions.contains(tx_hash) || self.ready_transactions.contains(tx_hash)
380    }
381
382    /// Remove the given transactions from the pool
383    fn remove_invalid(&mut self, tx_hashes: Vec<TxHash>) -> Vec<Arc<PoolTransaction<T>>> {
384        // early exit in case there is no invalid transactions.
385        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    /// Remove transactions by sender address
399    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    /// checks both pools for the matching transaction
420    ///
421    /// Returns `None` if the transaction does not exist in the pool
422    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 all markers are not satisfied import to future
446        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    /// Adds the transaction to the ready queue
455    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        // tracks whether we're processing the given `tx`
465        let mut is_new_tx = true;
466
467        // take first transaction from the list
468        while let Some(current_tx) = tx_queue.pop_front() {
469            // also add the transaction that the current transaction unlocks
470            tx_queue.extend(
471                self.pending_transactions.mark_and_unlock(&current_tx.transaction.provides),
472            );
473
474            let current_hash = current_tx.transaction.hash();
475            // try to add the transaction to the ready pool
476            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                    // tx removed from ready pool
482                    ready.removed.extend(replaced_transactions);
483                }
484                Err(err) => {
485                    // failed to add transaction
486                    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        // check for a cycle where importing a transaction resulted in pending transactions to be
498        // added while removing current transaction. in which case we move this transaction back to
499        // the pending queue
500        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    /// Prunes the transactions that provide the given markers
509    ///
510    /// This will effectively remove those transactions that satisfy the markers and transactions
511    /// from the pending queue might get promoted to if the markers unlock them.
512    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            // mark as satisfied and store the transactions that got unlocked
518            imports.extend(self.pending_transactions.mark_and_unlock(Some(&marker)));
519            // prune transactions
520            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
540/// Represents the outcome of a prune
541pub struct PruneResult<T> {
542    /// a list of added transactions that a pruned marker satisfied
543    pub promoted: Vec<AddedTransaction<T>>,
544    /// all transactions that  failed to be promoted and now are discarded
545    pub failed: Vec<TxHash>,
546    /// all transactions that were pruned from the ready pool
547    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    /// the hash of the submitted transaction
572    hash: TxHash,
573    /// transactions promoted to the ready queue
574    promoted: Vec<TxHash>,
575    /// transaction that failed and became discarded
576    discarded: Vec<TxHash>,
577    /// Transactions removed from the Ready pool
578    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    /// transaction was successfully added and being processed
595    Ready(ReadyTransaction<T>),
596    /// Transaction was successfully added but not yet queued for processing
597    Pending {
598        /// the hash of the submitted transaction
599        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}