Skip to main content

reth_network/transactions/
fetcher.rs

1//! Fetches transactions that peers announced with `NewPooledTransactionHashes`.
2//!
3//! The [`TransactionFetcher`] tracks every announced hash that is not known yet, together with the
4//! peers that announced it, and turns those announcements into `GetPooledTransactions` requests.
5//!
6//! # Model
7//!
8//! Every tracked hash is in exactly one of two states:
9//!
10//! - _pending_: waiting for one of the peers that announced it, its _candidates_, to become idle
11//! - _fetching_: part of exactly one inflight request
12//!
13//! A hash remembers the peers that announced it as candidates. Each peer keeps a FIFO queue of
14//! hashes to request, and the first peers to announce a hash get it queued right away. Later
15//! announcers are only remembered, so that a hash is not given up on when the first announcers
16//! fail to deliver it, without costing a queue entry per announcement. A request for an idle peer
17//! is built by draining its queue in the order the announcements were processed, skipping hashes
18//! that are being fetched from another peer or are not tracked anymore, until the request is
19//! full. Packing a request is therefore proportional to the request size and independent of the
20//! total number of pending hashes. FIFO refers to the fetcher's input order; callers may
21//! reorder or deduplicate hashes within a wire announcement.
22//!
23//! When a request resolves, delivered hashes are dropped from tracking. Undelivered hashes go back
24//! to pending and are queued for their remaining candidates, most recent announcers first, or are
25//! dropped if none remain. The responding peer is dropped as a candidate for every undelivered
26//! hash that precedes the last delivered hash of the request, treating it as skipped, and for
27//! all requested hashes if the response was empty or the request failed.
28//! Undelivered hashes after the last delivered hash keep the peer as a candidate only if it
29//! delivered at least half of the request, treating that tail as potentially truncated. This is
30//! a retry heuristic: transaction counts do not establish whether the response reached a byte
31//! limit. Requiring half the request limits repeated low-progress retries.
32//!
33//! Request timeouts are enforced by the peer's session, which resolves the request with
34//! [`RequestError::Timeout`], so the fetcher does not run any timers. A first timeout is retried
35//! once when the responding peer is the only remaining source; otherwise another source is tried.
36//!
37//! # Bounds
38//!
39//! - a configurable number of inflight requests per peer, one by default, and a global inflight
40//!   request limit
41//! - a per peer limit on the number of tracked hashes it is a candidate for, so a single peer
42//!   cannot flood the fetcher with announcements
43//! - a global limit on the number of tracked hashes, at which the peer tracking the most hashes
44//!   gives up its oldest pending hash, so a group of peers flooding the fetcher evicts its own
45//!   hashes rather than everyone else's
46//! - a fixed number of candidates per hash and a separate fetch-attempt limit, so remembering
47//!   recent fallback sources does not allow unlimited retries
48
49use 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
85/// Maximum live entries inspected in a global fallback eviction search.
86const MAX_EVICTION_ATTEMPTS: usize = 8;
87
88/// Fetches transactions that peers announced but that are not in the pool yet.
89///
90/// Announcements are recorded with [`Self::on_announcement`], requests are sent with
91/// [`Self::dispatch`] and resolved requests are yielded as [`FetchEvent`]s by the [`Stream`]
92/// implementation. See the [module docs](self) for how requests are scheduled.
93#[derive(Debug)]
94pub struct TransactionFetcher<N: NetworkPrimitives = EthNetworkPrimitives> {
95    /// All tracked hashes with their candidate peers and fetch state.
96    hashes: B256Map<TxEntry>,
97    /// Tracked hashes in the order they were added, used to evict the oldest pending hash once
98    /// the fetcher is at capacity and the peer tracking the most hashes has none pending. Entries
99    /// are removed lazily, so it may contain hashes that are not tracked anymore.
100    order: VecDeque<(TxHash, u64)>,
101    /// Fetch state of all peers that announced tracked hashes.
102    peers: HashMap<PeerKey, PeerState>,
103    /// Maps peer ids to their compact key.
104    peer_keys: HashMap<PeerId, PeerKey, FbBuildHasher<64>>,
105    /// Key assigned to the next new peer. Keys are never reused, so the candidates of a hash can
106    /// safely refer to peers that disconnected in the meantime.
107    next_peer_key: u32,
108    /// Distinguishes overlapping requests, including repeat requests for the same hash.
109    next_request_id: u64,
110    /// Distinguishes a re-announcement from stale queue entries for an earlier tracked hash.
111    next_generation: u64,
112    /// Idle peers with queued hashes, in the order they became ready.
113    ready: VecDeque<PeerKey>,
114    /// All inflight `GetPooledTransactions` requests.
115    inflight: FuturesUnordered<InflightRequest<N::PooledTransaction>>,
116    /// Number of tracked hashes that are part of an inflight request.
117    num_fetching: usize,
118    /// Reused when verifying responses, so no sets are allocated per response.
119    scratch_requested: B256Set,
120    /// Reused when verifying responses, so no sets are allocated per response.
121    scratch_delivered: B256Set,
122    /// Reused when processing announcements, so no vector is allocated per announcement.
123    scratch_queue: Vec<(TxHash, u64)>,
124    /// Configured limits.
125    config: TransactionFetcherConfig,
126    metrics: TransactionFetcherMetrics,
127}
128
129impl<N: NetworkPrimitives> TransactionFetcher<N> {
130    /// Creates a new fetcher with the given config.
131    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    /// Returns the fetcher's config.
155    pub const fn config(&self) -> &TransactionFetcherConfig {
156        &self.config
157    }
158
159    /// Returns the number of tracked hashes, pending and fetching.
160    pub fn num_hashes(&self) -> usize {
161        self.hashes.len()
162    }
163
164    /// Returns the number of hashes that are waiting for an idle candidate peer.
165    pub fn num_pending_hashes(&self) -> usize {
166        self.hashes.len().saturating_sub(self.num_fetching)
167    }
168
169    /// Returns the number of hashes that are part of an inflight request.
170    pub const fn num_fetching_hashes(&self) -> usize {
171        self.num_fetching
172    }
173
174    /// Returns the number of inflight requests.
175    pub fn num_inflight_requests(&self) -> usize {
176        self.inflight.len()
177    }
178
179    /// Returns `true` if there is no inflight request to the peer.
180    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    /// Records hashes announced by the peer and queues them for fetching.
188    ///
189    /// The metadata of a hash is the announced transaction type and size, if the announcement
190    /// carried it. The caller is expected to have filtered out hashes that are already known.
191    ///
192    /// Requests are only sent by [`Self::dispatch`].
193    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        // The peer's count of tracked hashes and the hashes to queue for it are kept locally and
207        // applied once at the end, so every announced hash costs a single map lookup.
208        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        // Entries at the front of `queue` already evicted or moved into the peer queue.
213        let mut queue_start = 0;
214        // Other peers can only lose tracked hashes during this announcement. Refresh their
215        // heap entries lazily instead of scanning every peer for each hash at capacity.
216        let mut eviction = None;
217        // Gossip commonly replaces the same fallback peer for every hash in a batch. Apply
218        // that peer's accounting once, flushing before eviction needs up-to-date counts.
219        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            // whether the hash is queued for the peer right away
226            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                        // announced before by this peer, keep the latest size
231                        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                        // Keep the first eager sources and the most recent fallback sources.
240                        // A fixed first-N list lets early co-announcers exclude every later source.
241                        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                    // The first announcers get the hash queued right away, later ones are only
253                    // remembered and asked once the earlier ones failed to deliver.
254                    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                        // the evicted hash may be one of this peer's, so its count is synced
275                        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    /// Sends `GetPooledTransactions` requests to idle peers that have queued hashes.
343    ///
344    /// Stops when the global inflight request limit is reached. `max_hashes_per_request` caps
345    /// each request independently of hashes inflight to other peers, so stalled peers cannot
346    /// shrink requests to responsive peers. Zero sends nothing. The caller must enforce its
347    /// import backpressure when consuming responses; decoded responses remain bounded by
348    /// the configured inflight request limit.
349    ///
350    /// Returns the number of requests sent. New requests are only polled by the next call to
351    /// [`Stream::poll_next`], so the caller must poll the fetcher again if any were sent.
352    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        // peers whose session channel is full, they get another chance on the next dispatch
363        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            // the session is gone if the manager doesn't know the peer anymore
374            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                    // the peer may be allowed more than one inflight request
412                    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    /// Stops tracking the given hashes because the transactions were received, e.g. over
435    /// broadcast.
436    ///
437    /// Hashes that are part of an inflight request are simply not rescheduled when that request
438    /// resolves.
439    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    /// Removes the peer as a candidate for all hashes it announced and drops pending hashes that
446    /// have no candidate left.
447    ///
448    /// An inflight request to the peer resolves with an error once its session is gone, which
449    /// reschedules the requested hashes.
450    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        // Every tracked hash is visited, since the peer is a candidate of hashes that are not in
457        // its queue as well. Disconnects are rare compared to announcements, so this is cheaper
458        // overall than checking for gone peers whenever candidates are counted.
459        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            // the hash must stay queued for at least one candidate
472            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 at the front in reverse announcement order.
480        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        // retried hashes go to the most recent announcers first
489        for key in requeued {
490            self.mark_ready(key, QueuePosition::Front);
491        }
492    }
493
494    /// Updates the fetcher's gauges.
495    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    /// Returns the key of the peer, registering the peer if it isn't known yet.
502    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    /// Marks the peer as ready for a request if it is idle and has hashes queued, either behind
512    /// the other ready peers or ahead of them.
513    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    /// Queues the hash for the peer, at the back or the front of its queue, and marks it as
528    /// queued in the hash's candidate entry.
529    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    /// Records a newly tracked hash in the eviction order.
541    ///
542    /// Stale lifetimes are removed lazily. Once the order reaches twice the capacity it is
543    /// compacted in place, retaining only the current generation of every tracked hash.
544    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    /// Evicts a pending hash to make room for one announced by `announcer` and returns whether
555    /// one was found.
556    ///
557    /// Prefer the oldest exclusively announced pending hash of the peer tracking the most
558    /// hashes, so a group of peers flooding the fetcher gives up its own hashes when possible. The
559    /// announcer gives way when it tracks as many hashes as the busiest peer. If the
560    /// chosen peer has no exclusively announced pending hash, the oldest pending hash overall
561    /// is evicted. Preferring exclusive hashes protects other peers' work when possible, while
562    /// the fallback prevents co-announcements from making the entire cache unevictable.
563    ///
564    /// `queued` are the hashes of the current announcement that are not in the announcer's queue
565    /// yet; `queued_start` marks entries already evicted or moved into the peer queue.
566    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                // Preserve shared entries skipped by the cursor in this announcement.
620                self.enqueue(key, hash, QueuePosition::Back);
621            }
622        }
623        self.evict_oldest_pending()
624    }
625
626    /// Evicts the oldest pending hash announced only by this peer, if one exists.
627    ///
628    /// Scan each selected peer's queue once per announcement. A bounded prefix can hide
629    /// exclusive victims behind shared hashes; caching avoids rescanning that prefix for every
630    /// eviction. Newly tracked hashes from the current announcement are handled by its cursor.
631    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    /// Evicts the oldest pending hash, including shared hashes, as a last resort.
667    ///
668    /// Shared hashes cannot be exempt from eviction: co-announcements could pin every slot.
669    /// Returns `false` if only fetching hashes were found within the search budget.
670    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    /// Drains the peer's queue into a request, in announcement order, until the request holds
693    /// `max_hashes` hashes or the expected response size reaches the configured soft limit. A
694    /// transaction that on its own exceeds the size limit is requested alone.
695    ///
696    /// The returned hashes are marked as being fetched by the peer.
697    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            // skip hashes that were delivered in the meantime, or that this peer is no longer a
707            // candidate for
708            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            // hashes that are being fetched elsewhere are queued again if that fetch fails
718            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    /// Reverts [`Self::pack_request`] for a request that could not be sent: the hashes are pending
745    /// again and are queued at the front of the peer's queue in their original order.
746    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    /// Queues a pending hash for the given candidates, the ones that don't have it queued, and
760    /// records those peers in `requeued`.
761    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    /// Stops tracking the hash.
776    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    /// Processes a resolved request and returns the corresponding event.
789    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                    // the peer only sent transactions we didn't ask for
849                    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    /// Settles the requested hashes of a resolved request: delivered hashes are dropped and
864    /// undelivered hashes are rescheduled for their remaining candidates.
865    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        // Position right after the last delivered hash. Drop the responder as a candidate for
874        // missing hashes before it. Keep it for the undelivered tail only if at least half of
875        // the requested hashes were delivered, treating that tail as potentially truncated.
876        // This is a retry heuristic: transaction counts do not establish whether the response
877        // reached a byte limit. Requiring half the request limits repeated low-progress retries.
878        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        // iterate in reverse so that queueing at the front of a queue preserves the request order
891        for (idx, hash) in requested.iter().enumerate().rev() {
892            if delivered.contains(hash) {
893                self.remove_hash(hash);
894                continue
895            }
896
897            // the hash was received elsewhere in the meantime, or was even announced and assigned
898            // to another peer again
899            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            // A short adaptive session timeout can expire before the first response from the
907            // only source arrives. Allow one retry; prefer another source whenever available.
908            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                // disconnected peers are pruned as well
912                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                // Rotating fallback candidates must not allow unlimited retries for one hash.
922                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        // retried hashes go to the most recent announcers first
937        for key in requeued {
938            self.mark_ready(key, QueuePosition::Front);
939        }
940    }
941}
942
943#[cfg(test)]
944impl<N: NetworkPrimitives> TransactionFetcher<N> {
945    /// Returns the connected peers that are candidates for the hash, in announcement order.
946    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    /// Returns the peer the hash is currently fetched from, if any.
961    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    /// Returns the hashes queued for the peer, oldest first. May include hashes that are not
967    /// tracked anymore or are being fetched elsewhere.
968    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    /// Returns the number of peers the fetcher tracks.
977    fn num_peers(&self) -> usize {
978        self.peers.len()
979    }
980
981    /// Panics if the internal bookkeeping is inconsistent.
982    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        // a hash flagged as queued is in the peer's queue, and a pending hash is queued for at
1025        // least one connected candidate, otherwise it could starve
1026        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    /// Advances all inflight requests and yields the next resolved request as an event.
1137    ///
1138    /// Never terminates, returns [`Poll::Pending`] while no request is inflight.
1139    fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
1140        // `FuturesUnordered` yields `None` while empty but keeps working once new requests are
1141        // pushed, so this is mapped to pending
1142        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/// Represents possible events from fetching transactions.
1150#[derive(Debug)]
1151pub enum FetchEvent<T = PooledTransaction> {
1152    /// Triggered when transactions are successfully fetched.
1153    TransactionsFetched {
1154        /// The ID of the peer from which transactions were fetched.
1155        peer_id: PeerId,
1156        /// The negotiated protocol of the session that served the response.
1157        version: EthVersion,
1158        /// The client version of the session that served the response.
1159        client_version: Arc<str>,
1160        /// The transactions that were fetched, if available.
1161        transactions: PooledTransactions<T>,
1162        /// Whether the peer should be penalized for sending unsolicited transactions or for
1163        /// misbehavior.
1164        report_peer: bool,
1165    },
1166    /// Triggered when there is an error in fetching transactions.
1167    FetchError {
1168        /// The ID of the peer from which an attempt to fetch transactions resulted in an error.
1169        peer_id: PeerId,
1170        /// The specific error that occurred while fetching.
1171        error: RequestError,
1172    },
1173    /// An empty response was received.
1174    EmptyResponse {
1175        /// The ID of the sender.
1176        peer_id: PeerId,
1177    },
1178}
1179
1180/// Compact identifier of a peer within the fetcher.
1181#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)]
1182struct PeerKey(u32);
1183
1184/// Cached eviction work for one announcement. Other peers can only lose candidates; newly
1185/// tracked hashes from this announcer are handled by the announcement's scratch cursor.
1186#[derive(Debug)]
1187struct EvictionState {
1188    peers: BinaryHeap<(usize, PeerKey)>,
1189    exclusive: HashMap<PeerKey, VecDeque<(TxHash, u64)>>,
1190}
1191
1192/// Position for newly queued hashes or ready peers.
1193#[derive(Debug, Clone, Copy)]
1194enum QueuePosition {
1195    Front,
1196    Back,
1197}
1198
1199/// Batches consecutive removals from the same candidate peer during an announcement.
1200#[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    /// Synchronizes tracked counts before capacity eviction or the end of an announcement.
1218    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/// State of a tracked hash.
1226#[derive(Debug)]
1227struct TxEntry {
1228    /// Peers that announced the hash and may be asked for it, oldest first.
1229    candidates: SmallVec<[Candidate; MAX_COUNT_CANDIDATE_PEERS_PER_HASH]>,
1230    /// Identity of this tracking lifetime, also stored in lazy queue entries.
1231    generation: u64,
1232    /// The peer and request currently fetching the transaction, if any.
1233    fetching_by: Option<(PeerKey, u64)>,
1234    /// Sent requests for this entry, bounded independently of candidate replacement.
1235    attempts: usize,
1236}
1237
1238impl TxEntry {
1239    /// A new entry for a hash that is queued for the announcing peer.
1240    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    /// Returns the candidates that don't have the hash queued, in the order they announced it.
1251    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/// A peer that announced a hash.
1265#[derive(Debug, Clone, Copy)]
1266struct Candidate {
1267    peer: PeerKey,
1268    /// The size the peer announced for the transaction, 0 if unknown, with
1269    /// [`Self::QUEUED`] set while the hash is in the peer's queue.
1270    size_and_queued: u32,
1271}
1272
1273impl Candidate {
1274    /// Marks the hash as queued for the peer. Announced sizes are capped well below this bit.
1275    const QUEUED: u32 = 1 << 31;
1276
1277    /// A candidate that has the hash queued.
1278    const fn queued(peer: PeerKey, size: u32) -> Self {
1279        Self { peer, size_and_queued: size | Self::QUEUED }
1280    }
1281
1282    /// A candidate that doesn't have the hash queued.
1283    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    /// Returns the size to account for when packing the hash into a request to this peer.
1304    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/// Fetch state of a peer.
1315#[derive(Debug)]
1316struct PeerState {
1317    peer_id: PeerId,
1318    /// Hashes the peer announced, oldest first. Entries are removed lazily, so the queue may
1319    /// contain hashes that are not tracked anymore or that the peer is no longer a candidate for.
1320    queue: VecDeque<(TxHash, u64)>,
1321    /// Number of tracked hashes that list this peer as a candidate.
1322    tracked: usize,
1323    /// Number of inflight requests to this peer.
1324    inflight: u8,
1325    /// Whether the peer is queued in the ready list.
1326    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    /// Queues the hash at the back or the front of the queue.
1335    ///
1336    /// Queues are bounded: once a queue holds twice `max_tracked` hashes, entries the peer is no
1337    /// longer a candidate for and duplicates are removed. A duplicate can occur when a hash is
1338    /// tracked again after it was delivered or evicted, while its old entry still lingers in the
1339    /// queue. Generations prevent those stale entries from selecting a newer tracking lifetime.
1340    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/// An inflight `GetPooledTransactions` request.
1364#[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    /// The requested hashes, in request order.
1372    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/// A resolved `GetPooledTransactions` request.
1396#[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
1407/// Returns the announced transaction size used for request packing, capped at the response soft
1408/// limit, or 0 if the announcement carried no size.
1409fn 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
1415/// Filters a response down to the transactions that were requested, dropping duplicates.
1416///
1417/// Records the delivered hashes in `delivered` and returns the number of unsolicited transactions
1418/// that were dropped.
1419fn 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    /// Counts how often a task is woken.
1474    #[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    /// A fetcher with mock peer sessions.
1506    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        /// Takes the next request the peer's session received.
1577        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        // no size metadata, so the count limit is what bounds the request
1681        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        // the peer is busy until the request resolves
1690        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        // the empty response dropped the peer as candidate for the first hash
1719        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            // A bad response keeps other hashes pending but drops the requested ones.
1744            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        // peer_a doesn't know the sizes, peer_b announces a huge first transaction
1764        rig.announce_unsized(peer_a, &hashes);
1765        rig.announce_with_sizes(peer_b, [(hashes[0], 1024 * KIB), (hashes[1], 100)]);
1766
1767        // peer_a's request is packed with the size estimate
1768        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        // peer_b's request honors the size it announced
1775        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        // absurd sizes must not break the size accounting
1788        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        // peer_b's queue entry was consumed without a request
1822        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        // peer_b skipped the hashes while they were inflight to peer_a
1844        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        // the hashes are pending again for peer_b only and queued in request order
1854        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        // no other peer announced the hashes, so they are dropped
1880        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        // the first half is delivered, the rest looks like truncation
1896        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        // Delivering fewer than half of the requested hashes drops peer_a for the rest,
1922        // limiting repeated low-progress retries regardless of transaction sizes.
1923        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        // the first hash was skipped and has no other candidate, the last one is retried
1944        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        // treated like a failed request, peer_b gets to try
1990        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        // one inflight and the pending hash arrive over broadcast
2029        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        // the response doesn't include the hash received over broadcast, which must not be
2035        // rescheduled, but includes the other one
2036        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        // dropping the session drops the pending request, which resolves it with an error
2106        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        // the limit is per peer, another peer can still announce the dropped hashes
2144        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        // delivering frees up the budget of the peer
2149        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        // the oldest pending hash makes room for a new one
2170        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        // hashes that are being fetched are not evicted, the announcement is dropped instead
2176        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        // the second hash is delivered, the third is pending again. Both peers then track one
2184        // hash and the announcing peer gives way, so its own older hash makes room.
2185        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        // Keep the eager sources and rotate the oldest fallback sources out for later ones.
2210        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        // the first announcer fetches the hash, retries go to the most recent announcers first
2223        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        // the hash is given up on once all candidates failed
2230        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        // The shared cursor entry must not hide newly queued exclusive victims.
2317        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        // the honest peer's hashes are the oldest when the flood fills the fetcher
2360        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        // the flooders gave up their oldest hashes
2372        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        // The global scan sees only fetching hashes. The last new hash is dropped after
2415        // the shared scratch entries have been moved into the announcer's real queue.
2416        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        // the same hashes are given up on and announced again over and over, every round leaves
2521        // stale entries in the eviction order behind
2522        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        // a new hash still evicts the oldest one rather than getting dropped
2533        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        // remembered candidates have no queue entry, they are forgotten nonetheless
2550        for peer_id in remembered {
2551            rig.disconnect(*peer_id);
2552        }
2553        assert_eq!(rig.fetcher.candidate_peers(&hash), queued);
2554
2555        // which frees their slots for later announcers
2556        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        // and a hash whose last connected candidate leaves is dropped
2562        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        // all peers that had the hash queued disconnect before fetching it
2583        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        // the deferred peer is served once a request slot frees up
2610        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        // Each peer can use the full request limit regardless of other inflight requests.
2632        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        // A smaller request limit is respected too.
2640        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        // All earlier peers keep their requests open. The newly ready peer still gets a
2658        // full request instead of being throttled by their inflight hashes.
2659        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        // occupy the only slot of the session channel
2674        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        // drain the blocker, the peer is retried on the next dispatch
2691        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        // the session task is gone, but the manager didn't process the disconnect yet
2710        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        // keep peer_a busy
2730        rig.announce(peer_a, &hashes(0..1));
2731        rig.dispatch();
2732
2733        // hashes announced by the busy peer are delivered by others over and over
2734        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        // arrives over broadcast, gets announced again and is assigned to another peer
2756        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        // the original request resolves without the hash, which must not touch the new fetch
2762        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        // more peers than a hash has candidate slots, so announcers get remembered and rejected
2777        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        // requests taken from the sessions that weren't answered yet
2792        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                        // drops the response senders of the peer's outstanding requests
2874                        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        // settle everything that is still inflight
2888        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        // the caller has to poll again after dispatching, only then the request is polled and
2933        // registers the waker that its response wakes
2934        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        // 24 of the 64 announced hashes were evicted, spread over the peers so that every
2959        // peer lost its oldest hashes and kept its newest
2960        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        // every peer still has hashes to request, the evicted ones are skipped
2971        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        // the last and the first hash are delivered in reverse order, the two in between were
2997        // skipped by peer_a
2998        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        // only the last hash of the first request is delivered, the rest was skipped and has no
3025        // other candidate
3026        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        // the peer reconnects and announces the hashes again while the old request is pending
3051        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        // the old request delivers one hash and skips the other, which is retried with the new
3059        // session
3060        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        // peer_b is busy while both announce the hashes, so they stay in its queue
3079        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}