1use alloy_consensus::{constants::EIP4844_TX_TYPE_ID, transaction::TxHashRef};
4use rayon::iter::{IntoParallelIterator, ParallelIterator};
5use smallvec::SmallVec;
6
7pub mod announcement;
9
10pub mod config;
12pub mod constants;
14pub mod fetcher;
16pub 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
98pub type PoolImportFuture =
102 Pin<Box<dyn Future<Output = Vec<PoolResult<AddedTransactionOutcome>>> + Send + 'static>>;
103
104#[derive(Debug, Clone)]
112pub struct TransactionsHandle<N: NetworkPrimitives = EthNetworkPrimitives> {
113 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 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 pub fn propagate(&self, hash: TxHash) {
134 self.send(TransactionsCommand::PropagateHash(hash))
135 }
136
137 pub fn propagate_hash_to(&self, hash: TxHash, peer: PeerId) {
141 self.propagate_hashes_to(Some(hash), peer)
142 }
143
144 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 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 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 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 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 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 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 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 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#[derive(Debug)]
320#[must_use = "Manager does nothing unless polled."]
321pub struct TransactionsManager<Pool, N: NetworkPrimitives = EthNetworkPrimitives> {
322 pool: Pool,
324 sender_recovery_cache: Option<SenderRecoveryCache>,
326 network: NetworkHandle<N>,
328 network_events: EventStream<NetworkEvent<PeerRequest<N>>>,
332 transaction_fetcher: TransactionFetcher<N>,
334 transactions_by_peers: B256Map<SmallVec<[PeerId; 1]>>,
339 announcement_hashes: B256Set,
341 pool_imports: FuturesUnordered<PoolImportFuture>,
353 pending_pool_imports_info: PendingPoolImportsInfo,
355 bad_imports: LruCache<TxHash, FbBuildHasher<32>>,
357 peers: HashMap<PeerId, PeerMetadata<N>, FbBuildHasher<64>>,
359 command_tx: mpsc::UnboundedSender<TransactionsCommand<N>>,
363 command_rx: UnboundedReceiverStream<TransactionsCommand<N>>,
368 pending_transactions: mpsc::Receiver<TxHash>,
377 transaction_events: MemoryBoundedReceiver<NetworkTransactionEvent<N>>,
379 config: TransactionsManagerConfig,
381 policies: NetworkPolicies<N>,
383 metrics: TransactionsManagerMetrics,
385 announced_tx_types_metrics: AnnouncedTxTypesMetrics,
387}
388
389impl<Pool: TransactionPool, N: NetworkPrimitives> TransactionsManager<Pool, N> {
390 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 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 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 pub fn handle(&self) -> TransactionsHandle<N> {
465 TransactionsHandle { manager_tx: self.command_tx.clone() }
466 }
467
468 pub fn with_sender_recovery_cache(mut self, cache: SenderRecoveryCache) -> Self {
470 self.sender_recovery_cache = Some(cache);
471 self
472 }
473
474 fn has_capacity_for_pending_pool_imports(&self) -> bool {
480 self.remaining_pool_import_capacity() > 0
481 }
482
483 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 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 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 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 fn on_good_import(&mut self, hash: TxHash) {
534 self.transactions_by_peers.remove(&hash);
535 }
536
537 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 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 !err.is_bad_transaction() || self.network.is_syncing() {
587 return
588 }
589 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 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 metrics.duration_poll_tx_manager.set(start.elapsed().as_secs_f64());
629 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 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 fn on_new_pooled_transaction_hashes(
657 &mut self,
658 peer_id: PeerId,
659 msg: NewPooledTransactionHashes,
660 ) {
661 if self.network.is_initially_syncing() {
663 return
664 }
665 if self.network.tx_gossip_disabled() {
666 return
667 }
668
669 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 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 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 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 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 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 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 fn on_new_pending_transactions(&mut self, hashes: Vec<TxHash>) {
823 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 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 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 full_transactions.extend(to_propagate);
856 } else {
857 for tx in to_propagate {
860 if !peer.seen_transactions.contains(tx.tx_hash()) {
861 full_transactions.push(&tx);
863 }
864 }
865 }
866
867 if full_transactions.is_empty() {
868 return None
870 }
871
872 let PropagateTransactions { pooled, full } = full_transactions.build();
873
874 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 peer.seen_transactions.insert(hash);
880 }
881
882 self.network.send_transactions_hashes(peer_id, new_pooled_hashes);
884 }
885
886 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 peer.seen_transactions.insert(*hash);
892 }
893
894 self.network.send_broadcast_pool_transactions(peer_id, new_full_transactions);
896 }
897
898 self.metrics.propagated_transactions.increment(propagated.len() as u64);
900
901 Some(propagated)
902 }
903
904 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 let propagated = {
918 let Some(peer) = self.peers.get_mut(&peer_id) else {
919 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 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 hashes.push(&tx);
938 }
939 }
940 }
941
942 let new_pooled_hashes = hashes.build();
943
944 if new_pooled_hashes.is_empty() {
945 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 self.network.send_transactions_hashes(peer_id, new_pooled_hashes);
958
959 self.metrics.propagated_transactions.increment(propagated.len() as u64);
961
962 propagated
963 };
964
965 self.pool.on_propagated(propagated);
967 }
968
969 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 let max_num_full = self.config.propagation_mode.full_peer_count(self.peers.len());
987
988 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 continue
994 }
995
996 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 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 for tx in &to_propagate {
1016 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 if let Some(mut new_pooled_hashes) = pooled {
1032 if new_pooled_hashes.len() >
1036 SOFT_LIMIT_COUNT_HASHES_IN_NEW_POOLED_TRANSACTIONS_BROADCAST_MESSAGE
1037 {
1038 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 self.network.send_transactions_hashes(*peer_id, new_pooled_hashes);
1059 }
1060
1061 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 self.network.send_broadcast_pool_transactions(*peer_id, new_full_transactions);
1071 }
1072 }
1073
1074 self.metrics.propagated_transactions.increment(propagated.len() as u64);
1076
1077 propagated
1078 }
1079
1080 fn propagate_all(&mut self, hashes: Vec<TxHash>) {
1085 if self.peers.is_empty() {
1086 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 self.pool.on_propagated(propagated);
1096 }
1097
1098 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 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 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 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 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 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 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 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 if !self.policies.propagation_policy().can_propagate(peer) {
1218 return
1219 }
1220
1221 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 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 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 self.handle_peer_session(info, messages);
1251 }
1252 NetworkEvent::Peer(PeerEvent::SessionEstablished(info)) => {
1253 let peer_id = info.peer_id;
1254 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 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 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 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 fn import_transactions(
1328 &mut self,
1329 peer_id: PeerId,
1330 transactions: PooledTransactions<N::PooledTransaction>,
1331 source: TransactionSource,
1332 ) {
1333 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 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 self.transaction_fetcher
1373 .on_transactions_received(transactions.iter().map(|tx| tx.tx_hash()));
1374
1375 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 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 let mut has_bad_transactions = false;
1407
1408 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 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 for tx in &new_txs {
1468 self.transactions_by_peers.insert(*tx.hash(), smallvec::smallvec![peer_id]);
1469 }
1470
1471 if !new_txs.is_empty() {
1474 let pool = self.pool.clone();
1475 let metric_pending_pool_imports = self.metrics.pending_pool_imports.clone();
1477 metric_pending_pool_imports.increment(new_txs.len() as f64);
1478
1479 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 metric_pending_pool_imports.decrement(added as f64);
1493 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 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 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
1552impl<
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 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 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 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 let fetcher_paused_at_capacity = !this.has_capacity_for_pending_pool_imports();
1633
1634 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 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 true
1667 } else {
1668 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 duration_metered_exec!(
1686 {
1687 let budget = this.remaining_pool_import_capacity();
1690 if budget > 0 && this.transaction_fetcher.dispatch(&this.peers, budget) > 0 {
1691 maybe_more_tx_fetch_events = true;
1693 }
1694 },
1695 poll_durations.acc_pending_fetch
1696 );
1697
1698 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 let resume_fetcher =
1712 fetcher_paused_at_capacity && this.has_capacity_for_pending_pool_imports();
1713
1714 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 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#[derive(Debug, Copy, Clone, Eq, PartialEq)]
1738enum PropagationMode {
1739 Basic,
1743 Forced,
1748}
1749
1750impl PropagationMode {
1751 const fn is_forced(self) -> bool {
1753 matches!(self, Self::Forced)
1754 }
1755}
1756
1757#[derive(Debug, Clone)]
1759struct PropagateTransaction {
1760 is_broadcastable_in_full: bool,
1761 propagation_size: usize,
1768 transaction: LazyEncodedTransaction,
1769}
1770
1771impl PropagateTransaction {
1772 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 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 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#[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#[derive(Debug, Clone)]
1866enum PropagateTransactionsBuilder {
1867 Pooled(PooledTransactionsHashesBuilder),
1868 Full(FullTransactionsBuilder),
1869}
1870
1871impl PropagateTransactionsBuilder {
1872 fn pooled(version: EthVersion, capacity: usize) -> Self {
1875 Self::Pooled(PooledTransactionsHashesBuilder::with_capacity(version, capacity))
1876 }
1877
1878 fn full(version: EthVersion, capacity: usize) -> Self {
1881 Self::Full(FullTransactionsBuilder::with_capacity(version, capacity))
1882 }
1883
1884 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 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 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
1913struct PropagateTransactions {
1915 pooled: Option<NewPooledTransactionHashes>,
1917 full: Option<BroadcastPoolTransactions>,
1919}
1920
1921#[derive(Debug, Clone)]
1926struct FullTransactionsBuilder {
1927 total_size: usize,
1929 transactions: Vec<LazyEncodedTransaction>,
1931 pooled: PooledTransactionsHashesBuilder,
1933}
1934
1935impl FullTransactionsBuilder {
1936 fn new(version: EthVersion) -> Self {
1938 Self {
1939 total_size: 0,
1940 pooled: PooledTransactionsHashesBuilder::new(version),
1941 transactions: vec![],
1942 }
1943 }
1944
1945 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 fn is_empty(&self) -> bool {
1959 self.transactions.is_empty() && self.pooled.is_empty()
1960 }
1961
1962 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 fn extend(&mut self, txs: impl IntoIterator<Item = PropagateTransaction>) {
1972 for tx in txs {
1973 self.push(&tx)
1974 }
1975 }
1976
1977 fn push(&mut self, transaction: &PropagateTransaction) {
1987 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 self.pooled.push(transaction);
2006 return
2007 }
2008
2009 self.total_size = new_size;
2010 self.transactions.push(transaction.shared());
2011 }
2012}
2013
2014#[derive(Debug, Clone)]
2017enum PooledTransactionsHashesBuilder {
2018 Eth66(NewPooledTransactionHashes66),
2019 Eth68(NewPooledTransactionHashes68),
2020 Eth72(NewPooledTransactionHashes72),
2021}
2022
2023impl PooledTransactionsHashesBuilder {
2026 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 msg.cell_mask = Some(NewPooledTransactionHashes72::ALL_CELLS_MASK);
2043 }
2044 }
2045 }
2046 }
2047
2048 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 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 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 msg.cell_mask = Some(NewPooledTransactionHashes72::ALL_CELLS_MASK);
2089 }
2090 }
2091 }
2092 }
2093
2094 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 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#[derive(Debug)]
2139enum TransactionSource {
2140 Broadcast,
2142 Response {
2144 version: EthVersion,
2146 client_version: Arc<str>,
2147 },
2148}
2149
2150impl TransactionSource {
2153 const fn is_broadcast(&self) -> bool {
2155 matches!(self, Self::Broadcast)
2156 }
2157}
2158
2159#[derive(Debug)]
2161pub struct PeerMetadata<N: NetworkPrimitives = EthNetworkPrimitives> {
2162 seen_transactions: LruCache<TxHash, FbBuildHasher<32>>,
2166 request_tx: PeerRequestSender<PeerRequest<N>>,
2168 version: EthVersion,
2170 client_version: Arc<str>,
2172 peer_kind: PeerKind,
2174}
2175
2176impl<N: NetworkPrimitives> PeerMetadata<N> {
2177 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 pub const fn request_tx(&self) -> &PeerRequestSender<PeerRequest<N>> {
2199 &self.request_tx
2200 }
2201
2202 pub const fn seen_transactions_mut(&mut self) -> &mut LruCache<TxHash, FbBuildHasher<32>> {
2204 &mut self.seen_transactions
2205 }
2206
2207 pub const fn version(&self) -> EthVersion {
2209 self.version
2210 }
2211
2212 pub fn client_version(&self) -> &str {
2214 &self.client_version
2215 }
2216
2217 pub const fn peer_kind(&self) -> PeerKind {
2219 self.peer_kind
2220 }
2221}
2222
2223#[derive(Debug)]
2225enum TransactionsCommand<N: NetworkPrimitives = EthNetworkPrimitives> {
2226 PropagateHash(B256),
2228 PropagateHashesTo(Vec<B256>, PeerId),
2230 GetActivePeers(oneshot::Sender<HashSet<PeerId>>),
2232 PropagateTransactionsTo(Vec<TxHash>, PeerId),
2234 PropagateTransactions(Vec<TxHash>),
2236 BroadcastTransactions(Vec<PropagateTransaction>),
2238 GetTransactionHashes { peers: Vec<PeerId>, tx: oneshot::Sender<HashMap<PeerId, B256Set>> },
2240 GetPeerSender {
2242 peer_id: PeerId,
2243 peer_request_sender: oneshot::Sender<Option<PeerRequestSender<PeerRequest<N>>>>,
2244 },
2245}
2246
2247#[derive(Debug)]
2249pub enum NetworkTransactionEvent<N: NetworkPrimitives = EthNetworkPrimitives> {
2250 IncomingTransactions {
2254 peer_id: PeerId,
2256 msg: Transactions<N::BroadcastedTransaction>,
2258 },
2259 IncomingPooledTransactionHashes {
2261 peer_id: PeerId,
2263 msg: NewPooledTransactionHashes,
2265 },
2266 GetPooledTransactions {
2268 peer_id: PeerId,
2270 request: GetPooledTransactions,
2272 response: oneshot::Sender<RequestResult<PooledTransactions<N::PooledTransaction>>>,
2274 },
2275 GetTransactionsHandle(oneshot::Sender<Option<TransactionsHandle<N>>>),
2277}
2278
2279#[derive(Debug)]
2281pub struct PendingPoolImportsInfo {
2282 pending_pool_imports: Arc<AtomicUsize>,
2284 max_pending_pool_imports: usize,
2286}
2287
2288impl PendingPoolImportsInfo {
2289 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 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 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 network_handle.update_sync_state(SyncState::Syncing);
2531 assert!(NetworkInfo::is_syncing(&network_handle));
2532 assert!(NetworkInfo::is_initially_syncing(&network_handle));
2533
2534 let mut established = listener0.take(2);
2536 while let Some(ev) = established.next().await {
2537 match ev {
2538 NetworkEvent::Peer(PeerEvent::SessionEstablished(info)) => {
2539 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 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 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 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 transactions.on_network_event(ev);
2608 }
2609 NetworkEvent::Peer(PeerEvent::PeerAdded(_peer_id)) => {}
2610 _ => {
2611 error!("unexpected event {ev:?}")
2612 }
2613 }
2614 }
2615 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 #[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 .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 poll_fn(|cx| {
2701 let _ = tx_manager.poll_unpin(cx);
2702 Poll::Ready(())
2703 })
2704 .await;
2705
2706 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 response
2724 .send(Ok(PooledTransactions(message)))
2725 .expect("should send peer_1 response to tx manager");
2726
2727 poll_fn(|cx| {
2729 let _ = tx_manager.poll_unpin(cx);
2730 Poll::Ready(())
2731 })
2732 .await;
2733
2734 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 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 transactions.on_network_event(ev);
2777 }
2778 NetworkEvent::Peer(PeerEvent::PeerAdded(_peer_id)) => {}
2779 ev => {
2780 error!("unexpected event {ev:?}")
2781 }
2782 }
2783 }
2784 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 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 assert!(!tx_manager.peers.contains_key(&peer_id));
2840 assert!(tx_manager.transaction_fetcher.queued_hashes(&peer_id).is_empty());
2841 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 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 #[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 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 assert!(tx_fetcher.is_idle(&peer_id_1));
3057 assert_eq!(tx_fetcher.num_inflight_requests(), 0);
3058
3059 assert_eq!(tx_fetcher.dispatch(&tx_manager.peers, usize::MAX), 1);
3061
3062 assert_eq!(tx_fetcher.num_pending_hashes(), 0);
3063 assert!(!tx_fetcher.is_idle(&peer_id_1));
3065 assert_eq!(tx_fetcher.num_inflight_requests(), 1);
3066
3067 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
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 assert!(tx_fetcher.is_idle(&peer_id));
3092 assert_eq!(tx_fetcher.num_inflight_requests(), 0);
3093 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 #[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 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 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 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 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 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 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 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 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 network.handle().update_sync_state(SyncState::Idle);
3291
3292 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 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 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 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 tx_manager.config.propagation_mode = TransactionPropagationMode::Max(0);
3372
3373 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 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 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 network.handle().update_sync_state(SyncState::Syncing);
3411 assert!(NetworkInfo::is_initially_syncing(&network.handle()));
3412
3413 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}