Skip to main content

reth_transaction_pool/pool/
mod.rs

1//! Transaction Pool internals.
2//!
3//! Incoming transactions are validated before they enter the pool first. The validation outcome can
4//! have 3 states:
5//!
6//!  1. Transaction can _never_ be valid
7//!  2. Transaction is _currently_ valid
8//!  3. Transaction is _currently_ invalid, but could potentially become valid in the future
9//!
10//! However, (2.) and (3.) of a transaction can only be determined on the basis of the current
11//! state, whereas (1.) holds indefinitely. This means once the state changes (2.) and (3.) the
12//! state of a transaction needs to be reevaluated again.
13//!
14//! The transaction pool is responsible for storing new, valid transactions and providing the next
15//! best transactions sorted by their priority. Where priority is determined by the transaction's
16//! score ([`TransactionOrdering`]).
17//!
18//! Furthermore, the following characteristics fall under (3.):
19//!
20//!  a) Nonce of a transaction is higher than the expected nonce for the next transaction of its
21//! sender. A distinction is made here whether multiple transactions from the same sender have
22//! gapless nonce increments.
23//!
24//!  a)(1) If _no_ transaction is missing in a chain of multiple
25//! transactions from the same sender (all nonce in row), all of them can in principle be executed
26//! on the current state one after the other.
27//!
28//!  a)(2) If there's a nonce gap, then all
29//! transactions after the missing transaction are blocked until the missing transaction arrives.
30//!
31//!  b) Transaction does not meet the dynamic fee cap requirement introduced by EIP-1559: The
32//! fee cap of the transaction needs to be no less than the base fee of block.
33//!
34//!
35//! In essence the transaction pool is made of three separate sub-pools:
36//!
37//!  - Pending Pool: Contains all transactions that are valid on the current state and satisfy (3.
38//!    a)(1): _No_ nonce gaps. A _pending_ transaction is considered _ready_ when it has the lowest
39//!    nonce of all transactions from the same sender. Once a _ready_ transaction with nonce `n` has
40//!    been executed, the next highest transaction from the same sender `n + 1` becomes ready.
41//!
42//!  - Queued Pool: Contains all transactions that are currently blocked by missing transactions:
43//!    (3. a)(2): _With_ nonce gaps or due to lack of funds.
44//!
45//!  - Basefee Pool: To account for the dynamic base fee requirement (3. b) which could render an
46//!    EIP-1559 and all subsequent transactions of the sender currently invalid.
47//!
48//! The classification of transactions is always dependent on the current state that is changed as
49//! soon as a new block is mined. Once a new block is mined, the account changeset must be applied
50//! to the transaction pool.
51//!
52//!
53//! Depending on the use case, consumers of the [`TransactionPool`](crate::traits::TransactionPool)
54//! are interested in (2.) and/or (3.).
55
56//! A generic [`TransactionPool`](crate::traits::TransactionPool) that only handles transactions.
57//!
58//! This Pool maintains two separate sub-pools for (2.) and (3.)
59//!
60//! ## Terminology
61//!
62//!  - _Pending_: pending transactions are transactions that fall under (2.). These transactions can
63//!    currently be executed and are stored in the pending sub-pool
64//!  - _Queued_: queued transactions are transactions that fall under category (3.). Those
65//!    transactions are _currently_ waiting for state changes that eventually move them into
66//!    category (2.) and become pending.
67
68use crate::{
69    blobstore::{BlobStore, PooledBlobSidecar},
70    error::{PoolError, PoolErrorKind, PoolResult},
71    identifier::{SenderId, SenderIdentifiers, TransactionId},
72    metrics::BlobStoreMetrics,
73    pool::{
74        listener::{
75            BlobTransactionSidecarListener, PendingTransactionHashListener, PoolEventBroadcast,
76            TransactionListener,
77        },
78        state::SubPool,
79        txpool::{SenderInfo, TxPool},
80        update::UpdateOutcome,
81    },
82    traits::{
83        AllPoolTransactions, BestTransactionsAttributes, BlockInfo, GetPooledTransactionLimit,
84        NewBlobSidecar, PoolSize, PoolTransaction, PropagatedTransactions, TransactionOrigin,
85    },
86    validate::{TransactionValidationOutcome, ValidPoolTransaction, ValidTransaction},
87    CanonicalStateUpdate, EthPoolTransaction, PoolConfig, TransactionOrdering,
88    TransactionValidator,
89};
90
91use alloy_primitives::{
92    map::{AddressSet, HashSet},
93    Address, TxHash, B256,
94};
95use parking_lot::{Mutex, RwLock, RwLockReadGuard, RwLockWriteGuard};
96use reth_eth_wire_types::HandleMempoolData;
97use reth_execution_types::ChangedAccount;
98
99use alloy_eips::{eip7594::BlobTransactionSidecarVariant, Typed2718};
100use reth_primitives_traits::Recovered;
101use rustc_hash::FxHashMap;
102use std::{
103    fmt,
104    sync::{
105        atomic::{AtomicBool, Ordering},
106        Arc,
107    },
108    time::Instant,
109};
110use tokio::sync::mpsc;
111use tracing::{debug, trace, warn};
112mod events;
113pub use best::{BestTransactionFilter, BestTransactionsWithPrioritizedSenders};
114pub use blob::{blob_tx_priority, fee_delta, BlobOrd, BlobTransactions};
115pub use events::{FullTransactionEvent, NewTransactionEvent, TransactionEvent};
116pub use listener::{AllTransactionsEvents, TransactionEvents, TransactionListenerKind};
117pub use parked::{BasefeeOrd, ParkedOrd, ParkedPool, QueuedOrd};
118pub use pending::PendingPool;
119
120mod best;
121pub use best::BestTransactions;
122
123mod blob;
124pub mod listener;
125mod parked;
126pub mod pending;
127pub mod size;
128pub(crate) mod state;
129pub mod txpool;
130mod update;
131
132/// Bound on number of pending transactions from `reth_network::TransactionsManager` to buffer.
133pub const PENDING_TX_LISTENER_BUFFER_SIZE: usize = 2048;
134/// Bound on number of new transactions from `reth_network::TransactionsManager` to buffer.
135pub const NEW_TX_LISTENER_BUFFER_SIZE: usize = 1024;
136
137const BLOB_SIDECAR_LISTENER_BUFFER_SIZE: usize = 512;
138
139/// Transaction pool internals.
140pub struct PoolInner<V, T, S>
141where
142    T: TransactionOrdering,
143{
144    /// Internal mapping of addresses to plain ints.
145    identifiers: RwLock<SenderIdentifiers>,
146    /// Transaction validator.
147    validator: V,
148    /// Storage for blob transactions
149    blob_store: S,
150    /// The internal pool that manages all transactions.
151    pool: RwLock<TxPool<T>>,
152    /// Pool settings.
153    config: PoolConfig,
154    /// Manages listeners for transaction state change events.
155    event_listener: RwLock<PoolEventBroadcast<T::Transaction>>,
156    /// Tracks whether any event listeners have ever been installed.
157    has_event_listeners: AtomicBool,
158    /// Listeners for new _full_ pending transactions.
159    pending_transaction_listener: RwLock<Vec<PendingTransactionHashListener>>,
160    /// Listeners for new transactions added to the pool.
161    transaction_listener: RwLock<Vec<TransactionListener<T::Transaction>>>,
162    /// Listener for new blob transaction sidecars added to the pool.
163    blob_transaction_sidecar_listener: Mutex<Vec<BlobTransactionSidecarListener>>,
164    /// Metrics for the blob store
165    blob_store_metrics: BlobStoreMetrics,
166}
167
168// === impl PoolInner ===
169
170impl<V, T, S> PoolInner<V, T, S>
171where
172    V: TransactionValidator,
173    T: TransactionOrdering<Transaction = <V as TransactionValidator>::Transaction>,
174    S: BlobStore,
175{
176    /// Create a new transaction pool instance.
177    pub fn new(validator: V, ordering: T, blob_store: S, config: PoolConfig) -> Self {
178        Self {
179            identifiers: Default::default(),
180            validator,
181            event_listener: Default::default(),
182            has_event_listeners: AtomicBool::new(false),
183            pool: RwLock::new(TxPool::new(ordering, config.clone())),
184            pending_transaction_listener: Default::default(),
185            transaction_listener: Default::default(),
186            blob_transaction_sidecar_listener: Default::default(),
187            config,
188            blob_store,
189            blob_store_metrics: Default::default(),
190        }
191    }
192
193    /// Returns the configured blob store.
194    pub const fn blob_store(&self) -> &S {
195        &self.blob_store
196    }
197
198    /// Returns stats about the size of the pool.
199    pub fn size(&self) -> PoolSize {
200        self.get_pool_data().size()
201    }
202
203    /// Returns the currently tracked block
204    pub fn block_info(&self) -> BlockInfo {
205        self.get_pool_data().block_info()
206    }
207    /// Sets the currently tracked block.
208    ///
209    /// This will also notify subscribers about any transactions that were promoted to the pending
210    /// pool due to fee changes.
211    pub fn set_block_info(&self, info: BlockInfo) {
212        let outcome = self.pool.write().set_block_info(info);
213
214        // Notify subscribers about promoted transactions due to fee changes
215        self.notify_on_transaction_updates(outcome.promoted, outcome.discarded);
216    }
217
218    /// Returns the internal [`SenderId`] for this address, allocating a new mapping when the
219    /// address is first observed.
220    ///
221    /// This must only be used on paths that intentionally begin tracking a sender, such as
222    /// transaction insertion. Read-only lookups should prefer [`Self::sender_id`] to avoid
223    /// growing the sender-id map for unknown addresses.
224    pub fn get_sender_id(&self, addr: Address) -> SenderId {
225        self.identifiers.write().sender_id_or_create(addr)
226    }
227
228    /// Returns the internal [`SenderId`] for this address if it is already tracked.
229    ///
230    /// Unlike [`Self::get_sender_id`], this never allocates a new sender mapping and is therefore
231    /// suitable for read-only queries or best-effort cleanup on unknown addresses.
232    pub fn sender_id(&self, addr: &Address) -> Option<SenderId> {
233        self.identifiers.read().sender_id(addr)
234    }
235
236    /// Returns the internal [`SenderId`]s for the given addresses.
237    pub fn get_sender_ids(&self, addrs: impl IntoIterator<Item = Address>) -> Vec<SenderId> {
238        self.identifiers.write().sender_ids_or_create(addrs)
239    }
240
241    /// Returns all senders in the pool
242    pub fn unique_senders(&self) -> AddressSet {
243        self.get_pool_data().unique_senders()
244    }
245
246    /// Converts the changed accounts to a map of sender ids to sender info (internal identifier
247    /// used for __tracked__ accounts)
248    fn changed_senders(
249        &self,
250        accs: impl Iterator<Item = ChangedAccount>,
251    ) -> FxHashMap<SenderId, SenderInfo> {
252        let identifiers = self.identifiers.read();
253        accs.into_iter()
254            .filter_map(|acc| {
255                let ChangedAccount { address, nonce, balance } = acc;
256                let sender_id = identifiers.sender_id(&address)?;
257                Some((sender_id, SenderInfo { state_nonce: nonce, balance }))
258            })
259            .collect()
260    }
261
262    /// Get the config the pool was configured with.
263    pub const fn config(&self) -> &PoolConfig {
264        &self.config
265    }
266
267    /// Get the validator reference.
268    pub const fn validator(&self) -> &V {
269        &self.validator
270    }
271
272    /// Adds a new transaction listener to the pool that gets notified about every new _pending_
273    /// transaction inserted into the pool
274    pub fn add_pending_listener(&self, kind: TransactionListenerKind) -> mpsc::Receiver<TxHash> {
275        let (sender, rx) = mpsc::channel(self.config.pending_tx_listener_buffer_size);
276        let listener = PendingTransactionHashListener { sender, kind };
277
278        let mut listeners = self.pending_transaction_listener.write();
279        // Clean up dead listeners before adding new one
280        listeners.retain(|l| !l.sender.is_closed());
281        listeners.push(listener);
282
283        rx
284    }
285
286    /// Adds a new transaction listener to the pool that gets notified about every new transaction.
287    pub fn add_new_transaction_listener(
288        &self,
289        kind: TransactionListenerKind,
290    ) -> mpsc::Receiver<NewTransactionEvent<T::Transaction>> {
291        let (sender, rx) = mpsc::channel(self.config.new_tx_listener_buffer_size);
292        let listener = TransactionListener { sender, kind };
293
294        let mut listeners = self.transaction_listener.write();
295        // Clean up dead listeners before adding new one
296        listeners.retain(|l| !l.sender.is_closed());
297        listeners.push(listener);
298
299        rx
300    }
301    /// Adds a new blob sidecar listener to the pool that gets notified about every new
302    /// eip4844 transaction's blob sidecar.
303    pub fn add_blob_sidecar_listener(&self) -> mpsc::Receiver<NewBlobSidecar> {
304        let (sender, rx) = mpsc::channel(BLOB_SIDECAR_LISTENER_BUFFER_SIZE);
305        let listener = BlobTransactionSidecarListener { sender };
306        self.blob_transaction_sidecar_listener.lock().push(listener);
307        rx
308    }
309
310    /// If the pool contains the transaction, this adds a new listener that gets notified about
311    /// transaction events.
312    pub fn add_transaction_event_listener(&self, tx_hash: TxHash) -> Option<TransactionEvents> {
313        if !self.get_pool_data().contains(&tx_hash) {
314            return None
315        }
316        let mut listener = self.event_listener.write();
317        let events = listener.subscribe(tx_hash);
318        self.mark_event_listener_installed();
319        Some(events)
320    }
321
322    /// Adds a listener for all transaction events.
323    pub fn add_all_transactions_event_listener(&self) -> AllTransactionsEvents<T::Transaction> {
324        let mut listener = self.event_listener.write();
325        let events = listener.subscribe_all();
326        self.mark_event_listener_installed();
327        events
328    }
329
330    #[inline]
331    fn has_event_listeners(&self) -> bool {
332        self.has_event_listeners.load(Ordering::Relaxed)
333    }
334
335    #[inline]
336    fn mark_event_listener_installed(&self) {
337        self.has_event_listeners.store(true, Ordering::Relaxed);
338    }
339
340    #[inline]
341    fn update_event_listener_state(&self, listener: &PoolEventBroadcast<T::Transaction>) {
342        if listener.is_empty() {
343            self.has_event_listeners.store(false, Ordering::Relaxed);
344        }
345    }
346
347    #[inline]
348    fn with_event_listener<F>(&self, emit: F)
349    where
350        F: FnOnce(&mut PoolEventBroadcast<T::Transaction>),
351    {
352        if !self.has_event_listeners() {
353            return
354        }
355        let mut listener = self.event_listener.write();
356        if !listener.is_empty() {
357            emit(&mut listener);
358        }
359        self.update_event_listener_state(&listener);
360    }
361
362    /// Returns a read lock to the pool's data.
363    pub fn get_pool_data(&self) -> RwLockReadGuard<'_, TxPool<T>> {
364        self.pool.read()
365    }
366
367    /// Returns transactions in the pool that can be propagated
368    pub fn pooled_transactions(&self) -> Vec<Arc<ValidPoolTransaction<T::Transaction>>> {
369        let mut out = Vec::new();
370        self.append_pooled_transactions(&mut out);
371        out
372    }
373
374    /// Returns hashes of transactions in the pool that can be propagated.
375    pub fn pooled_transactions_hashes(&self) -> Vec<TxHash> {
376        let mut out = Vec::new();
377        self.append_pooled_transactions_hashes(&mut out);
378        out
379    }
380
381    /// Returns only the first `max` transactions in the pool that can be propagated.
382    pub fn pooled_transactions_max(
383        &self,
384        max: usize,
385    ) -> Vec<Arc<ValidPoolTransaction<T::Transaction>>> {
386        if max == 0 {
387            return Vec::new()
388        }
389
390        let pool = self.get_pool_data();
391        let mut out = Vec::with_capacity(max.min(pool.all().len()));
392        out.extend(pool.all().transactions_iter().filter(|tx| tx.propagate).take(max).cloned());
393        out
394    }
395
396    /// Extends the given vector with all transactions in the pool that can be propagated.
397    pub fn append_pooled_transactions(
398        &self,
399        out: &mut Vec<Arc<ValidPoolTransaction<T::Transaction>>>,
400    ) {
401        out.extend(
402            self.get_pool_data().all().transactions_iter().filter(|tx| tx.propagate).cloned(),
403        );
404    }
405
406    /// Extends the given vector with pooled transactions for the given hashes that are allowed to
407    /// be propagated.
408    pub fn append_pooled_transaction_elements(
409        &self,
410        tx_hashes: &[TxHash],
411        limit: GetPooledTransactionLimit,
412        out: &mut Vec<<<V as TransactionValidator>::Transaction as PoolTransaction>::Pooled>,
413    ) where
414        <V as TransactionValidator>::Transaction: EthPoolTransaction,
415    {
416        let transactions = self.get_all_propagatable(tx_hashes);
417        let mut size = 0;
418        for transaction in transactions {
419            let encoded_len = transaction.encoded_length();
420            let Some(pooled) = self.to_pooled_transaction(transaction) else {
421                continue;
422            };
423
424            size += encoded_len;
425            out.push(pooled.into_inner());
426
427            if limit.exceeds(size) {
428                break
429            }
430        }
431    }
432
433    /// Extends the given vector with the hashes of all transactions in the pool that can be
434    /// propagated.
435    pub fn append_pooled_transactions_hashes(&self, out: &mut Vec<TxHash>) {
436        out.extend(
437            self.get_pool_data()
438                .all()
439                .transactions_iter()
440                .filter(|tx| tx.propagate)
441                .map(|tx| *tx.hash()),
442        );
443    }
444
445    /// Extends the given vector with only the first `max` transactions in the pool that can be
446    /// propagated.
447    pub fn append_pooled_transactions_max(
448        &self,
449        max: usize,
450        out: &mut Vec<Arc<ValidPoolTransaction<T::Transaction>>>,
451    ) {
452        out.extend(
453            self.get_pool_data()
454                .all()
455                .transactions_iter()
456                .filter(|tx| tx.propagate)
457                .take(max)
458                .cloned(),
459        );
460    }
461
462    /// Returns only the first `max` hashes of transactions in the pool that can be propagated.
463    pub fn pooled_transactions_hashes_max(&self, max: usize) -> Vec<TxHash> {
464        if max == 0 {
465            return Vec::new();
466        }
467
468        let pool = self.get_pool_data();
469        let mut out = Vec::with_capacity(max.min(pool.all().len()));
470        out.extend(
471            pool.all().transactions_iter().filter(|tx| tx.propagate).take(max).map(|tx| *tx.hash()),
472        );
473        out
474    }
475
476    /// Converts the internally tracked transaction to the pooled format.
477    ///
478    /// If the transaction is an EIP-4844 transaction, the blob sidecar is fetched from the blob
479    /// store and attached to the transaction.
480    fn to_pooled_transaction(
481        &self,
482        transaction: Arc<ValidPoolTransaction<T::Transaction>>,
483    ) -> Option<Recovered<<<V as TransactionValidator>::Transaction as PoolTransaction>::Pooled>>
484    where
485        <V as TransactionValidator>::Transaction: EthPoolTransaction,
486    {
487        if transaction.is_eip4844() {
488            let sidecar = self.blob_store.get(*transaction.hash()).ok()??;
489            transaction.transaction.clone().try_into_pooled_eip4844(sidecar)
490        } else {
491            transaction
492                .transaction
493                .clone_into_pooled()
494                .inspect_err(|err| {
495                    debug!(
496                        target: "txpool", %err,
497                        "failed to convert transaction to pooled element; skipping",
498                    );
499                })
500                .ok()
501        }
502    }
503
504    /// Returns pooled transactions for the given transaction hashes that are allowed to be
505    /// propagated.
506    pub fn get_pooled_transaction_elements(
507        &self,
508        tx_hashes: Vec<TxHash>,
509        limit: GetPooledTransactionLimit,
510    ) -> Vec<<<V as TransactionValidator>::Transaction as PoolTransaction>::Pooled>
511    where
512        <V as TransactionValidator>::Transaction: EthPoolTransaction,
513    {
514        let mut elements = Vec::new();
515        self.append_pooled_transaction_elements(&tx_hashes, limit, &mut elements);
516        elements.shrink_to_fit();
517        elements
518    }
519
520    /// Returns converted pooled transaction for the given transaction hash.
521    pub fn get_pooled_transaction_element(
522        &self,
523        tx_hash: TxHash,
524    ) -> Option<Recovered<<<V as TransactionValidator>::Transaction as PoolTransaction>::Pooled>>
525    where
526        <V as TransactionValidator>::Transaction: EthPoolTransaction,
527    {
528        self.get(&tx_hash).and_then(|tx| self.to_pooled_transaction(tx))
529    }
530
531    /// Updates the entire pool after a new block was executed.
532    pub fn on_canonical_state_change(&self, update: CanonicalStateUpdate<'_, V::Block>) {
533        trace!(target: "txpool", ?update, "updating pool on canonical state change");
534
535        let block_info = update.block_info();
536        let CanonicalStateUpdate {
537            new_tip, changed_accounts, mined_transactions, update_kind, ..
538        } = update;
539        self.validator.on_new_head_block(new_tip);
540
541        let changed_senders = self.changed_senders(changed_accounts.into_iter());
542
543        // update the pool
544        let outcome = self.pool.write().on_canonical_state_change(
545            block_info,
546            mined_transactions,
547            changed_senders,
548            update_kind,
549        );
550
551        // This will discard outdated transactions based on the account's nonce
552        self.delete_discarded_blobs(outcome.discarded.iter());
553
554        // notify listeners about updates
555        self.notify_on_new_state(outcome);
556    }
557
558    /// Performs account updates on the pool.
559    ///
560    /// This will either promote or discard transactions based on the new account state.
561    ///
562    /// This should be invoked when the pool drifted and accounts are updated manually
563    pub fn update_accounts(&self, accounts: Vec<ChangedAccount>) {
564        let changed_senders = self.changed_senders(accounts.into_iter());
565        let UpdateOutcome { promoted, discarded } =
566            self.pool.write().update_accounts(changed_senders);
567
568        self.notify_on_transaction_updates(promoted, discarded);
569    }
570
571    /// Add a single validated transaction into the pool.
572    ///
573    /// Returns the outcome and optionally metadata to be processed after the pool lock is
574    /// released.
575    ///
576    /// Note: this is only used internally by [`Self::add_transactions()`], all new transaction(s)
577    /// come in through that function, either as a batch or `std::iter::once`.
578    fn add_transaction(
579        &self,
580        pool: &mut RwLockWriteGuard<'_, TxPool<T>>,
581        origin: TransactionOrigin,
582        tx: TransactionValidationOutcome<T::Transaction>,
583    ) -> (PoolResult<AddedTransactionOutcome>, Option<AddedTransactionMeta<T::Transaction>>) {
584        match tx {
585            TransactionValidationOutcome::Valid {
586                balance,
587                state_nonce,
588                transaction,
589                propagate,
590                bytecode_hash,
591                authorities,
592            } => {
593                let sender_id = self.get_sender_id(transaction.sender());
594                let transaction_id = TransactionId::new(sender_id, transaction.nonce());
595
596                // split the valid transaction and the blob sidecar if it has any
597                let (transaction, blob_sidecar) = match transaction {
598                    ValidTransaction::Valid(tx) => (tx, None),
599                    ValidTransaction::ValidWithSidecar { transaction, sidecar } => {
600                        debug_assert!(
601                            transaction.is_eip4844(),
602                            "validator returned sidecar for non EIP-4844 transaction"
603                        );
604                        (transaction, Some(sidecar))
605                    }
606                };
607
608                let tx = ValidPoolTransaction {
609                    transaction,
610                    transaction_id,
611                    propagate,
612                    timestamp: Instant::now(),
613                    origin,
614                    authority_ids: authorities.map(|auths| self.get_sender_ids(auths)),
615                };
616
617                let added = match pool.add_transaction(tx, balance, state_nonce, bytecode_hash) {
618                    Ok(added) => added,
619                    Err(err) => return (Err(err), None),
620                };
621                let hash = *added.hash();
622                let state = added.transaction_state();
623
624                let meta = AddedTransactionMeta { added, blob_sidecar };
625
626                (Ok(AddedTransactionOutcome { hash, state }), Some(meta))
627            }
628            TransactionValidationOutcome::Invalid(tx, err) => {
629                self.with_event_listener(|listener| listener.invalid(tx.hash()));
630                (Err(PoolError::new(*tx.hash(), err)), None)
631            }
632            TransactionValidationOutcome::Error(tx_hash, err) => {
633                self.with_event_listener(|listener| listener.discarded(&tx_hash));
634                (Err(PoolError::other(tx_hash, err)), None)
635            }
636        }
637    }
638
639    /// Adds a transaction and returns the event stream.
640    pub fn add_transaction_and_subscribe(
641        &self,
642        origin: TransactionOrigin,
643        tx: TransactionValidationOutcome<T::Transaction>,
644    ) -> PoolResult<TransactionEvents> {
645        let listener = {
646            let mut listener = self.event_listener.write();
647            let events = listener.subscribe(tx.tx_hash());
648            self.mark_event_listener_installed();
649            events
650        };
651        let mut results = self.add_transactions(origin, std::iter::once(tx));
652        results.pop().expect("result length is the same as the input")?;
653        Ok(listener)
654    }
655
656    /// Adds all transactions in the iterator to the pool, returning a list of results.
657    ///
658    /// Convenience method that assigns the same origin to all transactions. Delegates to
659    /// [`Self::add_transactions_with_origins`].
660    pub fn add_transactions(
661        &self,
662        origin: TransactionOrigin,
663        transactions: impl IntoIterator<Item = TransactionValidationOutcome<T::Transaction>>,
664    ) -> Vec<PoolResult<AddedTransactionOutcome>> {
665        self.add_transactions_with_origins(transactions.into_iter().map(|tx| (origin, tx)))
666    }
667
668    /// Adds all transactions in the iterator to the pool, each with its own
669    /// [`TransactionOrigin`], returning a list of results.
670    pub fn add_transactions_with_origins(
671        &self,
672        transactions: impl IntoIterator<
673            Item = (TransactionOrigin, TransactionValidationOutcome<T::Transaction>),
674        >,
675    ) -> Vec<PoolResult<AddedTransactionOutcome>> {
676        // Collect results and metadata while holding the pool write lock
677        let (mut results, added_metas, discarded) = {
678            let mut pool = self.pool.write();
679            let mut added_metas = Vec::new();
680
681            let results = transactions
682                .into_iter()
683                .map(|(origin, tx)| {
684                    let (result, meta) = self.add_transaction(&mut pool, origin, tx);
685
686                    // Only collect metadata for successful insertions
687                    if result.is_ok() &&
688                        let Some(meta) = meta
689                    {
690                        added_metas.push(meta);
691                    }
692
693                    result
694                })
695                .collect::<Vec<_>>();
696
697            // Enforce the pool size limits if at least one transaction was added successfully
698            let discarded = if results.iter().any(Result::is_ok) {
699                let discarded = pool.discard_worst();
700                pool.update_size_metrics();
701                discarded
702            } else {
703                Default::default()
704            };
705
706            (results, added_metas, discarded)
707        };
708
709        for meta in added_metas {
710            self.on_added_transaction(meta);
711        }
712
713        if !discarded.is_empty() {
714            // Delete any blobs associated with discarded blob transactions
715            self.delete_discarded_blobs(discarded.iter());
716            self.with_event_listener(|listener| listener.discarded_many(&discarded));
717
718            // Linear search avoids allocating a hash set for small eviction batches.
719            const MAX_LINEAR_SEARCH_DISCARDS: usize = 4;
720            let discarded_hashes = (discarded.len() > MAX_LINEAR_SEARCH_DISCARDS)
721                .then(|| discarded.iter().map(|tx| *tx.hash()).collect::<HashSet<_>>());
722            let is_discarded = |hash: &TxHash| match &discarded_hashes {
723                Some(hashes) => hashes.contains(hash),
724                None => discarded.iter().any(|tx| tx.hash() == hash),
725            };
726
727            // A newly added transaction may be immediately discarded, so we need to
728            // adjust the result here
729            for res in &mut results {
730                if let Ok(AddedTransactionOutcome { hash, .. }) = res &&
731                    is_discarded(hash)
732                {
733                    *res = Err(PoolError::new(*hash, PoolErrorKind::DiscardedOnInsert))
734                }
735            }
736        };
737
738        results
739    }
740
741    /// Process a transaction that was added to the pool.
742    ///
743    /// Performs blob storage operations and sends all notifications. This should be called
744    /// after the pool write lock has been released to avoid blocking pool operations.
745    fn on_added_transaction(&self, meta: AddedTransactionMeta<T::Transaction>) {
746        // Handle blob sidecar storage and notifications for EIP-4844 transactions
747        if let Some(sidecar) = meta.blob_sidecar {
748            let hash = *meta.added.hash();
749            self.on_new_blob_sidecar(&hash, &sidecar);
750            self.insert_blob(hash, sidecar);
751        }
752
753        // Delete replaced blob sidecar if any
754        if let Some(replaced) = meta.added.replaced_blob_transaction() {
755            debug!(target: "txpool", "[{:?}] delete replaced blob sidecar", replaced);
756            self.delete_blob(replaced);
757        }
758
759        // Delete discarded blob sidecars if any, this doesnt do any IO.
760        if let Some(discarded) = meta.added.discarded_transactions() {
761            self.delete_discarded_blobs(discarded.iter());
762        }
763
764        // Notify pending transaction listeners
765        if let Some(pending) = meta.added.as_pending() {
766            self.on_new_pending_transaction(pending);
767        }
768
769        // Notify event listeners
770        self.notify_event_listeners(&meta.added);
771
772        // Notify new transaction listeners
773        self.on_new_transaction(meta.added.into_new_transaction_event());
774    }
775
776    /// Notify all listeners about a new pending transaction.
777    ///
778    /// See also [`Self::add_pending_listener`]
779    ///
780    /// CAUTION: This function is only intended to be used manually in order to use this type's
781    /// pending transaction receivers when manually implementing the
782    /// [`TransactionPool`](crate::TransactionPool) trait for a custom pool implementation
783    /// [`TransactionPool::pending_transactions_listener_for`](crate::TransactionPool).
784    pub fn on_new_pending_transaction(&self, pending: &AddedPendingTransaction<T::Transaction>) {
785        let mut needs_cleanup = false;
786
787        {
788            let listeners = self.pending_transaction_listener.read();
789            for listener in listeners.iter() {
790                if !listener.send_all(pending.pending_transactions(listener.kind)) {
791                    needs_cleanup = true;
792                }
793            }
794        }
795
796        // Clean up dead listeners if we detected any closed channels
797        if needs_cleanup {
798            self.pending_transaction_listener
799                .write()
800                .retain(|listener| !listener.sender.is_closed());
801        }
802    }
803
804    /// Notify all listeners about a newly inserted pending transaction.
805    ///
806    /// See also [`Self::add_new_transaction_listener`]
807    ///
808    /// CAUTION: This function is only intended to be used manually in order to use this type's
809    /// transaction receivers when manually implementing the
810    /// [`TransactionPool`](crate::TransactionPool) trait for a custom pool implementation
811    /// [`TransactionPool::new_transactions_listener_for`](crate::TransactionPool).
812    pub fn on_new_transaction(&self, event: NewTransactionEvent<T::Transaction>) {
813        let mut needs_cleanup = false;
814
815        {
816            let listeners = self.transaction_listener.read();
817            for listener in listeners.iter() {
818                if listener.kind.is_propagate_only() && !event.transaction.propagate {
819                    if listener.sender.is_closed() {
820                        needs_cleanup = true;
821                    }
822                    // Skip non-propagate transactions for propagate-only listeners
823                    continue
824                }
825
826                if !listener.send(event.clone()) {
827                    needs_cleanup = true;
828                }
829            }
830        }
831
832        // Clean up dead listeners if we detected any closed channels
833        if needs_cleanup {
834            self.transaction_listener.write().retain(|listener| !listener.sender.is_closed());
835        }
836    }
837
838    /// Notify all listeners about a blob sidecar for a newly inserted blob (eip4844) transaction.
839    fn on_new_blob_sidecar(&self, tx_hash: &TxHash, sidecar: &BlobTransactionSidecarVariant) {
840        let mut sidecar_listeners = self.blob_transaction_sidecar_listener.lock();
841        if sidecar_listeners.is_empty() {
842            return
843        }
844        let sidecar = Arc::new(sidecar.clone());
845        sidecar_listeners.retain_mut(|listener| {
846            let new_blob_event = NewBlobSidecar { tx_hash: *tx_hash, sidecar: sidecar.clone() };
847            match listener.sender.try_send(new_blob_event) {
848                Ok(()) => true,
849                Err(err) => {
850                    if matches!(err, mpsc::error::TrySendError::Full(_)) {
851                        debug!(
852                            target: "txpool",
853                            "[{:?}] failed to send blob sidecar; channel full",
854                            sidecar,
855                        );
856                        true
857                    } else {
858                        false
859                    }
860                }
861            }
862        })
863    }
864
865    /// Notifies transaction listeners about changes once a block was processed.
866    fn notify_on_new_state(&self, outcome: OnNewCanonicalStateOutcome<T::Transaction>) {
867        trace!(target: "txpool", promoted=outcome.promoted.len(), discarded= outcome.discarded.len() ,"notifying listeners on state change");
868
869        // notify about promoted pending transactions - emit hashes
870        let mut needs_pending_cleanup = false;
871        {
872            let listeners = self.pending_transaction_listener.read();
873            for listener in listeners.iter() {
874                if !listener.send_all(outcome.pending_transactions(listener.kind)) {
875                    needs_pending_cleanup = true;
876                }
877            }
878        }
879        if needs_pending_cleanup {
880            self.pending_transaction_listener.write().retain(|l| !l.sender.is_closed());
881        }
882
883        // emit full transactions
884        let mut needs_tx_cleanup = false;
885        {
886            let listeners = self.transaction_listener.read();
887            for listener in listeners.iter() {
888                if !listener.send_all(outcome.full_pending_transactions(listener.kind)) {
889                    needs_tx_cleanup = true;
890                }
891            }
892        }
893        if needs_tx_cleanup {
894            self.transaction_listener.write().retain(|l| !l.sender.is_closed());
895        }
896
897        let OnNewCanonicalStateOutcome { mined, promoted, discarded, block_hash } = outcome;
898
899        // broadcast specific transaction events
900        self.with_event_listener(|listener| {
901            for tx in &mined {
902                listener.mined(tx, block_hash);
903            }
904            for tx in &promoted {
905                listener.pending(tx.hash(), None);
906            }
907            for tx in &discarded {
908                listener.discarded(tx.hash());
909            }
910        })
911    }
912
913    /// Notifies all listeners about the transaction movements.
914    ///
915    /// This will emit events according to the provided changes.
916    ///
917    /// CAUTION: This function is only intended to be used manually in order to use this type's
918    /// [`TransactionEvents`] receivers when manually implementing the
919    /// [`TransactionPool`](crate::TransactionPool) trait for a custom pool implementation
920    /// [`TransactionPool::transaction_event_listener`](crate::TransactionPool).
921    pub fn notify_on_transaction_updates(
922        &self,
923        promoted: Vec<Arc<ValidPoolTransaction<T::Transaction>>>,
924        discarded: Vec<Arc<ValidPoolTransaction<T::Transaction>>>,
925    ) {
926        // Notify about promoted pending transactions (similar to notify_on_new_state)
927        if !promoted.is_empty() {
928            let mut needs_pending_cleanup = false;
929            {
930                let listeners = self.pending_transaction_listener.read();
931                for listener in listeners.iter() {
932                    let promoted_hashes = promoted.iter().filter_map(|tx| {
933                        if listener.kind.is_propagate_only() && !tx.propagate {
934                            None
935                        } else {
936                            Some(*tx.hash())
937                        }
938                    });
939                    if !listener.send_all(promoted_hashes) {
940                        needs_pending_cleanup = true;
941                    }
942                }
943            }
944            if needs_pending_cleanup {
945                self.pending_transaction_listener.write().retain(|l| !l.sender.is_closed());
946            }
947
948            // in this case we should also emit promoted transactions in full
949            let mut needs_tx_cleanup = false;
950            {
951                let listeners = self.transaction_listener.read();
952                for listener in listeners.iter() {
953                    let promoted_txs = promoted.iter().filter_map(|tx| {
954                        if listener.kind.is_propagate_only() && !tx.propagate {
955                            None
956                        } else {
957                            Some(NewTransactionEvent::pending(tx.clone()))
958                        }
959                    });
960                    if !listener.send_all(promoted_txs) {
961                        needs_tx_cleanup = true;
962                    }
963                }
964            }
965            if needs_tx_cleanup {
966                self.transaction_listener.write().retain(|l| !l.sender.is_closed());
967            }
968        }
969
970        self.with_event_listener(|listener| {
971            for tx in &promoted {
972                listener.pending(tx.hash(), None);
973            }
974            for tx in &discarded {
975                listener.discarded(tx.hash());
976            }
977        });
978
979        if !discarded.is_empty() {
980            // This deletes outdated blob txs from the blob store, based on the account's nonce.
981            // This is called during txpool maintenance when the pool drifted.
982            self.delete_discarded_blobs(discarded.iter());
983        }
984    }
985
986    /// Fire events for the newly added transaction if there are any.
987    ///
988    /// See also [`Self::add_transaction_event_listener`].
989    ///
990    /// CAUTION: This function is only intended to be used manually in order to use this type's
991    /// [`TransactionEvents`] receivers when manually implementing the
992    /// [`TransactionPool`](crate::TransactionPool) trait for a custom pool implementation
993    /// [`TransactionPool::transaction_event_listener`](crate::TransactionPool).
994    pub fn notify_event_listeners(&self, tx: &AddedTransaction<T::Transaction>) {
995        self.with_event_listener(|listener| match tx {
996            AddedTransaction::Pending(tx) => {
997                let AddedPendingTransaction { transaction, promoted, discarded, replaced } = tx;
998
999                listener.pending(transaction.hash(), replaced.clone());
1000                for tx in promoted {
1001                    listener.pending(tx.hash(), None);
1002                }
1003                for tx in discarded {
1004                    listener.discarded(tx.hash());
1005                }
1006            }
1007            AddedTransaction::Parked { transaction, replaced, queued_reason, .. } => {
1008                listener.queued(transaction.hash(), queued_reason.clone());
1009                if let Some(replaced) = replaced {
1010                    listener.replaced(replaced.clone(), *transaction.hash());
1011                }
1012            }
1013        });
1014    }
1015
1016    /// Returns an iterator that yields transactions that are ready to be included in the block.
1017    pub fn best_transactions(&self) -> BestTransactions<T> {
1018        self.get_pool_data().best_transactions()
1019    }
1020
1021    /// Returns an iterator that yields transactions that are ready to be included in the block with
1022    /// the given base fee and optional blob fee attributes.
1023    pub fn best_transactions_with_attributes(
1024        &self,
1025        best_transactions_attributes: BestTransactionsAttributes,
1026    ) -> Box<dyn crate::traits::BestTransactions<Item = Arc<ValidPoolTransaction<T::Transaction>>>>
1027    {
1028        self.get_pool_data().best_transactions_with_attributes(best_transactions_attributes)
1029    }
1030
1031    /// Returns only the first `max` transactions in the pending pool.
1032    pub fn pending_transactions_max(
1033        &self,
1034        max: usize,
1035    ) -> Vec<Arc<ValidPoolTransaction<T::Transaction>>> {
1036        self.get_pool_data().pending_transactions_iter().take(max).collect()
1037    }
1038
1039    /// Returns all transactions from the pending sub-pool
1040    pub fn pending_transactions(&self) -> Vec<Arc<ValidPoolTransaction<T::Transaction>>> {
1041        self.get_pool_data().pending_transactions()
1042    }
1043
1044    /// Returns all transactions from parked pools
1045    pub fn queued_transactions(&self) -> Vec<Arc<ValidPoolTransaction<T::Transaction>>> {
1046        self.get_pool_data().queued_transactions()
1047    }
1048
1049    /// Returns all transactions in the pool
1050    pub fn all_transactions(&self) -> AllPoolTransactions<T::Transaction> {
1051        let pool = self.get_pool_data();
1052        AllPoolTransactions {
1053            pending: pool.pending_transactions(),
1054            queued: pool.queued_transactions(),
1055        }
1056    }
1057
1058    /// Returns all transactions of the given sender, collected under a single read guard so that
1059    /// both sides reflect the same pool state.
1060    pub fn all_transactions_by_sender(
1061        &self,
1062        sender: Address,
1063    ) -> AllPoolTransactions<T::Transaction> {
1064        let Some(sender_id) = self.sender_id(&sender) else { return Default::default() };
1065        let pool = self.get_pool_data();
1066        AllPoolTransactions {
1067            pending: pool.pending_txs_by_sender(sender_id),
1068            queued: pool.queued_txs_by_sender(sender_id),
1069        }
1070    }
1071
1072    /// Returns _all_ transactions in the pool
1073    pub fn all_transaction_hashes(&self) -> Vec<TxHash> {
1074        self.get_pool_data().all().transactions_iter().map(|tx| *tx.hash()).collect()
1075    }
1076
1077    /// Removes and returns all matching transactions from the pool.
1078    ///
1079    /// This behaves as if the transactions got discarded (_not_ mined), effectively introducing a
1080    /// nonce gap for the given transactions.
1081    pub fn remove_transactions(
1082        &self,
1083        hashes: Vec<TxHash>,
1084    ) -> Vec<Arc<ValidPoolTransaction<T::Transaction>>> {
1085        if hashes.is_empty() {
1086            return Vec::new()
1087        }
1088        let removed = self.pool.write().remove_transactions(hashes);
1089
1090        self.with_event_listener(|listener| listener.discarded_many(&removed));
1091
1092        removed
1093    }
1094
1095    /// Removes and returns all matching transactions and their dependent transactions from the
1096    /// pool.
1097    pub fn remove_transactions_and_descendants(
1098        &self,
1099        hashes: Vec<TxHash>,
1100    ) -> Vec<Arc<ValidPoolTransaction<T::Transaction>>> {
1101        if hashes.is_empty() {
1102            return Vec::new()
1103        }
1104        let removed = self.pool.write().remove_transactions_and_descendants(hashes);
1105
1106        self.with_event_listener(|listener| {
1107            for tx in &removed {
1108                listener.discarded(tx.hash());
1109            }
1110        });
1111
1112        removed
1113    }
1114
1115    /// Removes and returns all transactions by the specified sender from the pool.
1116    pub fn remove_transactions_by_sender(
1117        &self,
1118        sender: Address,
1119    ) -> Vec<Arc<ValidPoolTransaction<T::Transaction>>> {
1120        let Some(sender_id) = self.sender_id(&sender) else { return Vec::new() };
1121        let removed = self.pool.write().remove_transactions_by_sender(sender_id);
1122
1123        self.with_event_listener(|listener| listener.discarded_many(&removed));
1124
1125        removed
1126    }
1127
1128    /// Prunes and returns all matching transactions from the pool.
1129    ///
1130    /// This removes the transactions as if they were mined: descendant transactions are **not**
1131    /// parked and remain eligible for inclusion.
1132    pub fn prune_transactions(
1133        &self,
1134        hashes: Vec<TxHash>,
1135    ) -> Vec<Arc<ValidPoolTransaction<T::Transaction>>> {
1136        if hashes.is_empty() {
1137            return Vec::new()
1138        }
1139
1140        self.pool.write().prune_transactions(hashes)
1141    }
1142
1143    /// Retains only transactions that are not present in the pool.
1144    pub fn retain_unknown<A>(&self, announcement: &mut A)
1145    where
1146        A: HandleMempoolData,
1147    {
1148        if announcement.is_empty() {
1149            return
1150        }
1151        let pool = self.get_pool_data();
1152        announcement.retain_by_hash(|tx| !pool.contains(tx))
1153    }
1154
1155    /// Retains only transactions that are present in the pool.
1156    pub fn retain_contains<A>(&self, announcement: &mut A)
1157    where
1158        A: HandleMempoolData,
1159    {
1160        if announcement.is_empty() {
1161            return
1162        }
1163        let pool = self.get_pool_data();
1164        announcement.retain_by_hash(|tx| pool.contains(tx))
1165    }
1166
1167    /// Returns the transaction by hash.
1168    pub fn get(&self, tx_hash: &TxHash) -> Option<Arc<ValidPoolTransaction<T::Transaction>>> {
1169        self.get_pool_data().get(tx_hash)
1170    }
1171
1172    /// Returns all transactions of the address
1173    pub fn get_transactions_by_sender(
1174        &self,
1175        sender: Address,
1176    ) -> Vec<Arc<ValidPoolTransaction<T::Transaction>>> {
1177        let Some(sender_id) = self.sender_id(&sender) else { return Vec::new() };
1178        self.get_pool_data().get_transactions_by_sender(sender_id)
1179    }
1180
1181    /// Returns a pending transaction sent by the given sender with the given nonce.
1182    pub fn get_pending_transaction_by_sender_and_nonce(
1183        &self,
1184        sender: Address,
1185        nonce: u64,
1186    ) -> Option<Arc<ValidPoolTransaction<T::Transaction>>> {
1187        let sender_id = self.sender_id(&sender)?;
1188        self.get_pool_data().get_pending_transaction_by_sender_and_nonce(sender_id, nonce)
1189    }
1190
1191    /// Returns all queued transactions of the address by sender
1192    pub fn get_queued_transactions_by_sender(
1193        &self,
1194        sender: Address,
1195    ) -> Vec<Arc<ValidPoolTransaction<T::Transaction>>> {
1196        let Some(sender_id) = self.sender_id(&sender) else { return Vec::new() };
1197        self.get_pool_data().queued_txs_by_sender(sender_id)
1198    }
1199
1200    /// Returns all pending transactions filtered by predicate
1201    pub fn pending_transactions_with_predicate(
1202        &self,
1203        predicate: impl FnMut(&ValidPoolTransaction<T::Transaction>) -> bool,
1204    ) -> Vec<Arc<ValidPoolTransaction<T::Transaction>>> {
1205        self.get_pool_data().pending_transactions_with_predicate(predicate)
1206    }
1207
1208    /// Returns all pending transactions of the address by sender
1209    pub fn get_pending_transactions_by_sender(
1210        &self,
1211        sender: Address,
1212    ) -> Vec<Arc<ValidPoolTransaction<T::Transaction>>> {
1213        let Some(sender_id) = self.sender_id(&sender) else { return Vec::new() };
1214        self.get_pool_data().pending_txs_by_sender(sender_id)
1215    }
1216
1217    /// Returns the highest transaction of the address
1218    pub fn get_highest_transaction_by_sender(
1219        &self,
1220        sender: Address,
1221    ) -> Option<Arc<ValidPoolTransaction<T::Transaction>>> {
1222        let sender_id = self.sender_id(&sender)?;
1223        self.get_pool_data().get_highest_transaction_by_sender(sender_id)
1224    }
1225
1226    /// Returns the transaction with the highest nonce that is executable given the on chain nonce.
1227    pub fn get_highest_consecutive_transaction_by_sender(
1228        &self,
1229        sender: Address,
1230        on_chain_nonce: u64,
1231    ) -> Option<Arc<ValidPoolTransaction<T::Transaction>>> {
1232        let sender_id = self.sender_id(&sender)?;
1233        self.get_pool_data().get_highest_consecutive_transaction_by_sender(
1234            sender_id.into_transaction_id(on_chain_nonce),
1235        )
1236    }
1237
1238    /// Returns the transaction given a [`TransactionId`]
1239    pub fn get_transaction_by_transaction_id(
1240        &self,
1241        transaction_id: &TransactionId,
1242    ) -> Option<Arc<ValidPoolTransaction<T::Transaction>>> {
1243        self.get_pool_data().all().get(transaction_id).map(|tx| tx.transaction.clone())
1244    }
1245
1246    /// Returns all transactions that where submitted with the given [`TransactionOrigin`]
1247    pub fn get_transactions_by_origin(
1248        &self,
1249        origin: TransactionOrigin,
1250    ) -> Vec<Arc<ValidPoolTransaction<T::Transaction>>> {
1251        self.get_pool_data()
1252            .all()
1253            .transactions_iter()
1254            .filter(|tx| tx.origin == origin)
1255            .cloned()
1256            .collect()
1257    }
1258
1259    /// Returns all pending transactions filtered by [`TransactionOrigin`]
1260    pub fn get_pending_transactions_by_origin(
1261        &self,
1262        origin: TransactionOrigin,
1263    ) -> Vec<Arc<ValidPoolTransaction<T::Transaction>>> {
1264        self.get_pool_data().pending_transactions_iter().filter(|tx| tx.origin == origin).collect()
1265    }
1266
1267    /// Returns all the transactions belonging to the hashes.
1268    ///
1269    /// If no transaction exists, it is skipped.
1270    pub fn get_all(&self, txs: Vec<TxHash>) -> Vec<Arc<ValidPoolTransaction<T::Transaction>>> {
1271        if txs.is_empty() {
1272            return Vec::new()
1273        }
1274        self.get_pool_data().get_all(txs).collect()
1275    }
1276
1277    /// Returns all the transactions belonging to the hashes that are propagatable.
1278    ///
1279    /// If no transaction exists, it is skipped.
1280    fn get_all_propagatable(
1281        &self,
1282        txs: &[TxHash],
1283    ) -> Vec<Arc<ValidPoolTransaction<T::Transaction>>> {
1284        if txs.is_empty() {
1285            return Vec::new()
1286        }
1287        let pool = self.get_pool_data();
1288        txs.iter().filter_map(|tx| pool.get(tx).filter(|tx| tx.propagate)).collect()
1289    }
1290
1291    /// Notify about propagated transactions.
1292    pub fn on_propagated(&self, txs: PropagatedTransactions) {
1293        if txs.is_empty() {
1294            return
1295        }
1296        self.with_event_listener(|listener| {
1297            txs.into_iter().for_each(|(hash, peers)| listener.propagated(&hash, peers));
1298        });
1299    }
1300
1301    /// Number of transactions in the entire pool
1302    pub fn len(&self) -> usize {
1303        self.get_pool_data().len()
1304    }
1305
1306    /// Whether the pool is empty
1307    pub fn is_empty(&self) -> bool {
1308        self.get_pool_data().is_empty()
1309    }
1310
1311    /// Returns whether or not the pool is over its configured size and transaction count limits.
1312    pub fn is_exceeded(&self) -> bool {
1313        self.pool.read().is_exceeded()
1314    }
1315
1316    /// Inserts a blob transaction into the blob store
1317    fn insert_blob(&self, hash: TxHash, blob: PooledBlobSidecar) {
1318        debug!(target: "txpool", "[{:?}] storing blob sidecar", hash);
1319        if let Err(err) = self.blob_store.insert(hash, blob) {
1320            warn!(target: "txpool", %err, "[{:?}] failed to insert blob", hash);
1321            self.blob_store_metrics.blobstore_failed_inserts.increment(1);
1322        }
1323        self.update_blob_store_metrics();
1324    }
1325
1326    /// Delete a blob from the blob store
1327    pub fn delete_blob(&self, blob: TxHash) {
1328        let _ = self.blob_store.delete(blob);
1329    }
1330
1331    /// Delete all blobs from the blob store
1332    pub fn delete_blobs(&self, txs: Vec<TxHash>) {
1333        let _ = self.blob_store.delete_all(txs);
1334    }
1335
1336    /// Cleans up the blob store
1337    pub fn cleanup_blobs(&self) {
1338        let stat = self.blob_store.cleanup();
1339        self.blob_store_metrics.blobstore_failed_deletes.increment(stat.delete_failed as u64);
1340        self.update_blob_store_metrics();
1341    }
1342
1343    fn update_blob_store_metrics(&self) {
1344        if let Some(data_size) = self.blob_store.data_size_hint() {
1345            self.blob_store_metrics.blobstore_byte_size.set(data_size as f64);
1346        }
1347        self.blob_store_metrics.blobstore_entries.set(self.blob_store.blobs_len() as f64);
1348    }
1349
1350    /// Deletes all blob transactions that were discarded.
1351    fn delete_discarded_blobs<'a>(
1352        &'a self,
1353        transactions: impl IntoIterator<Item = &'a Arc<ValidPoolTransaction<T::Transaction>>>,
1354    ) {
1355        let blob_txs = transactions
1356            .into_iter()
1357            .filter(|tx| tx.transaction.is_eip4844())
1358            .map(|tx| *tx.hash())
1359            .collect();
1360        self.delete_blobs(blob_txs);
1361    }
1362}
1363
1364impl<V, T: TransactionOrdering, S> fmt::Debug for PoolInner<V, T, S> {
1365    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1366        f.debug_struct("PoolInner").field("config", &self.config).finish_non_exhaustive()
1367    }
1368}
1369
1370/// Metadata for a transaction that was added to the pool.
1371///
1372/// This holds all the data needed to complete post-insertion operations (notifications,
1373/// blob storage).
1374#[derive(Debug)]
1375struct AddedTransactionMeta<T: PoolTransaction> {
1376    /// The transaction that was added to the pool
1377    added: AddedTransaction<T>,
1378    /// Optional blob sidecar for EIP-4844 transactions
1379    blob_sidecar: Option<PooledBlobSidecar>,
1380}
1381
1382/// Tracks an added transaction and all graph changes caused by adding it.
1383#[derive(Debug, Clone)]
1384pub struct AddedPendingTransaction<T: PoolTransaction> {
1385    /// Inserted transaction.
1386    pub transaction: Arc<ValidPoolTransaction<T>>,
1387    /// Replaced transaction.
1388    pub replaced: Option<Arc<ValidPoolTransaction<T>>>,
1389    /// transactions promoted to the pending queue
1390    pub promoted: Vec<Arc<ValidPoolTransaction<T>>>,
1391    /// transactions that failed and became discarded
1392    pub discarded: Vec<Arc<ValidPoolTransaction<T>>>,
1393}
1394
1395impl<T: PoolTransaction> AddedPendingTransaction<T> {
1396    /// Returns all transactions that were promoted to the pending pool and adhere to the given
1397    /// [`TransactionListenerKind`].
1398    ///
1399    /// If the kind is [`TransactionListenerKind::PropagateOnly`], then only transactions that
1400    /// are allowed to be propagated are returned.
1401    pub(crate) fn pending_transactions(
1402        &self,
1403        kind: TransactionListenerKind,
1404    ) -> impl Iterator<Item = B256> + '_ {
1405        let iter = std::iter::once(&self.transaction).chain(self.promoted.iter());
1406        PendingTransactionIter { kind, iter }
1407    }
1408}
1409
1410pub(crate) struct PendingTransactionIter<Iter> {
1411    kind: TransactionListenerKind,
1412    iter: Iter,
1413}
1414
1415impl<'a, Iter, T> Iterator for PendingTransactionIter<Iter>
1416where
1417    Iter: Iterator<Item = &'a Arc<ValidPoolTransaction<T>>>,
1418    T: PoolTransaction + 'a,
1419{
1420    type Item = B256;
1421
1422    fn next(&mut self) -> Option<Self::Item> {
1423        loop {
1424            let next = self.iter.next()?;
1425            if self.kind.is_propagate_only() && !next.propagate {
1426                continue
1427            }
1428            return Some(*next.hash())
1429        }
1430    }
1431}
1432
1433/// An iterator over full pending transactions
1434pub(crate) struct FullPendingTransactionIter<Iter> {
1435    kind: TransactionListenerKind,
1436    iter: Iter,
1437}
1438
1439impl<'a, Iter, T> Iterator for FullPendingTransactionIter<Iter>
1440where
1441    Iter: Iterator<Item = &'a Arc<ValidPoolTransaction<T>>>,
1442    T: PoolTransaction + 'a,
1443{
1444    type Item = NewTransactionEvent<T>;
1445
1446    fn next(&mut self) -> Option<Self::Item> {
1447        loop {
1448            let next = self.iter.next()?;
1449            if self.kind.is_propagate_only() && !next.propagate {
1450                continue
1451            }
1452            return Some(NewTransactionEvent {
1453                subpool: SubPool::Pending,
1454                transaction: next.clone(),
1455            })
1456        }
1457    }
1458}
1459
1460/// Represents a transaction that was added into the pool and its state
1461#[derive(Debug, Clone)]
1462pub enum AddedTransaction<T: PoolTransaction> {
1463    /// Transaction was successfully added and moved to the pending pool.
1464    Pending(AddedPendingTransaction<T>),
1465    /// Transaction was successfully added but not yet ready for processing and moved to a
1466    /// parked pool instead.
1467    Parked {
1468        /// Inserted transaction.
1469        transaction: Arc<ValidPoolTransaction<T>>,
1470        /// Replaced transaction.
1471        replaced: Option<Arc<ValidPoolTransaction<T>>>,
1472        /// The subpool it was moved to.
1473        subpool: SubPool,
1474        /// The specific reason why the transaction is queued (if applicable).
1475        queued_reason: Option<QueuedReason>,
1476    },
1477}
1478
1479impl<T: PoolTransaction> AddedTransaction<T> {
1480    /// Returns whether the transaction has been added to the pending pool.
1481    pub const fn as_pending(&self) -> Option<&AddedPendingTransaction<T>> {
1482        match self {
1483            Self::Pending(tx) => Some(tx),
1484            _ => None,
1485        }
1486    }
1487
1488    /// Returns the replaced transaction if there was one
1489    pub const fn replaced(&self) -> Option<&Arc<ValidPoolTransaction<T>>> {
1490        match self {
1491            Self::Pending(tx) => tx.replaced.as_ref(),
1492            Self::Parked { replaced, .. } => replaced.as_ref(),
1493        }
1494    }
1495
1496    /// Returns the discarded transactions if there were any
1497    pub(crate) fn discarded_transactions(&self) -> Option<&[Arc<ValidPoolTransaction<T>>]> {
1498        match self {
1499            Self::Pending(tx) => Some(&tx.discarded),
1500            Self::Parked { .. } => None,
1501        }
1502    }
1503
1504    /// Returns the hash of the replaced transaction if it is a blob transaction.
1505    pub(crate) fn replaced_blob_transaction(&self) -> Option<B256> {
1506        self.replaced().filter(|tx| tx.transaction.is_eip4844()).map(|tx| *tx.transaction.hash())
1507    }
1508
1509    /// Returns the hash of the transaction
1510    pub fn hash(&self) -> &TxHash {
1511        match self {
1512            Self::Pending(tx) => tx.transaction.hash(),
1513            Self::Parked { transaction, .. } => transaction.hash(),
1514        }
1515    }
1516
1517    /// Converts this type into the event type for listeners
1518    pub fn into_new_transaction_event(self) -> NewTransactionEvent<T> {
1519        match self {
1520            Self::Pending(tx) => {
1521                NewTransactionEvent { subpool: SubPool::Pending, transaction: tx.transaction }
1522            }
1523            Self::Parked { transaction, subpool, .. } => {
1524                NewTransactionEvent { transaction, subpool }
1525            }
1526        }
1527    }
1528
1529    /// Returns the subpool this transaction was added to
1530    pub(crate) const fn subpool(&self) -> SubPool {
1531        match self {
1532            Self::Pending(_) => SubPool::Pending,
1533            Self::Parked { subpool, .. } => *subpool,
1534        }
1535    }
1536
1537    /// Returns the [`TransactionId`] of the added transaction
1538    #[cfg(test)]
1539    pub(crate) fn id(&self) -> &TransactionId {
1540        match self {
1541            Self::Pending(added) => added.transaction.id(),
1542            Self::Parked { transaction, .. } => transaction.id(),
1543        }
1544    }
1545
1546    /// Returns the queued reason if the transaction is parked with a queued reason.
1547    pub const fn queued_reason(&self) -> Option<&QueuedReason> {
1548        match self {
1549            Self::Pending(_) => None,
1550            Self::Parked { queued_reason, .. } => queued_reason.as_ref(),
1551        }
1552    }
1553
1554    /// Returns the transaction state based on the subpool and queued reason.
1555    pub fn transaction_state(&self) -> AddedTransactionState {
1556        match self.subpool() {
1557            SubPool::Pending => AddedTransactionState::Pending,
1558            _ => {
1559                // For non-pending transactions, use the queued reason directly from the
1560                // AddedTransaction
1561                if let Some(reason) = self.queued_reason() {
1562                    AddedTransactionState::Queued(reason.clone())
1563                } else {
1564                    // Fallback - this shouldn't happen with the new implementation
1565                    AddedTransactionState::Queued(QueuedReason::NonceGap)
1566                }
1567            }
1568        }
1569    }
1570}
1571
1572/// The specific reason why a transaction is queued (not ready for execution)
1573#[derive(Debug, Clone, PartialEq, Eq)]
1574pub enum QueuedReason {
1575    /// Transaction has a nonce gap - missing prior transactions
1576    NonceGap,
1577    /// Transaction has parked ancestors - waiting for other transactions to be mined
1578    ParkedAncestors,
1579    /// Sender has insufficient balance to cover the transaction cost
1580    InsufficientBalance,
1581    /// Transaction exceeds the block gas limit
1582    TooMuchGas,
1583    /// Transaction doesn't meet the base fee requirement
1584    InsufficientBaseFee,
1585    /// Transaction doesn't meet the blob fee requirement (EIP-4844)
1586    InsufficientBlobFee,
1587}
1588
1589/// The state of a transaction when is was added to the pool
1590#[derive(Debug, Clone, PartialEq, Eq)]
1591pub enum AddedTransactionState {
1592    /// Ready for execution
1593    Pending,
1594    /// Not ready for execution due to a specific condition
1595    Queued(QueuedReason),
1596}
1597
1598impl AddedTransactionState {
1599    /// Returns whether the transaction was submitted as queued.
1600    pub const fn is_queued(&self) -> bool {
1601        matches!(self, Self::Queued(_))
1602    }
1603
1604    /// Returns whether the transaction was submitted as pending.
1605    pub const fn is_pending(&self) -> bool {
1606        matches!(self, Self::Pending)
1607    }
1608
1609    /// Returns the specific queued reason if the transaction is queued.
1610    pub const fn queued_reason(&self) -> Option<&QueuedReason> {
1611        match self {
1612            Self::Queued(reason) => Some(reason),
1613            Self::Pending => None,
1614        }
1615    }
1616}
1617
1618/// The outcome of a successful transaction addition
1619#[derive(Debug, Clone, PartialEq, Eq)]
1620pub struct AddedTransactionOutcome {
1621    /// The hash of the transaction
1622    pub hash: TxHash,
1623    /// The state of the transaction
1624    pub state: AddedTransactionState,
1625}
1626
1627impl AddedTransactionOutcome {
1628    /// Returns whether the transaction was submitted as queued.
1629    pub const fn is_queued(&self) -> bool {
1630        self.state.is_queued()
1631    }
1632
1633    /// Returns whether the transaction was submitted as pending.
1634    pub const fn is_pending(&self) -> bool {
1635        self.state.is_pending()
1636    }
1637}
1638
1639/// Contains all state changes after a [`CanonicalStateUpdate`] was processed
1640#[derive(Debug)]
1641pub(crate) struct OnNewCanonicalStateOutcome<T: PoolTransaction> {
1642    /// Hash of the block.
1643    pub(crate) block_hash: B256,
1644    /// All mined transactions.
1645    pub(crate) mined: Vec<TxHash>,
1646    /// Transactions promoted to the pending pool.
1647    pub(crate) promoted: Vec<Arc<ValidPoolTransaction<T>>>,
1648    /// transaction that were discarded during the update
1649    pub(crate) discarded: Vec<Arc<ValidPoolTransaction<T>>>,
1650}
1651
1652impl<T: PoolTransaction> OnNewCanonicalStateOutcome<T> {
1653    /// Returns all transactions that were promoted to the pending pool and adhere to the given
1654    /// [`TransactionListenerKind`].
1655    ///
1656    /// If the kind is [`TransactionListenerKind::PropagateOnly`], then only transactions that
1657    /// are allowed to be propagated are returned.
1658    pub(crate) fn pending_transactions(
1659        &self,
1660        kind: TransactionListenerKind,
1661    ) -> impl Iterator<Item = B256> + '_ {
1662        let iter = self.promoted.iter();
1663        PendingTransactionIter { kind, iter }
1664    }
1665
1666    /// Returns all FULL transactions that were promoted to the pending pool and adhere to the given
1667    /// [`TransactionListenerKind`].
1668    ///
1669    /// If the kind is [`TransactionListenerKind::PropagateOnly`], then only transactions that
1670    /// are allowed to be propagated are returned.
1671    pub(crate) fn full_pending_transactions(
1672        &self,
1673        kind: TransactionListenerKind,
1674    ) -> impl Iterator<Item = NewTransactionEvent<T>> + '_ {
1675        let iter = self.promoted.iter();
1676        FullPendingTransactionIter { kind, iter }
1677    }
1678}
1679
1680#[cfg(test)]
1681mod tests {
1682    use crate::{
1683        blobstore::{BlobStore, InMemoryBlobStore, PooledBlobSidecar},
1684        identifier::SenderId,
1685        test_utils::{testing_pool, MockTransaction, TestPoolBuilder},
1686        validate::ValidTransaction,
1687        BlockInfo, PoolConfig, SubPoolLimit, TransactionOrigin, TransactionPool,
1688        TransactionPoolExt, TransactionValidationOutcome, ValidPoolTransaction, U256,
1689    };
1690    use alloy_consensus::Transaction;
1691    use alloy_eips::{eip4844::BlobTransactionSidecar, eip7594::BlobTransactionSidecarVariant};
1692    use alloy_primitives::Address;
1693    use std::{fs, path::PathBuf, sync::Arc};
1694
1695    #[tokio::test]
1696    async fn all_transactions_by_sender_across_reclassification() {
1697        let pool = testing_pool();
1698        let sender = Address::with_last_byte(1);
1699        // nonces 0 and 1 are pending, the gapped nonce 9 is queued
1700        for nonce in [0, 1, 9] {
1701            let tx =
1702                MockTransaction::legacy().with_sender(sender).with_nonce(nonce).with_gas_price(100);
1703            pool.add_transaction(TransactionOrigin::External, tx).await.unwrap();
1704        }
1705        let nonces = |txs: &[Arc<ValidPoolTransaction<MockTransaction>>]| {
1706            let mut nonces: Vec<_> = txs.iter().map(|tx| tx.transaction.nonce()).collect();
1707            nonces.sort_unstable();
1708            nonces
1709        };
1710
1711        let txs = pool.all_transactions_by_sender(sender);
1712        assert_eq!(nonces(&txs.pending), [0, 1]);
1713        assert_eq!(nonces(&txs.queued), [9]);
1714
1715        // a higher base fee reclassifies the pending transactions as queued; the snapshot must
1716        // report each of them on exactly one side
1717        pool.set_block_info(BlockInfo {
1718            pending_basefee: 200,
1719            block_gas_limit: 30_000_000,
1720            ..Default::default()
1721        });
1722        let txs = pool.all_transactions_by_sender(sender);
1723        assert!(txs.pending.is_empty());
1724        assert_eq!(nonces(&txs.queued), [0, 1, 9]);
1725        assert_eq!(nonces(&pool.get_transactions_by_sender(sender)), [0, 1, 9]);
1726    }
1727
1728    #[test]
1729    fn test_discard_blobs_on_blob_tx_eviction() {
1730        let blobs = {
1731            // Read the contents of the JSON file into a string.
1732            let json_content = fs::read_to_string(
1733                PathBuf::from(env!("CARGO_MANIFEST_DIR")).join("test_data/blob1.json"),
1734            )
1735            .expect("Failed to read the blob data file");
1736
1737            // Parse the JSON contents into a serde_json::Value.
1738            let json_value: serde_json::Value =
1739                serde_json::from_str(&json_content).expect("Failed to deserialize JSON");
1740
1741            // Extract blob data from JSON and convert it to Blob.
1742            vec![
1743                // Extract the "data" field from the JSON and parse it as a string.
1744                json_value
1745                    .get("data")
1746                    .unwrap()
1747                    .as_str()
1748                    .expect("Data is not a valid string")
1749                    .to_string(),
1750            ]
1751        };
1752
1753        // Generate a BlobTransactionSidecar from the blobs.
1754        let sidecar = BlobTransactionSidecarVariant::Eip4844(
1755            BlobTransactionSidecar::try_from_blobs_hex(blobs).unwrap(),
1756        );
1757
1758        // Define the maximum limit for blobs in the sub-pool.
1759        let blob_limit = SubPoolLimit::new(1000, usize::MAX);
1760
1761        // Create a test pool with default configuration and the specified blob limit.
1762        let test_pool = &TestPoolBuilder::default()
1763            .with_config(PoolConfig { blob_limit, ..Default::default() })
1764            .pool;
1765
1766        // Set the block info for the pool, including a pending blob fee.
1767        test_pool
1768            .set_block_info(BlockInfo { pending_blob_fee: Some(10_000_000), ..Default::default() });
1769
1770        // Create an in-memory blob store.
1771        let blob_store = InMemoryBlobStore::default();
1772
1773        // Loop to add transactions to the pool and test blob eviction.
1774        for n in 0..blob_limit.max_txs + 10 {
1775            // Create a mock transaction with the generated blob sidecar.
1776            let mut tx = MockTransaction::eip4844_with_sidecar(sidecar.clone());
1777
1778            // Set non zero size
1779            tx.set_size(1844674407370951);
1780
1781            // Insert the sidecar into the blob store if the current index is within the blob limit.
1782            if n < blob_limit.max_txs {
1783                blob_store.insert(*tx.get_hash(), sidecar.clone().into()).unwrap();
1784            }
1785
1786            // Add the transaction to the pool with external origin and valid outcome.
1787            test_pool.add_transactions(
1788                TransactionOrigin::External,
1789                [TransactionValidationOutcome::Valid {
1790                    balance: U256::from(1_000),
1791                    state_nonce: 0,
1792                    bytecode_hash: None,
1793                    transaction: ValidTransaction::ValidWithSidecar {
1794                        transaction: tx,
1795                        sidecar: PooledBlobSidecar::from(sidecar.clone()),
1796                    },
1797                    propagate: true,
1798                    authorities: None,
1799                }],
1800            );
1801        }
1802
1803        // Assert that the size of the pool's blob component is equal to the maximum blob limit.
1804        assert_eq!(test_pool.size().blob, blob_limit.max_txs);
1805
1806        // Assert that the size of the pool's blob_size component matches the expected value.
1807        assert_eq!(test_pool.size().blob_size, 1844674407370951000);
1808
1809        // Assert that the pool's blob store matches the expected blob store.
1810        assert_eq!(*test_pool.blob_store(), blob_store);
1811    }
1812
1813    #[test]
1814    fn test_auths_stored_in_identifiers() {
1815        // Create a test pool with default configuration.
1816        let test_pool = &TestPoolBuilder::default().with_config(Default::default()).pool;
1817
1818        let auth = Address::new([1; 20]);
1819        let tx = MockTransaction::eip7702();
1820
1821        test_pool.add_transactions(
1822            TransactionOrigin::Local,
1823            [TransactionValidationOutcome::Valid {
1824                balance: U256::from(1_000),
1825                state_nonce: 0,
1826                bytecode_hash: None,
1827                transaction: ValidTransaction::Valid(tx),
1828                propagate: true,
1829                authorities: Some(vec![auth]),
1830            }],
1831        );
1832
1833        let identifiers = test_pool.identifiers.read();
1834        assert_eq!(identifiers.sender_id(&auth), Some(SenderId::from(1)));
1835    }
1836
1837    #[test]
1838    fn sender_queries_do_not_allocate_ids_for_unknown_addresses() {
1839        let test_pool = &TestPoolBuilder::default().with_config(Default::default()).pool;
1840        let sender = Address::new([9; 20]);
1841
1842        assert_eq!(test_pool.sender_id(&sender), None);
1843        assert!(test_pool.get_transactions_by_sender(sender).is_empty());
1844        assert!(test_pool.get_pending_transaction_by_sender_and_nonce(sender, 0).is_none());
1845        assert!(test_pool.get_queued_transactions_by_sender(sender).is_empty());
1846        assert!(test_pool.get_pending_transactions_by_sender(sender).is_empty());
1847        assert!(test_pool.get_highest_transaction_by_sender(sender).is_none());
1848        assert!(test_pool.get_highest_consecutive_transaction_by_sender(sender, 0).is_none());
1849        assert!(test_pool.remove_transactions_by_sender(sender).is_empty());
1850        assert_eq!(test_pool.sender_id(&sender), None);
1851    }
1852}