Skip to main content

reth_network/transactions/
mod.rs

1//! Transactions management for the p2p network.
2
3use alloy_consensus::{constants::EIP4844_TX_TYPE_ID, transaction::TxHashRef};
4use rayon::iter::{IntoParallelIterator, ParallelIterator};
5use smallvec::SmallVec;
6
7/// Normalized transaction announcements.
8pub mod announcement;
9
10/// Aggregation on configurable parameters for [`TransactionsManager`].
11pub mod config;
12/// Default and spec'd bounds.
13pub mod constants;
14/// Component responsible for fetching transactions from [`NewPooledTransactionHashes`].
15pub mod fetcher;
16/// Defines the traits for transaction-related policies.
17pub mod policy;
18#[cfg(test)]
19mod test_harness;
20
21pub use self::constants::{
22    tx_fetcher::DEFAULT_SOFT_LIMIT_BYTE_SIZE_POOLED_TRANSACTIONS_RESP_ON_PACK_GET_POOLED_TRANSACTIONS_REQ,
23    SOFT_LIMIT_BYTE_SIZE_POOLED_TRANSACTIONS_RESPONSE,
24};
25use announcement::TransactionAnnouncement;
26use config::AnnouncementAcceptance;
27pub use config::{
28    AnnouncementFilteringPolicy, TransactionFetcherConfig, TransactionIngressPolicy,
29    TransactionPropagationMode, TransactionPropagationPolicy, TransactionsManagerConfig,
30};
31use policy::NetworkPolicies;
32
33pub(crate) use fetcher::{FetchEvent, TransactionFetcher};
34
35use self::constants::{
36    tx_manager::*, DEFAULT_SOFT_LIMIT_BYTE_SIZE_TRANSACTIONS_BROADCAST_MESSAGE,
37    SOFT_LIMIT_COUNT_HASHES_IN_GET_POOLED_TRANSACTIONS_REQUEST,
38};
39use crate::{
40    budget::{
41        DEFAULT_BUDGET_TRY_DRAIN_NETWORK_TRANSACTION_EVENTS,
42        DEFAULT_BUDGET_TRY_DRAIN_PENDING_POOL_IMPORTS, DEFAULT_BUDGET_TRY_DRAIN_STREAM,
43    },
44    cache::LruCache,
45    duration_metered_exec, metered_poll_nested_stream_with_budget,
46    metrics::{AnnouncedTxTypesMetrics, TransactionsManagerMetrics},
47    transactions::config::{StrictEthAnnouncementFilter, TransactionPropagationKind},
48    NetworkHandle, TxTypesCounter,
49};
50use alloy_eips::eip2718::Typed2718;
51use alloy_primitives::{
52    bytes::BufMut,
53    map::{hash_map::Entry, B256Map, B256Set, FbBuildHasher, HashMap, HashSet},
54    TxHash, B256,
55};
56use alloy_rlp::Encodable;
57use constants::SOFT_LIMIT_COUNT_HASHES_IN_NEW_POOLED_TRANSACTIONS_BROADCAST_MESSAGE;
58use futures::{stream::FuturesUnordered, Future, StreamExt};
59use reth_eth_wire::{
60    BroadcastPoolTransactions, Cells, EthNetworkPrimitives, EthVersion, GetCells,
61    GetPooledTransactions, LazyEncoded, LazyEncodedTransaction, NetworkPrimitives,
62    NewPooledTransactionHashes, NewPooledTransactionHashes66, NewPooledTransactionHashes68,
63    NewPooledTransactionHashes72, PooledTransactions, Transactions,
64};
65use reth_ethereum_primitives::TxType;
66use reth_evm::SenderRecoveryCache;
67use reth_metrics::common::mpsc::MemoryBoundedReceiver;
68use reth_network_api::{
69    events::{PeerEvent, SessionInfo},
70    NetworkEvent, NetworkEventListenerProvider, PeerKind, PeerRequest, PeerRequestSender, Peers,
71};
72use reth_network_p2p::{
73    error::{RequestError, RequestResult},
74    sync::SyncStateProvider,
75};
76use reth_network_peers::PeerId;
77use reth_network_types::ReputationChangeKind;
78use reth_primitives_traits::{InMemorySize, SignedTransaction};
79use reth_tokio_util::EventStream;
80use reth_transaction_pool::{
81    error::{PoolError, PoolResult},
82    AddedTransactionOutcome, GetPooledTransactionLimit, PoolTransaction, PropagateKind,
83    PropagatedTransactions, TransactionPool, ValidPoolTransaction,
84};
85use std::{
86    pin::Pin,
87    sync::{
88        atomic::{AtomicUsize, Ordering},
89        Arc,
90    },
91    task::{Context, Poll},
92    time::{Duration, Instant},
93};
94use tokio::sync::{mpsc, oneshot, oneshot::error::RecvError};
95use tokio_stream::wrappers::UnboundedReceiverStream;
96use tracing::{debug, trace};
97
98/// The future for importing transactions into the pool.
99///
100/// Resolves with the result of each transaction import.
101pub type PoolImportFuture =
102    Pin<Box<dyn Future<Output = Vec<PoolResult<AddedTransactionOutcome>>> + Send + 'static>>;
103
104/// Api to interact with [`TransactionsManager`] task.
105///
106/// This can be obtained via [`TransactionsManager::handle`] and can be used to manually interact
107/// with the [`TransactionsManager`] task once it is spawned.
108///
109/// For example [`TransactionsHandle::get_peer_transaction_hashes`] returns the transaction hashes
110/// known by a specific peer.
111#[derive(Debug, Clone)]
112pub struct TransactionsHandle<N: NetworkPrimitives = EthNetworkPrimitives> {
113    /// Command channel to the [`TransactionsManager`]
114    manager_tx: mpsc::UnboundedSender<TransactionsCommand<N>>,
115}
116
117impl<N: NetworkPrimitives> TransactionsHandle<N> {
118    fn send(&self, cmd: TransactionsCommand<N>) {
119        let _ = self.manager_tx.send(cmd);
120    }
121
122    /// Fetch the [`PeerRequestSender`] for the given peer.
123    async fn peer_handle(
124        &self,
125        peer_id: PeerId,
126    ) -> Result<Option<PeerRequestSender<PeerRequest<N>>>, RecvError> {
127        let (tx, rx) = oneshot::channel();
128        self.send(TransactionsCommand::GetPeerSender { peer_id, peer_request_sender: tx });
129        rx.await
130    }
131
132    /// Manually propagate the transaction that belongs to the hash.
133    pub fn propagate(&self, hash: TxHash) {
134        self.send(TransactionsCommand::PropagateHash(hash))
135    }
136
137    /// Manually propagate the transaction hash to a specific peer.
138    ///
139    /// Note: this only propagates if the pool contains the transaction.
140    pub fn propagate_hash_to(&self, hash: TxHash, peer: PeerId) {
141        self.propagate_hashes_to(Some(hash), peer)
142    }
143
144    /// Manually propagate the transaction hashes to a specific peer.
145    ///
146    /// Note: this only propagates the transactions that are known to the pool.
147    pub fn propagate_hashes_to(&self, hash: impl IntoIterator<Item = TxHash>, peer: PeerId) {
148        let hashes = hash.into_iter().collect::<Vec<_>>();
149        if hashes.is_empty() {
150            return
151        }
152        self.send(TransactionsCommand::PropagateHashesTo(hashes, peer))
153    }
154
155    /// Request the active peer IDs from the [`TransactionsManager`].
156    pub async fn get_active_peers(&self) -> Result<HashSet<PeerId>, RecvError> {
157        let (tx, rx) = oneshot::channel();
158        self.send(TransactionsCommand::GetActivePeers(tx));
159        rx.await
160    }
161
162    /// Manually propagate full transaction hashes to a specific peer.
163    ///
164    /// Do nothing if transactions are empty.
165    pub fn propagate_transactions_to(&self, transactions: Vec<TxHash>, peer: PeerId) {
166        if transactions.is_empty() {
167            return
168        }
169        self.send(TransactionsCommand::PropagateTransactionsTo(transactions, peer))
170    }
171
172    /// Manually propagate the given transaction hashes to all peers.
173    ///
174    /// It's up to the [`TransactionsManager`] whether the transactions are sent as hashes or in
175    /// full.
176    pub fn propagate_transactions(&self, transactions: Vec<TxHash>) {
177        if transactions.is_empty() {
178            return
179        }
180        self.send(TransactionsCommand::PropagateTransactions(transactions))
181    }
182
183    /// Manually propagate the given transactions to all peers.
184    ///
185    /// It's up to the [`TransactionsManager`] whether the transactions are sent as hashes or in
186    /// full.
187    pub fn broadcast_transactions(
188        &self,
189        transactions: impl IntoIterator<Item = N::BroadcastedTransaction>,
190    ) {
191        let transactions =
192            transactions.into_iter().map(PropagateTransaction::new).collect::<Vec<_>>();
193        if transactions.is_empty() {
194            return
195        }
196        self.send(TransactionsCommand::BroadcastTransactions(transactions))
197    }
198
199    /// Request the transaction hashes known by specific peers.
200    pub async fn get_transaction_hashes(
201        &self,
202        peers: Vec<PeerId>,
203    ) -> Result<HashMap<PeerId, B256Set>, RecvError> {
204        if peers.is_empty() {
205            return Ok(Default::default())
206        }
207        let (tx, rx) = oneshot::channel();
208        self.send(TransactionsCommand::GetTransactionHashes { peers, tx });
209        rx.await
210    }
211
212    /// Request the transaction hashes known by a specific peer.
213    pub async fn get_peer_transaction_hashes(&self, peer: PeerId) -> Result<B256Set, RecvError> {
214        let res = self.get_transaction_hashes(vec![peer]).await?;
215        Ok(res.into_values().next().unwrap_or_default())
216    }
217
218    /// Requests the transactions directly from the given peer.
219    ///
220    /// Returns `None` if the peer is not connected.
221    ///
222    /// **Note**: this returns the response from the peer as received.
223    pub async fn get_pooled_transactions_from(
224        &self,
225        peer_id: PeerId,
226        hashes: Vec<B256>,
227    ) -> Result<Option<Vec<N::PooledTransaction>>, RequestError> {
228        let Some(peer) = self.peer_handle(peer_id).await? else { return Ok(None) };
229
230        let (tx, rx) = oneshot::channel();
231        let request = PeerRequest::GetPooledTransactions { request: hashes.into(), response: tx };
232        peer.try_send(request).ok();
233
234        rx.await?.map(|res| Some(res.0))
235    }
236
237    /// Requests cells from a specific eth/72 peer.
238    ///
239    /// This is the transport-level client primitive for EIP-8070. The caller supplies the
240    /// transaction hashes and the common cell mask for the request. Announcement-driven batching,
241    /// retries, cell verification, and blob reconstruction are layered above this method.
242    ///
243    /// Returns `None` if the peer is no longer connected.
244    pub async fn get_cells_from(
245        &self,
246        peer_id: PeerId,
247        hashes: Vec<B256>,
248        cell_mask: alloy_primitives::B128,
249    ) -> Result<Option<Cells>, RequestError> {
250        if hashes.is_empty() {
251            return Ok(Some(Cells { cell_mask, ..Default::default() }))
252        }
253
254        let Some(peer) = self.peer_handle(peer_id).await? else { return Ok(None) };
255
256        let (tx, rx) = oneshot::channel();
257        let request =
258            PeerRequest::GetCells { request: GetCells { hashes, cell_mask }, response: tx };
259        peer.try_send(request).map_err(|_| RequestError::ChannelClosed)?;
260
261        rx.await?.map(Some)
262    }
263}
264
265/// Manages transactions on top of the p2p network.
266///
267/// This can be spawned to another task and is supposed to be run as background service.
268/// [`TransactionsHandle`] can be used as frontend to programmatically send commands to it and
269/// interact with it.
270///
271/// The [`TransactionsManager`] is responsible for:
272///    - handling incoming eth messages for transactions.
273///    - serving transaction requests.
274///    - propagate transactions
275///
276/// This type communicates with the [`NetworkManager`](crate::NetworkManager) in both directions.
277///   - receives incoming network messages.
278///   - sends messages to dispatch (responses, propagate tx)
279///
280/// It is directly connected to the [`TransactionPool`] to retrieve requested transactions and
281/// propagate new transactions over the network.
282///
283/// It can be configured with different policies for transaction propagation and announcement
284/// filtering. See [`NetworkPolicies`] for more details.
285///
286/// ## Network Transaction Processing
287///
288/// ### Message Types
289///
290/// - **`Transactions`**: Full transaction broadcasts (rejects blob transactions)
291/// - **`NewPooledTransactionHashes`**: Hash announcements
292///
293/// ### Peer Tracking
294///
295/// - Maintains per-peer transaction cache (default: 10,240 entries)
296/// - Prevents duplicate imports and enables efficient propagation
297///
298/// ### Bad Transaction Handling
299///
300/// Caches and rejects transactions with consensus violations (gas, signature, chain ID).
301/// Penalizes peers sending invalid transactions.
302///
303/// ### Import Management
304///
305/// Limits concurrent pool imports and backs off when approaching capacity.
306///
307/// ### Transaction Fetching
308///
309/// For announced transactions: filters known → queues unknown → fetches → imports
310///
311/// ### Propagation Rules
312///
313/// Based on: origin (Local/External/Private), peer capabilities, and network state.
314/// Disabled during initial sync.
315///
316/// ### Security
317///
318/// Rate limiting via reputation, bad transaction isolation, peer scoring.
319#[derive(Debug)]
320#[must_use = "Manager does nothing unless polled."]
321pub struct TransactionsManager<Pool, N: NetworkPrimitives = EthNetworkPrimitives> {
322    /// Access to the transaction pool.
323    pool: Pool,
324    /// Cache of recovered transaction senders shared with payload execution, if enabled.
325    sender_recovery_cache: Option<SenderRecoveryCache>,
326    /// Network access.
327    network: NetworkHandle<N>,
328    /// Subscriptions to all network related events.
329    ///
330    /// From which we get all new incoming transaction related messages.
331    network_events: EventStream<NetworkEvent<PeerRequest<N>>>,
332    /// Transaction fetcher to handle inflight and missing transaction requests.
333    transaction_fetcher: TransactionFetcher<N>,
334    /// All currently pending transactions grouped by peers.
335    ///
336    /// This way we can track incoming transactions and prevent multiple pool imports for the same
337    /// transaction
338    transactions_by_peers: B256Map<SmallVec<[PeerId; 1]>>,
339    /// Scratch hashes reused when deduplicating announcements.
340    announcement_hashes: B256Set,
341    /// Transactions that are currently imported into the `Pool`.
342    ///
343    /// The import process includes:
344    ///  - validation of the transactions, e.g. transaction is well formed: valid tx type, fees are
345    ///    valid, or for 4844 transaction the blobs are valid. See also
346    ///    [`EthTransactionValidator`](reth_transaction_pool::validate::EthTransactionValidator)
347    /// - if the transaction is valid, it is added into the pool.
348    ///
349    /// Once the new transaction reaches the __pending__ state it will be emitted by the pool via
350    /// [`TransactionPool::pending_transactions_listener`] and arrive at the `pending_transactions`
351    /// receiver.
352    pool_imports: FuturesUnordered<PoolImportFuture>,
353    /// Stats on pending pool imports that help the node self-monitor.
354    pending_pool_imports_info: PendingPoolImportsInfo,
355    /// Bad imports.
356    bad_imports: LruCache<TxHash, FbBuildHasher<32>>,
357    /// All the connected peers.
358    peers: HashMap<PeerId, PeerMetadata<N>, FbBuildHasher<64>>,
359    /// Send half for the command channel.
360    ///
361    /// This is kept so that a new [`TransactionsHandle`] can be created at any time.
362    command_tx: mpsc::UnboundedSender<TransactionsCommand<N>>,
363    /// Incoming commands from [`TransactionsHandle`].
364    ///
365    /// This will only receive commands if a user manually sends a command to the manager through
366    /// the [`TransactionsHandle`] to interact with this type directly.
367    command_rx: UnboundedReceiverStream<TransactionsCommand<N>>,
368    /// A stream that yields new __pending__ transactions.
369    ///
370    /// A transaction is considered __pending__ if it is executable on the current state of the
371    /// chain. In other words, this only yields transactions that satisfy all consensus
372    /// requirements, these include:
373    ///   - no nonce gaps
374    ///   - all dynamic fee requirements are (currently) met
375    ///   - account has enough balance to cover the transaction's gas
376    pending_transactions: mpsc::Receiver<TxHash>,
377    /// Incoming events from the [`NetworkManager`](crate::NetworkManager).
378    transaction_events: MemoryBoundedReceiver<NetworkTransactionEvent<N>>,
379    /// How the `TransactionsManager` is configured.
380    config: TransactionsManagerConfig,
381    /// Network Policies
382    policies: NetworkPolicies<N>,
383    /// `TransactionsManager` metrics
384    metrics: TransactionsManagerMetrics,
385    /// `AnnouncedTxTypes` metrics
386    announced_tx_types_metrics: AnnouncedTxTypesMetrics,
387}
388
389impl<Pool: TransactionPool, N: NetworkPrimitives> TransactionsManager<Pool, N> {
390    /// Sets up a new instance.
391    ///
392    /// Note: This expects an existing [`NetworkManager`](crate::NetworkManager) instance.
393    pub fn new(
394        network: NetworkHandle<N>,
395        pool: Pool,
396        from_network: MemoryBoundedReceiver<NetworkTransactionEvent<N>>,
397        transactions_manager_config: TransactionsManagerConfig,
398    ) -> Self {
399        Self::with_policy(
400            network,
401            pool,
402            from_network,
403            transactions_manager_config,
404            NetworkPolicies::new(
405                TransactionPropagationKind::default(),
406                StrictEthAnnouncementFilter::default(),
407            ),
408        )
409    }
410}
411
412impl<Pool: TransactionPool, N: NetworkPrimitives> TransactionsManager<Pool, N> {
413    /// Sets up a new instance with given the settings.
414    ///
415    /// Note: This expects an existing [`NetworkManager`](crate::NetworkManager) instance.
416    pub fn with_policy(
417        network: NetworkHandle<N>,
418        pool: Pool,
419        from_network: MemoryBoundedReceiver<NetworkTransactionEvent<N>>,
420        transactions_manager_config: TransactionsManagerConfig,
421        policies: NetworkPolicies<N>,
422    ) -> Self {
423        let network_events = network.event_listener();
424
425        let (command_tx, command_rx) = mpsc::unbounded_channel();
426
427        let transaction_fetcher =
428            TransactionFetcher::new(transactions_manager_config.transaction_fetcher_config.clone());
429
430        // install a listener for new __pending__ transactions that are allowed to be propagated
431        // over the network
432        let pending = pool.pending_transactions_listener();
433        let pending_pool_imports_info =
434            PendingPoolImportsInfo::new(transactions_manager_config.max_pending_pool_imports);
435        let metrics = TransactionsManagerMetrics::default();
436        metrics
437            .capacity_pending_pool_imports
438            .increment(pending_pool_imports_info.max_pending_pool_imports as u64);
439
440        Self {
441            pool,
442            sender_recovery_cache: None,
443            network,
444            network_events,
445            transaction_fetcher,
446            transactions_by_peers: Default::default(),
447            announcement_hashes: B256Set::default(),
448            pool_imports: Default::default(),
449            pending_pool_imports_info,
450            bad_imports: LruCache::with_hasher(DEFAULT_MAX_COUNT_BAD_IMPORTS, Default::default()),
451            peers: Default::default(),
452            command_tx,
453            command_rx: UnboundedReceiverStream::new(command_rx),
454            pending_transactions: pending,
455            transaction_events: from_network,
456            config: transactions_manager_config,
457            policies,
458            metrics,
459            announced_tx_types_metrics: AnnouncedTxTypesMetrics::default(),
460        }
461    }
462
463    /// Returns a new handle that can send commands to this type.
464    pub fn handle(&self) -> TransactionsHandle<N> {
465        TransactionsHandle { manager_tx: self.command_tx.clone() }
466    }
467
468    /// Uses the provided sender recovery cache.
469    pub fn with_sender_recovery_cache(mut self, cache: SenderRecoveryCache) -> Self {
470        self.sender_recovery_cache = Some(cache);
471        self
472    }
473
474    /// Returns whether pending pool imports are below the soft limit.
475    ///
476    /// Fetch responses are admitted in full when capacity remains, so one response can exceed
477    /// the limit by at most 255 transactions (the 256-hash request cap minus one free slot).
478    /// Further responses and requests wait until imports fall below the limit again.
479    fn has_capacity_for_pending_pool_imports(&self) -> bool {
480        self.remaining_pool_import_capacity() > 0
481    }
482
483    /// Returns the remaining capacity below the soft limit for pending pool imports.
484    ///
485    /// Outstanding requests do not reserve capacity. A fetched response is admitted in full
486    /// if any capacity remains; while imports exceed the limit, this returns zero.
487    fn remaining_pool_import_capacity(&self) -> usize {
488        self.pending_pool_imports_info.max_pending_pool_imports.saturating_sub(
489            self.pending_pool_imports_info.pending_pool_imports.load(Ordering::Relaxed),
490        )
491    }
492
493    /// Returns import capacity available to broadcasts after reserving space for fetched hashes.
494    fn remaining_broadcast_import_capacity(&self) -> usize {
495        let capacity = self.remaining_pool_import_capacity();
496        if self.transaction_fetcher.num_hashes() == 0 {
497            return capacity
498        }
499        // Reserve one batch (or half a small import budget) for fetched transactions.
500        // Broadcasts are processed first and must not consume every newly freed slot.
501        let reserved = self
502            .pending_pool_imports_info
503            .max_pending_pool_imports
504            .div_ceil(2)
505            .min(SOFT_LIMIT_COUNT_HASHES_IN_GET_POOLED_TRANSACTIONS_REQUEST);
506        capacity.saturating_sub(reserved)
507    }
508
509    fn report_peer_bad_transactions(&self, peer_id: PeerId) {
510        self.report_peer(peer_id, ReputationChangeKind::BadTransactions);
511        self.metrics.reported_bad_transactions.increment(1);
512    }
513
514    fn report_peer(&self, peer_id: PeerId, kind: ReputationChangeKind) {
515        trace!(target: "net::tx", ?peer_id, ?kind, "reporting reputation change");
516        self.network.reputation_change(peer_id, kind);
517    }
518
519    fn report_already_seen(&self, peer_id: PeerId) {
520        trace!(target: "net::tx", ?peer_id, "Penalizing peer for already seen transaction");
521        self.network.reputation_change(peer_id, ReputationChangeKind::AlreadySeenTransaction);
522    }
523
524    /// Handles a closed peer session, removing the peer from transaction-local tracking state.
525    fn on_peer_session_closed(&mut self, peer_id: &PeerId) {
526        if let Some(mut peer) = self.peers.remove(peer_id) {
527            self.policies.propagation_policy_mut().on_session_closed(&mut peer);
528        }
529        self.transaction_fetcher.on_peer_disconnected(peer_id);
530    }
531
532    /// Clear the transaction
533    fn on_good_import(&mut self, hash: TxHash) {
534        self.transactions_by_peers.remove(&hash);
535    }
536
537    /// Handles a failed transaction import.
538    ///
539    /// Blob sidecar errors (e.g. invalid proof, missing sidecar) are penalized via
540    /// `report_peer_bad_transactions` but NOT cached in `bad_imports` — the transaction itself
541    /// may be valid when fetched from another peer with correct sidecar data.
542    ///
543    /// Other bad transactions are penalized and cached in `bad_imports` to avoid fetching or
544    /// importing them again.
545    ///
546    /// Errors that count as bad transactions are:
547    ///
548    /// - intrinsic gas too low
549    /// - exceeds gas limit
550    /// - gas uint overflow
551    /// - exceeds max init code size
552    /// - oversized data
553    /// - signer account has bytecode
554    /// - chain id mismatch
555    /// - old legacy chain id
556    /// - tx type not supported
557    ///
558    /// (and additionally for blobs txns...)
559    ///
560    /// - no blobs
561    /// - too many blobs
562    /// - invalid kzg proof
563    /// - kzg error
564    /// - not blob transaction (tx type mismatch)
565    /// - wrong versioned kzg commitment hash
566    fn on_bad_import(&mut self, err: PoolError) {
567        let peers = self.transactions_by_peers.remove(&err.hash);
568
569        if err.is_bad_blob_sidecar() {
570            // Blob sidecar errors: penalize but do NOT cache the hash as bad.
571            // The transaction may be valid — only the sidecar from this peer was wrong.
572            // Using regular penalties means repeated offenders still get disconnected.
573            //
574            // Eth/72 peers cannot reach this path with an elided sidecar: their blob
575            // transactions are dropped before import in `import_transactions`, so any sidecar
576            // failure here means actual blob data failed validation.
577            if let Some(peers) = peers {
578                for peer_id in peers {
579                    self.report_peer_bad_transactions(peer_id);
580                }
581            }
582            return
583        }
584
585        // if we're _currently_ syncing, we ignore a bad transaction
586        if !err.is_bad_transaction() || self.network.is_syncing() {
587            return
588        }
589        // otherwise we penalize the peer that sent the bad transaction, with the assumption that
590        // the peer should have known that this transaction is bad (e.g. violating consensus rules)
591        if let Some(peers) = peers {
592            for peer_id in peers {
593                self.report_peer_bad_transactions(peer_id);
594            }
595        }
596        self.metrics.bad_imports.increment(1);
597        self.bad_imports.insert(err.hash);
598    }
599
600    fn on_request_error(&self, peer_id: PeerId, req_err: RequestError) {
601        let kind = match req_err {
602            RequestError::UnsupportedCapability => ReputationChangeKind::BadProtocol,
603            RequestError::Timeout => ReputationChangeKind::Timeout,
604            RequestError::ChannelClosed | RequestError::ConnectionDropped => {
605                // peer is already disconnected
606                return
607            }
608            RequestError::BadResponse => return self.report_peer_bad_transactions(peer_id),
609        };
610        self.report_peer(peer_id, kind);
611    }
612
613    #[inline]
614    fn update_poll_metrics(&self, start: Instant, poll_durations: TxManagerPollDurations) {
615        let metrics = &self.metrics;
616
617        let TxManagerPollDurations {
618            acc_network_events,
619            acc_pending_imports,
620            acc_tx_events,
621            acc_imported_txns,
622            acc_fetch_events,
623            acc_pending_fetch,
624            acc_cmds,
625        } = poll_durations;
626
627        // update metrics for whole poll function
628        metrics.duration_poll_tx_manager.set(start.elapsed().as_secs_f64());
629        // update metrics for nested expressions
630        metrics.acc_duration_poll_network_events.set(acc_network_events.as_secs_f64());
631        metrics.acc_duration_poll_pending_pool_imports.set(acc_pending_imports.as_secs_f64());
632        metrics.acc_duration_poll_transaction_events.set(acc_tx_events.as_secs_f64());
633        metrics.acc_duration_poll_imported_transactions.set(acc_imported_txns.as_secs_f64());
634        metrics.acc_duration_poll_fetch_events.set(acc_fetch_events.as_secs_f64());
635        metrics.acc_duration_fetch_pending_hashes.set(acc_pending_fetch.as_secs_f64());
636        metrics.acc_duration_poll_commands.set(acc_cmds.as_secs_f64());
637    }
638}
639
640impl<Pool: TransactionPool, N: NetworkPrimitives> TransactionsManager<Pool, N> {
641    /// Processes a batch import results.
642    fn on_batch_import_result(&mut self, batch_results: Vec<PoolResult<AddedTransactionOutcome>>) {
643        for res in batch_results {
644            match res {
645                Ok(AddedTransactionOutcome { hash, .. }) => {
646                    self.on_good_import(hash);
647                }
648                Err(err) => {
649                    self.on_bad_import(err);
650                }
651            }
652        }
653    }
654
655    /// Request handler for an incoming `NewPooledTransactionHashes`
656    fn on_new_pooled_transaction_hashes(
657        &mut self,
658        peer_id: PeerId,
659        msg: NewPooledTransactionHashes,
660    ) {
661        // If the node is initially syncing, ignore transactions
662        if self.network.is_initially_syncing() {
663            return
664        }
665        if self.network.tx_gossip_disabled() {
666            return
667        }
668
669        // get handle to peer's session, if the session is still active
670        let Some(peer) = self.peers.get_mut(&peer_id) else {
671            trace!(
672                peer_id = format!("{peer_id:#}"),
673                ?msg,
674                "discarding announcement from inactive peer"
675            );
676
677            return
678        };
679        let client = peer.client_version.clone();
680
681        if msg.is_empty() {
682            self.report_peer(peer_id, ReputationChangeKind::BadAnnouncement);
683            return
684        }
685
686        let Ok(mut announcement) =
687            TransactionAnnouncement::from_message(&msg, &mut self.announcement_hashes)
688        else {
689            self.report_peer(peer_id, ReputationChangeKind::BadAnnouncement);
690            return
691        };
692        let has_duplicates = announcement.len() != msg.len();
693
694        // keep track of the transactions the peer knows
695        let count_txns_already_seen_by_peer =
696            msg.iter_hashes().filter(|hash| !peer.seen_transactions.insert(**hash)).count();
697        drop(msg);
698        if count_txns_already_seen_by_peer > 0 {
699            // this may occur if transactions are sent or announced to a peer, at the same time as
700            // the peer sends/announces those hashes to us. this is because, marking
701            // txns as seen by a peer is done optimistically upon sending them to the
702            // peer.
703            self.metrics.messages_with_hashes_already_seen_by_peer.increment(1);
704            self.metrics
705                .occurrences_hash_already_seen_by_peer
706                .increment(count_txns_already_seen_by_peer as u64);
707
708            trace!(target: "net::tx",
709                %count_txns_already_seen_by_peer,
710                peer_id=format!("{peer_id:#}"),
711                ?client,
712                "Peer sent hashes that have already been marked as seen by peer"
713            );
714
715            self.report_already_seen(peer_id);
716        }
717
718        if has_duplicates {
719            self.report_peer(peer_id, ReputationChangeKind::BadAnnouncement);
720        }
721
722        // Exclude pending imports and known bad transactions before counting types and applying
723        // policy. These cheap checks shrink the batch before acquiring the pool lock.
724        let mut should_report_peer = false;
725        let mut tx_types_counter = TxTypesCounter::default();
726        let has_eth68_metadata = announcement.version().has_eth68_metadata();
727
728        announcement.retain(|tx| {
729            if self.transactions_by_peers.contains_key(&tx.hash) ||
730                self.bad_imports.contains(&tx.hash)
731            {
732                return false
733            }
734            let (ty_byte, size_val) = tx.metadata.map_or((0, 0), |metadata| {
735                match TxType::try_from(metadata.tx_type) {
736                    Ok(tx_type) => tx_types_counter.increase_by_tx_type(tx_type),
737                    Err(_) => tx_types_counter.increase_other(),
738                }
739                (metadata.tx_type, metadata.size)
740            });
741
742            let decision = self
743                .policies
744                .announcement_filter()
745                .decide_on_announcement(ty_byte, &tx.hash, size_val);
746
747            match decision {
748                AnnouncementAcceptance::Accept => true,
749                AnnouncementAcceptance::Ignore => false,
750                AnnouncementAcceptance::Reject { penalize_peer } => {
751                    if penalize_peer {
752                        should_report_peer = true;
753                    }
754                    false
755                }
756            }
757        });
758
759        if has_eth68_metadata {
760            self.announced_tx_types_metrics.update_eth68_announcement_metrics(tx_types_counter);
761        }
762
763        if should_report_peer {
764            self.report_peer(peer_id, ReputationChangeKind::BadAnnouncement);
765        }
766
767        // Filter out hashes already in the pool.
768        //
769        // known txns have already been successfully fetched or received over gossip.
770        //
771        // most hashes will be filtered out here since the mempool protocol is a gossip
772        // protocol, healthy peers will send many of the same hashes.
773        //
774        let hashes_count_pre_pool_filter = announcement.len();
775        self.pool.retain_unknown(&mut announcement);
776        if hashes_count_pre_pool_filter > announcement.len() {
777            let already_known_hashes_count = hashes_count_pre_pool_filter - announcement.len();
778            self.metrics
779                .occurrences_hashes_already_in_pool
780                .increment(already_known_hashes_count as u64);
781        }
782
783        if announcement.is_empty() {
784            // nothing to request
785            return
786        }
787
788        trace!(target: "net::tx::propagation",
789            peer_id=format!("{peer_id:#}"),
790            hashes_len=announcement.len(),
791            msg_version=%announcement.version(),
792            client_version=%client,
793            "received unknown hashes in announcement from peer"
794        );
795
796        // Queue hashes for fetching; requests are sent when the manager is polled.
797        self.transaction_fetcher.on_announcement(peer_id, announcement);
798    }
799}
800
801impl<Pool, N> TransactionsManager<Pool, N>
802where
803    Pool: TransactionPool + Unpin + 'static,
804    N: NetworkPrimitives<
805            BroadcastedTransaction: SignedTransaction,
806            PooledTransaction: SignedTransaction,
807        > + Unpin,
808    Pool::Transaction:
809        PoolTransaction<Consensus = N::BroadcastedTransaction, Pooled = N::PooledTransaction>,
810{
811    /// Invoked when transactions in the local mempool are considered __pending__.
812    ///
813    /// When a transaction in the local mempool is moved to the pending pool, we propagate them to
814    /// connected peers over network using the `Transactions` and `NewPooledTransactionHashes`
815    /// messages. The Transactions message relays complete transaction objects and is typically
816    /// sent to a small, random fraction of connected peers.
817    ///
818    /// All other peers receive a notification of the transaction hash and can request the
819    /// complete transaction object if it is unknown to them. The dissemination of complete
820    /// transactions to a fraction of peers usually ensures that all nodes receive the transaction
821    /// and won't need to request it.
822    fn on_new_pending_transactions(&mut self, hashes: Vec<TxHash>) {
823        // We intentionally do not gate this on initial sync.
824        // During initial sync we skip importing tx announcements from peers in
825        // `on_new_pooled_transaction_hashes`, so transactions reaching this path are local.
826        if self.network.tx_gossip_disabled() {
827            return
828        }
829
830        trace!(target: "net::tx", num_hashes=?hashes.len(), "Start propagating transactions");
831
832        self.propagate_all(hashes);
833    }
834
835    /// Propagate the full transactions to a specific peer.
836    ///
837    /// Returns the propagated transactions.
838    fn propagate_full_transactions_to_peer(
839        &mut self,
840        txs: Vec<TxHash>,
841        peer_id: PeerId,
842        propagation_mode: PropagationMode,
843    ) -> Option<PropagatedTransactions> {
844        let peer = self.peers.get_mut(&peer_id)?;
845        trace!(target: "net::tx", ?peer_id, "Propagating transactions to peer");
846        let mut propagated = PropagatedTransactions::default();
847
848        // filter all transactions unknown to the peer
849        let mut full_transactions = FullTransactionsBuilder::new(peer.version);
850
851        let to_propagate = self.pool.get_all(txs).into_iter().map(PropagateTransaction::pool_tx);
852
853        if propagation_mode.is_forced() {
854            // skip cache check if forced
855            full_transactions.extend(to_propagate);
856        } else {
857            // Iterate through the transactions to propagate and fill the hashes and full
858            // transaction
859            for tx in to_propagate {
860                if !peer.seen_transactions.contains(tx.tx_hash()) {
861                    // Only include if the peer hasn't seen the transaction
862                    full_transactions.push(&tx);
863                }
864            }
865        }
866
867        if full_transactions.is_empty() {
868            // nothing to propagate
869            return None
870        }
871
872        let PropagateTransactions { pooled, full } = full_transactions.build();
873
874        // send hashes if any
875        if let Some(new_pooled_hashes) = pooled {
876            for hash in new_pooled_hashes.iter_hashes().copied() {
877                propagated.record(hash, PropagateKind::Hash(peer_id));
878                // mark transaction as seen by peer
879                peer.seen_transactions.insert(hash);
880            }
881
882            // send hashes of transactions
883            self.network.send_transactions_hashes(peer_id, new_pooled_hashes);
884        }
885
886        // send full transactions, if any
887        if let Some(new_full_transactions) = full {
888            for hash in new_full_transactions.iter_hashes() {
889                propagated.record(*hash, PropagateKind::Full(peer_id));
890                // mark transaction as seen by peer
891                peer.seen_transactions.insert(*hash);
892            }
893
894            // send full transactions
895            self.network.send_broadcast_pool_transactions(peer_id, new_full_transactions);
896        }
897
898        // Update propagated transactions metrics
899        self.metrics.propagated_transactions.increment(propagated.len() as u64);
900
901        Some(propagated)
902    }
903
904    /// Propagate the transaction hashes to the given peer
905    ///
906    /// Note: This will only send the hashes for transactions that exist in the pool.
907    fn propagate_hashes_to(
908        &mut self,
909        hashes: Vec<TxHash>,
910        peer_id: PeerId,
911        propagation_mode: PropagationMode,
912    ) {
913        trace!(target: "net::tx", "Start propagating transactions as hashes");
914
915        // This fetches a transactions from the pool, including the blob transactions, which are
916        // only ever sent as hashes.
917        let propagated = {
918            let Some(peer) = self.peers.get_mut(&peer_id) else {
919                // no such peer
920                return
921            };
922
923            let to_propagate =
924                self.pool.get_all(hashes).into_iter().map(PropagateTransaction::pool_tx);
925
926            let mut propagated = PropagatedTransactions::default();
927
928            // check if transaction is known to peer
929            let mut hashes = PooledTransactionsHashesBuilder::new(peer.version);
930
931            if propagation_mode.is_forced() {
932                hashes.extend(to_propagate)
933            } else {
934                for tx in to_propagate {
935                    if !peer.seen_transactions.contains(tx.tx_hash()) {
936                        // Include if the peer hasn't seen it
937                        hashes.push(&tx);
938                    }
939                }
940            }
941
942            let new_pooled_hashes = hashes.build();
943
944            if new_pooled_hashes.is_empty() {
945                // nothing to propagate
946                return
947            }
948
949            for hash in new_pooled_hashes.iter_hashes().copied() {
950                propagated.record(hash, PropagateKind::Hash(peer_id));
951                peer.seen_transactions.insert(hash);
952            }
953
954            trace!(target: "net::tx::propagation", ?peer_id, ?new_pooled_hashes, "Propagating transactions to peer");
955
956            // send hashes of transactions
957            self.network.send_transactions_hashes(peer_id, new_pooled_hashes);
958
959            // Update propagated transactions metrics
960            self.metrics.propagated_transactions.increment(propagated.len() as u64);
961
962            propagated
963        };
964
965        // notify pool so events get fired
966        self.pool.on_propagated(propagated);
967    }
968
969    /// Propagate the transactions to all connected peers either as full objects or hashes.
970    ///
971    /// The message for new pooled hashes depends on the negotiated version of the stream.
972    /// See [`NewPooledTransactionHashes`]
973    ///
974    /// Note: EIP-4844 are disallowed from being broadcast in full and are only ever sent as hashes, see also <https://eips.ethereum.org/EIPS/eip-4844#networking>.
975    fn propagate_transactions(
976        &mut self,
977        to_propagate: Vec<PropagateTransaction>,
978        propagation_mode: PropagationMode,
979    ) -> PropagatedTransactions {
980        let mut propagated = PropagatedTransactions::default();
981        if self.network.tx_gossip_disabled() {
982            return propagated
983        }
984
985        // send full transactions to a set of the connected peers based on the configured mode
986        let max_num_full = self.config.propagation_mode.full_peer_count(self.peers.len());
987
988        // Note: Assuming ~random~ order due to random state of the peers map hasher
989        let mut num_full_peers = 0;
990        for (peer_id, peer) in &mut self.peers {
991            if !self.policies.propagation_policy().can_propagate(peer) {
992                // skip peers we should not propagate to
993                continue
994            }
995
996            // determine whether to send full tx objects or hashes.
997            let mut builder = if num_full_peers < max_num_full {
998                num_full_peers += 1;
999                PropagateTransactionsBuilder::full(peer.version, to_propagate.len())
1000            } else {
1001                PropagateTransactionsBuilder::pooled(peer.version, to_propagate.len())
1002            };
1003
1004            // Transactions are optimistically marked as seen by the peer when included in the
1005            // message, see `PeerMetadata::seen_transactions`.
1006            if propagation_mode.is_forced() {
1007                for tx in &to_propagate {
1008                    peer.seen_transactions.insert(*tx.tx_hash());
1009                    builder.push(tx);
1010                }
1011            } else {
1012                // Iterate through the transactions to propagate and fill the hashes and full
1013                // transaction lists, before deciding whether or not to send full transactions to
1014                // the peer.
1015                for tx in &to_propagate {
1016                    // Only include the transaction if the peer hasn't seen it yet
1017                    if peer.seen_transactions.insert(*tx.tx_hash()) {
1018                        builder.push(tx);
1019                    }
1020                }
1021            }
1022
1023            if builder.is_empty() {
1024                trace!(target: "net::tx", ?peer_id, "Nothing to propagate to peer; has seen all transactions");
1025                continue
1026            }
1027
1028            let PropagateTransactions { pooled, full } = builder.build();
1029
1030            // send hashes if any
1031            if let Some(mut new_pooled_hashes) = pooled {
1032                // Unhappy path: too many hashes for a single message. This should not happen
1033                // during regular propagation, which is capped at the soft limit per batch, and
1034                // is only reachable via manual propagation commands with oversized batches.
1035                if new_pooled_hashes.len() >
1036                    SOFT_LIMIT_COUNT_HASHES_IN_NEW_POOLED_TRANSACTIONS_BROADCAST_MESSAGE
1037                {
1038                    // hashes that exceed the limit are not sent, so they must not be tracked as
1039                    // seen by the peer
1040                    for hash in new_pooled_hashes
1041                        .iter_hashes()
1042                        .skip(SOFT_LIMIT_COUNT_HASHES_IN_NEW_POOLED_TRANSACTIONS_BROADCAST_MESSAGE)
1043                    {
1044                        peer.seen_transactions.remove(hash);
1045                    }
1046                    new_pooled_hashes.truncate(
1047                        SOFT_LIMIT_COUNT_HASHES_IN_NEW_POOLED_TRANSACTIONS_BROADCAST_MESSAGE,
1048                    );
1049                }
1050
1051                for hash in new_pooled_hashes.iter_hashes().copied() {
1052                    propagated.record(hash, PropagateKind::Hash(*peer_id));
1053                }
1054
1055                trace!(target: "net::tx", ?peer_id, num_txs=?new_pooled_hashes.len(), "Propagating tx hashes to peer");
1056
1057                // send hashes of transactions
1058                self.network.send_transactions_hashes(*peer_id, new_pooled_hashes);
1059            }
1060
1061            // send full transactions, if any
1062            if let Some(new_full_transactions) = full {
1063                for hash in new_full_transactions.iter_hashes() {
1064                    propagated.record(*hash, PropagateKind::Full(*peer_id));
1065                }
1066
1067                trace!(target: "net::tx", ?peer_id, num_txs=?new_full_transactions.len(), "Propagating full transactions to peer");
1068
1069                // send full transactions
1070                self.network.send_broadcast_pool_transactions(*peer_id, new_full_transactions);
1071            }
1072        }
1073
1074        // Update propagated transactions metrics
1075        self.metrics.propagated_transactions.increment(propagated.len() as u64);
1076
1077        propagated
1078    }
1079
1080    /// Propagates the given transactions to the peers
1081    ///
1082    /// This fetches all transaction from the pool, including the 4844 blob transactions but
1083    /// __without__ their sidecar, because 4844 transactions are only ever announced as hashes.
1084    fn propagate_all(&mut self, hashes: Vec<TxHash>) {
1085        if self.peers.is_empty() {
1086            // nothing to propagate
1087            return
1088        }
1089        let propagated = self.propagate_transactions(
1090            self.pool.get_all(hashes).into_iter().map(PropagateTransaction::pool_tx).collect(),
1091            PropagationMode::Basic,
1092        );
1093
1094        // notify pool so events get fired
1095        self.pool.on_propagated(propagated);
1096    }
1097
1098    /// Request handler for an incoming request for transactions
1099    fn on_get_pooled_transactions(
1100        &mut self,
1101        peer_id: PeerId,
1102        request: GetPooledTransactions,
1103        response: oneshot::Sender<RequestResult<PooledTransactions<N::PooledTransaction>>>,
1104    ) {
1105        // fast exit if gossip is disabled
1106        if self.network.tx_gossip_disabled() {
1107            let _ = response.send(Ok(PooledTransactions::default()));
1108            return
1109        }
1110        if let Some(peer) = self.peers.get_mut(&peer_id) {
1111            // don't serve pending transactions to peers the policy doesn't propagate to
1112            if !self.policies.propagation_policy().can_propagate(peer) {
1113                let _ = response.send(Ok(PooledTransactions::default()));
1114                return
1115            }
1116
1117            let transactions = self.pool.get_pooled_transaction_elements(
1118                request.0,
1119                GetPooledTransactionLimit::ResponseSizeSoftLimit(
1120                    self.config
1121                        .transaction_fetcher_config
1122                        .soft_limit_byte_size_pooled_transactions_response,
1123                ),
1124            );
1125            trace!(target: "net::tx::propagation", sent_txs=?transactions.iter().map(|tx| tx.tx_hash()), "Sending requested transactions to peer");
1126
1127            // we sent a response at which point we assume that the peer is aware of the
1128            // transactions
1129            peer.seen_transactions.extend(transactions.iter().map(|tx| *tx.tx_hash()));
1130
1131            let resp = PooledTransactions(transactions);
1132            let _ = response.send(Ok(resp));
1133        }
1134    }
1135
1136    /// Handles a command received from a detached [`TransactionsHandle`]
1137    fn on_command(&mut self, cmd: TransactionsCommand<N>) {
1138        match cmd {
1139            TransactionsCommand::PropagateHash(hash) => {
1140                self.on_new_pending_transactions(vec![hash])
1141            }
1142            TransactionsCommand::PropagateHashesTo(hashes, peer) => {
1143                self.propagate_hashes_to(hashes, peer, PropagationMode::Forced)
1144            }
1145            TransactionsCommand::GetActivePeers(tx) => {
1146                let peers = self.peers.keys().copied().collect::<HashSet<_>>();
1147                tx.send(peers).ok();
1148            }
1149            TransactionsCommand::PropagateTransactionsTo(txs, peer) => {
1150                if let Some(propagated) =
1151                    self.propagate_full_transactions_to_peer(txs, peer, PropagationMode::Forced)
1152                {
1153                    self.pool.on_propagated(propagated);
1154                }
1155            }
1156            TransactionsCommand::PropagateTransactions(txs) => self.propagate_all(txs),
1157            TransactionsCommand::BroadcastTransactions(txs) => {
1158                let propagated = self.propagate_transactions(txs, PropagationMode::Forced);
1159                self.pool.on_propagated(propagated);
1160            }
1161            TransactionsCommand::GetTransactionHashes { peers, tx } => {
1162                let mut res = HashMap::with_capacity_and_hasher(peers.len(), Default::default());
1163                for peer_id in peers {
1164                    let hashes = self
1165                        .peers
1166                        .get(&peer_id)
1167                        .map(|peer| peer.seen_transactions.iter().copied().collect::<B256Set>())
1168                        .unwrap_or_default();
1169                    res.insert(peer_id, hashes);
1170                }
1171                tx.send(res).ok();
1172            }
1173            TransactionsCommand::GetPeerSender { peer_id, peer_request_sender } => {
1174                let sender = self.peers.get(&peer_id).map(|peer| peer.request_tx.clone());
1175                peer_request_sender.send(sender).ok();
1176            }
1177        }
1178    }
1179
1180    /// Handles session establishment and peer transactions initialization.
1181    ///
1182    /// This is invoked when a new session is established.
1183    fn handle_peer_session(
1184        &mut self,
1185        info: SessionInfo,
1186        messages: PeerRequestSender<PeerRequest<N>>,
1187    ) {
1188        let SessionInfo { peer_id, client_version, version, .. } = info;
1189
1190        // Insert a new peer into the peerset.
1191        let peer = PeerMetadata::<N>::new(
1192            messages,
1193            version,
1194            client_version,
1195            self.config.max_transactions_seen_by_peer_history,
1196            info.peer_kind,
1197        );
1198        let peer = match self.peers.entry(peer_id) {
1199            Entry::Occupied(mut entry) => {
1200                entry.insert(peer);
1201                entry.into_mut()
1202            }
1203            Entry::Vacant(entry) => entry.insert(peer),
1204        };
1205
1206        self.policies.propagation_policy_mut().on_session_established(peer);
1207
1208        // Send a `NewPooledTransactionHashes` to the peer with up to
1209        // `SOFT_LIMIT_COUNT_HASHES_IN_NEW_POOLED_TRANSACTIONS_BROADCAST_MESSAGE`
1210        // transactions in the pool.
1211        if self.network.is_initially_syncing() || self.network.tx_gossip_disabled() {
1212            trace!(target: "net::tx", ?peer_id, "Skipping transaction broadcast: node syncing or gossip disabled");
1213            return
1214        }
1215
1216        // skip peers we should not propagate to
1217        if !self.policies.propagation_policy().can_propagate(peer) {
1218            return
1219        }
1220
1221        // Get transactions to broadcast
1222        let pooled_txs = self.pool.pooled_transactions_max(
1223            SOFT_LIMIT_COUNT_HASHES_IN_NEW_POOLED_TRANSACTIONS_BROADCAST_MESSAGE,
1224        );
1225        if pooled_txs.is_empty() {
1226            trace!(target: "net::tx", ?peer_id, "No transactions in the pool to broadcast");
1227            return;
1228        }
1229
1230        // Build and send transaction hashes message
1231        let mut msg_builder = PooledTransactionsHashesBuilder::new(version);
1232        for pooled_tx in pooled_txs {
1233            peer.seen_transactions.insert(*pooled_tx.hash());
1234            msg_builder.push_pooled(pooled_tx);
1235        }
1236
1237        debug!(target: "net::tx", ?peer_id, tx_count = msg_builder.len(), "Broadcasting transaction hashes");
1238        let msg = msg_builder.build();
1239        self.network.send_transactions_hashes(peer_id, msg);
1240    }
1241
1242    /// Handles a received event related to common network events.
1243    fn on_network_event(&mut self, event_result: NetworkEvent<PeerRequest<N>>) {
1244        match event_result {
1245            NetworkEvent::Peer(PeerEvent::SessionClosed { peer_id, .. }) => {
1246                self.on_peer_session_closed(&peer_id);
1247            }
1248            NetworkEvent::ActivePeerSession { info, messages } => {
1249                // process active peer session and broadcast available transaction from the pool
1250                self.handle_peer_session(info, messages);
1251            }
1252            NetworkEvent::Peer(PeerEvent::SessionEstablished(info)) => {
1253                let peer_id = info.peer_id;
1254                // get messages from existing peer
1255                let messages = match self.peers.get(&peer_id) {
1256                    Some(p) => p.request_tx.clone(),
1257                    None => {
1258                        debug!(target: "net::tx", ?peer_id, "No peer request sender found");
1259                        return;
1260                    }
1261                };
1262                self.handle_peer_session(info, messages);
1263            }
1264            _ => {}
1265        }
1266    }
1267
1268    /// Returns true if the ingress policy allows processing messages from the given peer.
1269    fn accepts_incoming_from(&self, peer_id: &PeerId) -> bool {
1270        if self.config.ingress_policy.allows_all() {
1271            return true;
1272        }
1273        let Some(peer) = self.peers.get(peer_id) else {
1274            return false;
1275        };
1276        self.config.ingress_policy.allows(peer.peer_kind())
1277    }
1278
1279    /// Handles dedicated transaction events related to the `eth` protocol.
1280    fn on_network_tx_event(&mut self, event: NetworkTransactionEvent<N>) {
1281        match event {
1282            NetworkTransactionEvent::IncomingTransactions { peer_id, msg } => {
1283                if !self.accepts_incoming_from(&peer_id) {
1284                    trace!(target: "net::tx", peer_id=format!("{peer_id:#}"), policy=?self.config.ingress_policy, "Ignoring full transactions from peer blocked by ingress policy");
1285                    return;
1286                }
1287
1288                // ensure we didn't receive any blob transactions as these are disallowed to be
1289                // broadcasted in full
1290
1291                let has_blob_txs = msg.has_eip4844();
1292
1293                let non_blob_txs = msg
1294                    .into_iter()
1295                    .map(N::PooledTransaction::try_from)
1296                    .filter_map(Result::ok)
1297                    .collect();
1298
1299                self.import_transactions(peer_id, non_blob_txs, TransactionSource::Broadcast);
1300
1301                if has_blob_txs {
1302                    debug!(target: "net::tx", ?peer_id, "received bad full blob transaction broadcast");
1303                    self.report_peer_bad_transactions(peer_id);
1304                }
1305            }
1306            NetworkTransactionEvent::IncomingPooledTransactionHashes { peer_id, msg } => {
1307                if !self.accepts_incoming_from(&peer_id) {
1308                    trace!(target: "net::tx", peer_id=format!("{peer_id:#}"), policy=?self.config.ingress_policy, "Ignoring transaction hashes from peer blocked by ingress policy");
1309                    return;
1310                }
1311                self.on_new_pooled_transaction_hashes(peer_id, msg)
1312            }
1313            NetworkTransactionEvent::GetPooledTransactions { peer_id, request, response } => {
1314                self.on_get_pooled_transactions(peer_id, request, response)
1315            }
1316            NetworkTransactionEvent::GetTransactionsHandle(response) => {
1317                let _ = response.send(Some(self.handle()));
1318            }
1319        }
1320    }
1321
1322    /// Starts the import process for the given transactions.
1323    ///
1324    /// Callers must check import capacity before consuming a fetch response. Responses are
1325    /// admitted in full and may overshoot the soft limit by one bounded response; broadcasts
1326    /// are truncated to their available capacity.
1327    fn import_transactions(
1328        &mut self,
1329        peer_id: PeerId,
1330        transactions: PooledTransactions<N::PooledTransaction>,
1331        source: TransactionSource,
1332    ) {
1333        // If the node is pipeline syncing, ignore transactions
1334        if self.network.is_initially_syncing() {
1335            return
1336        }
1337        if self.network.tx_gossip_disabled() {
1338            return
1339        }
1340
1341        let (version, client_version) = match &source {
1342            TransactionSource::Broadcast => {
1343                let Some(peer) = self.peers.get(&peer_id) else { return };
1344                (peer.version(), peer.client_version.clone())
1345            }
1346            TransactionSource::Response { version, client_version } => {
1347                (*version, client_version.clone())
1348            }
1349        };
1350        let is_broadcast = source.is_broadcast();
1351        let mut transactions = transactions.0;
1352
1353        // Fetched batches may exceed the soft limit: their request size is bounded, and the
1354        // caller stops consuming responses until capacity becomes available again.
1355        if is_broadcast {
1356            let capacity = self.remaining_broadcast_import_capacity();
1357            if transactions.len() > capacity {
1358                self.metrics
1359                    .skipped_transactions_pending_pool_imports_at_capacity
1360                    .increment((transactions.len() - capacity) as u64);
1361                transactions.truncate(capacity);
1362            }
1363        }
1364
1365        if transactions.is_empty() {
1366            return
1367        }
1368
1369        let start = Instant::now();
1370
1371        // Stop fetching transactions received through either broadcasts or responses.
1372        self.transaction_fetcher
1373            .on_transactions_received(transactions.iter().map(|tx| tx.tx_hash()));
1374
1375        // track that the peer knows these transaction, but only if this is a new broadcast.
1376        // If we received the transactions as the response to our `GetPooledTransactions``
1377        // requests (based on received `NewPooledTransactionHashes`) then we already
1378        // recorded the hashes as seen by this peer in `Self::on_new_pooled_transaction_hashes`.
1379        let mut num_already_seen_by_peer = 0;
1380        if is_broadcast && let Some(peer) = self.peers.get_mut(&peer_id) {
1381            for tx in &transactions {
1382                if !peer.seen_transactions.insert(*tx.tx_hash()) {
1383                    num_already_seen_by_peer += 1;
1384                }
1385            }
1386        }
1387
1388        // Eth/72 `PooledTransactions` responses elide blob payloads from type 3 transactions per
1389        // EIP-8070, so their sidecars can never validate. Drop them before touching the pool;
1390        // geth equivalently diverts these bodies into a buffer that is completed with cells
1391        // fetched via `GetCells`, which is not implemented yet.
1392        if version == EthVersion::Eth72 {
1393            let len_before = transactions.len();
1394            transactions.retain(|tx| !tx.is_eip4844());
1395            let dropped = len_before - transactions.len();
1396            if dropped > 0 {
1397                trace!(target: "net::tx",
1398                    peer_id=format!("{peer_id:#}"),
1399                    dropped,
1400                    "dropped blob transactions from eth72 response, cell fetching not implemented"
1401                );
1402            }
1403        }
1404
1405        // tracks the quality of the given transactions
1406        let mut has_bad_transactions = false;
1407
1408        // 1. Remove known, already-tracked, and invalid transactions first since these are
1409        // cheap in-memory checks against local maps
1410        transactions.retain(|tx| {
1411            if let Entry::Occupied(mut entry) = self.transactions_by_peers.entry(*tx.tx_hash()) {
1412                let peers = entry.get_mut();
1413                if !peers.contains(&peer_id) {
1414                    peers.push(peer_id);
1415                }
1416                return false
1417            }
1418            if self.bad_imports.contains(tx.tx_hash()) {
1419                trace!(target: "net::tx",
1420                    peer_id=format!("{peer_id:#}"),
1421                    hash=%tx.tx_hash(),
1422                    %client_version,
1423                    "received a known bad transaction from peer"
1424                );
1425                has_bad_transactions = true;
1426                return false;
1427            }
1428            true
1429        });
1430
1431        // 2. filter out txns already inserted into pool
1432        let txns_count_pre_pool_filter = transactions.len();
1433        self.pool.retain_unknown(&mut transactions);
1434        if txns_count_pre_pool_filter > transactions.len() {
1435            let already_known_txns_count = txns_count_pre_pool_filter - transactions.len();
1436            self.metrics
1437                .occurrences_transactions_already_in_pool
1438                .increment(already_known_txns_count as u64);
1439        }
1440
1441        let txs_len = transactions.len();
1442
1443        let recover = |tx| {
1444            let recovered = Pool::Transaction::try_recover_with_cache_opt(
1445                tx,
1446                self.sender_recovery_cache.as_ref(),
1447            );
1448            match recovered {
1449                Ok(tx) => Some(tx),
1450                Err(badtx) => {
1451                    trace!(target: "net::tx",
1452                        peer_id=format!("{peer_id:#}"),
1453                        hash=%badtx.tx_hash(),
1454                        client_version=%client_version,
1455                        "failed ecrecovery for transaction"
1456                    );
1457                    None
1458                }
1459            }
1460        };
1461
1462        let new_txs = transactions.into_par_iter().filter_map(recover).collect::<Vec<_>>();
1463
1464        has_bad_transactions |= new_txs.len() != txs_len;
1465
1466        // Record the transactions as seen by the peer
1467        for tx in &new_txs {
1468            self.transactions_by_peers.insert(*tx.hash(), smallvec::smallvec![peer_id]);
1469        }
1470
1471        // 3. import new transactions as a batch to minimize lock contention on the underlying
1472        // pool
1473        if !new_txs.is_empty() {
1474            let pool = self.pool.clone();
1475            // update metrics
1476            let metric_pending_pool_imports = self.metrics.pending_pool_imports.clone();
1477            metric_pending_pool_imports.increment(new_txs.len() as f64);
1478
1479            // update self-monitoring info
1480            self.pending_pool_imports_info
1481                .pending_pool_imports
1482                .fetch_add(new_txs.len(), Ordering::Relaxed);
1483            let tx_manager_info_pending_pool_imports =
1484                self.pending_pool_imports_info.pending_pool_imports.clone();
1485
1486            trace!(target: "net::tx::propagation", new_txs_len=?new_txs.len(), "Importing new transactions");
1487            let import = Box::pin(async move {
1488                let added = new_txs.len();
1489                let res = pool.add_external_transactions(new_txs).await;
1490
1491                // update metrics
1492                metric_pending_pool_imports.decrement(added as f64);
1493                // update self-monitoring info
1494                tx_manager_info_pending_pool_imports.fetch_sub(added, Ordering::Relaxed);
1495
1496                res
1497            });
1498
1499            self.pool_imports.push(import);
1500        }
1501
1502        if num_already_seen_by_peer > 0 {
1503            self.metrics.messages_with_transactions_already_seen_by_peer.increment(1);
1504            self.metrics
1505                .occurrences_of_transaction_already_seen_by_peer
1506                .increment(num_already_seen_by_peer);
1507            trace!(target: "net::tx", num_txs=%num_already_seen_by_peer, ?peer_id, client=%client_version, "Peer sent already seen transactions");
1508        }
1509
1510        if has_bad_transactions {
1511            // peer sent us invalid transactions
1512            self.report_peer_bad_transactions(peer_id)
1513        }
1514
1515        if num_already_seen_by_peer > 0 {
1516            self.report_already_seen(peer_id);
1517        }
1518
1519        self.metrics.pool_import_prepare_duration.record(start.elapsed());
1520    }
1521
1522    /// Processes a [`FetchEvent`].
1523    fn on_fetch_event(&mut self, fetch_event: FetchEvent<N::PooledTransaction>) {
1524        match fetch_event {
1525            FetchEvent::TransactionsFetched {
1526                peer_id,
1527                transactions,
1528                report_peer,
1529                version,
1530                client_version,
1531            } => {
1532                self.import_transactions(
1533                    peer_id,
1534                    transactions,
1535                    TransactionSource::Response { version, client_version },
1536                );
1537                if report_peer {
1538                    self.report_peer(peer_id, ReputationChangeKind::BadTransactions);
1539                }
1540            }
1541            FetchEvent::FetchError { peer_id, error } => {
1542                trace!(target: "net::tx", ?peer_id, %error, "requesting transactions from peer failed");
1543                self.on_request_error(peer_id, error);
1544            }
1545            FetchEvent::EmptyResponse { peer_id } => {
1546                trace!(target: "net::tx", ?peer_id, "peer returned empty response");
1547            }
1548        }
1549    }
1550}
1551
1552/// An endless future. Preemption ensure that future is non-blocking, nonetheless. See
1553/// [`crate::NetworkManager`] for more context on the design pattern.
1554///
1555/// This should be spawned or used as part of `tokio::select!`.
1556//
1557// spawned in `NodeConfig::start_network`(reth_node_core::NodeConfig) and
1558// `NetworkConfig::start_network`(reth_network::NetworkConfig)
1559impl<
1560        Pool: TransactionPool + Unpin + 'static,
1561        N: NetworkPrimitives<
1562                BroadcastedTransaction: SignedTransaction,
1563                PooledTransaction: SignedTransaction,
1564            > + Unpin,
1565    > Future for TransactionsManager<Pool, N>
1566where
1567    Pool::Transaction:
1568        PoolTransaction<Consensus = N::BroadcastedTransaction, Pooled = N::PooledTransaction>,
1569{
1570    type Output = ();
1571
1572    fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
1573        let start = Instant::now();
1574        let mut poll_durations = TxManagerPollDurations::default();
1575
1576        let this = self.get_mut();
1577
1578        // All streams are polled until their corresponding budget is exhausted, then we manually
1579        // yield back control to tokio. See `NetworkManager` for more context on the design
1580        // pattern.
1581
1582        // Advance network/peer related events (update peers map).
1583        let maybe_more_network_events = metered_poll_nested_stream_with_budget!(
1584            poll_durations.acc_network_events,
1585            "net::tx",
1586            "Network events stream",
1587            DEFAULT_BUDGET_TRY_DRAIN_STREAM,
1588            this.network_events.poll_next_unpin(cx),
1589            |event| this.on_network_event(event)
1590        );
1591
1592        // Advance incoming transaction events (stream new txns/announcements from
1593        // network manager and queue for import to pool/fetch txns).
1594        //
1595        // Announcements are queued in the transaction fetcher, requests for them are sent
1596        // further below.
1597        //
1598        // Decoded broadcasts may exceed the available import capacity. Admission truncates
1599        // them before filtering and sender recovery, reserving capacity for fetched transactions.
1600        let maybe_more_tx_events = metered_poll_nested_stream_with_budget!(
1601            poll_durations.acc_tx_events,
1602            "net::tx",
1603            "Network transaction events stream",
1604            DEFAULT_BUDGET_TRY_DRAIN_NETWORK_TRANSACTION_EVENTS,
1605            this.transaction_events.poll_next_unpin(cx),
1606            |event: NetworkTransactionEvent<N>| this.on_network_tx_event(event),
1607        );
1608
1609        // Advance inflight fetch requests (flush transaction fetcher and queue for
1610        // import to pool).
1611        //
1612        // Pace response consumption with the soft import limit. Dispatch checks capacity but
1613        // does not reserve it for outstanding requests, so concurrent responses can exceed it.
1614        // Admit one complete response while capacity remains (at most 256 transactions), then
1615        // pause here if it fills or overshoots the limit. Existing requests continue in flight;
1616        // completed responses wait in their channels until pool imports free capacity.
1617        let mut maybe_more_tx_fetch_events = metered_poll_nested_stream_with_budget!(
1618            poll_durations.acc_fetch_events,
1619            "net::tx",
1620            "Transaction fetch events stream",
1621            DEFAULT_BUDGET_TRY_DRAIN_STREAM,
1622            if this.has_capacity_for_pending_pool_imports() {
1623                this.transaction_fetcher.poll_next_unpin(cx)
1624            } else {
1625                Poll::Pending
1626            },
1627            |event| this.on_fetch_event(event),
1628        );
1629
1630        // Remember whether response polling stopped at capacity so imports completing below
1631        // can trigger another poll even if the waiting response has not registered a waker yet.
1632        let fetcher_paused_at_capacity = !this.has_capacity_for_pending_pool_imports();
1633
1634        // Advance admitted batches and free capacity for more responses.
1635        let maybe_more_pool_imports = metered_poll_nested_stream_with_budget!(
1636            poll_durations.acc_pending_imports,
1637            "net::tx",
1638            "Batched pool imports stream",
1639            DEFAULT_BUDGET_TRY_DRAIN_PENDING_POOL_IMPORTS,
1640            this.pool_imports.poll_next_unpin(cx),
1641            |batch_results| this.on_batch_import_result(batch_results)
1642        );
1643
1644        // Advances new __pending__ transactions, transactions that were successfully inserted into
1645        // pending set in pool (are valid), and propagates them (inform peers which
1646        // transactions we have seen).
1647        //
1648        // This is polled after pool imports so transactions that became pending in this poll
1649        // iteration are propagated immediately, instead of waiting for the task to be woken
1650        // again.
1651        //
1652        // We try to drain this to batch the transactions in a single message.
1653        //
1654        // We don't expect this buffer to be large, since only pending transactions are
1655        // emitted here.
1656        let mut new_txs = Vec::new();
1657        let maybe_more_pending_txns = match this.pending_transactions.poll_recv_many(
1658            cx,
1659            &mut new_txs,
1660            SOFT_LIMIT_COUNT_HASHES_IN_NEW_POOLED_TRANSACTIONS_BROADCAST_MESSAGE,
1661        ) {
1662            Poll::Ready(count) => {
1663                if count == SOFT_LIMIT_COUNT_HASHES_IN_NEW_POOLED_TRANSACTIONS_BROADCAST_MESSAGE {
1664                    // we filled the entire buffer capacity and need to try again on the next poll
1665                    // immediately
1666                    true
1667                } else {
1668                    // try once more, because mostlikely the channel is now empty and the waker is
1669                    // registered if this is pending, if we filled additional hashes, we poll again
1670                    // on the next iteration
1671                    let limit =
1672                        SOFT_LIMIT_COUNT_HASHES_IN_NEW_POOLED_TRANSACTIONS_BROADCAST_MESSAGE -
1673                            new_txs.len();
1674                    this.pending_transactions.poll_recv_many(cx, &mut new_txs, limit).is_ready()
1675                }
1676            }
1677            Poll::Pending => false,
1678        };
1679        if !new_txs.is_empty() {
1680            this.on_new_pending_transactions(new_txs);
1681        }
1682
1683        // Dispatch only below the soft import limit. Each request is capped independently;
1684        // outstanding requests do not reserve capacity against other peers.
1685        duration_metered_exec!(
1686            {
1687                // Peers whose session channel was full are retried on the next poll, which any
1688                // other event triggers.
1689                let budget = this.remaining_pool_import_capacity();
1690                if budget > 0 && this.transaction_fetcher.dispatch(&this.peers, budget) > 0 {
1691                    // poll the fetcher again so the new inflight requests register their wakers
1692                    maybe_more_tx_fetch_events = true;
1693                }
1694            },
1695            poll_durations.acc_pending_fetch
1696        );
1697
1698        // Advance commands (propagate/fetch/serve txns).
1699        let maybe_more_commands = metered_poll_nested_stream_with_budget!(
1700            poll_durations.acc_cmds,
1701            "net::tx",
1702            "Commands channel",
1703            DEFAULT_BUDGET_TRY_DRAIN_STREAM,
1704            this.command_rx.poll_next_unpin(cx),
1705            |cmd| this.on_command(cmd)
1706        );
1707
1708        this.transaction_fetcher.update_metrics();
1709
1710        // Revisit the fetcher if imports freed capacity after response polling was paused.
1711        let resume_fetcher =
1712            fetcher_paused_at_capacity && this.has_capacity_for_pending_pool_imports();
1713
1714        // all channels are fully drained and import futures pending
1715        if maybe_more_network_events ||
1716            maybe_more_commands ||
1717            maybe_more_tx_events ||
1718            maybe_more_tx_fetch_events ||
1719            maybe_more_pool_imports ||
1720            maybe_more_pending_txns ||
1721            resume_fetcher
1722        {
1723            // make sure we're woken up again
1724            cx.waker().wake_by_ref();
1725            return Poll::Pending
1726        }
1727
1728        this.update_poll_metrics(start, poll_durations);
1729
1730        Poll::Pending
1731    }
1732}
1733
1734/// Represents the different modes of transaction propagation.
1735///
1736/// This enum is used to determine how transactions are propagated to peers in the network.
1737#[derive(Debug, Copy, Clone, Eq, PartialEq)]
1738enum PropagationMode {
1739    /// Default propagation mode.
1740    ///
1741    /// Transactions are only sent to peers that haven't seen them yet.
1742    Basic,
1743    /// Forced propagation mode.
1744    ///
1745    /// Transactions are sent to all peers regardless of whether they have been sent or received
1746    /// before.
1747    Forced,
1748}
1749
1750impl PropagationMode {
1751    /// Returns `true` if the propagation kind is `Forced`.
1752    const fn is_forced(self) -> bool {
1753        matches!(self, Self::Forced)
1754    }
1755}
1756
1757/// A transaction that's about to be propagated to multiple peers.
1758#[derive(Debug, Clone)]
1759struct PropagateTransaction {
1760    is_broadcastable_in_full: bool,
1761    /// Size advertised in `NewPooledTransactionHashes` metadata and used for full broadcast
1762    /// soft-limit accounting.
1763    ///
1764    /// This is the network encoded transaction size. For pool-backed blob transactions, this is
1765    /// the pool's cached encoded length, which includes the sidecar returned by
1766    /// `PooledTransactions`.
1767    propagation_size: usize,
1768    transaction: LazyEncodedTransaction,
1769}
1770
1771impl PropagateTransaction {
1772    /// Create a new instance from a transaction supplied directly for propagation.
1773    ///
1774    /// Direct transactions use their EIP-2718 encoded length so eth/68+ hash announcements carry
1775    /// the same size metadata as [`NewPooledTransactionHashes68::push`] and
1776    /// [`NewPooledTransactionHashes72::push`].
1777    fn new<T: SignedTransaction>(transaction: T) -> Self {
1778        let is_broadcastable_in_full = transaction.is_broadcastable_in_full();
1779        let propagation_size = transaction.encode_2718_len();
1780
1781        Self {
1782            is_broadcastable_in_full,
1783            propagation_size,
1784            transaction: LazyEncoded::new(transaction),
1785        }
1786    }
1787
1788    /// Create a new instance from a pooled transaction.
1789    ///
1790    /// Pool transactions already cache the network encoded size used by txpool admission and
1791    /// pooled hash announcements. For blob transactions, this includes the sidecar size expected in
1792    /// a `PooledTransactions` response.
1793    fn pool_tx<P: PoolTransaction>(tx: Arc<ValidPoolTransaction<P>>) -> Self {
1794        let is_broadcastable_in_full = tx.transaction.consensus_ref().is_broadcastable_in_full();
1795        let propagation_size = tx.encoded_length();
1796        Self {
1797            is_broadcastable_in_full,
1798            propagation_size,
1799            transaction: LazyEncoded::new(PropagatePooledTransactionEncoder::new(tx)),
1800        }
1801    }
1802
1803    fn tx_hash(&self) -> &TxHash {
1804        self.transaction.tx_hash()
1805    }
1806
1807    /// Returns the network encoded size used for propagation limits and hash metadata.
1808    const fn propagation_size(&self) -> usize {
1809        self.propagation_size
1810    }
1811
1812    fn tx_type(&self) -> u8 {
1813        self.transaction.ty()
1814    }
1815
1816    const fn is_broadcastable_in_full(&self) -> bool {
1817        self.is_broadcastable_in_full
1818    }
1819
1820    fn shared(&self) -> LazyEncodedTransaction {
1821        self.transaction.clone()
1822    }
1823}
1824
1825/// A pooled transaction encoder that avoids cloning into the consensus transaction for propagation.
1826#[derive(Debug)]
1827struct PropagatePooledTransactionEncoder<P: PoolTransaction> {
1828    transaction: Arc<ValidPoolTransaction<P>>,
1829}
1830
1831impl<P: PoolTransaction> PropagatePooledTransactionEncoder<P> {
1832    const fn new(transaction: Arc<ValidPoolTransaction<P>>) -> Self {
1833        Self { transaction }
1834    }
1835
1836    fn encode_uncached(&self, out: &mut dyn BufMut) {
1837        (*self.transaction.transaction.consensus_ref().inner()).encode(out);
1838    }
1839}
1840
1841impl<P: PoolTransaction> Encodable for PropagatePooledTransactionEncoder<P> {
1842    fn encode(&self, out: &mut dyn BufMut) {
1843        self.encode_uncached(out);
1844    }
1845
1846    fn length(&self) -> usize {
1847        (*self.transaction.transaction.consensus_ref().inner()).length()
1848    }
1849}
1850
1851impl<P: PoolTransaction> TxHashRef for PropagatePooledTransactionEncoder<P> {
1852    fn tx_hash(&self) -> &TxHash {
1853        self.transaction.hash()
1854    }
1855}
1856
1857impl<P: PoolTransaction> Typed2718 for PropagatePooledTransactionEncoder<P> {
1858    fn ty(&self) -> u8 {
1859        self.transaction.transaction.ty()
1860    }
1861}
1862
1863/// Helper type to construct the appropriate message to send to the peer based on whether the peer
1864/// should receive them in full or as pooled
1865#[derive(Debug, Clone)]
1866enum PropagateTransactionsBuilder {
1867    Pooled(PooledTransactionsHashesBuilder),
1868    Full(FullTransactionsBuilder),
1869}
1870
1871impl PropagateTransactionsBuilder {
1872    /// Create a builder for pooled transactions with capacity for the expected number of
1873    /// transactions.
1874    fn pooled(version: EthVersion, capacity: usize) -> Self {
1875        Self::Pooled(PooledTransactionsHashesBuilder::with_capacity(version, capacity))
1876    }
1877
1878    /// Create a builder that sends transactions in full and records transactions that don't fit,
1879    /// with capacity for the expected number of transactions.
1880    fn full(version: EthVersion, capacity: usize) -> Self {
1881        Self::Full(FullTransactionsBuilder::with_capacity(version, capacity))
1882    }
1883
1884    /// Returns true if no transactions are recorded.
1885    fn is_empty(&self) -> bool {
1886        match self {
1887            Self::Pooled(builder) => builder.is_empty(),
1888            Self::Full(builder) => builder.is_empty(),
1889        }
1890    }
1891
1892    /// Consumes the type and returns the built messages that should be sent to the peer.
1893    fn build(self) -> PropagateTransactions {
1894        match self {
1895            Self::Pooled(pooled) => {
1896                PropagateTransactions { pooled: Some(pooled.build()), full: None }
1897            }
1898            Self::Full(full) => full.build(),
1899        }
1900    }
1901}
1902
1903impl PropagateTransactionsBuilder {
1904    /// Appends a transaction to the list.
1905    fn push(&mut self, transaction: &PropagateTransaction) {
1906        match self {
1907            Self::Pooled(builder) => builder.push(transaction),
1908            Self::Full(builder) => builder.push(transaction),
1909        }
1910    }
1911}
1912
1913/// Represents how the transactions should be sent to a peer if any.
1914struct PropagateTransactions {
1915    /// The pooled transaction hashes to send.
1916    pooled: Option<NewPooledTransactionHashes>,
1917    /// The transactions to send in full.
1918    full: Option<BroadcastPoolTransactions>,
1919}
1920
1921/// Helper type for constructing the full transaction message that enforces the
1922/// [`DEFAULT_SOFT_LIMIT_BYTE_SIZE_TRANSACTIONS_BROADCAST_MESSAGE`] for full transaction broadcast
1923/// and enforces other propagation rules for EIP-4844 and tracks those transactions that can't be
1924/// broadcasted in full.
1925#[derive(Debug, Clone)]
1926struct FullTransactionsBuilder {
1927    /// The soft limit to enforce for a single broadcast message of full transactions.
1928    total_size: usize,
1929    /// All transactions to be broadcasted.
1930    transactions: Vec<LazyEncodedTransaction>,
1931    /// Transactions that didn't fit into the broadcast message
1932    pooled: PooledTransactionsHashesBuilder,
1933}
1934
1935impl FullTransactionsBuilder {
1936    /// Create a builder for the negotiated version of the peer's session
1937    fn new(version: EthVersion) -> Self {
1938        Self {
1939            total_size: 0,
1940            pooled: PooledTransactionsHashesBuilder::new(version),
1941            transactions: vec![],
1942        }
1943    }
1944
1945    /// Create a builder with capacity for the expected number of full transactions.
1946    ///
1947    /// The overflow hashes builder remains lazily allocated since most transactions are expected
1948    /// to be broadcast in full.
1949    fn with_capacity(version: EthVersion, capacity: usize) -> Self {
1950        Self {
1951            total_size: 0,
1952            pooled: PooledTransactionsHashesBuilder::new(version),
1953            transactions: Vec::with_capacity(capacity),
1954        }
1955    }
1956
1957    /// Returns whether or not any transactions are in the [`FullTransactionsBuilder`].
1958    fn is_empty(&self) -> bool {
1959        self.transactions.is_empty() && self.pooled.is_empty()
1960    }
1961
1962    /// Returns the messages that should be propagated to the peer.
1963    fn build(self) -> PropagateTransactions {
1964        let pooled = Some(self.pooled.build()).filter(|pooled| !pooled.is_empty());
1965        let full =
1966            (!self.transactions.is_empty()).then_some(BroadcastPoolTransactions(self.transactions));
1967        PropagateTransactions { pooled, full }
1968    }
1969
1970    /// Appends all transactions.
1971    fn extend(&mut self, txs: impl IntoIterator<Item = PropagateTransaction>) {
1972        for tx in txs {
1973            self.push(&tx)
1974        }
1975    }
1976
1977    /// Append a transaction to the list of full transaction if the total message bytes size doesn't
1978    /// exceed the soft maximum target byte size. The limit is soft, meaning if one single
1979    /// transaction goes over the limit, it will be broadcasted in its own [`Transactions`]
1980    /// message. The same pattern is followed in filling a [`GetPooledTransactions`] request in
1981    /// [`TransactionFetcher::dispatch`].
1982    ///
1983    /// If the transaction is unsuitable for broadcast or would exceed the softlimit, it is appended
1984    /// to list of pooled transactions, (e.g. 4844 transactions).
1985    /// See also [`SignedTransaction::is_broadcastable_in_full`].
1986    fn push(&mut self, transaction: &PropagateTransaction) {
1987        // Do not send full 4844 transaction hashes to peers.
1988        //
1989        //  Nodes MUST NOT automatically broadcast blob transactions to their peers.
1990        //  Instead, those transactions are only announced using
1991        //  `NewPooledTransactionHashes` messages, and can then be manually requested
1992        //  via `GetPooledTransactions`.
1993        //
1994        // From: <https://eips.ethereum.org/EIPS/eip-4844#networking>
1995        if !transaction.is_broadcastable_in_full() {
1996            self.pooled.push(transaction);
1997            return
1998        }
1999
2000        let new_size = self.total_size + transaction.propagation_size();
2001        if new_size > DEFAULT_SOFT_LIMIT_BYTE_SIZE_TRANSACTIONS_BROADCAST_MESSAGE &&
2002            self.total_size > 0
2003        {
2004            // transaction does not fit into the message
2005            self.pooled.push(transaction);
2006            return
2007        }
2008
2009        self.total_size = new_size;
2010        self.transactions.push(transaction.shared());
2011    }
2012}
2013
2014/// A helper type to create the pooled transactions message based on the negotiated version of the
2015/// session with the peer
2016#[derive(Debug, Clone)]
2017enum PooledTransactionsHashesBuilder {
2018    Eth66(NewPooledTransactionHashes66),
2019    Eth68(NewPooledTransactionHashes68),
2020    Eth72(NewPooledTransactionHashes72),
2021}
2022
2023// === impl PooledTransactionsHashesBuilder ===
2024
2025impl PooledTransactionsHashesBuilder {
2026    /// Push a transaction from the pool to the list.
2027    fn push_pooled<T: PoolTransaction>(&mut self, pooled_tx: Arc<ValidPoolTransaction<T>>) {
2028        match self {
2029            Self::Eth66(msg) => msg.push(*pooled_tx.hash()),
2030            Self::Eth68(msg) => {
2031                msg.hashes.push(*pooled_tx.hash());
2032                msg.sizes.push(pooled_tx.encoded_length());
2033                msg.types.push(pooled_tx.transaction.ty());
2034            }
2035            Self::Eth72(msg) => {
2036                msg.hashes.push(*pooled_tx.hash());
2037                msg.sizes.push(pooled_tx.encoded_length());
2038                let ty = pooled_tx.transaction.ty();
2039                msg.types.push(ty);
2040                if ty == EIP4844_TX_TYPE_ID {
2041                    // The pool holds the full sidecar, so every cell can be served.
2042                    msg.cell_mask = Some(NewPooledTransactionHashes72::ALL_CELLS_MASK);
2043                }
2044            }
2045        }
2046    }
2047
2048    /// Returns whether or not any transactions are in the [`PooledTransactionsHashesBuilder`].
2049    fn is_empty(&self) -> bool {
2050        match self {
2051            Self::Eth66(hashes) => hashes.is_empty(),
2052            Self::Eth68(hashes) => hashes.is_empty(),
2053            Self::Eth72(hashes) => hashes.is_empty(),
2054        }
2055    }
2056
2057    /// Returns the number of transactions in the builder.
2058    fn len(&self) -> usize {
2059        match self {
2060            Self::Eth66(hashes) => hashes.len(),
2061            Self::Eth68(hashes) => hashes.len(),
2062            Self::Eth72(hashes) => hashes.len(),
2063        }
2064    }
2065
2066    /// Appends all hashes
2067    fn extend(&mut self, txs: impl IntoIterator<Item = PropagateTransaction>) {
2068        for tx in txs {
2069            self.push(&tx);
2070        }
2071    }
2072
2073    fn push(&mut self, tx: &PropagateTransaction) {
2074        match self {
2075            Self::Eth66(msg) => msg.push(*tx.tx_hash()),
2076            Self::Eth68(msg) => {
2077                msg.hashes.push(*tx.tx_hash());
2078                msg.sizes.push(tx.propagation_size());
2079                msg.types.push(tx.tx_type());
2080            }
2081            Self::Eth72(msg) => {
2082                msg.hashes.push(*tx.tx_hash());
2083                msg.sizes.push(tx.propagation_size());
2084                let ty = tx.tx_type();
2085                msg.types.push(ty);
2086                if ty == EIP4844_TX_TYPE_ID {
2087                    // The pool holds the full sidecar, so every cell can be served.
2088                    msg.cell_mask = Some(NewPooledTransactionHashes72::ALL_CELLS_MASK);
2089                }
2090            }
2091        }
2092    }
2093
2094    /// Create a builder for the negotiated version of the peer's session
2095    fn new(version: EthVersion) -> Self {
2096        match version {
2097            EthVersion::Eth66 | EthVersion::Eth67 => Self::Eth66(Default::default()),
2098            EthVersion::Eth68 | EthVersion::Eth69 | EthVersion::Eth70 | EthVersion::Eth71 => {
2099                Self::Eth68(Default::default())
2100            }
2101            EthVersion::Eth72 => Self::Eth72(Default::default()),
2102        }
2103    }
2104
2105    /// Create a builder for the negotiated version of the peer's session with capacity for the
2106    /// expected number of hashes.
2107    fn with_capacity(version: EthVersion, capacity: usize) -> Self {
2108        match version {
2109            EthVersion::Eth66 | EthVersion::Eth67 => {
2110                Self::Eth66(NewPooledTransactionHashes66::with_capacity(capacity))
2111            }
2112            EthVersion::Eth68 | EthVersion::Eth69 | EthVersion::Eth70 | EthVersion::Eth71 => {
2113                Self::Eth68(NewPooledTransactionHashes68::with_capacity(capacity))
2114            }
2115            EthVersion::Eth72 => Self::Eth72(NewPooledTransactionHashes72::with_capacity(capacity)),
2116        }
2117    }
2118
2119    fn build(self) -> NewPooledTransactionHashes {
2120        match self {
2121            Self::Eth66(mut msg) => {
2122                msg.shrink_to_fit();
2123                msg.into()
2124            }
2125            Self::Eth68(mut msg) => {
2126                msg.shrink_to_fit();
2127                msg.into()
2128            }
2129            Self::Eth72(mut msg) => {
2130                msg.shrink_to_fit();
2131                msg.into()
2132            }
2133        }
2134    }
2135}
2136
2137/// How we received the transactions.
2138#[derive(Debug)]
2139enum TransactionSource {
2140    /// Transactions were broadcast to us via [`Transactions`] message.
2141    Broadcast,
2142    /// Transactions were sent as the response to a `GetPooledTransactions` request issued by us.
2143    Response {
2144        /// Session metadata remains available after the responding peer disconnects.
2145        version: EthVersion,
2146        client_version: Arc<str>,
2147    },
2148}
2149
2150// === impl TransactionSource ===
2151
2152impl TransactionSource {
2153    /// Whether the transaction were sent as broadcast.
2154    const fn is_broadcast(&self) -> bool {
2155        matches!(self, Self::Broadcast)
2156    }
2157}
2158
2159/// Tracks a single peer in the context of [`TransactionsManager`].
2160#[derive(Debug)]
2161pub struct PeerMetadata<N: NetworkPrimitives = EthNetworkPrimitives> {
2162    /// Optimistically keeps track of transactions that we know the peer has seen. Optimistic, in
2163    /// the sense that transactions are preemptively marked as seen by peer when they are sent to
2164    /// the peer.
2165    seen_transactions: LruCache<TxHash, FbBuildHasher<32>>,
2166    /// A communication channel directly to the peer's session task.
2167    request_tx: PeerRequestSender<PeerRequest<N>>,
2168    /// negotiated version of the session.
2169    version: EthVersion,
2170    /// The peer's client version.
2171    client_version: Arc<str>,
2172    /// The kind of peer.
2173    peer_kind: PeerKind,
2174}
2175
2176impl<N: NetworkPrimitives> PeerMetadata<N> {
2177    /// Returns a new instance of [`PeerMetadata`].
2178    pub fn new(
2179        request_tx: PeerRequestSender<PeerRequest<N>>,
2180        version: EthVersion,
2181        client_version: Arc<str>,
2182        max_transactions_seen_by_peer: u32,
2183        peer_kind: PeerKind,
2184    ) -> Self {
2185        Self {
2186            seen_transactions: LruCache::with_hasher(
2187                max_transactions_seen_by_peer,
2188                Default::default(),
2189            ),
2190            request_tx,
2191            version,
2192            client_version,
2193            peer_kind,
2194        }
2195    }
2196
2197    /// Returns a reference to the peer's request sender channel.
2198    pub const fn request_tx(&self) -> &PeerRequestSender<PeerRequest<N>> {
2199        &self.request_tx
2200    }
2201
2202    /// Returns a mutable reference to the seen transactions LRU cache.
2203    pub const fn seen_transactions_mut(&mut self) -> &mut LruCache<TxHash, FbBuildHasher<32>> {
2204        &mut self.seen_transactions
2205    }
2206
2207    /// Returns the negotiated `EthVersion` of the session.
2208    pub const fn version(&self) -> EthVersion {
2209        self.version
2210    }
2211
2212    /// Returns a reference to the peer's client version string.
2213    pub fn client_version(&self) -> &str {
2214        &self.client_version
2215    }
2216
2217    /// Returns the peer's kind.
2218    pub const fn peer_kind(&self) -> PeerKind {
2219        self.peer_kind
2220    }
2221}
2222
2223/// Commands to send to the [`TransactionsManager`]
2224#[derive(Debug)]
2225enum TransactionsCommand<N: NetworkPrimitives = EthNetworkPrimitives> {
2226    /// Propagate a transaction hash to the network.
2227    PropagateHash(B256),
2228    /// Propagate transaction hashes to a specific peer.
2229    PropagateHashesTo(Vec<B256>, PeerId),
2230    /// Request the list of active peer IDs from the [`TransactionsManager`].
2231    GetActivePeers(oneshot::Sender<HashSet<PeerId>>),
2232    /// Propagate a collection of full transactions to a specific peer.
2233    PropagateTransactionsTo(Vec<TxHash>, PeerId),
2234    /// Propagate a collection of hashes to all peers.
2235    PropagateTransactions(Vec<TxHash>),
2236    /// Propagate a collection of broadcastable transactions in full to all peers.
2237    BroadcastTransactions(Vec<PropagateTransaction>),
2238    /// Request transaction hashes known by specific peers from the [`TransactionsManager`].
2239    GetTransactionHashes { peers: Vec<PeerId>, tx: oneshot::Sender<HashMap<PeerId, B256Set>> },
2240    /// Requests a clone of the sender channel to the peer.
2241    GetPeerSender {
2242        peer_id: PeerId,
2243        peer_request_sender: oneshot::Sender<Option<PeerRequestSender<PeerRequest<N>>>>,
2244    },
2245}
2246
2247/// All events related to transactions emitted by the network.
2248#[derive(Debug)]
2249pub enum NetworkTransactionEvent<N: NetworkPrimitives = EthNetworkPrimitives> {
2250    /// Represents the event of receiving a list of transactions from a peer.
2251    ///
2252    /// This indicates transactions that were broadcasted to us from the peer.
2253    IncomingTransactions {
2254        /// The ID of the peer from which the transactions were received.
2255        peer_id: PeerId,
2256        /// The received transactions.
2257        msg: Transactions<N::BroadcastedTransaction>,
2258    },
2259    /// Represents the event of receiving a list of transaction hashes from a peer.
2260    IncomingPooledTransactionHashes {
2261        /// The ID of the peer from which the transaction hashes were received.
2262        peer_id: PeerId,
2263        /// The received new pooled transaction hashes.
2264        msg: NewPooledTransactionHashes,
2265    },
2266    /// Represents the event of receiving a `GetPooledTransactions` request from a peer.
2267    GetPooledTransactions {
2268        /// The ID of the peer from which the request was received.
2269        peer_id: PeerId,
2270        /// The received `GetPooledTransactions` request.
2271        request: GetPooledTransactions,
2272        /// The sender for responding to the request with a result of `PooledTransactions`.
2273        response: oneshot::Sender<RequestResult<PooledTransactions<N::PooledTransaction>>>,
2274    },
2275    /// Represents the event of receiving a `GetTransactionsHandle` request.
2276    GetTransactionsHandle(oneshot::Sender<Option<TransactionsHandle<N>>>),
2277}
2278
2279/// Tracks stats about the [`TransactionsManager`].
2280#[derive(Debug)]
2281pub struct PendingPoolImportsInfo {
2282    /// Number of transactions about to be inserted into the pool.
2283    pending_pool_imports: Arc<AtomicUsize>,
2284    /// Max number of transactions allowed to be imported concurrently.
2285    max_pending_pool_imports: usize,
2286}
2287
2288impl PendingPoolImportsInfo {
2289    /// Returns a new [`PendingPoolImportsInfo`].
2290    pub fn new(max_pending_pool_imports: usize) -> Self {
2291        Self { pending_pool_imports: Arc::new(AtomicUsize::default()), max_pending_pool_imports }
2292    }
2293
2294    /// Returns `true` if the number of pool imports is under a given tolerated max.
2295    pub fn has_capacity(&self, max_pending_pool_imports: usize) -> bool {
2296        self.pending_pool_imports.load(Ordering::Relaxed) < max_pending_pool_imports
2297    }
2298}
2299
2300impl Default for PendingPoolImportsInfo {
2301    fn default() -> Self {
2302        Self::new(DEFAULT_MAX_COUNT_PENDING_POOL_IMPORTS)
2303    }
2304}
2305
2306#[derive(Debug, Default)]
2307struct TxManagerPollDurations {
2308    acc_network_events: Duration,
2309    acc_pending_imports: Duration,
2310    acc_tx_events: Duration,
2311    acc_imported_txns: Duration,
2312    acc_fetch_events: Duration,
2313    acc_pending_fetch: Duration,
2314    acc_cmds: Duration,
2315}
2316
2317impl<N: NetworkPrimitives> InMemorySize for NetworkTransactionEvent<N> {
2318    // `N::BroadcastedTransaction` and `N::PooledTransaction` already implement
2319    // `InMemorySize` via `SignedTransaction: InMemorySize`, so no extra bound is needed.
2320    fn size(&self) -> usize {
2321        match self {
2322            Self::IncomingTransactions { peer_id, msg } => {
2323                core::mem::size_of_val(peer_id) +
2324                    msg.0.iter().map(InMemorySize::size).sum::<usize>()
2325            }
2326            Self::IncomingPooledTransactionHashes { peer_id, msg } => {
2327                core::mem::size_of_val(peer_id) + msg.size()
2328            }
2329            Self::GetPooledTransactions { peer_id, request, response } => {
2330                core::mem::size_of_val(peer_id) +
2331                    request.0.len() * core::mem::size_of::<TxHash>() +
2332                    core::mem::size_of_val(response)
2333            }
2334            Self::GetTransactionsHandle(_) => 0,
2335        }
2336    }
2337}
2338
2339#[cfg(test)]
2340mod tests {
2341    use super::*;
2342    use crate::{
2343        test_utils::{
2344            transactions::{buffer_hash_to_tx_fetcher, new_mock_session, new_tx_manager},
2345            Testnet,
2346        },
2347        transactions::config::RelaxedEthAnnouncementFilter,
2348        NetworkConfigBuilder, NetworkManager,
2349    };
2350    use alloy_consensus::{Transaction as _, TxEip1559, TxLegacy};
2351    use alloy_eips::{eip2718::Encodable2718, eip4844::BlobTransactionValidationError};
2352    use alloy_primitives::{hex, Signature, TxKind, B256, U256};
2353    use alloy_rlp::Decodable;
2354    use futures::FutureExt;
2355    use reth_chainspec::MIN_TRANSACTION_GAS;
2356    use reth_ethereum_primitives::{PooledTransactionVariant, Transaction, TransactionSigned};
2357    use reth_network_api::{NetworkInfo, PeerKind};
2358    use reth_network_p2p::{
2359        error::{RequestError, RequestResult},
2360        sync::{NetworkSyncUpdater, SyncState},
2361    };
2362    use reth_storage_api::noop::NoopProvider;
2363    use reth_tasks::Runtime;
2364    use reth_transaction_pool::{
2365        blobstore::InMemoryBlobStore,
2366        error::{Eip4844PoolTransactionError, InvalidPoolTransactionError, PoolError},
2367        identifier::SenderIdentifiers,
2368        test_utils::{
2369            testing_pool, MockTransaction, MockTransactionFactory, OkValidator, TestPool,
2370            TransactionGenerator,
2371        },
2372        CoinbaseTipOrdering, EthPooledTransaction, Pool, TransactionOrigin, ValidPoolTransaction,
2373    };
2374    use secp256k1::SecretKey;
2375    use std::{
2376        future::poll_fn,
2377        net::{IpAddr, Ipv4Addr, SocketAddr},
2378        str::FromStr,
2379        time::Instant,
2380    };
2381    use tracing::error;
2382
2383    type EthTestPool = Pool<
2384        OkValidator<EthPooledTransaction>,
2385        CoinbaseTipOrdering<EthPooledTransaction>,
2386        InMemoryBlobStore,
2387    >;
2388
2389    #[tokio::test]
2390    async fn announcement_policy_preserves_order_and_skips_pending_and_bad_imports() {
2391        #[derive(Debug, Clone, Default)]
2392        struct RecordingPolicy(Arc<parking_lot::Mutex<Vec<(u8, TxHash, usize)>>>);
2393
2394        impl AnnouncementFilteringPolicy<EthNetworkPrimitives> for RecordingPolicy {
2395            fn decide_on_announcement(
2396                &self,
2397                tx_type: u8,
2398                hash: &TxHash,
2399                size: usize,
2400            ) -> AnnouncementAcceptance {
2401                self.0.lock().push((tx_type, *hash, size));
2402                match tx_type {
2403                    0xff => AnnouncementAcceptance::Reject { penalize_peer: true },
2404                    1 => AnnouncementAcceptance::Ignore,
2405                    _ => AnnouncementAcceptance::Accept,
2406                }
2407            }
2408        }
2409
2410        let (mut manager, _network) = new_tx_manager().await;
2411        let policy = RecordingPolicy::default();
2412        manager.policies =
2413            NetworkPolicies::new(TransactionPropagationKind::default(), policy.clone());
2414        let peer_id = PeerId::new([1; 64]);
2415        let (peer, mut rx) = new_mock_session(peer_id, EthVersion::Eth68);
2416        manager.peers.insert(peer_id, peer);
2417        let [pending, accepted, ignored, rejected, later, bad] =
2418            [1, 2, 3, 4, 5, 6].map(B256::repeat_byte);
2419        manager.transactions_by_peers.insert(pending, Default::default());
2420        manager.bad_imports.insert(bad);
2421        manager.on_new_pooled_transaction_hashes(
2422            peer_id,
2423            NewPooledTransactionHashes68 {
2424                hashes: vec![pending, accepted, accepted, ignored, rejected, later, bad],
2425                types: vec![0xff, 2, 0xff, 1, 0xff, 2, 0xff],
2426                sizes: vec![100, 200, 300, 400, 500, 600, 700],
2427            }
2428            .into(),
2429        );
2430
2431        assert_eq!(
2432            *policy.0.lock(),
2433            [(2, accepted, 200), (1, ignored, 400), (0xff, rejected, 500), (2, later, 600)]
2434        );
2435        poll_fn(|cx| {
2436            let _ = manager.poll_unpin(cx);
2437            Poll::Ready(())
2438        })
2439        .await;
2440        let PeerRequest::GetPooledTransactions { request, .. } = rx.try_recv().unwrap() else {
2441            panic!("expected transaction request")
2442        };
2443        assert_eq!(request.0, [accepted, later]);
2444    }
2445
2446    async fn new_eth_tx_manager() -> (
2447        TransactionsManager<EthTestPool, EthNetworkPrimitives>,
2448        NetworkManager<EthNetworkPrimitives>,
2449    ) {
2450        let secret_key = SecretKey::new(&mut rand_08::thread_rng());
2451        let client = NoopProvider::default();
2452
2453        let config = NetworkConfigBuilder::new(secret_key, Runtime::test())
2454            .listener_port(0)
2455            .disable_discovery()
2456            .build(client);
2457
2458        let pool = Pool::new(
2459            OkValidator::default(),
2460            CoinbaseTipOrdering::default(),
2461            InMemoryBlobStore::default(),
2462            Default::default(),
2463        );
2464
2465        let transactions_manager_config = config.transactions_manager_config.clone();
2466        let (_network_handle, network, transactions, _) = NetworkManager::new(config)
2467            .await
2468            .unwrap()
2469            .into_builder()
2470            .transactions(pool.clone(), transactions_manager_config)
2471            .split_with_handle();
2472
2473        (transactions, network)
2474    }
2475
2476    fn valid_eth_pool_transaction(
2477        transaction: EthPooledTransaction,
2478    ) -> Arc<ValidPoolTransaction<EthPooledTransaction>> {
2479        let mut ids = SenderIdentifiers::default();
2480        let transaction_id =
2481            ids.sender_id_or_create(transaction.sender()).into_transaction_id(transaction.nonce());
2482
2483        Arc::new(ValidPoolTransaction {
2484            propagate: false,
2485            transaction_id,
2486            transaction,
2487            timestamp: Instant::now(),
2488            origin: TransactionOrigin::External,
2489            authority_ids: None,
2490        })
2491    }
2492
2493    fn gen_eip1559_pooled_with_nonce<R: rand::RngCore>(
2494        tx_gen: &mut TransactionGenerator<R>,
2495        nonce: u64,
2496    ) -> EthPooledTransaction {
2497        EthPooledTransaction::try_from_consensus(
2498            tx_gen.transaction().nonce(nonce).into_eip1559().try_into_recovered().unwrap(),
2499        )
2500        .unwrap()
2501    }
2502
2503    #[tokio::test(flavor = "multi_thread")]
2504    async fn test_ignored_tx_broadcasts_while_initially_syncing() {
2505        reth_tracing::init_test_tracing();
2506        let net = Testnet::create(3).await.spawn();
2507        let [peer0, peer1, _] = net.peers_array();
2508
2509        let listener0 = peer0.event_listener();
2510        peer0.add_peer(peer1);
2511        let secret_key = SecretKey::new(&mut rand_08::thread_rng());
2512
2513        let client = NoopProvider::default();
2514        let pool = testing_pool();
2515        let config = NetworkConfigBuilder::eth(secret_key, Runtime::test())
2516            .disable_discovery()
2517            .listener_port(0)
2518            .build(client);
2519        let transactions_manager_config = config.transactions_manager_config.clone();
2520        let (network_handle, network, mut transactions, _) = NetworkManager::new(config)
2521            .await
2522            .unwrap()
2523            .into_builder()
2524            .transactions(pool.clone(), transactions_manager_config)
2525            .split_with_handle();
2526
2527        tokio::task::spawn(network);
2528
2529        // go to syncing (pipeline sync)
2530        network_handle.update_sync_state(SyncState::Syncing);
2531        assert!(NetworkInfo::is_syncing(&network_handle));
2532        assert!(NetworkInfo::is_initially_syncing(&network_handle));
2533
2534        // wait for all initiator connections
2535        let mut established = listener0.take(2);
2536        while let Some(ev) = established.next().await {
2537            match ev {
2538                NetworkEvent::Peer(PeerEvent::SessionEstablished(info)) => {
2539                    // to insert a new peer in transactions peerset
2540                    transactions
2541                        .on_network_event(NetworkEvent::Peer(PeerEvent::SessionEstablished(info)))
2542                }
2543                NetworkEvent::Peer(PeerEvent::PeerAdded(_peer_id)) => {}
2544                ev => {
2545                    error!("unexpected event {ev:?}")
2546                }
2547            }
2548        }
2549        // random tx: <https://etherscan.io/getRawTx?tx=0x9448608d36e721ef403c53b00546068a6474d6cbab6816c3926de449898e7bce>
2550        let input = hex!(
2551            "02f871018302a90f808504890aef60826b6c94ddf4c5025d1a5742cf12f74eec246d4432c295e487e09c3bbcc12b2b80c080a0f21a4eacd0bf8fea9c5105c543be5a1d8c796516875710fafafdf16d16d8ee23a001280915021bb446d1973501a67f93d2b38894a514b976e7b46dc2fe54598d76"
2552        );
2553        let signed_tx = TransactionSigned::decode(&mut &input[..]).unwrap();
2554        transactions.on_network_tx_event(NetworkTransactionEvent::IncomingTransactions {
2555            peer_id: *peer1.peer_id(),
2556            msg: Transactions(vec![signed_tx.clone()]),
2557        });
2558        poll_fn(|cx| {
2559            let _ = transactions.poll_unpin(cx);
2560            Poll::Ready(())
2561        })
2562        .await;
2563        assert!(pool.is_empty());
2564    }
2565
2566    #[tokio::test(flavor = "multi_thread")]
2567    async fn test_tx_broadcasts_through_two_syncs() {
2568        reth_tracing::init_test_tracing();
2569        let net = Testnet::create(3).await.spawn();
2570        let [peer0, peer1, _] = net.peers_array();
2571
2572        let listener0 = peer0.event_listener();
2573        peer0.add_peer(peer1);
2574        let secret_key = SecretKey::new(&mut rand_08::thread_rng());
2575
2576        let client = NoopProvider::default();
2577        let pool = testing_pool();
2578        let config = NetworkConfigBuilder::new(secret_key, Runtime::test())
2579            .disable_discovery()
2580            .listener_port(0)
2581            .build(client);
2582        let transactions_manager_config = config.transactions_manager_config.clone();
2583        let (network_handle, network, mut transactions, _) = NetworkManager::new(config)
2584            .await
2585            .unwrap()
2586            .into_builder()
2587            .transactions(pool.clone(), transactions_manager_config)
2588            .split_with_handle();
2589
2590        tokio::task::spawn(network);
2591
2592        // go to syncing (pipeline sync) to idle and then to syncing (live)
2593        network_handle.update_sync_state(SyncState::Syncing);
2594        assert!(NetworkInfo::is_syncing(&network_handle));
2595        network_handle.update_sync_state(SyncState::Idle);
2596        assert!(!NetworkInfo::is_syncing(&network_handle));
2597        network_handle.update_sync_state(SyncState::Syncing);
2598        assert!(NetworkInfo::is_syncing(&network_handle));
2599
2600        // wait for all initiator connections
2601        let mut established = listener0.take(2);
2602        while let Some(ev) = established.next().await {
2603            match ev {
2604                NetworkEvent::ActivePeerSession { .. } |
2605                NetworkEvent::Peer(PeerEvent::SessionEstablished(_)) => {
2606                    // to insert a new peer in transactions peerset
2607                    transactions.on_network_event(ev);
2608                }
2609                NetworkEvent::Peer(PeerEvent::PeerAdded(_peer_id)) => {}
2610                _ => {
2611                    error!("unexpected event {ev:?}")
2612                }
2613            }
2614        }
2615        // random tx: <https://etherscan.io/getRawTx?tx=0x9448608d36e721ef403c53b00546068a6474d6cbab6816c3926de449898e7bce>
2616        let input = hex!(
2617            "02f871018302a90f808504890aef60826b6c94ddf4c5025d1a5742cf12f74eec246d4432c295e487e09c3bbcc12b2b80c080a0f21a4eacd0bf8fea9c5105c543be5a1d8c796516875710fafafdf16d16d8ee23a001280915021bb446d1973501a67f93d2b38894a514b976e7b46dc2fe54598d76"
2618        );
2619        let signed_tx = TransactionSigned::decode(&mut &input[..]).unwrap();
2620        transactions.on_network_tx_event(NetworkTransactionEvent::IncomingTransactions {
2621            peer_id: *peer1.peer_id(),
2622            msg: Transactions(vec![signed_tx.clone()]),
2623        });
2624        poll_fn(|cx| {
2625            let _ = transactions.poll_unpin(cx);
2626            Poll::Ready(())
2627        })
2628        .await;
2629        assert!(!NetworkInfo::is_initially_syncing(&network_handle));
2630        assert!(NetworkInfo::is_syncing(&network_handle));
2631        assert!(!pool.is_empty());
2632    }
2633
2634    // Ensure that the transaction manager correctly handles the `IncomingPooledTransactionHashes`
2635    // event and is able to retrieve the corresponding transactions.
2636    #[tokio::test(flavor = "multi_thread")]
2637    async fn test_handle_incoming_transactions_hashes() {
2638        reth_tracing::init_test_tracing();
2639
2640        let secret_key = SecretKey::new(&mut rand_08::thread_rng());
2641        let client = NoopProvider::default();
2642
2643        let config = NetworkConfigBuilder::new(secret_key, Runtime::test())
2644            // let OS choose port
2645            .listener_port(0)
2646            .disable_discovery()
2647            .build(client);
2648
2649        let pool = testing_pool();
2650
2651        let transactions_manager_config = config.transactions_manager_config.clone();
2652        let (_network_handle, _network, mut tx_manager, _) = NetworkManager::new(config)
2653            .await
2654            .unwrap()
2655            .into_builder()
2656            .transactions(pool.clone(), transactions_manager_config)
2657            .split_with_handle();
2658
2659        let peer_id_1 = PeerId::new([1; 64]);
2660        let eth_version = EthVersion::Eth66;
2661
2662        let txs = vec![TransactionSigned::new_unhashed(
2663            Transaction::Legacy(TxLegacy {
2664                chain_id: Some(4),
2665                nonce: 15u64,
2666                gas_price: 2200000000,
2667                gas_limit: 34811,
2668                to: TxKind::Call(hex!("cf7f9e66af820a19257a2108375b180b0ec49167").into()),
2669                value: U256::from(1234u64),
2670                input: Default::default(),
2671            }),
2672            Signature::new(
2673                U256::from_str(
2674                    "0x35b7bfeb9ad9ece2cbafaaf8e202e706b4cfaeb233f46198f00b44d4a566a981",
2675                )
2676                .unwrap(),
2677                U256::from_str(
2678                    "0x612638fb29427ca33b9a3be2a0a561beecfe0269655be160d35e72d366a6a860",
2679                )
2680                .unwrap(),
2681                true,
2682            ),
2683        )];
2684
2685        let txs_hashes: Vec<B256> = txs.iter().map(|tx| *tx.hash()).collect();
2686
2687        let (peer_1, mut to_mock_session_rx) = new_mock_session(peer_id_1, eth_version);
2688        tx_manager.peers.insert(peer_id_1, peer_1);
2689
2690        assert!(pool.is_empty());
2691
2692        tx_manager.on_network_tx_event(NetworkTransactionEvent::IncomingPooledTransactionHashes {
2693            peer_id: peer_id_1,
2694            msg: NewPooledTransactionHashes::from(NewPooledTransactionHashes66::from(
2695                txs_hashes.clone(),
2696            )),
2697        });
2698
2699        // advance the transaction manager future to send the request
2700        poll_fn(|cx| {
2701            let _ = tx_manager.poll_unpin(cx);
2702            Poll::Ready(())
2703        })
2704        .await;
2705
2706        // mock session of peer_1 receives request
2707        let req = to_mock_session_rx
2708            .recv()
2709            .await
2710            .expect("peer_1 session should receive request with buffered hashes");
2711        let PeerRequest::GetPooledTransactions { request, response } = req else { unreachable!() };
2712        assert_eq!(request, GetPooledTransactions::from(txs_hashes.clone()));
2713
2714        let message: Vec<PooledTransactionVariant> = txs
2715            .into_iter()
2716            .map(|tx| {
2717                PooledTransactionVariant::try_from(tx)
2718                    .expect("Failed to convert MockTransaction to PooledTransaction")
2719            })
2720            .collect();
2721
2722        // return the transactions corresponding to the transaction hashes.
2723        response
2724            .send(Ok(PooledTransactions(message)))
2725            .expect("should send peer_1 response to tx manager");
2726
2727        // adance the transaction manager future
2728        poll_fn(|cx| {
2729            let _ = tx_manager.poll_unpin(cx);
2730            Poll::Ready(())
2731        })
2732        .await;
2733
2734        // ensure that the transactions corresponding to the transaction hashes have been
2735        // successfully retrieved and stored in the Pool.
2736        assert_eq!(pool.get_all(txs_hashes.clone()).len(), txs_hashes.len());
2737    }
2738
2739    #[tokio::test(flavor = "multi_thread")]
2740    async fn test_handle_incoming_transactions() {
2741        reth_tracing::init_test_tracing();
2742        let net = Testnet::create(3).await.spawn();
2743        let [peer0, peer1, _] = net.peers_array();
2744
2745        let listener0 = peer0.event_listener();
2746
2747        peer0.add_peer(peer1);
2748        let secret_key = SecretKey::new(&mut rand_08::thread_rng());
2749
2750        let client = NoopProvider::default();
2751        let pool = testing_pool();
2752        let config = NetworkConfigBuilder::new(secret_key, Runtime::test())
2753            .disable_discovery()
2754            .listener_port(0)
2755            .build(client);
2756        let transactions_manager_config = config.transactions_manager_config.clone();
2757        let (network_handle, network, mut transactions, _) = NetworkManager::new(config)
2758            .await
2759            .unwrap()
2760            .into_builder()
2761            .transactions(pool.clone(), transactions_manager_config)
2762            .split_with_handle();
2763        tokio::task::spawn(network);
2764
2765        network_handle.update_sync_state(SyncState::Idle);
2766
2767        assert!(!NetworkInfo::is_syncing(&network_handle));
2768
2769        // wait for all initiator connections
2770        let mut established = listener0.take(2);
2771        while let Some(ev) = established.next().await {
2772            match ev {
2773                NetworkEvent::ActivePeerSession { .. } |
2774                NetworkEvent::Peer(PeerEvent::SessionEstablished(_)) => {
2775                    // to insert a new peer in transactions peerset
2776                    transactions.on_network_event(ev);
2777                }
2778                NetworkEvent::Peer(PeerEvent::PeerAdded(_peer_id)) => {}
2779                ev => {
2780                    error!("unexpected event {ev:?}")
2781                }
2782            }
2783        }
2784        // random tx: <https://etherscan.io/getRawTx?tx=0x9448608d36e721ef403c53b00546068a6474d6cbab6816c3926de449898e7bce>
2785        let input = hex!(
2786            "02f871018302a90f808504890aef60826b6c94ddf4c5025d1a5742cf12f74eec246d4432c295e487e09c3bbcc12b2b80c080a0f21a4eacd0bf8fea9c5105c543be5a1d8c796516875710fafafdf16d16d8ee23a001280915021bb446d1973501a67f93d2b38894a514b976e7b46dc2fe54598d76"
2787        );
2788        let signed_tx = TransactionSigned::decode(&mut &input[..]).unwrap();
2789        transactions.on_network_tx_event(NetworkTransactionEvent::IncomingTransactions {
2790            peer_id: *peer1.peer_id(),
2791            msg: Transactions(vec![signed_tx.clone()]),
2792        });
2793        assert!(transactions
2794            .transactions_by_peers
2795            .get(signed_tx.tx_hash())
2796            .unwrap()
2797            .contains(peer1.peer_id()));
2798
2799        // advance the transaction manager future
2800        poll_fn(|cx| {
2801            let _ = transactions.poll_unpin(cx);
2802            Poll::Ready(())
2803        })
2804        .await;
2805
2806        assert!(!pool.is_empty());
2807        assert!(pool.get(signed_tx.tx_hash()).is_some());
2808    }
2809
2810    #[tokio::test(flavor = "multi_thread")]
2811    async fn test_session_closed_cleans_transaction_peer_state() {
2812        let (mut tx_manager, _network) = new_tx_manager().await;
2813        let peer_id = PeerId::new([1; 64]);
2814        let fallback_peer = PeerId::new([2; 64]);
2815        let (peer, _) = new_mock_session(peer_id, EthVersion::Eth66);
2816        let (fallback, _) = new_mock_session(fallback_peer, EthVersion::Eth66);
2817        let hash_shared = B256::from_slice(&[1; 32]);
2818
2819        tx_manager.peers.insert(peer_id, peer);
2820        tx_manager.peers.insert(fallback_peer, fallback);
2821        buffer_hash_to_tx_fetcher(&mut tx_manager.transaction_fetcher, hash_shared, peer_id, None);
2822        buffer_hash_to_tx_fetcher(
2823            &mut tx_manager.transaction_fetcher,
2824            hash_shared,
2825            fallback_peer,
2826            None,
2827        );
2828        assert_eq!(
2829            tx_manager.transaction_fetcher.candidate_peers(&hash_shared),
2830            vec![peer_id, fallback_peer]
2831        );
2832
2833        tx_manager.on_network_event(NetworkEvent::Peer(PeerEvent::SessionClosed {
2834            peer_id,
2835            reason: None,
2836        }));
2837
2838        // peer removed from peers map and from the fetcher
2839        assert!(!tx_manager.peers.contains_key(&peer_id));
2840        assert!(tx_manager.transaction_fetcher.queued_hashes(&peer_id).is_empty());
2841        // fallback peer is still available for the hash
2842        assert_eq!(
2843            tx_manager.transaction_fetcher.candidate_peers(&hash_shared),
2844            vec![fallback_peer]
2845        );
2846    }
2847
2848    #[tokio::test(flavor = "multi_thread")]
2849    async fn test_bad_blob_sidecar_not_cached_as_bad_import() {
2850        let (mut tx_manager, _network) = new_tx_manager().await;
2851        let peer_id = PeerId::new([1; 64]);
2852        let hash = B256::from_slice(&[1; 32]);
2853
2854        tx_manager.network.update_sync_state(SyncState::Idle);
2855        tx_manager.transactions_by_peers.insert(hash, smallvec::smallvec![peer_id]);
2856
2857        let err = PoolError::new(
2858            hash,
2859            InvalidPoolTransactionError::Eip4844(Eip4844PoolTransactionError::InvalidEip4844Blob(
2860                BlobTransactionValidationError::InvalidProof,
2861            )),
2862        );
2863
2864        tx_manager.on_bad_import(err);
2865
2866        assert!(!tx_manager.bad_imports.contains(&hash));
2867    }
2868
2869    #[tokio::test(flavor = "multi_thread")]
2870    async fn test_missing_blob_sidecar_not_cached_as_bad_import() {
2871        let (mut tx_manager, _network) = new_tx_manager().await;
2872        let peer_id = PeerId::new([1; 64]);
2873        let hash = B256::from_slice(&[3; 32]);
2874
2875        tx_manager.network.update_sync_state(SyncState::Idle);
2876        tx_manager.transactions_by_peers.insert(hash, smallvec::smallvec![peer_id]);
2877
2878        let err = PoolError::new(
2879            hash,
2880            InvalidPoolTransactionError::Eip4844(
2881                Eip4844PoolTransactionError::MissingEip4844BlobSidecar,
2882            ),
2883        );
2884
2885        tx_manager.on_bad_import(err);
2886
2887        assert!(!tx_manager.bad_imports.contains(&hash));
2888    }
2889
2890    #[tokio::test(flavor = "multi_thread")]
2891    async fn test_non_blob_sidecar_error_still_cached_as_bad_import() {
2892        let (mut tx_manager, _network) = new_tx_manager().await;
2893        let peer_id = PeerId::new([1; 64]);
2894        let hash = B256::from_slice(&[2; 32]);
2895
2896        tx_manager.network.update_sync_state(SyncState::Idle);
2897        tx_manager.transactions_by_peers.insert(hash, smallvec::smallvec![peer_id]);
2898
2899        let err = PoolError::new(
2900            hash,
2901            InvalidPoolTransactionError::Eip4844(Eip4844PoolTransactionError::NoEip4844Blobs),
2902        );
2903
2904        tx_manager.on_bad_import(err);
2905
2906        assert!(tx_manager.bad_imports.contains(&hash));
2907    }
2908
2909    #[tokio::test(flavor = "multi_thread")]
2910    async fn test_on_get_pooled_transactions_network() {
2911        reth_tracing::init_test_tracing();
2912        let net = Testnet::create(2).await.spawn();
2913        let [peer0, peer1] = net.peers_array();
2914
2915        let listener0 = peer0.event_listener();
2916
2917        peer0.add_peer(peer1);
2918        let secret_key = SecretKey::new(&mut rand_08::thread_rng());
2919
2920        let client = NoopProvider::default();
2921        let pool = testing_pool();
2922        let config = NetworkConfigBuilder::new(secret_key, Runtime::test())
2923            .disable_discovery()
2924            .listener_port(0)
2925            .build(client);
2926        let transactions_manager_config = config.transactions_manager_config.clone();
2927        let (network_handle, network, mut transactions, _) = NetworkManager::new(config)
2928            .await
2929            .unwrap()
2930            .into_builder()
2931            .transactions(pool.clone(), transactions_manager_config)
2932            .split_with_handle();
2933        tokio::task::spawn(network);
2934
2935        network_handle.update_sync_state(SyncState::Idle);
2936
2937        assert!(!NetworkInfo::is_syncing(&network_handle));
2938
2939        // wait for all initiator connections
2940        let mut established = listener0.take(2);
2941        while let Some(ev) = established.next().await {
2942            match ev {
2943                NetworkEvent::ActivePeerSession { .. } |
2944                NetworkEvent::Peer(PeerEvent::SessionEstablished(_)) => {
2945                    transactions.on_network_event(ev);
2946                }
2947                NetworkEvent::Peer(PeerEvent::PeerAdded(_peer_id)) => {}
2948                ev => {
2949                    error!("unexpected event {ev:?}")
2950                }
2951            }
2952        }
2953
2954        let tx = MockTransaction::eip1559();
2955        let _ = transactions
2956            .pool
2957            .add_transaction(reth_transaction_pool::TransactionOrigin::External, tx.clone())
2958            .await;
2959
2960        let request = GetPooledTransactions(vec![*tx.get_hash()]);
2961
2962        let (send, receive) =
2963            oneshot::channel::<RequestResult<PooledTransactions<PooledTransactionVariant>>>();
2964
2965        transactions.on_network_tx_event(NetworkTransactionEvent::GetPooledTransactions {
2966            peer_id: *peer1.peer_id(),
2967            request,
2968            response: send,
2969        });
2970
2971        match receive.await.unwrap() {
2972            Ok(PooledTransactions(transactions)) => {
2973                assert_eq!(transactions.len(), 1);
2974            }
2975            Err(e) => {
2976                panic!("error: {e:?}");
2977            }
2978        }
2979    }
2980
2981    // Ensure that when the remote peer only returns part of the requested transactions, the
2982    // replied transactions are removed from the `tx_fetcher`, while the unresponsive ones are
2983    // re-buffered.
2984    #[tokio::test]
2985    async fn test_partially_tx_response() {
2986        reth_tracing::init_test_tracing();
2987
2988        let mut tx_manager = new_tx_manager().await.0;
2989        let tx_fetcher = &mut tx_manager.transaction_fetcher;
2990
2991        let peer_id_1 = PeerId::new([1; 64]);
2992        let eth_version = EthVersion::Eth66;
2993
2994        let txs = vec![
2995            TransactionSigned::new_unhashed(
2996                Transaction::Legacy(TxLegacy {
2997                    chain_id: Some(4),
2998                    nonce: 15u64,
2999                    gas_price: 2200000000,
3000                    gas_limit: 34811,
3001                    to: TxKind::Call(hex!("cf7f9e66af820a19257a2108375b180b0ec49167").into()),
3002                    value: U256::from(1234u64),
3003                    input: Default::default(),
3004                }),
3005                Signature::new(
3006                    U256::from_str(
3007                        "0x35b7bfeb9ad9ece2cbafaaf8e202e706b4cfaeb233f46198f00b44d4a566a981",
3008                    )
3009                    .unwrap(),
3010                    U256::from_str(
3011                        "0x612638fb29427ca33b9a3be2a0a561beecfe0269655be160d35e72d366a6a860",
3012                    )
3013                    .unwrap(),
3014                    true,
3015                ),
3016            ),
3017            TransactionSigned::new_unhashed(
3018                Transaction::Eip1559(TxEip1559 {
3019                    chain_id: 4,
3020                    nonce: 26u64,
3021                    max_priority_fee_per_gas: 1500000000,
3022                    max_fee_per_gas: 1500000013,
3023                    gas_limit: MIN_TRANSACTION_GAS,
3024                    to: TxKind::Call(hex!("61815774383099e24810ab832a5b2a5425c154d5").into()),
3025                    value: U256::from(3000000000000000000u64),
3026                    input: Default::default(),
3027                    access_list: Default::default(),
3028                }),
3029                Signature::new(
3030                    U256::from_str(
3031                        "0x59e6b67f48fb32e7e570dfb11e042b5ad2e55e3ce3ce9cd989c7e06e07feeafd",
3032                    )
3033                    .unwrap(),
3034                    U256::from_str(
3035                        "0x016b83f4f980694ed2eee4d10667242b1f40dc406901b34125b008d334d47469",
3036                    )
3037                    .unwrap(),
3038                    true,
3039                ),
3040            ),
3041        ];
3042
3043        let txs_hashes: Vec<B256> = txs.iter().map(|tx| *tx.hash()).collect();
3044
3045        let (mut peer_1, mut to_mock_session_rx) = new_mock_session(peer_id_1, eth_version);
3046        // mark hashes as seen by peer so it can fish them out from the cache for hashes pending
3047        // fetch
3048        peer_1.seen_transactions.insert(txs_hashes[0]);
3049        peer_1.seen_transactions.insert(txs_hashes[1]);
3050        tx_manager.peers.insert(peer_id_1, peer_1);
3051
3052        buffer_hash_to_tx_fetcher(tx_fetcher, txs_hashes[0], peer_id_1, None);
3053        buffer_hash_to_tx_fetcher(tx_fetcher, txs_hashes[1], peer_id_1, None);
3054
3055        // peer_1 is idle
3056        assert!(tx_fetcher.is_idle(&peer_id_1));
3057        assert_eq!(tx_fetcher.num_inflight_requests(), 0);
3058
3059        // sends requests for buffered hashes to peer_1
3060        assert_eq!(tx_fetcher.dispatch(&tx_manager.peers, usize::MAX), 1);
3061
3062        assert_eq!(tx_fetcher.num_pending_hashes(), 0);
3063        // as long as request is in flight peer_1 is not idle
3064        assert!(!tx_fetcher.is_idle(&peer_id_1));
3065        assert_eq!(tx_fetcher.num_inflight_requests(), 1);
3066
3067        // mock session of peer_1 receives request
3068        let req = to_mock_session_rx
3069            .recv()
3070            .await
3071            .expect("peer_1 session should receive request with buffered hashes");
3072        let PeerRequest::GetPooledTransactions { response, .. } = req else { unreachable!() };
3073
3074        let message: Vec<PooledTransactionVariant> = txs
3075            .into_iter()
3076            .take(1)
3077            .map(|tx| {
3078                PooledTransactionVariant::try_from(tx)
3079                    .expect("Failed to convert MockTransaction to PooledTransaction")
3080            })
3081            .collect();
3082        // response partial request
3083        response
3084            .send(Ok(PooledTransactions(message)))
3085            .expect("should send peer_1 response to tx manager");
3086        let Some(FetchEvent::TransactionsFetched { peer_id, .. }) = tx_fetcher.next().await else {
3087            unreachable!()
3088        };
3089
3090        // request has resolved, peer_1 is idle again
3091        assert!(tx_fetcher.is_idle(&peer_id));
3092        assert_eq!(tx_fetcher.num_inflight_requests(), 0);
3093        // the undelivered hash at the end of the request is retried with peer_1
3094        assert_eq!(tx_fetcher.num_pending_hashes(), 1);
3095        assert_eq!(tx_fetcher.candidate_peers(&txs_hashes[1]), vec![peer_id_1]);
3096    }
3097
3098    /// Tests that hashes of a failed request are retried with the alternate peer that announced
3099    /// them, and are given up on when no peer is left to fetch them from.
3100    #[tokio::test]
3101    async fn test_failed_request_retries_on_alternate_peer() {
3102        reth_tracing::init_test_tracing();
3103
3104        let mut tx_manager = new_tx_manager().await.0;
3105        let tx_fetcher = &mut tx_manager.transaction_fetcher;
3106
3107        let peer_id_1 = PeerId::new([1; 64]);
3108        let peer_id_2 = PeerId::new([2; 64]);
3109        let eth_version = EthVersion::Eth66;
3110        let seen_hashes = [B256::from_slice(&[1; 32]), B256::from_slice(&[2; 32])];
3111
3112        let (peer_1, mut to_mock_session_rx_1) = new_mock_session(peer_id_1, eth_version);
3113        let (peer_2, mut to_mock_session_rx_2) = new_mock_session(peer_id_2, eth_version);
3114        tx_manager.peers.insert(peer_id_1, peer_1);
3115        tx_manager.peers.insert(peer_id_2, peer_2);
3116
3117        // both peers announce the hashes, peer_1 first
3118        for hash in seen_hashes {
3119            buffer_hash_to_tx_fetcher(tx_fetcher, hash, peer_id_1, None);
3120            buffer_hash_to_tx_fetcher(tx_fetcher, hash, peer_id_2, None);
3121        }
3122        assert!(tx_fetcher.is_idle(&peer_id_1));
3123        assert!(tx_fetcher.is_idle(&peer_id_2));
3124
3125        // the hashes are requested from peer_1 only
3126        assert_eq!(tx_fetcher.dispatch(&tx_manager.peers, usize::MAX), 1);
3127        assert_eq!(tx_fetcher.num_pending_hashes(), 0);
3128        assert!(!tx_fetcher.is_idle(&peer_id_1));
3129        assert!(tx_fetcher.is_idle(&peer_id_2));
3130
3131        let req = to_mock_session_rx_1
3132            .recv()
3133            .await
3134            .expect("peer_1 session should receive request with buffered hashes");
3135        let PeerRequest::GetPooledTransactions { request, response } = req else { unreachable!() };
3136        let GetPooledTransactions(hashes) = request;
3137        assert_eq!(
3138            hashes.into_iter().collect::<B256Set>(),
3139            seen_hashes.into_iter().collect::<B256Set>()
3140        );
3141
3142        // fail request to peer_1
3143        response
3144            .send(Err(RequestError::BadResponse))
3145            .expect("should send peer_1 response to tx manager");
3146        let Some(FetchEvent::FetchError { peer_id, .. }) = tx_fetcher.next().await else {
3147            unreachable!()
3148        };
3149        assert_eq!(peer_id, peer_id_1);
3150
3151        // request has resolved, peer_1 is idle again and the hashes are pending for peer_2 only
3152        assert!(tx_fetcher.is_idle(&peer_id_1));
3153        assert_eq!(tx_fetcher.num_inflight_requests(), 0);
3154        assert_eq!(tx_fetcher.num_pending_hashes(), 2);
3155        assert_eq!(tx_fetcher.candidate_peers(&seen_hashes[0]), vec![peer_id_2]);
3156
3157        // the hashes are retried with peer_2
3158        assert_eq!(tx_fetcher.dispatch(&tx_manager.peers, usize::MAX), 1);
3159        assert_eq!(tx_fetcher.num_pending_hashes(), 0);
3160        let req = to_mock_session_rx_2
3161            .recv()
3162            .await
3163            .expect("peer_2 session should receive request with buffered hashes");
3164        let PeerRequest::GetPooledTransactions { response, .. } = req else { unreachable!() };
3165
3166        // fail request to peer_2 as well
3167        response
3168            .send(Err(RequestError::BadResponse))
3169            .expect("should send peer_2 response to tx manager");
3170        let Some(FetchEvent::FetchError { .. }) = tx_fetcher.next().await else { unreachable!() };
3171
3172        // no peer is left to fetch the hashes from, they are dropped
3173        assert_eq!(tx_fetcher.num_hashes(), 0);
3174        assert_eq!(tx_fetcher.num_inflight_requests(), 0);
3175    }
3176
3177    #[test]
3178    fn test_direct_propagation_transaction_uses_2718_size() {
3179        let mut tx_gen = TransactionGenerator::new(rand::rng());
3180        let tx = tx_gen.gen_eip1559();
3181        let expected_size = tx.encode_2718_len();
3182
3183        let tx = PropagateTransaction::new(tx);
3184
3185        assert_eq!(tx.propagation_size(), expected_size);
3186    }
3187
3188    #[test]
3189    fn test_transaction_builder_empty() {
3190        let mut builder = PropagateTransactionsBuilder::pooled(EthVersion::Eth68, 0);
3191        assert!(builder.is_empty());
3192
3193        let mut tx_gen = TransactionGenerator::new(rand::rng());
3194        let tx =
3195            PropagateTransaction::pool_tx(valid_eth_pool_transaction(tx_gen.gen_eip1559_pooled()));
3196        builder.push(&tx);
3197        assert!(!builder.is_empty());
3198
3199        let txs = builder.build();
3200        assert!(txs.full.is_none());
3201        let txs = txs.pooled.unwrap();
3202        assert_eq!(txs.len(), 1);
3203    }
3204
3205    #[test]
3206    fn test_pooled_propagation_transaction_encoder_length_matches_network_encoding() {
3207        let mut tx_gen = TransactionGenerator::new(rand::rng());
3208        let tx = valid_eth_pool_transaction(tx_gen.gen_eip1559_pooled());
3209        let pooled = PropagatePooledTransactionEncoder::new(tx);
3210
3211        let mut pooled_encoded = Vec::new();
3212        pooled.encode(&mut pooled_encoded);
3213        assert_eq!(pooled.length(), pooled_encoded.len());
3214
3215        let broadcast = BroadcastPoolTransactions(vec![LazyEncoded::new(pooled)]);
3216        let mut first_encoded = Vec::new();
3217        broadcast.encode(&mut first_encoded);
3218        let mut second_encoded = Vec::new();
3219        broadcast.encode(&mut second_encoded);
3220        assert_eq!(first_encoded, second_encoded);
3221
3222        let mut encoded = first_encoded.as_slice();
3223        let decoded = Transactions::<TransactionSigned>::decode(&mut encoded).unwrap();
3224        assert_eq!(decoded.len(), 1);
3225        assert!(encoded.is_empty());
3226    }
3227
3228    #[test]
3229    fn test_transaction_builder_large() {
3230        let mut builder = PropagateTransactionsBuilder::full(EthVersion::Eth68, 0);
3231        assert!(builder.is_empty());
3232
3233        let mut tx_gen = TransactionGenerator::new(rand::rng());
3234        let mut tx = tx_gen.gen_eip1559_pooled();
3235        // create a transaction that still fits
3236        tx.encoded_length = DEFAULT_SOFT_LIMIT_BYTE_SIZE_TRANSACTIONS_BROADCAST_MESSAGE + 1;
3237        let tx = PropagateTransaction::pool_tx(valid_eth_pool_transaction(tx));
3238        builder.push(&tx);
3239        assert!(!builder.is_empty());
3240
3241        let txs = builder.clone().build();
3242        assert!(txs.pooled.is_none());
3243        let txs = txs.full.unwrap();
3244        assert_eq!(txs.len(), 1);
3245
3246        builder.push(&tx);
3247
3248        let txs = builder.clone().build();
3249        let pooled = txs.pooled.unwrap();
3250        assert_eq!(pooled.len(), 1);
3251        let txs = txs.full.unwrap();
3252        assert_eq!(txs.len(), 1);
3253    }
3254
3255    #[test]
3256    fn test_transaction_builder_eip4844() {
3257        let mut builder = PropagateTransactionsBuilder::full(EthVersion::Eth68, 0);
3258        assert!(builder.is_empty());
3259
3260        let mut tx_gen = TransactionGenerator::new(rand::rng());
3261        let tx =
3262            PropagateTransaction::pool_tx(valid_eth_pool_transaction(tx_gen.gen_eip4844_pooled()));
3263        builder.push(&tx);
3264        assert!(!builder.is_empty());
3265
3266        let txs = builder.clone().build();
3267        assert!(txs.full.is_none());
3268        let txs = txs.pooled.unwrap();
3269        assert_eq!(txs.len(), 1);
3270
3271        let tx =
3272            PropagateTransaction::pool_tx(valid_eth_pool_transaction(tx_gen.gen_eip1559_pooled()));
3273        builder.push(&tx);
3274
3275        let txs = builder.clone().build();
3276        let pooled = txs.pooled.unwrap();
3277        assert_eq!(pooled.len(), 1);
3278        let txs = txs.full.unwrap();
3279        assert_eq!(txs.len(), 1);
3280    }
3281
3282    #[tokio::test]
3283    async fn test_propagate_full() {
3284        reth_tracing::init_test_tracing();
3285
3286        let (mut tx_manager, network) = new_eth_tx_manager().await;
3287        let peer_id = PeerId::random();
3288
3289        // ensure not syncing
3290        network.handle().update_sync_state(SyncState::Idle);
3291
3292        // mock a peer
3293        let (tx, _rx) = mpsc::channel::<PeerRequest>(1);
3294
3295        let session_info = SessionInfo {
3296            peer_id,
3297            remote_addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 0),
3298            client_version: Arc::from(""),
3299            capabilities: Arc::new(vec![].into()),
3300            status: Arc::new(Default::default()),
3301            version: EthVersion::Eth68,
3302            peer_kind: PeerKind::Basic,
3303        };
3304        let messages: PeerRequestSender<PeerRequest> = PeerRequestSender::new(peer_id, tx);
3305        tx_manager
3306            .on_network_event(NetworkEvent::ActivePeerSession { info: session_info, messages });
3307        let mut propagate = vec![];
3308        let mut tx_gen = TransactionGenerator::new(rand::rng());
3309        let eip1559_tx = valid_eth_pool_transaction(tx_gen.gen_eip1559_pooled());
3310        propagate.push(eip1559_tx.clone());
3311        let eip4844_tx = valid_eth_pool_transaction(tx_gen.gen_eip4844_pooled());
3312        propagate.push(eip4844_tx.clone());
3313
3314        let propagated = tx_manager.propagate_transactions(
3315            propagate.clone().into_iter().map(PropagateTransaction::pool_tx).collect(),
3316            PropagationMode::Basic,
3317        );
3318        assert_eq!(propagated.len(), 2);
3319        let prop_txs = propagated.get(eip1559_tx.transaction.hash()).unwrap();
3320        assert_eq!(prop_txs.len(), 1);
3321        assert!(prop_txs[0].is_full());
3322
3323        let prop_txs = propagated.get(eip4844_tx.transaction.hash()).unwrap();
3324        assert_eq!(prop_txs.len(), 1);
3325        assert!(prop_txs[0].is_hash());
3326
3327        let peer = tx_manager.peers.get(&peer_id).unwrap();
3328        assert!(peer.seen_transactions.contains(eip1559_tx.transaction.hash()));
3329        assert!(peer.seen_transactions.contains(eip1559_tx.transaction.hash()));
3330        peer.seen_transactions.contains(eip4844_tx.transaction.hash());
3331
3332        // propagate again
3333        let propagated = tx_manager.propagate_transactions(
3334            propagate.into_iter().map(PropagateTransaction::pool_tx).collect(),
3335            PropagationMode::Basic,
3336        );
3337        assert!(propagated.is_empty());
3338    }
3339
3340    #[test]
3341    fn test_eth72_hashes_builder_sets_cell_mask_for_blob_txs() {
3342        let mut tx_gen = TransactionGenerator::new(rand::rng());
3343
3344        // no blob transactions: the mask stays unset and encodes as a zero mask
3345        let mut builder = PooledTransactionsHashesBuilder::new(EthVersion::Eth72);
3346        builder.push(&PropagateTransaction::pool_tx(valid_eth_pool_transaction(
3347            tx_gen.gen_eip1559_pooled(),
3348        )));
3349        let msg = builder.build();
3350        assert_eq!(msg.as_eth72().unwrap().cell_mask, None);
3351
3352        // announcing a blob transaction advertises every cell as available
3353        let mut builder = PooledTransactionsHashesBuilder::new(EthVersion::Eth72);
3354        builder.push(&PropagateTransaction::pool_tx(valid_eth_pool_transaction(
3355            tx_gen.gen_eip1559_pooled(),
3356        )));
3357        builder.push_pooled(valid_eth_pool_transaction(tx_gen.gen_eip4844_pooled()));
3358        let msg = builder.build();
3359        assert_eq!(
3360            msg.as_eth72().unwrap().cell_mask,
3361            Some(NewPooledTransactionHashes72::ALL_CELLS_MASK)
3362        );
3363    }
3364
3365    #[tokio::test]
3366    async fn test_truncated_hash_announcement_not_marked_seen() {
3367        reth_tracing::init_test_tracing();
3368
3369        let (mut tx_manager, network) = new_eth_tx_manager().await;
3370        // all peers receive hash announcements only
3371        tx_manager.config.propagation_mode = TransactionPropagationMode::Max(0);
3372
3373        // ensure not syncing
3374        network.handle().update_sync_state(SyncState::Idle);
3375
3376        let peer_id = PeerId::random();
3377        let (peer, _rx) = new_mock_session(peer_id, EthVersion::Eth68);
3378        tx_manager.peers.insert(peer_id, peer);
3379
3380        // one more transaction than fits into a single hashes broadcast message
3381        let mut tx_gen = TransactionGenerator::new(rand::rng());
3382        let txs = (0..=SOFT_LIMIT_COUNT_HASHES_IN_NEW_POOLED_TRANSACTIONS_BROADCAST_MESSAGE)
3383            .map(|nonce| {
3384                valid_eth_pool_transaction(gen_eip1559_pooled_with_nonce(&mut tx_gen, nonce as u64))
3385            })
3386            .collect::<Vec<_>>();
3387        let last_sent = *txs[txs.len() - 2].hash();
3388        let truncated = *txs[txs.len() - 1].hash();
3389
3390        let propagated = tx_manager.propagate_transactions(
3391            txs.into_iter().map(PropagateTransaction::pool_tx).collect(),
3392            PropagationMode::Basic,
3393        );
3394
3395        // the truncated hash was not sent, so it must not be tracked as seen by the peer
3396        assert!(propagated.get(&truncated).is_none());
3397        let peer = tx_manager.peers.get(&peer_id).unwrap();
3398        assert!(!peer.seen_transactions.contains(&truncated));
3399        assert!(peer.seen_transactions.contains(&last_sent));
3400    }
3401
3402    #[tokio::test]
3403    async fn test_propagate_pending_txs_while_initially_syncing() {
3404        reth_tracing::init_test_tracing();
3405
3406        let (mut tx_manager, network) = new_eth_tx_manager().await;
3407        let peer_id = PeerId::random();
3408
3409        // Keep the node in initial sync mode.
3410        network.handle().update_sync_state(SyncState::Syncing);
3411        assert!(NetworkInfo::is_initially_syncing(&network.handle()));
3412
3413        // Add a peer so propagation has a destination.
3414        let (peer, _rx) = new_mock_session(peer_id, EthVersion::Eth68);
3415        tx_manager.peers.insert(peer_id, peer);
3416
3417        let mut tx_gen = TransactionGenerator::new(rand::rng());
3418        let tx = gen_eip1559_pooled_with_nonce(&mut tx_gen, 0);
3419        let tx_hash = *tx.hash();
3420        tx_manager
3421            .pool
3422            .add_transaction(reth_transaction_pool::TransactionOrigin::External, tx.clone())
3423            .await
3424            .expect("transaction should be accepted into the pool");
3425
3426        tx_manager.on_new_pending_transactions(vec![tx_hash]);
3427
3428        let peer = tx_manager.peers.get(&peer_id).expect("peer should exist");
3429        assert!(peer.seen_transactions.contains(&tx_hash));
3430    }
3431
3432    #[tokio::test]
3433    async fn test_relaxed_filter_ignores_unknown_tx_types() {
3434        reth_tracing::init_test_tracing();
3435
3436        let transactions_manager_config = TransactionsManagerConfig::default();
3437
3438        let propagation_policy = TransactionPropagationKind::default();
3439        let announcement_policy = RelaxedEthAnnouncementFilter::default();
3440
3441        let policy_bundle = NetworkPolicies::new(propagation_policy, announcement_policy);
3442
3443        let pool = testing_pool();
3444        let secret_key = SecretKey::new(&mut rand_08::thread_rng());
3445        let client = NoopProvider::default();
3446
3447        let network_config = NetworkConfigBuilder::new(secret_key, Runtime::test())
3448            .listener_port(0)
3449            .disable_discovery()
3450            .build(client.clone());
3451
3452        let mut network_manager = NetworkManager::new(network_config).await.unwrap();
3453        let (to_tx_manager_tx, from_network_rx) =
3454            reth_metrics::common::mpsc::memory_bounded_channel::<
3455                NetworkTransactionEvent<EthNetworkPrimitives>,
3456            >(
3457                crate::transactions::constants::tx_manager::DEFAULT_TX_MANAGER_CHANNEL_MEMORY_LIMIT_BYTES,
3458                "test_tx_channel",
3459            );
3460        network_manager.set_transactions(to_tx_manager_tx);
3461        let network_handle = network_manager.handle().clone();
3462        let network_service_handle = tokio::spawn(network_manager);
3463
3464        let mut tx_manager = TransactionsManager::<TestPool, EthNetworkPrimitives>::with_policy(
3465            network_handle.clone(),
3466            pool.clone(),
3467            from_network_rx,
3468            transactions_manager_config,
3469            policy_bundle,
3470        );
3471
3472        let peer_id = PeerId::random();
3473        let eth_version = EthVersion::Eth68;
3474        let (mock_peer_metadata, mut mock_session_rx) = new_mock_session(peer_id, eth_version);
3475        tx_manager.peers.insert(peer_id, mock_peer_metadata);
3476
3477        let mut tx_factory = MockTransactionFactory::default();
3478
3479        let valid_known_tx = tx_factory.create_eip1559();
3480        let known_tx_signed: Arc<ValidPoolTransaction<MockTransaction>> = Arc::new(valid_known_tx);
3481
3482        let known_tx_hash = *known_tx_signed.hash();
3483        let known_tx_type_byte = known_tx_signed.transaction.tx_type();
3484        let known_tx_size = known_tx_signed.encoded_length();
3485
3486        let unknown_tx_hash = B256::random();
3487        let unknown_tx_type_byte = 0xff_u8;
3488        let unknown_tx_size = 150;
3489
3490        let announcement_msg = NewPooledTransactionHashes::Eth68(NewPooledTransactionHashes68 {
3491            types: vec![known_tx_type_byte, unknown_tx_type_byte],
3492            sizes: vec![known_tx_size, unknown_tx_size],
3493            hashes: vec![known_tx_hash, unknown_tx_hash],
3494        });
3495
3496        tx_manager.on_new_pooled_transaction_hashes(peer_id, announcement_msg);
3497
3498        poll_fn(|cx| {
3499            let _ = tx_manager.poll_unpin(cx);
3500            Poll::Ready(())
3501        })
3502        .await;
3503
3504        let mut requested_hashes_in_getpooled = B256Set::default();
3505        let mut unexpected_request_received = false;
3506
3507        match tokio::time::timeout(std::time::Duration::from_millis(200), mock_session_rx.recv())
3508            .await
3509        {
3510            Ok(Some(PeerRequest::GetPooledTransactions { request, response: tx_response_ch })) => {
3511                let GetPooledTransactions(hashes) = request;
3512                for hash in hashes {
3513                    requested_hashes_in_getpooled.insert(hash);
3514                }
3515                let _ = tx_response_ch.send(Ok(PooledTransactions(vec![])));
3516            }
3517            Ok(Some(other_request)) => {
3518                tracing::error!(?other_request, "Received unexpected PeerRequest type");
3519                unexpected_request_received = true;
3520            }
3521            Ok(None) => tracing::info!("Mock session channel closed or no request received."),
3522            Err(_timeout_err) => {
3523                tracing::info!("Timeout: No GetPooledTransactions request received.")
3524            }
3525        }
3526
3527        assert!(
3528            requested_hashes_in_getpooled.contains(&known_tx_hash),
3529            "Should have requested the known EIP-1559 transaction. Requested: {requested_hashes_in_getpooled:?}"
3530        );
3531        assert!(
3532            !requested_hashes_in_getpooled.contains(&unknown_tx_hash),
3533            "Should NOT have requested the unknown transaction type. Requested: {requested_hashes_in_getpooled:?}"
3534        );
3535        assert!(
3536            !unexpected_request_received,
3537            "An unexpected P2P request was received by the mock peer."
3538        );
3539
3540        network_service_handle.abort();
3541    }
3542}