1use super::{
50 announcement::{AnnouncedTransaction, TransactionMetadata},
51 config::TransactionFetcherConfig,
52 constants::{
53 tx_fetcher::{
54 AVERAGE_BYTE_SIZE_TX_ENCODED, MAX_COUNT_CANDIDATE_PEERS_PER_HASH,
55 MAX_COUNT_EAGER_CANDIDATE_PEERS_PER_HASH, MAX_FETCH_ATTEMPTS_PER_HASH,
56 },
57 SOFT_LIMIT_BYTE_SIZE_POOLED_TRANSACTIONS_RESPONSE,
58 SOFT_LIMIT_COUNT_HASHES_IN_GET_POOLED_TRANSACTIONS_REQUEST,
59 },
60 PeerMetadata,
61};
62use crate::metrics::TransactionFetcherMetrics;
63use alloy_consensus::transaction::PooledTransaction;
64use alloy_primitives::{
65 map::{B256Map, B256Set, Entry, FbBuildHasher, HashMap},
66 TxHash,
67};
68use futures::{stream::FuturesUnordered, Future, FutureExt, Stream, StreamExt};
69use reth_eth_wire::{EthVersion, GetPooledTransactions, PooledTransactions};
70use reth_eth_wire_types::{EthNetworkPrimitives, NetworkPrimitives};
71use reth_network_api::PeerRequest;
72use reth_network_p2p::error::{RequestError, RequestResult};
73use reth_network_peers::PeerId;
74use reth_primitives_traits::SignedTransaction;
75use smallvec::SmallVec;
76use std::{
77 collections::{BinaryHeap, VecDeque},
78 pin::Pin,
79 sync::Arc,
80 task::{ready, Context, Poll},
81};
82use tokio::sync::{mpsc::error::TrySendError, oneshot};
83use tracing::trace;
84
85const MAX_EVICTION_ATTEMPTS: usize = 8;
87
88#[derive(Debug)]
94pub struct TransactionFetcher<N: NetworkPrimitives = EthNetworkPrimitives> {
95 hashes: B256Map<TxEntry>,
97 order: VecDeque<(TxHash, u64)>,
101 peers: HashMap<PeerKey, PeerState>,
103 peer_keys: HashMap<PeerId, PeerKey, FbBuildHasher<64>>,
105 next_peer_key: u32,
108 next_request_id: u64,
110 next_generation: u64,
112 ready: VecDeque<PeerKey>,
114 inflight: FuturesUnordered<InflightRequest<N::PooledTransaction>>,
116 num_fetching: usize,
118 scratch_requested: B256Set,
120 scratch_delivered: B256Set,
122 scratch_queue: Vec<(TxHash, u64)>,
124 config: TransactionFetcherConfig,
126 metrics: TransactionFetcherMetrics,
127}
128
129impl<N: NetworkPrimitives> TransactionFetcher<N> {
130 pub fn new(config: TransactionFetcherConfig) -> Self {
132 let metrics = TransactionFetcherMetrics::default();
133 metrics.capacity_inflight_requests.increment(config.max_inflight_requests as u64);
134
135 Self {
136 hashes: Default::default(),
137 order: Default::default(),
138 peers: Default::default(),
139 peer_keys: Default::default(),
140 next_peer_key: 0,
141 next_request_id: 0,
142 next_generation: 0,
143 ready: Default::default(),
144 inflight: Default::default(),
145 num_fetching: 0,
146 scratch_requested: Default::default(),
147 scratch_delivered: Default::default(),
148 scratch_queue: Default::default(),
149 config,
150 metrics,
151 }
152 }
153
154 pub const fn config(&self) -> &TransactionFetcherConfig {
156 &self.config
157 }
158
159 pub fn num_hashes(&self) -> usize {
161 self.hashes.len()
162 }
163
164 pub fn num_pending_hashes(&self) -> usize {
166 self.hashes.len().saturating_sub(self.num_fetching)
167 }
168
169 pub const fn num_fetching_hashes(&self) -> usize {
171 self.num_fetching
172 }
173
174 pub fn num_inflight_requests(&self) -> usize {
176 self.inflight.len()
177 }
178
179 pub fn is_idle(&self, peer_id: &PeerId) -> bool {
181 self.peer_keys
182 .get(peer_id)
183 .and_then(|key| self.peers.get(key))
184 .is_none_or(|peer| peer.inflight == 0)
185 }
186
187 pub fn on_announcement(
194 &mut self,
195 peer_id: PeerId,
196 announcement: impl IntoIterator<Item = AnnouncedTransaction>,
197 ) {
198 let key = self.peer_key(peer_id);
199 let max_per_peer = self.config.max_announced_hashes_per_peer as usize;
200 let max_total = self.config.max_capacity_cache_txns_pending_fetch as usize;
201
202 let mut dropped_peer_limit = 0u64;
203 let mut evicted_at_capacity = 0u64;
204 let mut dropped_at_capacity = 0u64;
205
206 let Some(peer) = self.peers.get(&key) else { return };
209 let mut tracked = peer.tracked;
210 let mut queue = std::mem::take(&mut self.scratch_queue);
211 queue.clear();
212 let mut queue_start = 0;
214 let mut eviction = None;
217 let mut retired_candidates = DeferredCandidateRemovals::default();
220
221 for AnnouncedTransaction { hash, metadata } in announcement {
222 let size = announced_size(metadata);
223 let at_capacity = self.hashes.len() >= max_total;
224
225 let (eager, generation) = match self.hashes.entry(hash) {
227 Entry::Occupied(mut occupied) => {
228 let entry = occupied.get_mut();
229 if let Some(candidate) = entry.candidate_mut(key) {
230 candidate.set_size(size);
232 continue
233 }
234 if tracked >= max_per_peer {
235 dropped_peer_limit += 1;
236 continue
237 }
238 if entry.candidates.len() == MAX_COUNT_CANDIDATE_PEERS_PER_HASH {
239 let replace = (MAX_COUNT_EAGER_CANDIDATE_PEERS_PER_HASH..
242 entry.candidates.len())
243 .find(|&idx| {
244 entry.fetching_by.is_none_or(|(fetching, _)| {
245 entry.candidates[idx].peer != fetching
246 })
247 })
248 .expect("only one candidate can be fetching");
249 let removed = entry.candidates.remove(replace);
250 retired_candidates.record_removed(removed.peer, &mut self.peers);
251 }
252 let eager = entry.candidates.len() < MAX_COUNT_EAGER_CANDIDATE_PEERS_PER_HASH ||
255 (entry.fetching_by.is_none() &&
256 !entry.candidates.iter().any(Candidate::is_queued));
257 let candidate = if eager {
258 Candidate::queued(key, size)
259 } else {
260 Candidate::unqueued(key, size)
261 };
262 entry.candidates.push(candidate);
263 (eager, entry.generation)
264 }
265 Entry::Vacant(vacant) => {
266 let generation = self.next_generation;
267 self.next_generation += 1;
268 if tracked >= max_per_peer {
269 dropped_peer_limit += 1;
270 continue
271 }
272 if at_capacity {
273 retired_candidates.flush(&mut self.peers);
274 if let Some(peer) = self.peers.get_mut(&key) {
276 peer.tracked = tracked;
277 }
278 if self.evict_pending(key, &queue, &mut queue_start, &mut eviction) {
279 evicted_at_capacity += 1;
280 tracked = self.peers.get(&key).map_or(tracked, |peer| peer.tracked);
281 } else {
282 dropped_at_capacity += 1;
283 continue
284 }
285 self.hashes.insert(hash, TxEntry::new(key, size, generation));
286 } else {
287 vacant.insert(TxEntry::new(key, size, generation));
288 }
289 self.record_order(hash, generation);
290 (true, generation)
291 }
292 };
293
294 tracked += 1;
295 if eager {
296 queue.push((hash, generation));
297 }
298 }
299
300 retired_candidates.flush(&mut self.peers);
301 let queued = queue.len() - queue_start;
302 if let Some(peer) = self.peers.get_mut(&key) {
303 peer.tracked = tracked;
304 for &(hash, generation) in &queue[queue_start..] {
305 peer.push_queue(
306 &self.hashes,
307 key,
308 (hash, generation),
309 QueuePosition::Back,
310 max_per_peer,
311 );
312 }
313 }
314 queue.clear();
315 self.scratch_queue = queue;
316
317 if dropped_peer_limit > 0 {
318 self.metrics.announced_hashes_dropped_peer_limit.increment(dropped_peer_limit);
319 }
320 if evicted_at_capacity > 0 {
321 self.metrics.hashes_evicted_at_capacity.increment(evicted_at_capacity);
322 }
323 if dropped_at_capacity > 0 {
324 self.metrics.announced_hashes_dropped_at_capacity.increment(dropped_at_capacity);
325 }
326
327 if queued > 0 {
328 trace!(target: "net::tx",
329 peer_id=format!("{peer_id:#}"),
330 queued,
331 dropped_peer_limit,
332 evicted_at_capacity,
333 dropped_at_capacity,
334 "queued announced hashes"
335 );
336 }
337 if queued > 0 || queue_start > 0 {
338 self.mark_ready(key, QueuePosition::Back);
339 }
340 }
341
342 pub fn dispatch(
353 &mut self,
354 peers: &HashMap<PeerId, PeerMetadata<N>, FbBuildHasher<64>>,
355 max_hashes_per_request: usize,
356 ) -> usize {
357 if max_hashes_per_request == 0 {
358 return 0
359 }
360 let max_inflight = self.config.max_inflight_requests as usize;
361 let mut sent = 0;
362 let mut retry = SmallVec::<[PeerKey; 4]>::new();
364
365 while self.inflight.len() < max_inflight {
366 let Some(key) = self.ready.pop_front() else { break };
367 let Some(peer) = self.peers.get_mut(&key) else { continue };
368 peer.ready = false;
369 if peer.inflight >= self.config.max_inflight_requests_per_peer {
370 continue
371 }
372 let peer_id = peer.peer_id;
373 let Some(session) = peers.get(&peer_id) else { continue };
375
376 let limit = max_hashes_per_request
377 .min(SOFT_LIMIT_COUNT_HASHES_IN_GET_POOLED_TRANSACTIONS_REQUEST);
378 let request_id = self.next_request_id;
379 self.next_request_id += 1;
380 let hashes = self.pack_request(key, request_id, limit);
381 if hashes.is_empty() {
382 continue
383 }
384
385 let (response, rx) = oneshot::channel();
386 let request = PeerRequest::GetPooledTransactions {
387 request: GetPooledTransactions(hashes.clone()),
388 response,
389 };
390
391 match session.request_tx().try_send(request) {
392 Ok(()) => {
393 trace!(target: "net::tx",
394 peer_id=format!("{peer_id:#}"),
395 hashes=hashes.len(),
396 "sending `GetPooledTransactions` request to peer's session"
397 );
398 if let Some(peer) = self.peers.get_mut(&key) {
399 peer.inflight += 1;
400 }
401 self.inflight.push(InflightRequest {
402 peer: key,
403 request_id,
404 peer_id,
405 version: session.version(),
406 client_version: session.client_version.clone(),
407 hashes,
408 response: rx,
409 });
410 sent += 1;
411 self.mark_ready(key, QueuePosition::Back);
413 }
414 Err(err) => {
415 self.unpack_request(key, request_id, hashes);
416 match err {
417 TrySendError::Full(_) => {
418 self.metrics.egress_peer_channel_full.increment(1);
419 retry.push(key);
420 }
421 TrySendError::Closed(_) => self.on_peer_disconnected(&peer_id),
422 }
423 }
424 }
425 }
426
427 for key in retry {
428 self.mark_ready(key, QueuePosition::Back);
429 }
430
431 sent
432 }
433
434 pub fn on_transactions_received<'a>(&mut self, hashes: impl IntoIterator<Item = &'a TxHash>) {
440 for hash in hashes {
441 self.remove_hash(hash);
442 }
443 }
444
445 pub fn on_peer_disconnected(&mut self, peer_id: &PeerId) {
451 let Some(key) = self.peer_keys.remove(peer_id) else { return };
452 if self.peers.remove(&key).is_none() {
453 return
454 }
455
456 let mut dropped = 0u64;
460 let mut requeue = Vec::new();
461 self.hashes.retain(|hash, entry| {
462 let before = entry.candidates.len();
463 entry.candidates.retain(|candidate| candidate.peer != key);
464 if entry.candidates.len() == before || entry.fetching_by.is_some() {
465 return true
466 }
467 if entry.candidates.is_empty() {
468 dropped += 1;
469 return false
470 }
471 if !entry.candidates.iter().any(|candidate| candidate.is_queued()) {
473 requeue.push((entry.generation, *hash, entry.unqueued_candidates()));
474 }
475 true
476 });
477
478 let mut requeued = SmallVec::<[PeerKey; MAX_COUNT_CANDIDATE_PEERS_PER_HASH]>::new();
479 requeue.sort_unstable_by_key(|(generation, _, _)| std::cmp::Reverse(*generation));
481 for (_, hash, targets) in requeue {
482 self.requeue(hash, targets, &mut requeued);
483 }
484
485 if dropped > 0 {
486 self.metrics.hashes_dropped_no_candidate_peers.increment(dropped);
487 }
488 for key in requeued {
490 self.mark_ready(key, QueuePosition::Front);
491 }
492 }
493
494 pub fn update_metrics(&self) {
496 self.metrics.inflight_transaction_requests.set(self.inflight.len() as f64);
497 self.metrics.hashes_inflight_transaction_requests.set(self.num_fetching as f64);
498 self.metrics.hashes_pending_fetch.set(self.num_pending_hashes() as f64);
499 }
500
501 fn peer_key(&mut self, peer_id: PeerId) -> PeerKey {
503 *self.peer_keys.entry(peer_id).or_insert_with(|| {
504 let key = PeerKey(self.next_peer_key);
505 self.next_peer_key += 1;
506 self.peers.insert(key, PeerState::new(peer_id));
507 key
508 })
509 }
510
511 fn mark_ready(&mut self, key: PeerKey, position: QueuePosition) {
514 if let Some(peer) = self.peers.get_mut(&key) &&
515 !peer.ready &&
516 !peer.queue.is_empty() &&
517 peer.inflight < self.config.max_inflight_requests_per_peer
518 {
519 peer.ready = true;
520 match position {
521 QueuePosition::Front => self.ready.push_front(key),
522 QueuePosition::Back => self.ready.push_back(key),
523 }
524 }
525 }
526
527 fn enqueue(&mut self, key: PeerKey, hash: TxHash, position: QueuePosition) {
530 let max_per_peer = self.config.max_announced_hashes_per_peer as usize;
531 let Some(peer) = self.peers.get_mut(&key) else { return };
532 let Some(entry) = self.hashes.get_mut(&hash) else { return };
533 let generation = entry.generation;
534 if let Some(candidate) = entry.candidate_mut(key) {
535 candidate.set_queued(true);
536 }
537 peer.push_queue(&self.hashes, key, (hash, generation), position, max_per_peer);
538 }
539
540 fn record_order(&mut self, hash: TxHash, generation: u64) {
545 let max_len = 2 * self.config.max_capacity_cache_txns_pending_fetch as usize;
546 if self.order.len() >= max_len {
547 self.order.retain(|(hash, generation)| {
548 self.hashes.get(hash).is_some_and(|entry| entry.generation == *generation)
549 });
550 }
551 self.order.push_back((hash, generation));
552 }
553
554 fn evict_pending(
567 &mut self,
568 announcer: PeerKey,
569 queued: &[(TxHash, u64)],
570 queued_start: &mut usize,
571 eviction: &mut Option<EvictionState>,
572 ) -> bool {
573 let eviction = eviction.get_or_insert_with(|| EvictionState {
574 peers: self
575 .peers
576 .iter()
577 .filter(|(key, peer)| **key != announcer && peer.tracked > 0)
578 .map(|(key, peer)| (peer.tracked, *key))
579 .collect(),
580 exclusive: Default::default(),
581 });
582 let heap = &mut eviction.peers;
583 let (mut tracked, mut key) = loop {
584 let Some(&(previous, key)) = heap.peek() else { break (0, announcer) };
585 let current = self.peers.get(&key).map_or(0, |peer| peer.tracked);
586 if current == previous {
587 break (current, key)
588 }
589 heap.pop();
590 if current > 0 {
591 heap.push((current, key));
592 }
593 };
594 if let Some(peer) = self.peers.get(&announcer) &&
595 peer.tracked >= tracked
596 {
597 key = announcer;
598 tracked = peer.tracked;
599 }
600 if tracked == 0 {
601 return self.evict_oldest_pending()
602 }
603 if self.evict_oldest_pending_of(key, &mut eviction.exclusive) {
604 return true
605 }
606 if key == announcer {
607 while *queued_start < queued.len() {
608 let (hash, generation) = queued[*queued_start];
609 *queued_start += 1;
610 let Some(entry) =
611 self.hashes.get(&hash).filter(|entry| entry.generation == generation)
612 else {
613 continue
614 };
615 if entry.candidates.len() == 1 && entry.fetching_by.is_none() {
616 self.remove_hash(&hash);
617 return true
618 }
619 self.enqueue(key, hash, QueuePosition::Back);
621 }
622 }
623 self.evict_oldest_pending()
624 }
625
626 fn evict_oldest_pending_of(
632 &mut self,
633 key: PeerKey,
634 exclusive: &mut HashMap<PeerKey, VecDeque<(TxHash, u64)>>,
635 ) -> bool {
636 let candidates = exclusive.entry(key).or_insert_with(|| {
637 self.peers
638 .get(&key)
639 .into_iter()
640 .flat_map(|peer| &peer.queue)
641 .filter(|(hash, generation)| {
642 self.hashes.get(hash).is_some_and(|entry| {
643 entry.generation == *generation &&
644 entry.candidates.len() == 1 &&
645 entry.has_candidate(key) &&
646 entry.fetching_by.is_none()
647 })
648 })
649 .copied()
650 .collect()
651 });
652 while let Some((hash, generation)) = candidates.pop_front() {
653 if self.hashes.get(&hash).is_some_and(|entry| {
654 entry.generation == generation &&
655 entry.candidates.len() == 1 &&
656 entry.has_candidate(key) &&
657 entry.fetching_by.is_none()
658 }) {
659 self.remove_hash(&hash);
660 return true
661 }
662 }
663 false
664 }
665
666 fn evict_oldest_pending(&mut self) -> bool {
671 let mut index = 0;
672 let mut attempts = 0;
673 while attempts < MAX_EVICTION_ATTEMPTS {
674 let Some(&(hash, generation)) = self.order.get(index) else { break };
675 let Some(entry) = self.hashes.get(&hash).filter(|entry| entry.generation == generation)
676 else {
677 self.order.remove(index);
678 continue
679 };
680 attempts += 1;
681 if entry.fetching_by.is_some() {
682 index += 1;
683 continue
684 }
685 self.order.remove(index);
686 self.remove_hash(&hash);
687 return true
688 }
689 false
690 }
691
692 fn pack_request(&mut self, key: PeerKey, request_id: u64, max_hashes: usize) -> Vec<TxHash> {
698 let max_bytes =
699 self.config.soft_limit_byte_size_pooled_transactions_response_on_pack_request;
700 let Some(peer) = self.peers.get_mut(&key) else { return Vec::new() };
701
702 let mut hashes = Vec::with_capacity(peer.queue.len().min(max_hashes));
703 let mut bytes = 0usize;
704
705 while let Some((hash, generation)) = peer.queue.pop_front() {
706 let Some(entry) =
709 self.hashes.get_mut(&hash).filter(|entry| entry.generation == generation)
710 else {
711 continue
712 };
713 let Some(candidate) = entry.candidate_mut(key) else { continue };
714 candidate.set_queued(false);
715 let size = candidate.request_size();
716
717 if entry.fetching_by.is_some() {
719 continue
720 }
721
722 if !hashes.is_empty() && bytes.saturating_add(size) > max_bytes {
723 if let Some(candidate) = entry.candidate_mut(key) {
724 candidate.set_queued(true);
725 }
726 peer.queue.push_front((hash, generation));
727 break
728 }
729
730 entry.fetching_by = Some((key, request_id));
731 entry.attempts += 1;
732 self.num_fetching += 1;
733 bytes = bytes.saturating_add(size);
734 hashes.push(hash);
735
736 if hashes.len() >= max_hashes {
737 break
738 }
739 }
740
741 hashes
742 }
743
744 fn unpack_request(&mut self, key: PeerKey, request_id: u64, hashes: Vec<TxHash>) {
747 for hash in hashes.into_iter().rev() {
748 if let Some(entry) = self.hashes.get_mut(&hash) &&
749 entry.fetching_by == Some((key, request_id))
750 {
751 entry.fetching_by = None;
752 entry.attempts -= 1;
753 self.num_fetching -= 1;
754 self.enqueue(key, hash, QueuePosition::Front);
755 }
756 }
757 }
758
759 fn requeue(
762 &mut self,
763 hash: TxHash,
764 targets: SmallVec<[PeerKey; MAX_COUNT_CANDIDATE_PEERS_PER_HASH]>,
765 requeued: &mut SmallVec<[PeerKey; MAX_COUNT_CANDIDATE_PEERS_PER_HASH]>,
766 ) {
767 for key in targets {
768 self.enqueue(key, hash, QueuePosition::Front);
769 if !requeued.contains(&key) {
770 requeued.push(key);
771 }
772 }
773 }
774
775 fn remove_hash(&mut self, hash: &TxHash) {
777 let Some(entry) = self.hashes.remove(hash) else { return };
778 if entry.fetching_by.is_some() {
779 self.num_fetching -= 1;
780 }
781 for candidate in &entry.candidates {
782 if let Some(peer) = self.peers.get_mut(&candidate.peer) {
783 peer.tracked = peer.tracked.saturating_sub(1);
784 }
785 }
786 }
787
788 fn on_resolved(
790 &mut self,
791 resolved: ResolvedRequest<N::PooledTransaction>,
792 ) -> FetchEvent<N::PooledTransaction> {
793 let ResolvedRequest {
794 peer: key,
795 request_id,
796 peer_id,
797 version,
798 client_version,
799 hashes: requested,
800 result,
801 } = resolved;
802
803 if let Some(peer) = self.peers.get_mut(&key) {
804 peer.inflight = peer.inflight.saturating_sub(1);
805 }
806
807 let mut delivered = std::mem::take(&mut self.scratch_delivered);
808 delivered.clear();
809
810 let outcome = match result {
811 Ok(mut transactions) => {
812 let mut requested_set = std::mem::take(&mut self.scratch_requested);
813 requested_set.clear();
814 requested_set.extend(requested.iter().copied());
815 let unsolicited =
816 verify_response(&mut transactions, &requested_set, &mut delivered);
817 self.scratch_requested = requested_set;
818 Ok((transactions, unsolicited))
819 }
820 Err(error) => Err(error),
821 };
822
823 let timed_out = matches!(&outcome, Err(RequestError::Timeout));
824 self.on_delivery(key, request_id, &requested, &delivered, timed_out);
825 self.scratch_delivered = delivered;
826 self.mark_ready(key, QueuePosition::Back);
827
828 match outcome {
829 Ok((transactions, unsolicited)) => {
830 if unsolicited > 0 {
831 self.metrics.unsolicited_transactions.increment(unsolicited as u64);
832 trace!(target: "net::tx",
833 peer_id=format!("{peer_id:#}"),
834 unsolicited,
835 "received transactions in `PooledTransactions` response that weren't requested"
836 );
837 }
838 if !transactions.is_empty() {
839 self.metrics.fetched_transactions.increment(transactions.len() as u64);
840 FetchEvent::TransactionsFetched {
841 peer_id,
842 version,
843 client_version,
844 transactions,
845 report_peer: unsolicited > 0,
846 }
847 } else if unsolicited > 0 {
848 FetchEvent::FetchError { peer_id, error: RequestError::BadResponse }
850 } else {
851 trace!(target: "net::tx",
852 peer_id=format!("{peer_id:#}"),
853 requested=requested.len(),
854 "received empty `PooledTransactions` response, peer failed to serve hashes it announced"
855 );
856 FetchEvent::EmptyResponse { peer_id }
857 }
858 }
859 Err(error) => FetchEvent::FetchError { peer_id, error },
860 }
861 }
862
863 fn on_delivery(
866 &mut self,
867 key: PeerKey,
868 request_id: u64,
869 requested: &[TxHash],
870 delivered: &B256Set,
871 timed_out: bool,
872 ) {
873 let cutoff = if 2 * delivered.len() >= requested.len() {
879 requested
880 .iter()
881 .rposition(|hash| delivered.contains(hash))
882 .map_or(requested.len(), |idx| idx + 1)
883 } else {
884 requested.len()
885 };
886
887 let mut dropped = 0u64;
888 let mut requeued = SmallVec::<[PeerKey; MAX_COUNT_CANDIDATE_PEERS_PER_HASH]>::new();
889
890 for (idx, hash) in requested.iter().enumerate().rev() {
892 if delivered.contains(hash) {
893 self.remove_hash(hash);
894 continue
895 }
896
897 let Some(entry) = self.hashes.get_mut(hash) else { continue };
900 if entry.fetching_by != Some((key, request_id)) {
901 continue
902 }
903 entry.fetching_by = None;
904 self.num_fetching -= 1;
905
906 let retry_only_source = timed_out && entry.attempts == 1 && entry.candidates.len() == 1;
909 let peers = &mut self.peers;
910 entry.candidates.retain(|candidate| {
911 let Some(peer) = peers.get_mut(&candidate.peer) else { return false };
913 if idx < cutoff && candidate.peer == key && !retry_only_source {
914 peer.tracked = peer.tracked.saturating_sub(1);
915 return false
916 }
917 true
918 });
919
920 if entry.candidates.is_empty() || entry.attempts >= MAX_FETCH_ATTEMPTS_PER_HASH {
921 self.remove_hash(hash);
923 dropped += 1;
924 continue
925 }
926
927 let targets = entry.unqueued_candidates();
928 self.requeue(*hash, targets, &mut requeued);
929 }
930
931 if dropped > 0 {
932 self.metrics.hashes_dropped_no_candidate_peers.increment(dropped);
933 trace!(target: "net::tx", dropped, "dropped hashes with no candidates or exhausted fetch attempts");
934 }
935
936 for key in requeued {
938 self.mark_ready(key, QueuePosition::Front);
939 }
940 }
941}
942
943#[cfg(test)]
944impl<N: NetworkPrimitives> TransactionFetcher<N> {
945 pub(super) fn candidate_peers(&self, hash: &TxHash) -> Vec<PeerId> {
947 self.hashes
948 .get(hash)
949 .map(|entry| {
950 entry
951 .candidates
952 .iter()
953 .filter_map(|candidate| self.peers.get(&candidate.peer))
954 .map(|peer| peer.peer_id)
955 .collect()
956 })
957 .unwrap_or_default()
958 }
959
960 fn fetching_peer(&self, hash: &TxHash) -> Option<PeerId> {
962 let (key, _) = self.hashes.get(hash)?.fetching_by?;
963 self.peers.get(&key).map(|peer| peer.peer_id)
964 }
965
966 pub(super) fn queued_hashes(&self, peer_id: &PeerId) -> Vec<TxHash> {
969 self.peer_keys
970 .get(peer_id)
971 .and_then(|key| self.peers.get(key))
972 .map(|peer| peer.queue.iter().map(|(hash, _)| *hash).collect())
973 .unwrap_or_default()
974 }
975
976 fn num_peers(&self) -> usize {
978 self.peers.len()
979 }
980
981 fn assert_invariants(&self) {
983 let fetching = self.hashes.values().filter(|entry| entry.fetching_by.is_some()).count();
984 assert_eq!(fetching, self.num_fetching, "fetching counter out of sync");
985
986 let ordered = self
987 .order
988 .iter()
989 .filter(|(hash, generation)| {
990 self.hashes.get(hash).is_some_and(|entry| entry.generation == *generation)
991 })
992 .map(|(hash, _)| *hash)
993 .collect::<B256Set>();
994 assert!(
995 self.order.len() <= 2 * self.config.max_capacity_cache_txns_pending_fetch as usize,
996 "eviction order grew beyond its bound"
997 );
998
999 for (hash, entry) in &self.hashes {
1000 assert!(
1001 !entry.candidates.is_empty() || entry.fetching_by.is_some(),
1002 "{hash} is pending without candidates"
1003 );
1004 assert!(
1005 entry.candidates.len() <= MAX_COUNT_CANDIDATE_PEERS_PER_HASH,
1006 "{hash} has too many candidates"
1007 );
1008 let mut unique = entry.candidates.iter().map(|c| c.peer).collect::<Vec<_>>();
1009 unique.sort_unstable();
1010 unique.dedup();
1011 assert_eq!(unique.len(), entry.candidates.len(), "{hash} has duplicate candidates");
1012 assert!(ordered.contains(hash), "{hash} is missing from the eviction order");
1013 if let Some((peer, request_id)) = entry.fetching_by {
1014 assert!(
1015 self.inflight.iter().any(|request| request.peer == peer &&
1016 request.request_id == request_id &&
1017 request.hashes.contains(hash)),
1018 "fetching hash has no matching request"
1019 );
1020 }
1021 assert!(entry.attempts <= MAX_FETCH_ATTEMPTS_PER_HASH);
1022 }
1023
1024 let queued = self
1027 .peers
1028 .iter()
1029 .map(|(key, peer)| {
1030 assert!(
1031 peer.queue.len() <= 2 * self.config.max_announced_hashes_per_peer as usize,
1032 "queue of {:#} grew beyond its bound",
1033 peer.peer_id
1034 );
1035 (
1036 *key,
1037 peer.queue
1038 .iter()
1039 .filter(|(hash, generation)| {
1040 self.hashes
1041 .get(hash)
1042 .is_some_and(|entry| entry.generation == *generation)
1043 })
1044 .map(|(hash, _)| *hash)
1045 .collect::<B256Set>(),
1046 )
1047 })
1048 .collect::<HashMap<_, _>>();
1049 for (hash, entry) in &self.hashes {
1050 let mut fetchable = entry.fetching_by.is_some();
1051 for candidate in &entry.candidates {
1052 let queue =
1053 queued.get(&candidate.peer).expect("candidate refers to a disconnected peer");
1054 if candidate.is_queued() {
1055 assert!(
1056 queue.contains(hash),
1057 "{hash} is flagged but not queued for {:?}",
1058 candidate.peer
1059 );
1060 fetchable = true;
1061 }
1062 }
1063 assert!(fetchable, "pending {hash} is not queued for any connected candidate");
1064 }
1065
1066 for (key, peer) in &self.peers {
1067 assert_eq!(
1068 peer.inflight as usize,
1069 self.inflight.iter().filter(|request| request.peer == *key).count()
1070 );
1071 assert!(peer.inflight <= self.config.max_inflight_requests_per_peer);
1072 if peer.inflight < self.config.max_inflight_requests_per_peer &&
1073 peer.queue.iter().any(|(hash, generation)| {
1074 self.hashes.get(hash).is_some_and(|entry| {
1075 entry.generation == *generation &&
1076 entry.fetching_by.is_none() &&
1077 entry
1078 .candidates
1079 .iter()
1080 .any(|candidate| candidate.peer == *key && candidate.is_queued())
1081 })
1082 })
1083 {
1084 assert!(peer.ready, "idle peer has queued work but is not ready");
1085 }
1086 let tracked = self.hashes.values().filter(|entry| entry.has_candidate(*key)).count();
1087 assert_eq!(tracked, peer.tracked, "tracked counter of {:#} out of sync", peer.peer_id);
1088 assert!(
1089 peer.tracked <= self.config.max_announced_hashes_per_peer as usize,
1090 "{:#} exceeds the per peer limit",
1091 peer.peer_id
1092 );
1093 assert_eq!(
1094 peer.ready,
1095 self.ready.contains(key),
1096 "ready flag of {:#} out of sync",
1097 peer.peer_id
1098 );
1099 if peer.ready {
1100 assert!(
1101 !peer.queue.is_empty(),
1102 "{:#} is ready without queued hashes",
1103 peer.peer_id
1104 );
1105 assert!(
1106 peer.inflight < self.config.max_inflight_requests_per_peer,
1107 "{:#} is ready but busy",
1108 peer.peer_id
1109 );
1110 }
1111 assert_eq!(
1112 self.peer_keys.get(&peer.peer_id),
1113 Some(key),
1114 "peer key mapping out of sync"
1115 );
1116 }
1117
1118 for (peer_id, key) in &self.peer_keys {
1119 assert!(self.peers.contains_key(key), "{peer_id:#} maps to an unknown key");
1120 }
1121
1122 let inflight = self.peers.values().map(|peer| peer.inflight as usize).sum::<usize>();
1123 assert!(inflight <= self.inflight.len(), "peers claim more inflight requests than exist");
1124 }
1125}
1126
1127impl<N: NetworkPrimitives> Default for TransactionFetcher<N> {
1128 fn default() -> Self {
1129 Self::new(TransactionFetcherConfig::default())
1130 }
1131}
1132
1133impl<N: NetworkPrimitives> Stream for TransactionFetcher<N> {
1134 type Item = FetchEvent<N::PooledTransaction>;
1135
1136 fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
1140 match self.inflight.poll_next_unpin(cx) {
1143 Poll::Ready(Some(resolved)) => Poll::Ready(Some(self.on_resolved(resolved))),
1144 Poll::Ready(None) | Poll::Pending => Poll::Pending,
1145 }
1146 }
1147}
1148
1149#[derive(Debug)]
1151pub enum FetchEvent<T = PooledTransaction> {
1152 TransactionsFetched {
1154 peer_id: PeerId,
1156 version: EthVersion,
1158 client_version: Arc<str>,
1160 transactions: PooledTransactions<T>,
1162 report_peer: bool,
1165 },
1166 FetchError {
1168 peer_id: PeerId,
1170 error: RequestError,
1172 },
1173 EmptyResponse {
1175 peer_id: PeerId,
1177 },
1178}
1179
1180#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)]
1182struct PeerKey(u32);
1183
1184#[derive(Debug)]
1187struct EvictionState {
1188 peers: BinaryHeap<(usize, PeerKey)>,
1189 exclusive: HashMap<PeerKey, VecDeque<(TxHash, u64)>>,
1190}
1191
1192#[derive(Debug, Clone, Copy)]
1194enum QueuePosition {
1195 Front,
1196 Back,
1197}
1198
1199#[derive(Debug, Default)]
1201struct DeferredCandidateRemovals {
1202 pending: Option<(PeerKey, usize)>,
1203}
1204
1205impl DeferredCandidateRemovals {
1206 fn record_removed(&mut self, key: PeerKey, peers: &mut HashMap<PeerKey, PeerState>) {
1207 if let Some((pending_key, count)) = &mut self.pending &&
1208 *pending_key == key
1209 {
1210 *count += 1;
1211 } else {
1212 self.flush(peers);
1213 self.pending = Some((key, 1));
1214 }
1215 }
1216
1217 fn flush(&mut self, peers: &mut HashMap<PeerKey, PeerState>) {
1219 if let Some((key, count)) = self.pending.take() {
1220 peers.get_mut(&key).expect("candidate peer is connected").tracked -= count;
1221 }
1222 }
1223}
1224
1225#[derive(Debug)]
1227struct TxEntry {
1228 candidates: SmallVec<[Candidate; MAX_COUNT_CANDIDATE_PEERS_PER_HASH]>,
1230 generation: u64,
1232 fetching_by: Option<(PeerKey, u64)>,
1234 attempts: usize,
1236}
1237
1238impl TxEntry {
1239 fn new(peer: PeerKey, size: u32, generation: u64) -> Self {
1241 let mut candidates = SmallVec::new();
1242 candidates.push(Candidate::queued(peer, size));
1243 Self { candidates, generation, fetching_by: None, attempts: 0 }
1244 }
1245
1246 fn has_candidate(&self, key: PeerKey) -> bool {
1247 self.candidates.iter().any(|candidate| candidate.peer == key)
1248 }
1249
1250 fn unqueued_candidates(&self) -> SmallVec<[PeerKey; MAX_COUNT_CANDIDATE_PEERS_PER_HASH]> {
1252 self.candidates
1253 .iter()
1254 .filter(|candidate| !candidate.is_queued())
1255 .map(|candidate| candidate.peer)
1256 .collect()
1257 }
1258
1259 fn candidate_mut(&mut self, key: PeerKey) -> Option<&mut Candidate> {
1260 self.candidates.iter_mut().find(|candidate| candidate.peer == key)
1261 }
1262}
1263
1264#[derive(Debug, Clone, Copy)]
1266struct Candidate {
1267 peer: PeerKey,
1268 size_and_queued: u32,
1271}
1272
1273impl Candidate {
1274 const QUEUED: u32 = 1 << 31;
1276
1277 const fn queued(peer: PeerKey, size: u32) -> Self {
1279 Self { peer, size_and_queued: size | Self::QUEUED }
1280 }
1281
1282 const fn unqueued(peer: PeerKey, size: u32) -> Self {
1284 Self { peer, size_and_queued: size }
1285 }
1286
1287 const fn is_queued(&self) -> bool {
1288 self.size_and_queued & Self::QUEUED != 0
1289 }
1290
1291 const fn set_queued(&mut self, queued: bool) {
1292 if queued {
1293 self.size_and_queued |= Self::QUEUED;
1294 } else {
1295 self.size_and_queued &= !Self::QUEUED;
1296 }
1297 }
1298
1299 const fn set_size(&mut self, size: u32) {
1300 self.size_and_queued = (self.size_and_queued & Self::QUEUED) | size;
1301 }
1302
1303 const fn request_size(&self) -> usize {
1305 let size = self.size_and_queued & !Self::QUEUED;
1306 if size == 0 {
1307 AVERAGE_BYTE_SIZE_TX_ENCODED
1308 } else {
1309 size as usize
1310 }
1311 }
1312}
1313
1314#[derive(Debug)]
1316struct PeerState {
1317 peer_id: PeerId,
1318 queue: VecDeque<(TxHash, u64)>,
1321 tracked: usize,
1323 inflight: u8,
1325 ready: bool,
1327}
1328
1329impl PeerState {
1330 const fn new(peer_id: PeerId) -> Self {
1331 Self { peer_id, queue: VecDeque::new(), tracked: 0, inflight: 0, ready: false }
1332 }
1333
1334 fn push_queue(
1341 &mut self,
1342 hashes: &B256Map<TxEntry>,
1343 key: PeerKey,
1344 queued: (TxHash, u64),
1345 position: QueuePosition,
1346 max_tracked: usize,
1347 ) {
1348 if self.queue.len() >= 2 * max_tracked {
1349 let mut seen = B256Set::with_capacity_and_hasher(self.tracked, Default::default());
1350 self.queue.retain(|(hash, generation)| {
1351 hashes.get(hash).is_some_and(|entry| {
1352 entry.generation == *generation && entry.has_candidate(key)
1353 }) && seen.insert(*hash)
1354 });
1355 }
1356 match position {
1357 QueuePosition::Front => self.queue.push_front(queued),
1358 QueuePosition::Back => self.queue.push_back(queued),
1359 }
1360 }
1361}
1362
1363#[derive(Debug)]
1365struct InflightRequest<T> {
1366 peer: PeerKey,
1367 request_id: u64,
1368 peer_id: PeerId,
1369 version: EthVersion,
1370 client_version: Arc<str>,
1371 hashes: Vec<TxHash>,
1373 response: oneshot::Receiver<RequestResult<PooledTransactions<T>>>,
1374}
1375
1376impl<T> Future for InflightRequest<T> {
1377 type Output = ResolvedRequest<T>;
1378
1379 fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
1380 let this = self.get_mut();
1381 let result =
1382 ready!(this.response.poll_unpin(cx)).unwrap_or(Err(RequestError::ChannelClosed));
1383 Poll::Ready(ResolvedRequest {
1384 peer: this.peer,
1385 request_id: this.request_id,
1386 peer_id: this.peer_id,
1387 version: this.version,
1388 client_version: this.client_version.clone(),
1389 hashes: std::mem::take(&mut this.hashes),
1390 result,
1391 })
1392 }
1393}
1394
1395#[derive(Debug)]
1397struct ResolvedRequest<T> {
1398 peer: PeerKey,
1399 request_id: u64,
1400 peer_id: PeerId,
1401 version: EthVersion,
1402 client_version: Arc<str>,
1403 hashes: Vec<TxHash>,
1404 result: RequestResult<PooledTransactions<T>>,
1405}
1406
1407fn announced_size(metadata: Option<TransactionMetadata>) -> u32 {
1410 metadata.map_or(0, |metadata| {
1411 metadata.size.min(SOFT_LIMIT_BYTE_SIZE_POOLED_TRANSACTIONS_RESPONSE) as u32
1412 })
1413}
1414
1415fn verify_response<T: SignedTransaction>(
1420 transactions: &mut PooledTransactions<T>,
1421 requested: &B256Set,
1422 delivered: &mut B256Set,
1423) -> usize {
1424 let mut unsolicited = 0;
1425 transactions.0.retain(|tx| {
1426 let hash = *tx.tx_hash();
1427 if !requested.contains(&hash) {
1428 unsolicited += 1;
1429 return false
1430 }
1431 delivered.insert(hash)
1432 });
1433 unsolicited
1434}
1435
1436#[cfg(test)]
1437mod tests {
1438 use super::*;
1439 use crate::test_utils::transactions::new_mock_session_with_capacity;
1440 use alloy_consensus::transaction::Recovered;
1441 use alloy_primitives::B256;
1442 use futures::task::{noop_waker_ref, waker, ArcWake};
1443 use rand::{rngs::StdRng, seq::IndexedRandom, Rng, SeedableRng};
1444 use reth_eth_wire::EthVersion;
1445 use reth_ethereum_primitives::{PooledTransactionVariant, TransactionSigned};
1446 use reth_transaction_pool::test_utils::MockTransactionFactory;
1447 use std::sync::{
1448 atomic::{AtomicUsize, Ordering},
1449 Arc,
1450 };
1451 use tokio::sync::mpsc;
1452
1453 type Fetcher = TransactionFetcher<EthNetworkPrimitives>;
1454 type ResponseSender =
1455 oneshot::Sender<RequestResult<PooledTransactions<PooledTransactionVariant>>>;
1456
1457 const KIB: usize = 1024;
1458
1459 fn peer(n: u8) -> PeerId {
1460 PeerId::new([n; 64])
1461 }
1462
1463 fn hash(n: u64) -> TxHash {
1464 let mut bytes = [0u8; 32];
1465 bytes[24..].copy_from_slice(&n.to_be_bytes());
1466 B256::from(bytes)
1467 }
1468
1469 fn hashes(range: std::ops::Range<u64>) -> Vec<TxHash> {
1470 range.map(hash).collect()
1471 }
1472
1473 #[derive(Default)]
1475 struct WakeCounter(AtomicUsize);
1476
1477 impl ArcWake for WakeCounter {
1478 fn wake_by_ref(arc_self: &Arc<Self>) {
1479 arc_self.0.fetch_add(1, Ordering::SeqCst);
1480 }
1481 }
1482
1483 impl WakeCounter {
1484 fn wakes(&self) -> usize {
1485 self.0.load(Ordering::SeqCst)
1486 }
1487 }
1488
1489 fn pooled_txs(count: usize) -> Vec<PooledTransactionVariant> {
1490 let mut factory = MockTransactionFactory::default();
1491 (0..count)
1492 .map(|_| {
1493 let recovered: Recovered<TransactionSigned> =
1494 factory.create_eip1559().transaction.into();
1495 PooledTransactionVariant::try_from(recovered.into_inner())
1496 .expect("eip1559 transaction converts to pooled transaction")
1497 })
1498 .collect()
1499 }
1500
1501 fn hashes_of(txs: &[PooledTransactionVariant]) -> Vec<TxHash> {
1502 txs.iter().map(|tx| *tx.tx_hash()).collect()
1503 }
1504
1505 struct Rig {
1507 fetcher: Fetcher,
1508 peers: HashMap<PeerId, PeerMetadata<EthNetworkPrimitives>, FbBuildHasher<64>>,
1509 sessions: HashMap<PeerId, mpsc::Receiver<PeerRequest>, FbBuildHasher<64>>,
1510 }
1511
1512 impl Rig {
1513 fn new() -> Self {
1514 Self::with_config(TransactionFetcherConfig::default())
1515 }
1516
1517 fn with_config(config: TransactionFetcherConfig) -> Self {
1518 Self {
1519 fetcher: Fetcher::new(config),
1520 peers: Default::default(),
1521 sessions: Default::default(),
1522 }
1523 }
1524
1525 fn verify(&self) {
1526 self.fetcher.assert_invariants();
1527 }
1528
1529 fn add_peer(&mut self, peer_id: PeerId) {
1530 self.add_peer_with_capacity(peer_id, 8);
1531 }
1532
1533 fn add_peer_with_capacity(&mut self, peer_id: PeerId, capacity: usize) {
1534 let (peer, rx) = new_mock_session_with_capacity(peer_id, EthVersion::Eth68, capacity);
1535 self.peers.insert(peer_id, peer);
1536 self.sessions.insert(peer_id, rx);
1537 }
1538
1539 fn announce(&mut self, peer_id: PeerId, hashes: &[TxHash]) {
1540 self.announce_with_sizes(peer_id, hashes.iter().map(|hash| (*hash, 512)));
1541 }
1542
1543 fn announce_with_sizes(
1544 &mut self,
1545 peer_id: PeerId,
1546 entries: impl IntoIterator<Item = (TxHash, usize)>,
1547 ) {
1548 self.fetcher.on_announcement(
1549 peer_id,
1550 entries.into_iter().map(|(hash, size)| AnnouncedTransaction {
1551 hash,
1552 metadata: Some(TransactionMetadata { tx_type: 2, size }),
1553 }),
1554 );
1555 self.verify();
1556 }
1557
1558 fn announce_unsized(&mut self, peer_id: PeerId, hashes: &[TxHash]) {
1559 self.fetcher.on_announcement(
1560 peer_id,
1561 hashes.iter().map(|&hash| AnnouncedTransaction { hash, metadata: None }),
1562 );
1563 self.verify();
1564 }
1565
1566 fn dispatch(&mut self) -> usize {
1567 self.dispatch_with_budget(usize::MAX)
1568 }
1569
1570 fn dispatch_with_budget(&mut self, max_hashes_per_request: usize) -> usize {
1571 let sent = self.fetcher.dispatch(&self.peers, max_hashes_per_request);
1572 self.verify();
1573 sent
1574 }
1575
1576 fn take_request(&mut self, peer_id: PeerId) -> Option<(Vec<TxHash>, ResponseSender)> {
1578 match self.sessions.get_mut(&peer_id)?.try_recv().ok()? {
1579 PeerRequest::GetPooledTransactions { request, response } => {
1580 Some((request.0, response))
1581 }
1582 _ => unreachable!("the fetcher only sends `GetPooledTransactions` requests"),
1583 }
1584 }
1585
1586 fn next_event(&mut self) -> Option<FetchEvent<PooledTransactionVariant>> {
1587 let mut cx = Context::from_waker(noop_waker_ref());
1588 let event = match self.fetcher.poll_next_unpin(&mut cx) {
1589 Poll::Ready(event) => event,
1590 Poll::Pending => None,
1591 };
1592 self.verify();
1593 event
1594 }
1595
1596 fn drain_events(&mut self) -> Vec<FetchEvent<PooledTransactionVariant>> {
1597 std::iter::from_fn(|| self.next_event()).collect()
1598 }
1599
1600 fn respond(
1601 &mut self,
1602 peer_id: PeerId,
1603 txs: Vec<PooledTransactionVariant>,
1604 ) -> FetchEvent<PooledTransactionVariant> {
1605 let (_, response) = self.take_request(peer_id).expect("request inflight");
1606 response.send(Ok(PooledTransactions(txs))).unwrap();
1607 self.next_event().expect("response yields an event")
1608 }
1609
1610 fn fail(
1611 &mut self,
1612 peer_id: PeerId,
1613 error: RequestError,
1614 ) -> FetchEvent<PooledTransactionVariant> {
1615 let (_, response) = self.take_request(peer_id).expect("request inflight");
1616 response.send(Err(error)).unwrap();
1617 self.next_event().expect("error yields an event")
1618 }
1619
1620 fn disconnect(&mut self, peer_id: PeerId) {
1621 self.peers.remove(&peer_id);
1622 self.sessions.remove(&peer_id);
1623 self.fetcher.on_peer_disconnected(&peer_id);
1624 self.verify();
1625 }
1626 }
1627
1628 #[test]
1629 fn announced_hashes_are_requested_from_announcing_peer() {
1630 let mut rig = Rig::new();
1631 let peer_a = peer(1);
1632 rig.add_peer(peer_a);
1633 let mut txs = pooled_txs(2);
1634 txs.sort_unstable_by_key(|tx| std::cmp::Reverse(*tx.tx_hash()));
1635 let hashes = hashes_of(&txs);
1636
1637 let mut cx = Context::from_waker(noop_waker_ref());
1638 assert!(rig.fetcher.poll_next_unpin(&mut cx).is_pending());
1639 rig.announce(peer_a, &hashes[..1]);
1640 rig.announce(peer_a, &hashes[1..]);
1641 rig.announce(peer_a, &hashes);
1642 assert_eq!(rig.fetcher.queued_hashes(&peer_a), hashes);
1643 assert_eq!(rig.fetcher.candidate_peers(&hashes[0]), vec![peer_a]);
1644 assert_eq!(rig.fetcher.num_pending_hashes(), 2);
1645 assert!(rig.fetcher.is_idle(&peer_a));
1646
1647 assert_eq!(rig.dispatch(), 1);
1648 assert!(!rig.fetcher.is_idle(&peer_a));
1649 assert_eq!(rig.fetcher.num_inflight_requests(), 1);
1650 assert_eq!(rig.fetcher.num_fetching_hashes(), 2);
1651 assert_eq!(rig.fetcher.num_pending_hashes(), 0);
1652 assert_eq!(rig.fetcher.fetching_peer(&hashes[0]), Some(peer_a));
1653
1654 let (requested, response) = rig.take_request(peer_a).unwrap();
1655 assert_eq!(requested, hashes);
1656 response.send(Ok(PooledTransactions(txs))).unwrap();
1657
1658 let FetchEvent::TransactionsFetched { peer_id, transactions, report_peer, .. } =
1659 rig.next_event().unwrap()
1660 else {
1661 panic!("expected fetched transactions")
1662 };
1663 assert_eq!(peer_id, peer_a);
1664 assert_eq!(transactions.len(), 2);
1665 assert!(!report_peer);
1666
1667 assert_eq!(rig.fetcher.num_hashes(), 0);
1668 assert_eq!(rig.fetcher.num_inflight_requests(), 0);
1669 assert!(rig.fetcher.is_idle(&peer_a));
1670 assert!(rig.fetcher.poll_next_unpin(&mut cx).is_pending());
1671 }
1672
1673 #[test]
1674 fn request_is_capped_by_hash_count() {
1675 let mut rig = Rig::new();
1676 let peer_a = peer(1);
1677 rig.add_peer(peer_a);
1678 let hashes = hashes(0..300);
1679
1680 rig.announce_unsized(peer_a, &hashes);
1682 assert_eq!(rig.dispatch(), 1);
1683
1684 let (requested, response) = rig.take_request(peer_a).unwrap();
1685 assert_eq!(requested.len(), SOFT_LIMIT_COUNT_HASHES_IN_GET_POOLED_TRANSACTIONS_REQUEST);
1686 assert_eq!(requested, hashes[..256]);
1687 assert_eq!(rig.fetcher.num_pending_hashes(), 44);
1688
1689 assert_eq!(rig.dispatch(), 0);
1691 response.send(Err(RequestError::BadResponse)).unwrap();
1692 rig.next_event().unwrap();
1693 assert_eq!(rig.fetcher.num_pending_hashes(), 44);
1694
1695 assert_eq!(rig.dispatch(), 1);
1696 let (requested, _) = rig.take_request(peer_a).unwrap();
1697 assert_eq!(requested, hashes[256..]);
1698 }
1699
1700 #[test]
1701 fn request_is_capped_by_expected_response_size() {
1702 let mut rig = Rig::new();
1703 let peer_a = peer(1);
1704 rig.add_peer(peer_a);
1705 let hashes = hashes(0..3);
1706
1707 rig.announce_with_sizes(
1708 peer_a,
1709 [(hashes[0], 100 * KIB), (hashes[1], 100 * KIB), (hashes[2], 100)],
1710 );
1711
1712 rig.dispatch();
1713 let (requested, response) = rig.take_request(peer_a).unwrap();
1714 assert_eq!(requested, hashes[..1], "second transaction doesn't fit in 128 KiB");
1715
1716 response.send(Ok(PooledTransactions(vec![]))).unwrap();
1717 rig.next_event().unwrap();
1718 assert_eq!(rig.fetcher.num_hashes(), 2);
1720
1721 rig.dispatch();
1722 let (requested, _) = rig.take_request(peer_a).unwrap();
1723 assert_eq!(requested, hashes[1..]);
1724 }
1725
1726 #[test]
1727 fn oversized_transaction_is_requested_alone() {
1728 let mut rig = Rig::new();
1729 let peer_a = peer(1);
1730 rig.add_peer(peer_a);
1731 let hashes = hashes(0..3);
1732
1733 rig.announce_with_sizes(
1734 peer_a,
1735 [(hashes[0], 100), (hashes[1], 1024 * KIB), (hashes[2], 100)],
1736 );
1737
1738 let mut requests = Vec::new();
1739 for _ in 0..3 {
1740 assert_eq!(rig.dispatch(), 1);
1741 let (requested, response) = rig.take_request(peer_a).unwrap();
1742 requests.push(requested);
1743 response.send(Err(RequestError::BadResponse)).unwrap();
1745 rig.next_event().unwrap();
1746 }
1747 assert_eq!(
1748 requests,
1749 vec![hashes[..1].to_vec(), hashes[1..2].to_vec(), hashes[2..].to_vec()]
1750 );
1751 assert_eq!(rig.fetcher.num_hashes(), 0);
1752 }
1753
1754 #[test]
1755 fn announced_size_is_per_peer() {
1756 let mut rig = Rig::new();
1757 let peer_a = peer(1);
1758 let peer_b = peer(2);
1759 rig.add_peer(peer_a);
1760 rig.add_peer(peer_b);
1761 let hashes = hashes(0..2);
1762
1763 rig.announce_unsized(peer_a, &hashes);
1765 rig.announce_with_sizes(peer_b, [(hashes[0], 1024 * KIB), (hashes[1], 100)]);
1766
1767 rig.dispatch();
1769 let (requested, response) = rig.take_request(peer_a).unwrap();
1770 assert_eq!(requested, hashes);
1771 response.send(Ok(PooledTransactions(vec![]))).unwrap();
1772 rig.next_event().unwrap();
1773
1774 rig.dispatch();
1776 let (requested, _) = rig.take_request(peer_b).unwrap();
1777 assert_eq!(requested, hashes[..1]);
1778 }
1779
1780 #[test]
1781 fn announced_sizes_are_capped() {
1782 let mut rig = Rig::new();
1783 let peer_a = peer(1);
1784 rig.add_peer(peer_a);
1785 let hashes = hashes(0..3);
1786
1787 rig.announce_with_sizes(
1789 peer_a,
1790 [(hashes[0], usize::MAX - 1), (hashes[1], usize::MAX), (hashes[2], 100)],
1791 );
1792
1793 rig.dispatch();
1794 let (requested, response) = rig.take_request(peer_a).unwrap();
1795 assert_eq!(requested, hashes[..1], "an oversized transaction is requested alone");
1796 response.send(Err(RequestError::BadResponse)).unwrap();
1797 rig.next_event().unwrap();
1798
1799 rig.dispatch();
1800 let (requested, _) = rig.take_request(peer_a).unwrap();
1801 assert_eq!(requested, hashes[1..2]);
1802 }
1803
1804 #[test]
1805 fn hash_is_fetched_from_one_peer_at_a_time() {
1806 let mut rig = Rig::new();
1807 let peer_a = peer(1);
1808 let peer_b = peer(2);
1809 rig.add_peer(peer_a);
1810 rig.add_peer(peer_b);
1811 let txs = pooled_txs(1);
1812 let hashes = hashes_of(&txs);
1813
1814 rig.announce(peer_a, &hashes);
1815 rig.announce(peer_b, &hashes);
1816 assert_eq!(rig.fetcher.candidate_peers(&hashes[0]), vec![peer_a, peer_b]);
1817
1818 assert_eq!(rig.dispatch(), 1);
1819 assert_eq!(rig.fetcher.fetching_peer(&hashes[0]), Some(peer_a));
1820 assert!(rig.take_request(peer_b).is_none());
1821 assert!(rig.fetcher.queued_hashes(&peer_b).is_empty());
1823 assert_eq!(rig.fetcher.candidate_peers(&hashes[0]), vec![peer_a, peer_b]);
1824
1825 rig.respond(peer_a, txs);
1826 assert_eq!(rig.fetcher.num_hashes(), 0);
1827 assert_eq!(rig.dispatch(), 0);
1828 }
1829
1830 #[test]
1831 fn failed_request_is_retried_with_alternate_peer() {
1832 let mut rig = Rig::new();
1833 let peer_a = peer(1);
1834 let peer_b = peer(2);
1835 rig.add_peer(peer_a);
1836 rig.add_peer(peer_b);
1837 let txs = pooled_txs(2);
1838 let hashes = hashes_of(&txs);
1839
1840 rig.announce(peer_a, &hashes);
1841 rig.announce(peer_b, &hashes);
1842 rig.dispatch();
1843 assert!(rig.fetcher.queued_hashes(&peer_b).is_empty());
1845
1846 let FetchEvent::FetchError { peer_id, error } = rig.fail(peer_a, RequestError::Timeout)
1847 else {
1848 panic!("expected fetch error")
1849 };
1850 assert_eq!(peer_id, peer_a);
1851 assert!(matches!(error, RequestError::Timeout));
1852
1853 assert_eq!(rig.fetcher.num_pending_hashes(), 2);
1855 assert_eq!(rig.fetcher.candidate_peers(&hashes[0]), vec![peer_b]);
1856 assert_eq!(rig.fetcher.queued_hashes(&peer_b), hashes);
1857
1858 assert_eq!(rig.dispatch(), 1);
1859 assert!(rig.take_request(peer_a).is_none());
1860 let event = rig.respond(peer_b, txs);
1861 assert!(
1862 matches!(event, FetchEvent::TransactionsFetched { peer_id, .. } if peer_id == peer_b)
1863 );
1864 assert_eq!(rig.fetcher.num_hashes(), 0);
1865 }
1866
1867 #[test]
1868 fn empty_response_drops_peer_as_candidate() {
1869 let mut rig = Rig::new();
1870 let peer_a = peer(1);
1871 rig.add_peer(peer_a);
1872 let hashes = hashes(0..2);
1873
1874 rig.announce(peer_a, &hashes);
1875 rig.dispatch();
1876 let event = rig.respond(peer_a, vec![]);
1877 assert!(matches!(event, FetchEvent::EmptyResponse { peer_id } if peer_id == peer_a));
1878
1879 assert_eq!(rig.fetcher.num_hashes(), 0);
1881 assert_eq!(rig.dispatch(), 0);
1882 }
1883
1884 #[test]
1885 fn partial_response_keeps_peer_for_truncated_tail() {
1886 let mut rig = Rig::new();
1887 let peer_a = peer(1);
1888 rig.add_peer(peer_a);
1889 let txs = pooled_txs(4);
1890 let hashes = hashes_of(&txs);
1891
1892 rig.announce(peer_a, &hashes);
1893 rig.dispatch();
1894
1895 rig.respond(peer_a, txs[..2].to_vec());
1897 assert_eq!(rig.fetcher.num_pending_hashes(), 2);
1898 assert_eq!(rig.fetcher.candidate_peers(&hashes[2]), vec![peer_a]);
1899 assert_eq!(rig.fetcher.queued_hashes(&peer_a), hashes[2..]);
1900
1901 rig.dispatch();
1902 let (requested, _) = rig.take_request(peer_a).unwrap();
1903 assert_eq!(requested, hashes[2..]);
1904 }
1905
1906 #[test]
1907 fn small_partial_response_drops_peer_for_all_undelivered_hashes() {
1908 let mut rig = Rig::new();
1909 let peer_a = peer(1);
1910 let peer_b = peer(2);
1911 rig.add_peer(peer_a);
1912 rig.add_peer(peer_b);
1913 let txs = pooled_txs(10);
1914 let hashes = hashes_of(&txs);
1915
1916 rig.announce(peer_a, &hashes);
1917 rig.announce(peer_b, &hashes[..5]);
1918 rig.dispatch();
1919 assert_eq!(rig.fetcher.fetching_peer(&hashes[0]), Some(peer_a));
1920
1921 rig.respond(peer_a, txs[..1].to_vec());
1924 assert_eq!(rig.fetcher.num_hashes(), 4, "hashes without another candidate are dropped");
1925 for hash in &hashes[1..5] {
1926 assert_eq!(rig.fetcher.candidate_peers(hash), vec![peer_b]);
1927 }
1928 assert!(rig.fetcher.queued_hashes(&peer_a).is_empty());
1929 }
1930
1931 #[test]
1932 fn skipped_hashes_without_alternates_are_dropped() {
1933 let mut rig = Rig::new();
1934 let peer_a = peer(1);
1935 rig.add_peer(peer_a);
1936 let txs = pooled_txs(4);
1937 let hashes = hashes_of(&txs);
1938
1939 rig.announce(peer_a, &hashes);
1940 rig.dispatch();
1941 rig.respond(peer_a, txs[1..3].to_vec());
1942
1943 assert_eq!(rig.fetcher.num_hashes(), 1);
1945 assert_eq!(rig.fetcher.candidate_peers(&hashes[3]), vec![peer_a]);
1946 assert!(rig.fetcher.candidate_peers(&hashes[0]).is_empty());
1947 }
1948
1949 #[test]
1950 fn unsolicited_and_duplicate_transactions_are_filtered_and_reported() {
1951 let mut rig = Rig::new();
1952 let peer_a = peer(1);
1953 rig.add_peer(peer_a);
1954 let txs = pooled_txs(2);
1955 let hashes = hashes_of(&txs);
1956
1957 rig.announce(peer_a, &hashes[..1]);
1958 rig.dispatch();
1959
1960 let FetchEvent::TransactionsFetched { transactions, report_peer, .. } =
1961 rig.respond(peer_a, vec![txs[0].clone(), txs[1].clone(), txs[0].clone()])
1962 else {
1963 panic!("expected fetched transactions")
1964 };
1965 assert!(report_peer);
1966 assert_eq!(transactions.0, txs[..1]);
1967 assert_eq!(rig.fetcher.num_hashes(), 0);
1968 }
1969
1970 #[test]
1971 fn response_with_only_unsolicited_transactions_is_a_bad_response() {
1972 let mut rig = Rig::new();
1973 let peer_a = peer(1);
1974 let peer_b = peer(2);
1975 rig.add_peer(peer_a);
1976 rig.add_peer(peer_b);
1977 let txs = pooled_txs(2);
1978 let hashes = hashes_of(&txs);
1979
1980 rig.announce(peer_a, &hashes[..1]);
1981 rig.announce(peer_b, &hashes[..1]);
1982 rig.dispatch();
1983
1984 let event = rig.respond(peer_a, txs[1..].to_vec());
1985 assert!(matches!(
1986 event,
1987 FetchEvent::FetchError { peer_id, error: RequestError::BadResponse } if peer_id == peer_a
1988 ));
1989 assert_eq!(rig.fetcher.candidate_peers(&hashes[0]), vec![peer_b]);
1991 rig.dispatch();
1992 assert!(rig.take_request(peer_b).is_some());
1993 }
1994
1995 #[test]
1996 fn duplicate_transactions_in_response_are_deduplicated() {
1997 let mut rig = Rig::new();
1998 let peer_a = peer(1);
1999 rig.add_peer(peer_a);
2000 let txs = pooled_txs(1);
2001
2002 rig.announce(peer_a, &hashes_of(&txs));
2003 rig.dispatch();
2004
2005 let FetchEvent::TransactionsFetched { transactions, report_peer, .. } =
2006 rig.respond(peer_a, vec![txs[0].clone(), txs[0].clone()])
2007 else {
2008 panic!("expected fetched transactions")
2009 };
2010 assert_eq!(transactions.len(), 1);
2011 assert!(!report_peer);
2012 }
2013
2014 #[test]
2015 fn received_transactions_stop_tracking() {
2016 let mut rig = Rig::new();
2017 let peer_a = peer(1);
2018 rig.add_peer(peer_a);
2019 let txs = pooled_txs(3);
2020 let hashes = hashes_of(&txs);
2021
2022 rig.announce(peer_a, &hashes[..2]);
2023 rig.dispatch();
2024 rig.announce(peer_a, &hashes[2..]);
2025 assert_eq!(rig.fetcher.num_fetching_hashes(), 2);
2026 assert_eq!(rig.fetcher.num_pending_hashes(), 1);
2027
2028 rig.fetcher.on_transactions_received([&hashes[0], &hashes[2]]);
2030 rig.fetcher.assert_invariants();
2031 assert_eq!(rig.fetcher.num_hashes(), 1);
2032 assert_eq!(rig.fetcher.num_fetching_hashes(), 1);
2033
2034 let FetchEvent::TransactionsFetched { transactions, .. } =
2037 rig.respond(peer_a, txs[1..2].to_vec())
2038 else {
2039 panic!("expected fetched transactions")
2040 };
2041 assert_eq!(transactions.len(), 1);
2042 assert_eq!(rig.fetcher.num_hashes(), 0);
2043 assert_eq!(rig.dispatch(), 0);
2044 }
2045
2046 #[test]
2047 fn response_delivering_hash_received_over_broadcast_is_still_returned() {
2048 let mut rig = Rig::new();
2049 let peer_a = peer(1);
2050 rig.add_peer(peer_a);
2051 let txs = pooled_txs(1);
2052 let hashes = hashes_of(&txs);
2053
2054 rig.announce(peer_a, &hashes);
2055 rig.dispatch();
2056 rig.fetcher.on_transactions_received(&hashes);
2057
2058 let FetchEvent::TransactionsFetched { transactions, report_peer, .. } =
2059 rig.respond(peer_a, txs)
2060 else {
2061 panic!("expected fetched transactions")
2062 };
2063 assert_eq!(transactions.len(), 1);
2064 assert!(!report_peer, "the transaction was requested, even if it arrived elsewhere first");
2065 assert_eq!(rig.fetcher.num_hashes(), 0);
2066 }
2067
2068 #[test]
2069 fn peer_disconnect_drops_hashes_without_alternates() {
2070 let mut rig = Rig::new();
2071 let peer_a = peer(1);
2072 let peer_b = peer(2);
2073 rig.add_peer(peer_a);
2074 rig.add_peer(peer_b);
2075 let hashes = hashes(0..2);
2076
2077 rig.announce(peer_a, &hashes);
2078 rig.announce(peer_b, &hashes[1..]);
2079
2080 rig.disconnect(peer_a);
2081 assert_eq!(rig.fetcher.num_hashes(), 1);
2082 assert_eq!(rig.fetcher.candidate_peers(&hashes[1]), vec![peer_b]);
2083 assert!(rig.fetcher.is_idle(&peer_a));
2084
2085 rig.dispatch();
2086 let (requested, _) = rig.take_request(peer_b).unwrap();
2087 assert_eq!(requested, hashes[1..]);
2088 }
2089
2090 #[test]
2091 fn inflight_request_of_disconnected_peer_is_rescheduled() {
2092 let mut rig = Rig::new();
2093 let peer_a = peer(1);
2094 let peer_b = peer(2);
2095 rig.add_peer(peer_a);
2096 rig.add_peer(peer_b);
2097 let txs = pooled_txs(2);
2098 let hashes = hashes_of(&txs);
2099
2100 rig.announce(peer_a, &hashes);
2101 rig.announce(peer_b, &hashes);
2102 rig.dispatch();
2103 assert_eq!(rig.fetcher.fetching_peer(&hashes[0]), Some(peer_a));
2104
2105 rig.disconnect(peer_a);
2107 assert_eq!(
2108 rig.fetcher.num_fetching_hashes(),
2109 2,
2110 "hashes stay inflight until the request resolves"
2111 );
2112
2113 let event = rig.next_event().unwrap();
2114 assert!(matches!(
2115 event,
2116 FetchEvent::FetchError { peer_id, error: RequestError::ChannelClosed } if peer_id == peer_a
2117 ));
2118 assert_eq!(rig.fetcher.num_pending_hashes(), 2);
2119 assert_eq!(rig.fetcher.candidate_peers(&hashes[0]), vec![peer_b]);
2120
2121 rig.dispatch();
2122 let event = rig.respond(peer_b, txs);
2123 assert!(matches!(event, FetchEvent::TransactionsFetched { .. }));
2124 assert_eq!(rig.fetcher.num_hashes(), 0);
2125 }
2126
2127 #[test]
2128 fn per_peer_announcement_limit_is_enforced() {
2129 let config =
2130 TransactionFetcherConfig { max_announced_hashes_per_peer: 3, ..Default::default() };
2131 let mut rig = Rig::with_config(config);
2132 let peer_a = peer(1);
2133 let peer_b = peer(2);
2134 rig.add_peer(peer_a);
2135 rig.add_peer(peer_b);
2136 let txs = pooled_txs(5);
2137 let hashes = hashes_of(&txs);
2138
2139 rig.announce(peer_a, &hashes);
2140 assert_eq!(rig.fetcher.num_hashes(), 3);
2141 assert_eq!(rig.fetcher.queued_hashes(&peer_a), hashes[..3]);
2142
2143 rig.announce(peer_b, &hashes[3..]);
2145 assert_eq!(rig.fetcher.num_hashes(), 5);
2146 assert_eq!(rig.fetcher.candidate_peers(&hashes[4]), vec![peer_b]);
2147
2148 rig.dispatch();
2150 rig.respond(peer_a, txs[..3].to_vec());
2151 rig.announce(peer_a, &hashes[3..]);
2152 assert_eq!(rig.fetcher.candidate_peers(&hashes[3]), vec![peer_b, peer_a]);
2153 }
2154
2155 #[test]
2156 fn global_capacity_evicts_oldest_pending_hash() {
2157 let config = TransactionFetcherConfig {
2158 max_capacity_cache_txns_pending_fetch: 2,
2159 ..Default::default()
2160 };
2161 let mut rig = Rig::with_config(config);
2162 let peer_a = peer(1);
2163 let peer_b = peer(2);
2164 rig.add_peer(peer_a);
2165 rig.add_peer(peer_b);
2166 let txs = pooled_txs(5);
2167 let hashes = hashes_of(&txs);
2168
2169 rig.announce(peer_a, &hashes[..3]);
2171 assert_eq!(rig.fetcher.num_hashes(), 2);
2172 assert!(rig.fetcher.candidate_peers(&hashes[0]).is_empty());
2173 assert_eq!(rig.fetcher.candidate_peers(&hashes[2]), vec![peer_a]);
2174
2175 rig.dispatch();
2177 let (requested, response) = rig.take_request(peer_a).unwrap();
2178 assert_eq!(requested, hashes[1..3]);
2179 rig.announce(peer_b, &hashes[3..4]);
2180 assert_eq!(rig.fetcher.num_hashes(), 2);
2181 assert!(rig.fetcher.candidate_peers(&hashes[3]).is_empty());
2182
2183 response.send(Ok(PooledTransactions(txs[1..2].to_vec()))).unwrap();
2186 rig.next_event().unwrap();
2187 assert_eq!(rig.fetcher.num_hashes(), 1);
2188 rig.announce(peer_b, &hashes[3..4]);
2189 assert_eq!(rig.fetcher.num_hashes(), 2);
2190 rig.announce(peer_b, &hashes[4..5]);
2191 assert_eq!(rig.fetcher.num_hashes(), 2);
2192 assert_eq!(rig.fetcher.candidate_peers(&hashes[2]), vec![peer_a]);
2193 assert!(rig.fetcher.candidate_peers(&hashes[3]).is_empty());
2194 assert_eq!(rig.fetcher.candidate_peers(&hashes[4]), vec![peer_b]);
2195 }
2196
2197 #[test]
2198 fn later_announcers_are_asked_after_the_first_ones_failed() {
2199 let mut rig = Rig::new();
2200 let batch = hashes(0..3);
2201 let hash = batch[0];
2202 let peers =
2203 (1..=MAX_COUNT_CANDIDATE_PEERS_PER_HASH as u8 + 2).map(peer).collect::<Vec<_>>();
2204 for peer_id in &peers {
2205 rig.add_peer(*peer_id);
2206 rig.announce(*peer_id, &batch);
2207 }
2208
2209 let candidates = peers[..MAX_COUNT_EAGER_CANDIDATE_PEERS_PER_HASH]
2211 .iter()
2212 .chain(&peers[MAX_COUNT_EAGER_CANDIDATE_PEERS_PER_HASH + 2..])
2213 .copied()
2214 .collect::<Vec<_>>();
2215 assert_eq!(rig.fetcher.candidate_peers(&hash), candidates);
2216 for (i, peer_id) in candidates.iter().enumerate() {
2217 let eager = i < MAX_COUNT_EAGER_CANDIDATE_PEERS_PER_HASH;
2218 assert_eq!(rig.fetcher.queued_hashes(peer_id).contains(&hash), eager, "peer {i}");
2219 }
2220 assert!(rig.fetcher.queued_hashes(&peers[MAX_COUNT_CANDIDATE_PEERS_PER_HASH]).is_empty());
2221
2222 assert_eq!(rig.dispatch(), 1);
2224 assert_eq!(rig.fetcher.fetching_peer(&hash), Some(peers[0]));
2225 rig.fail(peers[0], RequestError::Timeout);
2226 assert_eq!(rig.dispatch(), 1);
2227 assert_eq!(rig.fetcher.fetching_peer(&hash), Some(*candidates.last().unwrap()));
2228
2229 let mut order = vec![peers[0]];
2231 while let Some(fetching) = rig.fetcher.fetching_peer(&hash) {
2232 rig.fail(fetching, RequestError::Timeout);
2233 order.push(fetching);
2234 rig.dispatch();
2235 }
2236 let expected = std::iter::once(candidates[0])
2237 .chain(candidates[1..].iter().rev().copied())
2238 .collect::<Vec<_>>();
2239 assert_eq!(order, expected);
2240 assert_eq!(rig.fetcher.num_hashes(), 0);
2241 }
2242
2243 #[test]
2244 fn old_request_error_does_not_settle_a_new_request_to_the_same_peer() {
2245 let config =
2246 TransactionFetcherConfig { max_inflight_requests_per_peer: 2, ..Default::default() };
2247 let mut rig = Rig::with_config(config);
2248 let p = peer(1);
2249 let h = hash(1);
2250 rig.add_peer(p);
2251 rig.announce(p, &[h]);
2252 rig.dispatch();
2253 let (_, first) = rig.take_request(p).unwrap();
2254 rig.fetcher.on_transactions_received([&h]);
2255 rig.announce(p, &[h]);
2256 rig.dispatch();
2257 let (_, second) = rig.take_request(p).unwrap();
2258 first.send(Err(RequestError::Timeout)).unwrap();
2259 rig.next_event().unwrap();
2260 assert_eq!(rig.fetcher.num_fetching_hashes(), 1);
2261 assert_eq!(rig.fetcher.fetching_peer(&h), Some(p));
2262 assert_eq!(rig.fetcher.candidate_peers(&h), vec![p]);
2263 second.send(Ok(PooledTransactions(vec![]))).unwrap();
2264 rig.next_event().unwrap();
2265 assert_eq!(rig.fetcher.num_hashes(), 0);
2266 }
2267
2268 #[test]
2269 fn first_timeout_retries_the_only_source_once() {
2270 let mut rig = Rig::new();
2271 let p = peer(1);
2272 let h = hash(1);
2273 rig.add_peer(p);
2274 rig.announce(p, &[h]);
2275 rig.dispatch();
2276 rig.fail(p, RequestError::Timeout);
2277 assert_eq!(rig.fetcher.candidate_peers(&h), vec![p]);
2278 assert_eq!(rig.dispatch(), 1);
2279 rig.fail(p, RequestError::Timeout);
2280 assert_eq!(rig.fetcher.num_hashes(), 0);
2281 assert_eq!(rig.dispatch(), 0);
2282 }
2283
2284 #[test]
2285 fn replacing_candidates_cannot_extend_the_fetch_attempt_budget() {
2286 let mut rig = Rig::new();
2287 let h = hash(1);
2288 for n in 1..=MAX_COUNT_CANDIDATE_PEERS_PER_HASH as u8 {
2289 rig.add_peer(peer(n));
2290 rig.announce(peer(n), &[h]);
2291 }
2292 for n in 0..MAX_FETCH_ATTEMPTS_PER_HASH {
2293 let newcomer = peer((n + MAX_COUNT_CANDIDATE_PEERS_PER_HASH + 1) as u8);
2294 rig.add_peer(newcomer);
2295 rig.announce(newcomer, &[h]);
2296 assert_eq!(rig.dispatch(), 1);
2297 let fetching = rig.fetcher.fetching_peer(&h).unwrap();
2298 rig.fail(fetching, RequestError::Timeout);
2299 }
2300 assert_eq!(rig.fetcher.num_hashes(), 0);
2301 }
2302
2303 #[test]
2304 fn self_eviction_skips_shared_and_stale_announcement_entries() {
2305 let mut rig = Rig::with_config(TransactionFetcherConfig {
2306 max_capacity_cache_txns_pending_fetch: 4,
2307 ..Default::default()
2308 });
2309 let honest = peer(1);
2310 let flooder = peer(2);
2311 for p in [honest, flooder] {
2312 rig.add_peer(p);
2313 }
2314 rig.announce(honest, &hashes(0..2));
2315 rig.announce(flooder, &hashes(10..12));
2316 rig.announce(flooder, &[hash(0), hash(12), hash(13), hash(14), hash(15), hash(16)]);
2318 assert_eq!(rig.fetcher.candidate_peers(&hash(1)), vec![honest]);
2319 assert_eq!(rig.fetcher.candidate_peers(&hash(0)), vec![honest, flooder]);
2320 assert_eq!(rig.fetcher.num_hashes(), 4);
2321 }
2322
2323 #[test]
2324 fn stale_queue_entries_do_not_hide_exclusive_eviction_victims() {
2325 let mut rig = Rig::with_config(TransactionFetcherConfig {
2326 max_capacity_cache_txns_pending_fetch: 16,
2327 ..Default::default()
2328 });
2329 let honest = peer(1);
2330 let flooder = peer(2);
2331 let newcomer = peer(3);
2332 for p in [honest, flooder, newcomer] {
2333 rig.add_peer(p);
2334 }
2335 rig.announce(honest, &[hash(0)]);
2336 rig.announce(flooder, &hashes(1..16));
2337 rig.fetcher.on_transactions_received(hashes(1..9).iter());
2338 rig.announce(flooder, &hashes(16..24));
2339 rig.announce(newcomer, &[hash(24)]);
2340 assert_eq!(rig.fetcher.candidate_peers(&hash(0)), vec![honest]);
2341 assert!(rig.fetcher.candidate_peers(&hash(9)).is_empty());
2342 }
2343
2344 #[test]
2345 fn capacity_eviction_targets_the_peer_with_the_most_hashes() {
2346 let config = TransactionFetcherConfig {
2347 max_capacity_cache_txns_pending_fetch: 100,
2348 max_announced_hashes_per_peer: 60,
2349 ..Default::default()
2350 };
2351 let mut rig = Rig::with_config(config);
2352 let honest = peer(1);
2353 let flooders = [peer(2), peer(3)];
2354 rig.add_peer(honest);
2355 for peer_id in flooders {
2356 rig.add_peer(peer_id);
2357 }
2358
2359 let honest_hashes = hashes(0..10);
2361 rig.announce(honest, &honest_hashes);
2362 rig.announce(flooders[0], &hashes(100..160));
2363 rig.announce(flooders[1], &hashes(200..260));
2364
2365 assert_eq!(rig.fetcher.num_hashes(), 100);
2366 for hash in &honest_hashes {
2367 assert_eq!(rig.fetcher.candidate_peers(hash), vec![honest], "{hash} was evicted");
2368 }
2369 assert_eq!(rig.fetcher.queued_hashes(&honest), honest_hashes);
2370
2371 assert!(rig.fetcher.candidate_peers(&hashes(100..101)[0]).is_empty());
2373 assert_eq!(rig.fetcher.candidate_peers(&hashes(259..260)[0]), vec![flooders[1]]);
2374 }
2375
2376 #[test]
2377 fn shared_queue_prefix_does_not_hide_exclusive_eviction_victim() {
2378 let mut rig = Rig::with_config(TransactionFetcherConfig {
2379 max_capacity_cache_txns_pending_fetch: 10,
2380 ..Default::default()
2381 });
2382 let honest = peer(1);
2383 let flooder = peer(2);
2384 let newcomer = peer(3);
2385 for p in [honest, flooder, newcomer] {
2386 rig.add_peer(p);
2387 }
2388 rig.announce(honest, &hashes(0..9));
2389 rig.announce(flooder, &hashes(0..8));
2390 rig.announce(flooder, &[hash(9)]);
2391 rig.announce(newcomer, &[hash(10)]);
2392 for h in hashes(0..8) {
2393 assert_eq!(rig.fetcher.candidate_peers(&h), vec![honest, flooder]);
2394 }
2395 assert!(rig.fetcher.candidate_peers(&hash(9)).is_empty());
2396 assert_eq!(rig.fetcher.candidate_peers(&hash(10)), vec![newcomer]);
2397 }
2398
2399 #[test]
2400 fn capacity_fallback_keeps_the_announcer_ready() {
2401 let mut rig = Rig::with_config(TransactionFetcherConfig {
2402 max_capacity_cache_txns_pending_fetch: 9,
2403 ..Default::default()
2404 });
2405 let stalled = peer(1);
2406 let source = peer(2);
2407 let announcer = peer(3);
2408 for p in [stalled, source, announcer] {
2409 rig.add_peer(p);
2410 }
2411 rig.announce(stalled, &hashes(0..8));
2412 rig.dispatch();
2413 rig.announce(source, &[hash(8)]);
2414 rig.announce(announcer, &hashes(0..10));
2417 rig.disconnect(source);
2418 assert_eq!(rig.dispatch(), 1);
2419 assert_eq!(rig.take_request(announcer).unwrap().0, vec![hash(8)]);
2420 }
2421
2422 #[test]
2423 fn capacity_eviction_preserves_coannounced_hashes() {
2424 let config = TransactionFetcherConfig {
2425 max_capacity_cache_txns_pending_fetch: 4,
2426 ..Default::default()
2427 };
2428 let mut rig = Rig::with_config(config);
2429 let honest = peer(1);
2430 let flooder = peer(2);
2431 let newcomer = peer(3);
2432 for peer_id in [honest, flooder, newcomer] {
2433 rig.add_peer(peer_id);
2434 }
2435 let shared = hash(0);
2436 rig.announce(honest, &[shared]);
2437 rig.announce(flooder, &hashes(0..4));
2438 rig.announce(newcomer, &[hash(4)]);
2439 assert_eq!(rig.fetcher.num_hashes(), 4);
2440 assert_eq!(rig.fetcher.candidate_peers(&shared), vec![honest, flooder]);
2441 assert!(rig.fetcher.candidate_peers(&hash(1)).is_empty());
2442 assert_eq!(rig.fetcher.candidate_peers(&hash(4)), vec![newcomer]);
2443 rig.dispatch();
2444 assert_eq!(rig.fetcher.fetching_peer(&shared), Some(honest));
2445 assert_eq!(rig.take_request(flooder).unwrap().0, vec![hash(2), hash(3)]);
2446 }
2447
2448 #[test]
2449 fn capacity_with_only_shared_hashes_admits_new_announcements() {
2450 let config = TransactionFetcherConfig {
2451 max_capacity_cache_txns_pending_fetch: 2,
2452 ..Default::default()
2453 };
2454 let mut rig = Rig::with_config(config);
2455 for peer_id in [peer(1), peer(2), peer(3)] {
2456 rig.add_peer(peer_id);
2457 }
2458 for peer_id in [peer(1), peer(2)] {
2459 rig.announce(peer_id, &hashes(0..2));
2460 }
2461 rig.announce(peer(3), &[hash(2)]);
2462 assert_eq!(rig.fetcher.num_hashes(), 2);
2463 assert_eq!(rig.fetcher.candidate_peers(&hash(2)), vec![peer(3)]);
2464 assert!(rig.fetcher.candidate_peers(&hash(0)).is_empty());
2465 assert_eq!(rig.fetcher.candidate_peers(&hash(1)), vec![peer(1), peer(2)]);
2466 rig.dispatch();
2467 assert_eq!(rig.fetcher.num_fetching_hashes(), 2);
2468 }
2469
2470 #[test]
2471 fn retracked_hash_is_not_evicted_at_its_stale_queue_position() {
2472 let mut rig = Rig::with_config(TransactionFetcherConfig {
2473 max_capacity_cache_txns_pending_fetch: 3,
2474 ..Default::default()
2475 });
2476 let p = peer(1);
2477 rig.add_peer(p);
2478 rig.announce(p, &hashes(0..3));
2479 rig.fetcher.on_transactions_received([&hash(0)]);
2480 rig.announce(p, &[hash(0)]);
2481 rig.announce(p, &[hash(3)]);
2482 assert_eq!(rig.fetcher.candidate_peers(&hash(0)), vec![p]);
2483 assert!(rig.fetcher.candidate_peers(&hash(1)).is_empty());
2484 rig.dispatch();
2485 assert_eq!(rig.take_request(p).unwrap().0, vec![hash(2), hash(0), hash(3)]);
2486 }
2487
2488 #[test]
2489 fn shared_fallback_ignores_stale_tracking_generations() {
2490 let mut rig = Rig::with_config(TransactionFetcherConfig {
2491 max_capacity_cache_txns_pending_fetch: 3,
2492 ..Default::default()
2493 });
2494 for p in [peer(1), peer(2), peer(3)] {
2495 rig.add_peer(p);
2496 }
2497 for p in [peer(1), peer(2)] {
2498 rig.announce(p, &hashes(0..3));
2499 }
2500 rig.fetcher.on_transactions_received([&hash(0)]);
2501 for p in [peer(1), peer(2)] {
2502 rig.announce(p, &[hash(0)]);
2503 }
2504 rig.announce(peer(3), &[hash(3)]);
2505 assert_eq!(rig.fetcher.candidate_peers(&hash(0)), vec![peer(1), peer(2)]);
2506 assert!(rig.fetcher.candidate_peers(&hash(1)).is_empty());
2507 }
2508
2509 #[test]
2510 fn eviction_order_stays_bounded_when_hashes_are_tracked_again() {
2511 let config = TransactionFetcherConfig {
2512 max_capacity_cache_txns_pending_fetch: 8,
2513 ..Default::default()
2514 };
2515 let mut rig = Rig::with_config(config);
2516 let peer_a = peer(1);
2517 rig.add_peer(peer_a);
2518 let batch = hashes(0..8);
2519
2520 for _ in 0..10 {
2523 rig.announce(peer_a, &batch);
2524 assert_eq!(rig.fetcher.num_hashes(), 8);
2525 rig.dispatch();
2526 rig.fail(peer_a, RequestError::BadResponse);
2527 assert_eq!(rig.fetcher.num_hashes(), 0);
2528 }
2529 rig.announce(peer_a, &batch);
2530 assert_eq!(rig.fetcher.num_hashes(), 8);
2531
2532 rig.announce(peer_a, &hashes(8..9));
2534 assert_eq!(rig.fetcher.num_hashes(), 8);
2535 assert!(rig.fetcher.candidate_peers(&batch[0]).is_empty());
2536 }
2537
2538 #[test]
2539 fn disconnected_remembered_candidates_are_forgotten() {
2540 let mut rig = Rig::new();
2541 let hash = hash(1);
2542 let peers = (1..=MAX_COUNT_CANDIDATE_PEERS_PER_HASH as u8).map(peer).collect::<Vec<_>>();
2543 for peer_id in &peers {
2544 rig.add_peer(*peer_id);
2545 rig.announce(*peer_id, &[hash]);
2546 }
2547 let (queued, remembered) = peers.split_at(MAX_COUNT_EAGER_CANDIDATE_PEERS_PER_HASH);
2548
2549 for peer_id in remembered {
2551 rig.disconnect(*peer_id);
2552 }
2553 assert_eq!(rig.fetcher.candidate_peers(&hash), queued);
2554
2555 let newcomer = peer(42);
2557 rig.add_peer(newcomer);
2558 rig.announce(newcomer, &[hash]);
2559 assert!(rig.fetcher.candidate_peers(&hash).contains(&newcomer));
2560
2561 for peer_id in queued {
2563 rig.disconnect(*peer_id);
2564 }
2565 rig.disconnect(newcomer);
2566 assert_eq!(rig.fetcher.num_hashes(), 0);
2567 }
2568
2569 #[test]
2570 fn remembered_candidates_take_over_when_queued_ones_disconnect() {
2571 let mut rig = Rig::new();
2572 let hash = hash(1);
2573 let peers =
2574 (1..=MAX_COUNT_EAGER_CANDIDATE_PEERS_PER_HASH as u8 + 1).map(peer).collect::<Vec<_>>();
2575 for peer_id in &peers {
2576 rig.add_peer(*peer_id);
2577 rig.announce(*peer_id, &[hash]);
2578 }
2579 let last = *peers.last().unwrap();
2580 assert!(rig.fetcher.queued_hashes(&last).is_empty());
2581
2582 for peer_id in &peers[..peers.len() - 1] {
2584 rig.disconnect(*peer_id);
2585 }
2586 assert_eq!(rig.fetcher.candidate_peers(&hash), vec![last]);
2587 assert_eq!(rig.fetcher.queued_hashes(&last), vec![hash]);
2588 assert_eq!(rig.dispatch(), 1);
2589 assert_eq!(rig.fetcher.fetching_peer(&hash), Some(last));
2590 }
2591
2592 #[test]
2593 fn global_inflight_limit_defers_ready_peers() {
2594 let config = TransactionFetcherConfig { max_inflight_requests: 1, ..Default::default() };
2595 let mut rig = Rig::with_config(config);
2596 let peer_a = peer(1);
2597 let peer_b = peer(2);
2598 rig.add_peer(peer_a);
2599 rig.add_peer(peer_b);
2600
2601 rig.announce(peer_a, &hashes(0..2));
2602 rig.announce(peer_b, &hashes(2..4));
2603
2604 assert_eq!(rig.dispatch(), 1);
2605 let (_, response) = rig.take_request(peer_a).unwrap();
2606 assert!(rig.take_request(peer_b).is_none());
2607 assert_eq!(rig.dispatch(), 0);
2608
2609 rig.announce(peer_a, &hashes(4..5));
2611 drop(response);
2612 rig.next_event().unwrap();
2613 assert_eq!(rig.dispatch(), 1);
2614 assert!(rig.take_request(peer_b).is_some());
2615 assert!(rig.take_request(peer_a).is_none(), "peer_b was ready first");
2616 }
2617
2618 #[test]
2619 fn dispatch_applies_the_request_limit_independently_to_each_peer() {
2620 let mut rig = Rig::new();
2621 let peer_a = peer(1);
2622 let peer_b = peer(2);
2623 rig.add_peer(peer_a);
2624 rig.add_peer(peer_b);
2625
2626 rig.announce(peer_a, &hashes(0..100));
2627 rig.announce(peer_b, &hashes(100..200));
2628
2629 assert_eq!(rig.dispatch_with_budget(0), 0);
2630
2631 assert_eq!(rig.dispatch_with_budget(40), 2);
2633 let (requested, response_a) = rig.take_request(peer_a).unwrap();
2634 assert_eq!(requested, hashes(0..40));
2635 let (requested, _) = rig.take_request(peer_b).unwrap();
2636 assert_eq!(requested, hashes(100..140));
2637 assert_eq!(rig.fetcher.num_fetching_hashes(), 80);
2638
2639 response_a.send(Err(RequestError::BadResponse)).unwrap();
2641 rig.next_event().unwrap();
2642 assert_eq!(rig.dispatch_with_budget(10), 1);
2643 let (requested, _) = rig.take_request(peer_a).unwrap();
2644 assert_eq!(requested, hashes(40..50));
2645 }
2646
2647 #[test]
2648 fn stalled_peers_do_not_shrink_requests_to_an_idle_peer() {
2649 let mut rig = Rig::new();
2650 for n in 0..16 {
2651 let p = peer(n + 1);
2652 rig.add_peer(p);
2653 rig.announce(p, &hashes(u64::from(n) * 256..u64::from(n + 1) * 256));
2654 }
2655 assert_eq!(rig.dispatch_with_budget(4096), 16);
2656 assert_eq!(rig.fetcher.num_fetching_hashes(), 4096);
2657 let honest = peer(17);
2660 rig.add_peer(honest);
2661 rig.announce(honest, &hashes(4096..4352));
2662 assert_eq!(rig.dispatch_with_budget(4096), 1);
2663 assert_eq!(rig.take_request(honest).unwrap().0, hashes(4096..4352));
2664 }
2665
2666 #[test]
2667 fn full_session_channel_rolls_back_request() {
2668 let mut rig = Rig::new();
2669 let peer_a = peer(1);
2670 rig.add_peer_with_capacity(peer_a, 1);
2671 let hashes = hashes(0..3);
2672
2673 let (blocker, _rx) = oneshot::channel();
2675 rig.peers[&peer_a]
2676 .request_tx()
2677 .try_send(PeerRequest::GetPooledTransactions {
2678 request: GetPooledTransactions(vec![]),
2679 response: blocker,
2680 })
2681 .unwrap();
2682
2683 rig.announce(peer_a, &hashes);
2684 assert_eq!(rig.dispatch(), 0);
2685 assert_eq!(rig.fetcher.num_pending_hashes(), 3);
2686 assert_eq!(rig.fetcher.num_inflight_requests(), 0);
2687 assert!(rig.fetcher.is_idle(&peer_a));
2688 assert_eq!(rig.fetcher.queued_hashes(&peer_a), hashes);
2689
2690 let (blocked, _) = rig.take_request(peer_a).unwrap();
2692 assert!(blocked.is_empty());
2693 assert_eq!(rig.dispatch(), 1);
2694 let (requested, _) = rig.take_request(peer_a).unwrap();
2695 assert_eq!(requested, hashes);
2696 }
2697
2698 #[test]
2699 fn closed_session_channel_disconnects_peer() {
2700 let mut rig = Rig::new();
2701 let peer_a = peer(1);
2702 let peer_b = peer(2);
2703 rig.add_peer(peer_a);
2704 rig.add_peer(peer_b);
2705 let hashes = hashes(0..2);
2706
2707 rig.announce(peer_a, &hashes);
2708 rig.announce(peer_b, &hashes[..1]);
2709 rig.sessions.remove(&peer_a);
2711
2712 assert_eq!(rig.dispatch(), 1);
2713 assert!(rig.fetcher.candidate_peers(&hashes[0]).contains(&peer_b));
2714 assert_eq!(rig.fetcher.candidate_peers(&hashes[0]), vec![peer_b]);
2715 assert_eq!(rig.fetcher.num_hashes(), 1, "hash only announced by the gone peer is dropped");
2716 assert_eq!(rig.fetcher.fetching_peer(&hashes[0]), Some(peer_b));
2717 }
2718
2719 #[test]
2720 fn busy_peer_queue_is_compacted() {
2721 let config =
2722 TransactionFetcherConfig { max_announced_hashes_per_peer: 4, ..Default::default() };
2723 let mut rig = Rig::with_config(config);
2724 let peer_a = peer(1);
2725 let peer_b = peer(2);
2726 rig.add_peer(peer_a);
2727 rig.add_peer(peer_b);
2728
2729 rig.announce(peer_a, &hashes(0..1));
2731 rig.dispatch();
2732
2733 for round in 1..20u64 {
2735 let hashes = hashes(round * 10..round * 10 + 4);
2736 rig.announce(peer_a, &hashes);
2737 rig.fetcher.on_transactions_received(&hashes);
2738 rig.fetcher.assert_invariants();
2739 }
2740 assert!(rig.fetcher.queued_hashes(&peer_a).len() <= 8);
2741 }
2742
2743 #[test]
2744 fn late_response_for_reassigned_hash_is_ignored() {
2745 let mut rig = Rig::new();
2746 let peer_a = peer(1);
2747 let peer_b = peer(2);
2748 rig.add_peer(peer_a);
2749 rig.add_peer(peer_b);
2750 let txs = pooled_txs(1);
2751 let hashes = hashes_of(&txs);
2752
2753 rig.announce(peer_a, &hashes);
2754 rig.dispatch();
2755 rig.fetcher.on_transactions_received(&hashes);
2757 rig.announce(peer_b, &hashes);
2758 rig.dispatch();
2759 assert_eq!(rig.fetcher.fetching_peer(&hashes[0]), Some(peer_b));
2760
2761 rig.respond(peer_a, vec![]);
2763 assert_eq!(rig.fetcher.fetching_peer(&hashes[0]), Some(peer_b));
2764 assert_eq!(rig.fetcher.candidate_peers(&hashes[0]), vec![peer_b]);
2765
2766 rig.respond(peer_b, txs);
2767 assert_eq!(rig.fetcher.num_hashes(), 0);
2768 }
2769
2770 #[test]
2771 fn random_operations_keep_invariants() {
2772 let mut rng = StdRng::seed_from_u64(0x5eed);
2773 let txs = pooled_txs(150);
2774 let all_hashes = hashes_of(&txs);
2775 let by_hash = txs.iter().map(|tx| (*tx.tx_hash(), tx.clone())).collect::<B256Map<_>>();
2776 let peer_ids = (1..=24).map(peer).collect::<Vec<_>>();
2778
2779 let config = TransactionFetcherConfig {
2780 max_inflight_requests: 8,
2781 max_inflight_requests_per_peer: 2,
2782 max_capacity_cache_txns_pending_fetch: 120,
2783 max_announced_hashes_per_peer: 40,
2784 ..Default::default()
2785 };
2786 let mut rig = Rig::with_config(config);
2787 for peer_id in &peer_ids {
2788 rig.add_peer_with_capacity(*peer_id, 4);
2789 }
2790 let mut connected = peer_ids.clone();
2791 let mut outstanding: Vec<(PeerId, Vec<TxHash>, ResponseSender)> = Vec::new();
2793
2794 for _ in 0..5000 {
2795 match rng.random_range(0..100u32) {
2796 0..=39 => {
2797 let Some(&peer_id) = connected.choose(&mut rng) else { continue };
2798 let count = rng.random_range(1..=30);
2799 let entries = all_hashes
2800 .choose_multiple(&mut rng, count)
2801 .map(|hash| {
2802 let size = match rng.random_range(0..10) {
2803 0 => 0,
2804 1 => 200 * KIB,
2805 _ => rng.random_range(100..1500),
2806 };
2807 (*hash, size)
2808 })
2809 .collect::<Vec<_>>();
2810 rig.announce_with_sizes(peer_id, entries);
2811 }
2812 40..=59 => {
2813 let budget =
2814 if rng.random_bool(0.2) { rng.random_range(0..60) } else { usize::MAX };
2815 rig.dispatch_with_budget(budget);
2816 for peer_id in &connected {
2817 while let Some((requested, response)) = rig.take_request(*peer_id) {
2818 assert!(!requested.is_empty());
2819 assert!(
2820 requested.len() <=
2821 SOFT_LIMIT_COUNT_HASHES_IN_GET_POOLED_TRANSACTIONS_REQUEST
2822 );
2823 let unique = requested.iter().copied().collect::<B256Set>();
2824 assert_eq!(unique.len(), requested.len(), "request has duplicates");
2825 outstanding.push((*peer_id, requested, response));
2826 }
2827 }
2828 }
2829 60..=84 => {
2830 if outstanding.is_empty() {
2831 continue
2832 }
2833 let (_, requested, response) =
2834 outstanding.swap_remove(rng.random_range(0..outstanding.len()));
2835 match rng.random_range(0..10) {
2836 0..=5 => {
2837 let mut delivered = requested
2838 .iter()
2839 .filter(|_| rng.random_bool(0.7))
2840 .map(|hash| by_hash[hash].clone())
2841 .collect::<Vec<_>>();
2842 if rng.random_bool(0.1) &&
2843 let Some(duplicate) = delivered.first().cloned()
2844 {
2845 delivered.push(duplicate);
2846 }
2847 if rng.random_bool(0.1) {
2848 delivered.push(txs.choose(&mut rng).unwrap().clone());
2849 }
2850 let _ = response.send(Ok(PooledTransactions(delivered)));
2851 }
2852 6..=7 => {
2853 let _ = response.send(Ok(PooledTransactions(vec![])));
2854 }
2855 8 => {
2856 let _ = response.send(Err(RequestError::Timeout));
2857 }
2858 _ => drop(response),
2859 }
2860 rig.drain_events();
2861 }
2862 85..=89 => {
2863 let count = rng.random_range(1..=10);
2864 let received =
2865 all_hashes.choose_multiple(&mut rng, count).copied().collect::<Vec<_>>();
2866 rig.fetcher.on_transactions_received(&received);
2867 rig.fetcher.assert_invariants();
2868 }
2869 90..=94 => {
2870 if connected.len() > 1 {
2871 let peer_id = connected.swap_remove(rng.random_range(0..connected.len()));
2872 rig.disconnect(peer_id);
2873 outstanding.retain(|(id, ..)| *id != peer_id);
2875 rig.drain_events();
2876 }
2877 }
2878 _ => {
2879 if let Some(&peer_id) = peer_ids.iter().find(|id| !connected.contains(id)) {
2880 rig.add_peer_with_capacity(peer_id, 4);
2881 connected.push(peer_id);
2882 }
2883 }
2884 }
2885 }
2886
2887 drop(outstanding);
2889 for peer_id in &connected {
2890 rig.disconnect(*peer_id);
2891 }
2892 rig.drain_events();
2893 assert_eq!(rig.fetcher.num_inflight_requests(), 0);
2894 assert_eq!(rig.fetcher.num_fetching_hashes(), 0);
2895 assert_eq!(rig.fetcher.num_hashes(), 0);
2896 assert_eq!(rig.fetcher.num_peers(), 0);
2897 }
2898 #[test]
2899 fn resolved_requests_are_yielded_one_per_poll() {
2900 let mut rig = Rig::new();
2901 let peers = (1..=3).map(peer).collect::<Vec<_>>();
2902 for (i, peer_id) in peers.iter().enumerate() {
2903 rig.add_peer(*peer_id);
2904 rig.announce(*peer_id, &hashes(i as u64..i as u64 + 1));
2905 }
2906 assert_eq!(rig.dispatch(), 3);
2907 for peer_id in &peers {
2908 let (_, response) = rig.take_request(*peer_id).unwrap();
2909 response.send(Err(RequestError::Timeout)).unwrap();
2910 }
2911
2912 let events = rig.drain_events();
2913 assert_eq!(events.len(), 3);
2914 assert!(events.iter().all(|event| matches!(event, FetchEvent::FetchError { .. })));
2915 assert!(rig.next_event().is_none());
2916 assert_eq!(rig.fetcher.num_inflight_requests(), 0);
2917 }
2918
2919 #[test]
2920 fn requests_sent_after_a_poll_register_wakers_on_the_next_poll() {
2921 let mut rig = Rig::new();
2922 let peer_a = peer(1);
2923 rig.add_peer(peer_a);
2924 let counter = Arc::new(WakeCounter::default());
2925 let waker = waker(counter.clone());
2926 let mut cx = Context::from_waker(&waker);
2927
2928 assert!(rig.fetcher.poll_next_unpin(&mut cx).is_pending());
2929 rig.announce(peer_a, &hashes(0..1));
2930 assert_eq!(rig.dispatch(), 1);
2931
2932 assert!(rig.fetcher.poll_next_unpin(&mut cx).is_pending());
2935 let wakes = counter.wakes();
2936 let (_, response) = rig.take_request(peer_a).unwrap();
2937 response.send(Err(RequestError::Timeout)).unwrap();
2938 assert_eq!(counter.wakes(), wakes + 1);
2939 assert!(matches!(
2940 rig.fetcher.poll_next_unpin(&mut cx),
2941 Poll::Ready(Some(FetchEvent::FetchError { .. }))
2942 ));
2943 }
2944
2945 #[test]
2946 fn global_capacity_bounds_tracked_hashes_across_peers() {
2947 let config = TransactionFetcherConfig {
2948 max_capacity_cache_txns_pending_fetch: 40,
2949 ..Default::default()
2950 };
2951 let mut rig = Rig::with_config(config);
2952 let peers = (1..=4).map(peer).collect::<Vec<_>>();
2953 for (i, peer_id) in peers.iter().enumerate() {
2954 rig.add_peer(*peer_id);
2955 rig.announce_unsized(*peer_id, &hashes(i as u64 * 16..(i as u64 + 1) * 16));
2956 }
2957
2958 assert_eq!(rig.fetcher.num_hashes(), 40);
2961 for (i, peer_id) in peers.iter().enumerate() {
2962 let first = i as u64 * 16;
2963 assert!(rig.fetcher.candidate_peers(&hash(first)).is_empty());
2964 assert!(rig.fetcher.candidate_peers(&hash(first + 3)).is_empty());
2965 assert_eq!(rig.fetcher.candidate_peers(&hash(first + 8)), vec![*peer_id]);
2966 assert_eq!(rig.fetcher.candidate_peers(&hash(first + 15)), vec![*peer_id]);
2967 }
2968 rig.fetcher.assert_invariants();
2969
2970 assert_eq!(rig.dispatch(), 4);
2972 let mut requested_total = 0;
2973 for peer_id in &peers {
2974 let (requested, _) = rig.take_request(*peer_id).unwrap();
2975 assert!(requested.len() <= 256);
2976 requested_total += requested.len();
2977 }
2978 assert_eq!(requested_total, 40, "every remaining hash is requested");
2979 rig.fetcher.assert_invariants();
2980 }
2981
2982 #[test]
2983 fn delivery_order_does_not_matter() {
2984 let mut rig = Rig::new();
2985 let peer_a = peer(1);
2986 let peer_b = peer(2);
2987 rig.add_peer(peer_a);
2988 rig.add_peer(peer_b);
2989 let txs = pooled_txs(4);
2990 let hashes = hashes_of(&txs);
2991
2992 rig.announce(peer_a, &hashes);
2993 rig.announce(peer_b, &hashes);
2994 rig.dispatch();
2995
2996 let event = rig.respond(peer_a, vec![txs[3].clone(), txs[0].clone()]);
2999 let FetchEvent::TransactionsFetched { transactions, .. } = event else { panic!() };
3000 assert_eq!(transactions.len(), 2);
3001 assert_eq!(rig.fetcher.num_hashes(), 2);
3002 assert_eq!(rig.fetcher.candidate_peers(&hashes[1]), vec![peer_b]);
3003 assert_eq!(rig.fetcher.candidate_peers(&hashes[2]), vec![peer_b]);
3004 assert_eq!(rig.fetcher.queued_hashes(&peer_b), hashes[1..3]);
3005 }
3006
3007 #[test]
3008 fn partial_delivery_with_concurrent_requests_per_peer() {
3009 let config =
3010 TransactionFetcherConfig { max_inflight_requests_per_peer: 2, ..Default::default() };
3011 let mut rig = Rig::with_config(config);
3012 let peer_a = peer(1);
3013 rig.add_peer(peer_a);
3014 let txs = pooled_txs(300);
3015 let hashes = hashes_of(&txs);
3016
3017 rig.announce(peer_a, &hashes);
3018 assert_eq!(rig.dispatch(), 2);
3019 let (first, response_first) = rig.take_request(peer_a).unwrap();
3020 let (second, response_second) = rig.take_request(peer_a).unwrap();
3021 assert_eq!(first, hashes[..256]);
3022 assert_eq!(second, hashes[256..]);
3023
3024 response_first.send(Ok(PooledTransactions(txs[255..256].to_vec()))).unwrap();
3027 rig.next_event().unwrap();
3028 assert_eq!(rig.fetcher.num_hashes(), 44);
3029 assert_eq!(rig.fetcher.num_fetching_hashes(), 44);
3030 assert!(!rig.fetcher.is_idle(&peer_a));
3031
3032 response_second.send(Ok(PooledTransactions(txs[256..].to_vec()))).unwrap();
3033 rig.next_event().unwrap();
3034 assert_eq!(rig.fetcher.num_hashes(), 0);
3035 assert!(rig.fetcher.is_idle(&peer_a));
3036 }
3037
3038 #[test]
3039 fn old_request_resolves_after_peer_reconnected_and_reannounced() {
3040 let mut rig = Rig::new();
3041 let peer_a = peer(1);
3042 rig.add_peer(peer_a);
3043 let txs = pooled_txs(2);
3044 let hashes = hashes_of(&txs);
3045
3046 rig.announce(peer_a, &hashes);
3047 rig.dispatch();
3048 let (_, response) = rig.take_request(peer_a).unwrap();
3049
3050 rig.disconnect(peer_a);
3052 rig.add_peer(peer_a);
3053 rig.announce(peer_a, &hashes);
3054 assert_eq!(rig.fetcher.candidate_peers(&hashes[0]), vec![peer_a]);
3055 assert_eq!(rig.fetcher.num_fetching_hashes(), 2);
3056 assert_eq!(rig.dispatch(), 0, "the hashes are still being fetched");
3057
3058 response.send(Ok(PooledTransactions(txs[1..].to_vec()))).unwrap();
3061 rig.next_event().unwrap();
3062 assert_eq!(rig.fetcher.num_hashes(), 1);
3063 assert_eq!(rig.fetcher.candidate_peers(&hashes[0]), vec![peer_a]);
3064 assert_eq!(rig.dispatch(), 1);
3065 let (requested, _) = rig.take_request(peer_a).unwrap();
3066 assert_eq!(requested, hashes[..1]);
3067 }
3068
3069 #[test]
3070 fn requeued_hashes_are_not_queued_twice() {
3071 let mut rig = Rig::new();
3072 let peer_a = peer(1);
3073 let peer_b = peer(2);
3074 rig.add_peer(peer_a);
3075 rig.add_peer(peer_b);
3076 let hashes = hashes(0..3);
3077
3078 rig.announce(peer_b, &hashes[..1]);
3080 rig.dispatch();
3081 rig.announce(peer_a, &hashes[1..]);
3082 rig.announce(peer_b, &hashes[1..]);
3083 rig.dispatch();
3084 assert_eq!(rig.fetcher.fetching_peer(&hashes[1]), Some(peer_a));
3085 assert_eq!(rig.fetcher.queued_hashes(&peer_b), hashes[1..]);
3086
3087 rig.fail(peer_a, RequestError::Timeout);
3088 assert_eq!(rig.fetcher.queued_hashes(&peer_b), hashes[1..], "no duplicates");
3089 }
3090}