Skip to main content

anvil/eth/pool/
transactions.rs

1use crate::eth::{error::PoolError, util::hex_fmt_many};
2use alloy_consensus::{
3    Transaction, Typed2718,
4    crypto::RecoveryError,
5    transaction::{SignerRecoverable, TxHashRef},
6};
7use alloy_network::AnyRpcTransaction;
8use alloy_primitives::{
9    Address, TxHash,
10    map::{HashMap, HashSet},
11};
12use alloy_rlp::Encodable;
13use anvil_core::eth::transaction::PendingTransaction;
14use parking_lot::RwLock;
15use std::{cmp::Ordering, collections::BTreeSet, fmt, str::FromStr, sync::Arc, time::Instant};
16
17/// A unique identifying marker for a transaction
18pub type TxMarker = Vec<u8>;
19
20/// Result type for replaced transactions: the replaced pool transactions and the hashes they
21/// unlock.
22type ReplacedTransactions<T> = (Vec<Arc<PoolTransaction<T>>>, Vec<TxHash>);
23
24/// Modes that determine the transaction ordering of the mempool
25///
26/// This type controls the transaction order via the priority metric of a transaction
27#[derive(Clone, Copy, Debug, Default, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
28pub enum TransactionOrder {
29    /// Keep the pool transaction transactions sorted in the order they arrive.
30    ///
31    /// This will essentially assign every transaction the exact priority so the order is
32    /// determined by their internal id
33    Fifo,
34    /// This means that it prioritizes transactions based on the fees paid to the miner.
35    #[default]
36    Fees,
37}
38
39impl TransactionOrder {
40    /// Returns the priority of the transactions
41    pub fn priority<T: Transaction>(&self, tx: &T) -> TransactionPriority {
42        match self {
43            Self::Fifo => TransactionPriority::default(),
44            Self::Fees => TransactionPriority(tx.max_fee_per_gas()),
45        }
46    }
47}
48
49impl FromStr for TransactionOrder {
50    type Err = String;
51
52    fn from_str(s: &str) -> Result<Self, Self::Err> {
53        let s = s.to_lowercase();
54        let order = match s.as_str() {
55            "fees" => Self::Fees,
56            "fifo" => Self::Fifo,
57            _ => return Err(format!("Unknown TransactionOrder: `{s}`")),
58        };
59        Ok(order)
60    }
61}
62
63/// Metric value for the priority of a transaction.
64///
65/// The `TransactionPriority` determines the ordering of two transactions that have all their
66/// markers satisfied.
67#[derive(Clone, Copy, Debug, Default, PartialEq, Eq, PartialOrd, Ord)]
68pub struct TransactionPriority(pub u128);
69
70/// Internal Transaction type
71#[derive(Clone, PartialEq, Eq)]
72pub struct PoolTransaction<T> {
73    /// the pending eth transaction
74    pub pending_transaction: PendingTransaction<T>,
75    /// Markers required by the transaction
76    pub requires: Vec<TxMarker>,
77    /// Markers that this transaction provides
78    pub provides: Vec<TxMarker>,
79    /// priority of the transaction
80    pub priority: TransactionPriority,
81    /// Whether this transaction is being replayed from chain history.
82    pub is_replay: bool,
83}
84
85// == impl PoolTransaction ==
86
87impl<T> PoolTransaction<T> {
88    pub const fn new(transaction: PendingTransaction<T>) -> Self {
89        Self {
90            pending_transaction: transaction,
91            requires: vec![],
92            provides: vec![],
93            priority: TransactionPriority(0),
94            is_replay: false,
95        }
96    }
97
98    /// Marks this transaction as a historical replay.
99    pub const fn with_replay(mut self) -> Self {
100        self.is_replay = true;
101        self
102    }
103
104    /// Returns the hash of this transaction
105    pub const fn hash(&self) -> TxHash {
106        *self.pending_transaction.hash()
107    }
108}
109
110impl<T: Transaction> PoolTransaction<T> {
111    /// Returns the max fee per gas of this transaction
112    pub fn max_fee_per_gas(&self) -> u128 {
113        self.pending_transaction.transaction.max_fee_per_gas()
114    }
115}
116
117impl<T: Typed2718> PoolTransaction<T> {
118    /// Returns the type of the transaction
119    pub fn tx_type(&self) -> u8 {
120        self.pending_transaction.transaction.ty()
121    }
122}
123
124impl<T: fmt::Debug> fmt::Debug for PoolTransaction<T> {
125    fn fmt(&self, fmt: &mut fmt::Formatter<'_>) -> fmt::Result {
126        write!(fmt, "Transaction {{ ")?;
127        write!(fmt, "hash: {:?}, ", self.pending_transaction.hash())?;
128        write!(fmt, "requires: [{}], ", hex_fmt_many(self.requires.iter()))?;
129        write!(fmt, "provides: [{}], ", hex_fmt_many(self.provides.iter()))?;
130        write!(fmt, "raw tx: {:?}", self.pending_transaction)?;
131        write!(fmt, "}}")?;
132        Ok(())
133    }
134}
135
136impl<T> TryFrom<AnyRpcTransaction> for PoolTransaction<T>
137where
138    T: SignerRecoverable + TxHashRef + Encodable + TryFrom<AnyRpcTransaction>,
139    <T as TryFrom<AnyRpcTransaction>>::Error: Into<eyre::Error>,
140    RecoveryError: Into<eyre::Error>,
141{
142    type Error = eyre::Error;
143    fn try_from(value: AnyRpcTransaction) -> Result<Self, Self::Error> {
144        let typed_transaction = T::try_from(value).map_err(Into::into)?;
145        let pending_transaction = PendingTransaction::new(typed_transaction)?;
146        Ok(Self {
147            pending_transaction,
148            requires: vec![],
149            provides: vec![],
150            priority: TransactionPriority(0),
151            is_replay: false,
152        })
153    }
154}
155
156/// A waiting pool of transaction that are pending, but not yet ready to be included in a new block.
157///
158/// Keeps a set of transactions that are waiting for other transactions
159#[derive(Clone, Debug)]
160pub struct PendingTransactions<T> {
161    /// markers that aren't yet provided by any transaction
162    required_markers: HashMap<TxMarker, HashSet<TxHash>>,
163    /// mapping of the markers of a transaction to the hash of the transaction
164    waiting_markers: HashMap<Vec<TxMarker>, TxHash>,
165    /// the transactions that are not ready yet are waiting for another tx to finish
166    waiting_queue: HashMap<TxHash, PendingPoolTransaction<T>>,
167}
168
169impl<T> Default for PendingTransactions<T> {
170    fn default() -> Self {
171        Self {
172            required_markers: Default::default(),
173            waiting_markers: Default::default(),
174            waiting_queue: Default::default(),
175        }
176    }
177}
178
179impl<T> PendingTransactions<T> {
180    /// Returns an independent snapshot of the pending transactions.
181    pub(super) fn snapshot(&self) -> Self {
182        Self {
183            required_markers: self.required_markers.clone(),
184            waiting_markers: self.waiting_markers.clone(),
185            waiting_queue: self.waiting_queue.clone(),
186        }
187    }
188
189    /// Returns the number of transactions that are currently waiting
190    pub fn len(&self) -> usize {
191        self.waiting_queue.len()
192    }
193
194    pub fn is_empty(&self) -> bool {
195        self.waiting_queue.is_empty()
196    }
197
198    /// Clears internal state
199    pub fn clear(&mut self) {
200        self.required_markers.clear();
201        self.waiting_markers.clear();
202        self.waiting_queue.clear();
203    }
204
205    /// Returns an iterator over all transactions in the waiting pool
206    pub fn transactions(&self) -> impl Iterator<Item = Arc<PoolTransaction<T>>> + '_ {
207        self.waiting_queue.values().map(|tx| tx.transaction.clone())
208    }
209
210    /// Returns true if given transaction is part of the queue
211    pub fn contains(&self, hash: &TxHash) -> bool {
212        self.waiting_queue.contains_key(hash)
213    }
214
215    /// Returns the transaction for the hash if it's pending
216    pub fn get(&self, hash: &TxHash) -> Option<&PendingPoolTransaction<T>> {
217        self.waiting_queue.get(hash)
218    }
219
220    /// This will check off the markers of pending transactions.
221    ///
222    /// Returns the those transactions that become unlocked (all markers checked) and can be moved
223    /// to the ready queue.
224    pub fn mark_and_unlock(
225        &mut self,
226        markers: impl IntoIterator<Item = impl AsRef<TxMarker>>,
227    ) -> Vec<PendingPoolTransaction<T>> {
228        let mut unlocked_ready = Vec::new();
229        for mark in markers {
230            let mark = mark.as_ref();
231            if let Some(tx_hashes) = self.required_markers.remove(mark) {
232                for hash in tx_hashes {
233                    let tx = self.waiting_queue.get_mut(&hash).expect("tx is included;");
234                    tx.mark(mark);
235
236                    if tx.is_ready() {
237                        let tx = self.waiting_queue.remove(&hash).expect("tx is included;");
238                        self.waiting_markers.remove(&tx.transaction.provides);
239
240                        unlocked_ready.push(tx);
241                    }
242                }
243            }
244        }
245
246        unlocked_ready
247    }
248
249    /// Removes the transactions associated with the given hashes
250    ///
251    /// Returns all removed transactions.
252    pub fn remove(&mut self, hashes: Vec<TxHash>) -> Vec<Arc<PoolTransaction<T>>> {
253        let mut removed = vec![];
254        for hash in hashes {
255            if let Some(waiting_tx) = self.waiting_queue.remove(&hash) {
256                self.waiting_markers.remove(&waiting_tx.transaction.provides);
257                for marker in waiting_tx.missing_markers {
258                    let remove = if let Some(required) = self.required_markers.get_mut(&marker) {
259                        required.remove(&hash);
260                        required.is_empty()
261                    } else {
262                        false
263                    };
264                    if remove {
265                        self.required_markers.remove(&marker);
266                    }
267                }
268                removed.push(waiting_tx.transaction)
269            }
270        }
271        removed
272    }
273
274    /// Removes transactions and their transitive dependents from the waiting pool.
275    pub fn remove_with_dependents(
276        &mut self,
277        hashes: Vec<TxHash>,
278        invalidated: impl IntoIterator<Item = TxMarker>,
279    ) -> Vec<Arc<PoolTransaction<T>>> {
280        let mut required_by = HashMap::<TxMarker, Vec<TxHash>>::default();
281        for (hash, tx) in &self.waiting_queue {
282            for marker in &tx.transaction.requires {
283                required_by.entry(marker.clone()).or_default().push(*hash);
284            }
285        }
286
287        let mut to_remove = HashSet::<TxHash>::default();
288        let mut markers = invalidated.into_iter().collect::<Vec<_>>();
289        for hash in hashes {
290            if to_remove.insert(hash)
291                && let Some(tx) = self.waiting_queue.get(&hash)
292            {
293                markers.extend(tx.transaction.provides.iter().cloned());
294            }
295        }
296
297        while let Some(marker) = markers.pop() {
298            if let Some(dependents) = required_by.remove(&marker) {
299                for hash in dependents {
300                    if to_remove.insert(hash)
301                        && let Some(tx) = self.waiting_queue.get(&hash)
302                    {
303                        markers.extend(tx.transaction.provides.iter().cloned());
304                    }
305                }
306            }
307        }
308
309        self.remove(to_remove.into_iter().collect())
310    }
311}
312
313impl<T: Transaction> PendingTransactions<T> {
314    /// Adds a transaction to Pending queue of transactions
315    pub fn add_transaction(&mut self, tx: PendingPoolTransaction<T>) -> Result<(), PoolError> {
316        assert!(!tx.is_ready(), "transaction must not be ready");
317        assert!(
318            !self.waiting_queue.contains_key(&tx.transaction.hash()),
319            "transaction is already added"
320        );
321
322        let replaced_hash = if let Some(replace) = self
323            .waiting_markers
324            .get(&tx.transaction.provides)
325            .and_then(|hash| self.waiting_queue.get(hash))
326        {
327            // check if underpriced
328            if tx.transaction.max_fee_per_gas() <= replace.transaction.max_fee_per_gas() {
329                warn!(target: "txpool", "pending replacement transaction underpriced [{:?}]", tx.transaction.hash());
330                return Err(PoolError::ReplacementUnderpriced(tx.transaction.hash()));
331            }
332            Some(replace.transaction.hash())
333        } else {
334            None
335        };
336        // Remove old markers before inserting the replacement, which shares their keys.
337        if let Some(replaced_hash) = replaced_hash {
338            self.remove(vec![replaced_hash]);
339        }
340
341        // add all missing markers
342        for marker in &tx.missing_markers {
343            self.required_markers.entry(marker.clone()).or_default().insert(tx.transaction.hash());
344        }
345
346        // also track identifying markers
347        self.waiting_markers.insert(tx.transaction.provides.clone(), tx.transaction.hash());
348        // add tx to the queue
349        self.waiting_queue.insert(tx.transaction.hash(), tx);
350
351        Ok(())
352    }
353}
354
355/// A transaction in the pool
356pub struct PendingPoolTransaction<T> {
357    pub transaction: Arc<PoolTransaction<T>>,
358    /// markers required and have not been satisfied yet by other transactions in the pool
359    pub missing_markers: HashSet<TxMarker>,
360    /// timestamp when the tx was added
361    pub added_at: Instant,
362}
363
364impl<T> Clone for PendingPoolTransaction<T> {
365    fn clone(&self) -> Self {
366        Self {
367            transaction: Arc::clone(&self.transaction),
368            missing_markers: self.missing_markers.clone(),
369            added_at: self.added_at,
370        }
371    }
372}
373
374impl<T> PendingPoolTransaction<T> {
375    /// Creates a new `PendingPoolTransaction`.
376    ///
377    /// Determines the markers that are still missing before this transaction can be moved to the
378    /// ready queue.
379    pub fn new(transaction: PoolTransaction<T>, provided: &HashMap<TxMarker, TxHash>) -> Self {
380        let missing_markers = transaction
381            .requires
382            .iter()
383            .filter(|marker| {
384                // is true if the marker is already satisfied either via transaction in the pool
385                !provided.contains_key(&**marker)
386            })
387            .cloned()
388            .collect();
389
390        Self { transaction: Arc::new(transaction), missing_markers, added_at: Instant::now() }
391    }
392
393    /// Removes the required marker
394    pub fn mark(&mut self, marker: &TxMarker) {
395        self.missing_markers.remove(marker);
396    }
397
398    /// Returns true if transaction has all requirements satisfied.
399    pub fn is_ready(&self) -> bool {
400        self.missing_markers.is_empty()
401    }
402}
403
404impl<T: fmt::Debug> fmt::Debug for PendingPoolTransaction<T> {
405    fn fmt(&self, fmt: &mut fmt::Formatter<'_>) -> fmt::Result {
406        write!(fmt, "PendingTransaction {{ ")?;
407        write!(fmt, "added_at: {:?}, ", self.added_at)?;
408        write!(fmt, "tx: {:?}, ", self.transaction)?;
409        write!(fmt, "missing_markers: {{{}}}", hex_fmt_many(self.missing_markers.iter()))?;
410        write!(fmt, "}}")
411    }
412}
413
414pub struct TransactionsIterator<T> {
415    all: HashMap<TxHash, ReadyTransaction<T>>,
416    awaiting: HashMap<TxHash, (usize, PoolTransactionRef<T>)>,
417    independent: BTreeSet<PoolTransactionRef<T>>,
418    _invalid: HashSet<TxHash>,
419}
420
421impl<T> TransactionsIterator<T> {
422    /// Depending on number of satisfied requirements insert given ref
423    /// either to awaiting set or to best set.
424    fn independent_or_awaiting(&mut self, satisfied: usize, tx_ref: PoolTransactionRef<T>) {
425        if satisfied >= tx_ref.transaction.requires.len() {
426            // If we have satisfied all deps insert to best
427            self.independent.insert(tx_ref);
428        } else {
429            // otherwise we're still awaiting for some deps
430            self.awaiting.insert(tx_ref.transaction.hash(), (satisfied, tx_ref));
431        }
432    }
433}
434
435impl<T> Iterator for TransactionsIterator<T> {
436    type Item = Arc<PoolTransaction<T>>;
437
438    fn next(&mut self) -> Option<Self::Item> {
439        loop {
440            let best = self.independent.iter().next_back()?.clone();
441            let best = self.independent.take(&best)?;
442            let hash = best.transaction.hash();
443
444            let ready =
445                if let Some(ready) = self.all.get(&hash).cloned() { ready } else { continue };
446
447            // Insert transactions that just got unlocked.
448            for hash in &ready.unlocks {
449                // first check local awaiting transactions
450                let res = if let Some((mut satisfied, tx_ref)) = self.awaiting.remove(hash) {
451                    satisfied += 1;
452                    Some((satisfied, tx_ref))
453                    // then get from the pool
454                } else {
455                    self.all
456                        .get(hash)
457                        .map(|next| (next.requires_offset + 1, next.transaction.clone()))
458                };
459                if let Some((satisfied, tx_ref)) = res {
460                    self.independent_or_awaiting(satisfied, tx_ref)
461                }
462            }
463
464            return Some(best.transaction);
465        }
466    }
467}
468
469/// transactions that are ready to be included in a block.
470#[derive(Clone, Debug)]
471pub struct ReadyTransactions<T> {
472    /// keeps track of transactions inserted in the pool
473    ///
474    /// this way we can determine when transactions where submitted to the pool
475    id: u64,
476    /// markers that are provided by `ReadyTransaction`s
477    provided_markers: HashMap<TxMarker, TxHash>,
478    /// transactions that are ready
479    ready_tx: Arc<RwLock<HashMap<TxHash, ReadyTransaction<T>>>>,
480    /// independent transactions that can be included directly and don't require other transactions
481    /// Sorted by their id
482    independent_transactions: BTreeSet<PoolTransactionRef<T>>,
483}
484
485impl<T> Default for ReadyTransactions<T> {
486    fn default() -> Self {
487        Self {
488            id: 0,
489            provided_markers: Default::default(),
490            ready_tx: Default::default(),
491            independent_transactions: Default::default(),
492        }
493    }
494}
495
496impl<T> ReadyTransactions<T> {
497    /// Returns an independent snapshot of the ready transactions.
498    pub(super) fn snapshot(&self) -> Self {
499        Self {
500            id: self.id,
501            provided_markers: self.provided_markers.clone(),
502            ready_tx: Arc::new(RwLock::new(self.ready_tx.read().clone())),
503            independent_transactions: self.independent_transactions.clone(),
504        }
505    }
506
507    /// Returns an iterator over all transactions
508    pub fn get_transactions(&self) -> TransactionsIterator<T> {
509        TransactionsIterator {
510            all: self.ready_tx.read().clone(),
511            independent: self.independent_transactions.clone(),
512            awaiting: Default::default(),
513            _invalid: Default::default(),
514        }
515    }
516
517    /// Clears the internal state
518    pub fn clear(&mut self) {
519        self.provided_markers.clear();
520        self.ready_tx.write().clear();
521        self.independent_transactions.clear();
522    }
523
524    /// Returns true if the transaction is part of the queue.
525    pub fn contains(&self, hash: &TxHash) -> bool {
526        self.ready_tx.read().contains_key(hash)
527    }
528
529    /// Returns the number of ready transactions without cloning the snapshot
530    pub fn len(&self) -> usize {
531        self.ready_tx.read().len()
532    }
533
534    /// Returns true if there are no ready transactions
535    pub fn is_empty(&self) -> bool {
536        self.ready_tx.read().is_empty()
537    }
538
539    /// Returns the transaction for the hash if it's in the ready pool but not yet mined
540    pub fn get(&self, hash: &TxHash) -> Option<ReadyTransaction<T>> {
541        self.ready_tx.read().get(hash).cloned()
542    }
543
544    pub const fn provided_markers(&self) -> &HashMap<TxMarker, TxHash> {
545        &self.provided_markers
546    }
547
548    const fn next_id(&mut self) -> u64 {
549        let id = self.id;
550        self.id = self.id.wrapping_add(1);
551        id
552    }
553
554    /// Removes the transactions from the ready queue and returns the removed transactions.
555    /// This will also remove all transactions that depend on those.
556    pub fn clear_transactions(&mut self, tx_hashes: &[TxHash]) -> Vec<Arc<PoolTransaction<T>>> {
557        self.remove_with_markers(tx_hashes.to_vec(), None)
558    }
559
560    /// Removes the transactions that provide the marker
561    ///
562    /// This will also remove all transactions that lead to the transaction that provides the
563    /// marker.
564    pub fn prune_tags(&mut self, marker: TxMarker) -> Vec<Arc<PoolTransaction<T>>> {
565        let mut removed_tx = vec![];
566
567        // the markers to remove
568        let mut remove = vec![marker];
569
570        while let Some(marker) = remove.pop() {
571            let res = self
572                .provided_markers
573                .remove(&marker)
574                .and_then(|hash| self.ready_tx.write().remove(&hash));
575
576            if let Some(tx) = res {
577                let unlocks = tx.unlocks;
578                self.independent_transactions.remove(&tx.transaction);
579                let tx = tx.transaction.transaction;
580
581                // also prune previous transactions
582                {
583                    let hash = tx.hash();
584                    let mut ready = self.ready_tx.write();
585
586                    let mut previous_markers = |marker| -> Option<Vec<TxMarker>> {
587                        let prev_hash = self.provided_markers.get(marker)?;
588                        let tx2 = ready.get_mut(prev_hash)?;
589                        // remove hash
590                        if let Some(idx) = tx2.unlocks.iter().position(|i| i == &hash) {
591                            tx2.unlocks.swap_remove(idx);
592                        }
593                        tx2.unlocks.is_empty().then(|| tx2.transaction.transaction.provides.clone())
594                    };
595
596                    // find previous transactions
597                    for marker in &tx.requires {
598                        if let Some(mut tags_to_remove) = previous_markers(marker) {
599                            remove.append(&mut tags_to_remove);
600                        }
601                    }
602                }
603
604                // add the transactions that just got unlocked to independent set
605                for hash in unlocks {
606                    if let Some(tx) = self.ready_tx.write().get_mut(&hash) {
607                        tx.requires_offset += 1;
608                        if tx.requires_offset == tx.transaction.transaction.requires.len() {
609                            self.independent_transactions.insert(tx.transaction.clone());
610                        }
611                    }
612                }
613                // finally, remove the markers that this transaction provides
614                let current_marker = &marker;
615                for marker in &tx.provides {
616                    let removed = self.provided_markers.remove(marker);
617                    assert_eq!(
618                        removed,
619                        if current_marker == marker { None } else { Some(tx.hash()) },
620                        "The pool contains exactly one transaction providing given tag; the removed transaction
621						claims to provide that tag, so it has to be mapped to it's hash; qed"
622                    );
623                }
624                removed_tx.push(tx);
625            }
626        }
627
628        removed_tx
629    }
630
631    /// Removes transactions and those that depend on them and satisfy at least one marker in the
632    /// given filter set.
633    pub fn remove_with_markers(
634        &mut self,
635        mut tx_hashes: Vec<TxHash>,
636        marker_filter: Option<HashSet<TxMarker>>,
637    ) -> Vec<Arc<PoolTransaction<T>>> {
638        let mut removed = Vec::new();
639        let mut ready = self.ready_tx.write();
640
641        while let Some(hash) = tx_hashes.pop() {
642            if let Some(mut tx) = ready.remove(&hash) {
643                let invalidated = tx.transaction.transaction.provides.iter().filter(|mark| {
644                    marker_filter.as_ref().is_none_or(|filter| !filter.contains(&**mark))
645                });
646
647                let mut removed_some_marks = false;
648                // remove entries from provided_markers
649                for mark in invalidated {
650                    removed_some_marks = true;
651                    self.provided_markers.remove(mark);
652                }
653
654                // remove from unlocks
655                for mark in &tx.transaction.transaction.requires {
656                    if let Some(provider_hash) = self.provided_markers.get(mark)
657                        && let Some(provider_tx) = ready.get_mut(provider_hash)
658                        && let Some(idx) = provider_tx.unlocks.iter().position(|i| i == &hash)
659                    {
660                        provider_tx.unlocks.swap_remove(idx);
661                    }
662                }
663
664                // remove from the independent set
665                self.independent_transactions.remove(&tx.transaction);
666
667                if removed_some_marks {
668                    // remove all transactions that the current one unlocks
669                    tx_hashes.append(&mut tx.unlocks);
670                }
671
672                // remove transaction
673                removed.push(tx.transaction.transaction);
674            }
675        }
676
677        removed
678    }
679}
680
681impl<T: Transaction> ReadyTransactions<T> {
682    /// Adds a new transactions to the ready queue.
683    ///
684    /// # Panics
685    ///
686    /// If the pending transaction is not ready ([`PendingPoolTransaction::is_ready`])
687    /// or the transaction is already included.
688    pub fn add_transaction(
689        &mut self,
690        tx: PendingPoolTransaction<T>,
691    ) -> Result<Vec<Arc<PoolTransaction<T>>>, PoolError> {
692        assert!(tx.is_ready(), "transaction must be ready",);
693        assert!(
694            !self.ready_tx.read().contains_key(&tx.transaction.hash()),
695            "transaction already included"
696        );
697
698        let (replaced_tx, unlocks) = self.replaced_transactions(&tx.transaction)?;
699
700        let id = self.next_id();
701        let hash = tx.transaction.hash();
702
703        let mut independent = true;
704        let mut requires_offset = 0;
705        let mut ready = self.ready_tx.write();
706        // Add links to transactions that unlock the current one
707        for mark in &tx.transaction.requires {
708            // Check if the transaction that satisfies the mark is still in the queue.
709            if let Some(other) = self.provided_markers.get(mark) {
710                let tx = ready.get_mut(other).expect("hash included;");
711                tx.unlocks.push(hash);
712                // tx still depends on other tx
713                independent = false;
714            } else {
715                requires_offset += 1;
716            }
717        }
718
719        // update markers
720        for mark in tx.transaction.provides.iter().cloned() {
721            self.provided_markers.insert(mark, hash);
722        }
723
724        let transaction = PoolTransactionRef { id, transaction: tx.transaction };
725
726        // add to the independent set
727        if independent {
728            self.independent_transactions.insert(transaction.clone());
729        }
730
731        // insert to ready queue
732        ready.insert(hash, ReadyTransaction { transaction, unlocks, requires_offset });
733
734        Ok(replaced_tx)
735    }
736
737    /// Removes and returns those transactions that got replaced by the `tx`
738    fn replaced_transactions(
739        &mut self,
740        tx: &PoolTransaction<T>,
741    ) -> Result<ReplacedTransactions<T>, PoolError> {
742        // check if we are replacing transactions
743        let remove_hashes: HashSet<_> =
744            tx.provides.iter().filter_map(|mark| self.provided_markers.get(mark)).collect();
745
746        // early exit if we are not replacing anything.
747        if remove_hashes.is_empty() {
748            return Ok((Vec::new(), Vec::new()));
749        }
750
751        // check if we're replacing the same transaction and if it can be replaced
752        let mut unlocked_tx = Vec::new();
753        {
754            // construct a list of unlocked transactions
755            // also check for transactions that shouldn't be replaced because underpriced
756            let ready = self.ready_tx.read();
757            for to_remove in remove_hashes.iter().filter_map(|hash| ready.get(*hash)) {
758                // if we're attempting to replace a transaction that provides the exact same markers
759                // (addr + nonce) then we check for gas price
760                if to_remove.provides() == tx.provides {
761                    // check if underpriced
762                    if tx.pending_transaction.transaction.max_fee_per_gas()
763                        <= to_remove.max_fee_per_gas()
764                    {
765                        warn!(target: "txpool", "ready replacement transaction underpriced [{:?}]", tx.hash());
766                        return Err(PoolError::ReplacementUnderpriced(tx.hash()));
767                    }
768                    trace!(target: "txpool", "replacing ready transaction [{:?}] with higher gas price [{:?}]", to_remove.transaction.transaction.hash(), tx.hash());
769                }
770
771                unlocked_tx.extend(to_remove.unlocks.iter().copied())
772            }
773        }
774
775        let remove_hashes = remove_hashes.into_iter().copied().collect::<Vec<_>>();
776
777        let new_provides = tx.provides.iter().cloned().collect::<HashSet<_>>();
778        let removed_tx = self.remove_with_markers(remove_hashes, Some(new_provides));
779
780        Ok((removed_tx, unlocked_tx))
781    }
782}
783
784/// A reference to a transaction in the pool
785#[derive(Debug)]
786pub struct PoolTransactionRef<T> {
787    /// actual transaction
788    pub transaction: Arc<PoolTransaction<T>>,
789    /// identifier used to internally compare the transaction in the pool
790    pub id: u64,
791}
792
793impl<T> Clone for PoolTransactionRef<T> {
794    fn clone(&self) -> Self {
795        Self { transaction: Arc::clone(&self.transaction), id: self.id }
796    }
797}
798
799impl<T> Eq for PoolTransactionRef<T> {}
800
801impl<T> PartialEq<Self> for PoolTransactionRef<T> {
802    fn eq(&self, other: &Self) -> bool {
803        self.cmp(other) == Ordering::Equal
804    }
805}
806
807impl<T> PartialOrd<Self> for PoolTransactionRef<T> {
808    fn partial_cmp(&self, other: &Self) -> Option<Ordering> {
809        Some(self.cmp(other))
810    }
811}
812
813impl<T> Ord for PoolTransactionRef<T> {
814    fn cmp(&self, other: &Self) -> Ordering {
815        self.transaction
816            .priority
817            .cmp(&other.transaction.priority)
818            .then_with(|| other.id.cmp(&self.id))
819    }
820}
821
822#[derive(Debug)]
823pub struct ReadyTransaction<T> {
824    /// ref to the actual transaction
825    pub transaction: PoolTransactionRef<T>,
826    /// tracks the transactions that get unlocked by this transaction
827    pub unlocks: Vec<TxHash>,
828    /// amount of required markers that are inherently provided
829    pub requires_offset: usize,
830}
831
832impl<T> Clone for ReadyTransaction<T> {
833    fn clone(&self) -> Self {
834        Self {
835            transaction: self.transaction.clone(),
836            unlocks: self.unlocks.clone(),
837            requires_offset: self.requires_offset,
838        }
839    }
840}
841
842impl<T> ReadyTransaction<T> {
843    pub fn provides(&self) -> &[TxMarker] {
844        &self.transaction.transaction.provides
845    }
846}
847
848impl<T: Transaction> ReadyTransaction<T> {
849    pub fn max_fee_per_gas(&self) -> u128 {
850        self.transaction.transaction.max_fee_per_gas()
851    }
852}
853
854/// creates an unique identifier for aan (`nonce` + `Address`) combo
855pub fn to_marker(nonce: u64, from: Address) -> TxMarker {
856    let mut data = [0u8; 28];
857    data[..8].copy_from_slice(&nonce.to_le_bytes()[..]);
858    data[8..].copy_from_slice(&from.0[..]);
859    data.to_vec()
860}
861
862#[cfg(test)]
863mod tests {
864    use super::*;
865
866    #[test]
867    fn can_id_txs() {
868        let addr = Address::random();
869        assert_eq!(to_marker(1, addr), to_marker(1, addr));
870        assert_ne!(to_marker(2, addr), to_marker(1, addr));
871    }
872}