Skip to main content

reth_network/fetch/
mod.rs

1//! Fetch data from the network.
2
3mod client;
4
5pub use client::FetchClient;
6
7use crate::{message::BlockRequest, session::BlockRangeInfo};
8use alloy_primitives::B256;
9use futures::StreamExt;
10use reth_eth_wire::{
11    snap::SnapProtocolMessage, BlockAccessLists, Capabilities, EthNetworkPrimitives, EthVersion,
12    GetBlockAccessLists, GetBlockBodies, GetBlockHeaders, GetReceipts, NetworkPrimitives,
13};
14use reth_network_api::test_utils::PeersHandle;
15use reth_network_p2p::{
16    block_access_lists::client::BalRequirement,
17    error::{EthResponseValidator, PeerRequestResult, RequestError, RequestResult},
18    headers::client::HeadersRequest,
19    priority::Priority,
20    receipts::client::ReceiptsResponse,
21    snap::client::SnapResponse,
22};
23use reth_network_peers::PeerId;
24use reth_network_types::ReputationChangeKind;
25use std::{
26    collections::{HashMap, VecDeque},
27    ops::RangeInclusive,
28    sync::{
29        atomic::{AtomicU64, AtomicUsize, Ordering},
30        Arc,
31    },
32    task::{Context, Poll},
33};
34use tokio::sync::{mpsc, mpsc::UnboundedSender, oneshot};
35use tokio_stream::wrappers::UnboundedReceiverStream;
36
37type InflightHeadersRequest<H> = Request<HeadersRequest, PeerRequestResult<Vec<H>>>;
38type InflightBodiesRequest<B> = Request<(), PeerRequestResult<Vec<B>>>;
39type InflightReceiptsRequest<R> = Request<(), PeerRequestResult<ReceiptsResponse<R>>>;
40type InflightBlockAccessListsRequest = Request<(), PeerRequestResult<BlockAccessLists>>;
41type InflightSnapRequest = Request<(), PeerRequestResult<SnapResponse>>;
42
43/// Manages data fetching operations.
44///
45/// This type is hooked into the staged sync pipeline and delegates download request to available
46/// peers and sends the response once ready.
47///
48/// This type maintains a list of connected peers that are available for requests.
49#[derive(Debug)]
50pub struct StateFetcher<N: NetworkPrimitives = EthNetworkPrimitives> {
51    /// Currently active [`GetBlockHeaders`] requests
52    inflight_headers_requests: HashMap<PeerId, InflightHeadersRequest<N::BlockHeader>>,
53    /// Currently active [`GetBlockBodies`] requests
54    inflight_bodies_requests: HashMap<PeerId, InflightBodiesRequest<N::BlockBody>>,
55    /// Currently active [`GetBlockAccessLists`] requests
56    inflight_bals_requests: HashMap<PeerId, InflightBlockAccessListsRequest>,
57    /// Currently active `GetReceipts` requests
58    inflight_receipts_requests: HashMap<PeerId, InflightReceiptsRequest<N::Receipt>>,
59    /// Currently active `snap/2` requests
60    inflight_snap_requests: HashMap<PeerId, InflightSnapRequest>,
61    /// The list of _available_ peers for requests.
62    peers: HashMap<PeerId, Peer>,
63    /// The handle to the peers manager
64    peers_handle: PeersHandle,
65    /// Number of active peer sessions the node's currently handling.
66    num_active_peers: Arc<AtomicUsize>,
67    /// Requests queued for processing
68    queued_requests: VecDeque<DownloadRequest<N>>,
69    /// Receiver for new incoming download requests
70    download_requests_rx: UnboundedReceiverStream<DownloadRequest<N>>,
71    /// Sender for download requests, used to detach a [`FetchClient`]
72    download_requests_tx: UnboundedSender<DownloadRequest<N>>,
73}
74
75// === impl StateSyncer ===
76
77impl<N: NetworkPrimitives> StateFetcher<N> {
78    pub(crate) fn new(peers_handle: PeersHandle, num_active_peers: Arc<AtomicUsize>) -> Self {
79        let (download_requests_tx, download_requests_rx) = mpsc::unbounded_channel();
80        Self {
81            inflight_headers_requests: Default::default(),
82            inflight_bodies_requests: Default::default(),
83            inflight_bals_requests: Default::default(),
84            inflight_receipts_requests: Default::default(),
85            inflight_snap_requests: Default::default(),
86            peers: Default::default(),
87            peers_handle,
88            num_active_peers,
89            queued_requests: Default::default(),
90            download_requests_rx: UnboundedReceiverStream::new(download_requests_rx),
91            download_requests_tx,
92        }
93    }
94
95    /// Invoked when connected to a new peer.
96    pub(crate) fn new_active_peer(&mut self, peer: NewPeerInfo) {
97        let NewPeerInfo {
98            peer_id,
99            best_hash,
100            best_number,
101            capabilities,
102            timeout,
103            range_info,
104            supports_snap,
105        } = peer;
106        self.peers.insert(
107            peer_id,
108            Peer {
109                state: PeerState::Idle,
110                best_hash,
111                best_number,
112                capabilities,
113                timeout,
114                last_response_likely_bad: false,
115                range_info,
116                supports_snap,
117            },
118        );
119    }
120
121    /// Removes the peer from the peer list, after which it is no longer available for future
122    /// requests.
123    ///
124    /// Invoked when an active session was closed.
125    ///
126    /// This cancels also inflight request and sends an error to the receiver.
127    pub(crate) fn on_session_closed(&mut self, peer: &PeerId) {
128        self.peers.remove(peer);
129        if let Some(req) = self.inflight_headers_requests.remove(peer) {
130            let _ = req.response.send(Err(RequestError::ConnectionDropped));
131        }
132        if let Some(req) = self.inflight_bodies_requests.remove(peer) {
133            let _ = req.response.send(Err(RequestError::ConnectionDropped));
134        }
135        if let Some(req) = self.inflight_bals_requests.remove(peer) {
136            let _ = req.response.send(Err(RequestError::ConnectionDropped));
137        }
138        if let Some(req) = self.inflight_receipts_requests.remove(peer) {
139            let _ = req.response.send(Err(RequestError::ConnectionDropped));
140        }
141        if let Some(req) = self.inflight_snap_requests.remove(peer) {
142            let _ = req.response.send(Err(RequestError::ConnectionDropped));
143        }
144    }
145
146    /// Updates the block information for the peer.
147    ///
148    /// Returns `true` if this a newer block
149    pub(crate) fn update_peer_block(&mut self, peer_id: &PeerId, hash: B256, number: u64) -> bool {
150        if let Some(peer) = self.peers.get_mut(peer_id) &&
151            number > peer.best_number
152        {
153            peer.best_hash = hash;
154            peer.best_number = number;
155            return true
156        }
157        false
158    }
159
160    /// Invoked when an active session is about to be disconnected.
161    pub(crate) fn on_pending_disconnect(&mut self, peer_id: &PeerId) {
162        if let Some(peer) = self.peers.get_mut(peer_id) {
163            peer.state = PeerState::Closing;
164        }
165    }
166
167    /// Returns the _next_ idle peer that's ready to accept a request,
168    /// prioritizing those with the lowest timeout/latency and those that recently responded with
169    /// adequate data. Additionally, if full blocks are required this prioritizes peers that have
170    /// full history available
171    fn next_best_peer(&self, requirement: BestPeerRequirements) -> Option<PeerId> {
172        // filter out peers that aren't idle or don't meet the requirement
173        let mut idle = self
174            .peers
175            .iter()
176            .filter(|(_, peer)| peer.state.is_idle() && peer.satisfies(&requirement));
177
178        let mut best_peer = idle.next()?;
179
180        for maybe_better in idle {
181            // replace best peer if our current best peer sent us a bad response last time
182            if best_peer.1.last_response_likely_bad && !maybe_better.1.last_response_likely_bad {
183                best_peer = maybe_better;
184                continue
185            }
186
187            // replace best peer if this peer meets the requirements better
188            if maybe_better.1.is_better(best_peer.1, &requirement) {
189                best_peer = maybe_better;
190                continue
191            }
192
193            // replace best peer if this peer has better rtt and both have same range quality
194            if maybe_better.1.timeout() < best_peer.1.timeout() &&
195                !maybe_better.1.last_response_likely_bad
196            {
197                best_peer = maybe_better;
198            }
199        }
200
201        Some(*best_peer.0)
202    }
203
204    /// Returns whether any connected peer can serve BAL requests.
205    fn has_eth71_peer(&self) -> bool {
206        self.peers.values().any(|peer| {
207            !matches!(peer.state, PeerState::Closing) &&
208                peer.capabilities.supports_eth_at_least(&EthVersion::Eth71)
209        })
210    }
211
212    /// Returns the next action to return
213    fn poll_action(&mut self) -> PollAction {
214        // we only check and not pop here since we don't know yet whether a peer is available.
215        if self.queued_requests.is_empty() {
216            return PollAction::NoRequests
217        }
218
219        let request = self.queued_requests.pop_front().expect("not empty");
220        let Some(peer_id) = self.next_best_peer(request.best_peer_requirements()) else {
221            // Optional BAL/snap requests can lose their capable peer while queued; complete them
222            // instead of waiting for future peer churn.
223            if self.should_fail_fast(&request) {
224                request.send_err_response(RequestError::UnsupportedCapability);
225            } else {
226                // no peer matches this request's requirements; requeue at the back so other
227                // queued requests get a chance on the next poll instead of head-of-line blocking.
228                self.queued_requests.push_back(request);
229            }
230            return PollAction::NoPeersAvailable
231        };
232
233        let request = self.prepare_block_request(peer_id, request);
234
235        PollAction::Ready(FetchAction::BlockRequest { peer_id, request })
236    }
237
238    /// Advance the state the syncer
239    pub(crate) fn poll(&mut self, cx: &mut Context<'_>) -> Poll<FetchAction> {
240        // drain buffered actions first
241        loop {
242            let no_peers_available = match self.poll_action() {
243                PollAction::Ready(action) => return Poll::Ready(action),
244                PollAction::NoRequests => false,
245                PollAction::NoPeersAvailable => true,
246            };
247
248            loop {
249                // poll incoming requests
250                match self.download_requests_rx.poll_next_unpin(cx) {
251                    Poll::Ready(Some(request)) => {
252                        // Optional BAL/snap requests should not wait for future peer churn if no
253                        // connected peer can serve them right now.
254                        if self.should_fail_fast(&request) {
255                            request.send_err_response(RequestError::UnsupportedCapability);
256                            continue
257                        }
258
259                        match request.get_priority() {
260                            Priority::High => {
261                                // find first normal request and queue before it; add this request
262                                // to the back of the high-priority queue
263                                let pos = self
264                                    .queued_requests
265                                    .iter()
266                                    .position(|req| req.is_normal_priority())
267                                    .unwrap_or(0);
268                                self.queued_requests.insert(pos, request);
269                            }
270                            Priority::Normal => {
271                                self.queued_requests.push_back(request);
272                            }
273                        }
274                    }
275                    Poll::Ready(None) => {
276                        unreachable!("channel can't close")
277                    }
278                    Poll::Pending => break,
279                }
280            }
281
282            if self.queued_requests.is_empty() || no_peers_available {
283                return Poll::Pending
284            }
285        }
286    }
287
288    /// Returns whether any connected peer negotiated `snap/2`.
289    fn has_snap_peer(&self) -> bool {
290        self.peers
291            .values()
292            .any(|peer| !matches!(peer.state, PeerState::Closing) && peer.supports_snap)
293    }
294
295    /// Returns `true` if `request` cannot be served by any currently connected peer and should
296    /// fail immediately instead of waiting for future peer churn.
297    fn should_fail_fast(&self, request: &DownloadRequest<N>) -> bool {
298        (request.is_optional_bal() && !self.has_eth71_peer()) ||
299            (request.is_snap() && !self.has_snap_peer())
300    }
301
302    /// Handles a new request to a peer.
303    ///
304    /// Caution: this assumes the peer exists and is idle
305    fn prepare_block_request(&mut self, peer_id: PeerId, req: DownloadRequest<N>) -> BlockRequest {
306        // update the peer's state
307        if let Some(peer) = self.peers.get_mut(&peer_id) {
308            peer.state = req.peer_state();
309        }
310
311        self.prepare_inflight_block_request(peer_id, req)
312    }
313
314    /// Tracks an inflight request and converts it into a peer request.
315    fn prepare_inflight_block_request(
316        &mut self,
317        peer_id: PeerId,
318        req: DownloadRequest<N>,
319    ) -> BlockRequest {
320        match req {
321            DownloadRequest::GetBlockHeaders { request, response, .. } => {
322                let inflight = Request { request: request.clone(), response };
323                self.inflight_headers_requests.insert(peer_id, inflight);
324                let HeadersRequest { start, limit, direction } = request;
325                BlockRequest::GetBlockHeaders(GetBlockHeaders {
326                    start_block: start,
327                    limit,
328                    skip: 0,
329                    direction,
330                })
331            }
332            DownloadRequest::GetBlockBodies { request, response, .. } => {
333                let inflight = Request { request: (), response };
334                self.inflight_bodies_requests.insert(peer_id, inflight);
335                BlockRequest::GetBlockBodies(GetBlockBodies(request))
336            }
337            DownloadRequest::GetBlockAccessLists { request, response, .. } => {
338                let inflight = Request { request: (), response };
339                self.inflight_bals_requests.insert(peer_id, inflight);
340                BlockRequest::GetBlockAccessLists(GetBlockAccessLists(request))
341            }
342            DownloadRequest::GetReceipts { request, response, .. } => {
343                let inflight = Request { request: (), response };
344                self.inflight_receipts_requests.insert(peer_id, inflight);
345                BlockRequest::GetReceipts(GetReceipts(request))
346            }
347            DownloadRequest::GetSnap { request, response, .. } => {
348                let inflight = Request { request: (), response };
349                self.inflight_snap_requests.insert(peer_id, inflight);
350                BlockRequest::GetSnap(Box::new(request))
351            }
352        }
353    }
354
355    /// Returns a queued followup request the peer can serve.
356    ///
357    /// This is an immediate scheduling shortcut after a successful response. It skips queued
358    /// requests whose hard requirements do not match this peer, leaving them for the regular peer
359    /// selection path.
360    ///
361    /// Caution: this expects that the peer is _not_ closed.
362    fn followup_request(&mut self, peer_id: PeerId) -> Option<BlockResponseOutcome> {
363        let peer = self.peers.get_mut(&peer_id)?;
364        let req_idx = self.queued_requests.iter().position(|req| {
365            // Find the first queued request this peer can serve.
366            peer.satisfies(&req.best_peer_requirements())
367        })?;
368        let req = self.queued_requests.remove(req_idx).expect("valid request index");
369
370        peer.state = req.peer_state();
371        let req = self.prepare_inflight_block_request(peer_id, req);
372        Some(BlockResponseOutcome::Request(peer_id, req))
373    }
374
375    /// Called on a `GetBlockHeaders` response from a peer.
376    ///
377    /// This delegates the response and returns a [`BlockResponseOutcome`] to either queue in a
378    /// direct followup request or get the peer reported if the response was a
379    /// [`EthResponseValidator::reputation_change_err`]
380    pub(crate) fn on_block_headers_response(
381        &mut self,
382        peer_id: PeerId,
383        res: RequestResult<Vec<N::BlockHeader>>,
384    ) -> Option<BlockResponseOutcome> {
385        let is_error = res.is_err();
386        let maybe_reputation_change = res.reputation_change_err();
387
388        let resp = self.inflight_headers_requests.remove(&peer_id);
389
390        let is_likely_bad_response =
391            resp.as_ref().is_some_and(|r| res.is_likely_bad_headers_response(&r.request));
392
393        if let Some(resp) = resp {
394            // delegate the response
395            let _ = resp.response.send(res.map(|h| (peer_id, h).into()));
396        }
397
398        if let Some(peer) = self.peers.get_mut(&peer_id) {
399            // update the peer's response state
400            peer.last_response_likely_bad = is_likely_bad_response;
401
402            // If the peer is still ready to accept new requests, we try to send a followup
403            // request immediately.
404            if peer.state.on_request_finished() && !is_error && !is_likely_bad_response {
405                return self.followup_request(peer_id)
406            }
407        }
408
409        // if the response was an `Err` worth reporting the peer for then we return a `BadResponse`
410        // outcome
411        maybe_reputation_change
412            .map(|reputation_change| BlockResponseOutcome::BadResponse(peer_id, reputation_change))
413    }
414
415    /// Called on a `GetBlockBodies` response from a peer
416    pub(crate) fn on_block_bodies_response(
417        &mut self,
418        peer_id: PeerId,
419        res: RequestResult<Vec<N::BlockBody>>,
420    ) -> Option<BlockResponseOutcome> {
421        let is_likely_bad_response = res.as_ref().map_or(true, |bodies| bodies.is_empty());
422
423        if let Some(resp) = self.inflight_bodies_requests.remove(&peer_id) {
424            let _ = resp.response.send(res.map(|b| (peer_id, b).into()));
425        }
426        if let Some(peer) = self.peers.get_mut(&peer_id) {
427            // update the peer's response state
428            peer.last_response_likely_bad = is_likely_bad_response;
429
430            if peer.state.on_request_finished() && !is_likely_bad_response {
431                return self.followup_request(peer_id)
432            }
433        }
434        None
435    }
436
437    /// Called on a `GetBlockAccessLists` response from a peer
438    pub(crate) fn on_block_access_lists_response(
439        &mut self,
440        peer_id: PeerId,
441        res: RequestResult<BlockAccessLists>,
442    ) -> Option<BlockResponseOutcome> {
443        let is_likely_bad_response = res.is_err();
444
445        if let Some(resp) = self.inflight_bals_requests.remove(&peer_id) {
446            let _ = resp.response.send(res.map(|b| (peer_id, b).into()));
447        }
448        if let Some(peer) = self.peers.get_mut(&peer_id) {
449            peer.last_response_likely_bad = is_likely_bad_response;
450
451            if peer.state.on_request_finished() && !is_likely_bad_response {
452                return self.followup_request(peer_id)
453            }
454        }
455        None
456    }
457
458    /// Called on a `GetReceipts` response from a peer.
459    ///
460    /// All receipt variants (legacy with bloom, eth/69, eth/70) are expected to be normalized
461    /// to [`ReceiptsResponse`] by the caller before invoking this method.
462    pub(crate) fn on_receipts_response(
463        &mut self,
464        peer_id: PeerId,
465        res: RequestResult<ReceiptsResponse<N::Receipt>>,
466    ) -> Option<BlockResponseOutcome> {
467        let is_likely_bad_response = res.as_ref().map_or(true, |resp| resp.receipts.is_empty());
468
469        if let Some(resp) = self.inflight_receipts_requests.remove(&peer_id) {
470            let _ = resp.response.send(res.map(|r| (peer_id, r).into()));
471        }
472        if let Some(peer) = self.peers.get_mut(&peer_id) {
473            peer.last_response_likely_bad = is_likely_bad_response;
474
475            if peer.state.on_request_finished() && !is_likely_bad_response {
476                return self.followup_request(peer_id)
477            }
478        }
479        None
480    }
481
482    /// Called on a `snap/2` response from a peer.
483    pub(crate) fn on_snap_response(
484        &mut self,
485        peer_id: PeerId,
486        res: RequestResult<SnapResponse>,
487    ) -> Option<BlockResponseOutcome> {
488        let is_likely_bad_response = res.is_err();
489
490        if let Some(resp) = self.inflight_snap_requests.remove(&peer_id) {
491            let _ = resp.response.send(res.map(|r| (peer_id, r).into()));
492        }
493        if let Some(peer) = self.peers.get_mut(&peer_id) {
494            peer.last_response_likely_bad = is_likely_bad_response;
495
496            if peer.state.on_request_finished() && !is_likely_bad_response {
497                return self.followup_request(peer_id)
498            }
499        }
500        None
501    }
502
503    /// Returns a new [`FetchClient`] that can send requests to this type.
504    pub(crate) fn client(&self) -> FetchClient<N> {
505        FetchClient {
506            request_tx: self.download_requests_tx.clone(),
507            peers_handle: self.peers_handle.clone(),
508            num_active_peers: Arc::clone(&self.num_active_peers),
509        }
510    }
511}
512
513/// The outcome of [`StateFetcher::poll_action`]
514enum PollAction {
515    Ready(FetchAction),
516    NoRequests,
517    NoPeersAvailable,
518}
519
520/// Everything [`StateFetcher::new_active_peer`] needs to register a newly connected peer.
521#[derive(Debug)]
522pub(crate) struct NewPeerInfo {
523    /// The remote peer's identifier.
524    pub(crate) peer_id: PeerId,
525    /// Best known hash that the peer has.
526    pub(crate) best_hash: B256,
527    /// The best block number of the peer.
528    pub(crate) best_number: u64,
529    /// Capabilities announced by the peer.
530    pub(crate) capabilities: Arc<Capabilities>,
531    /// The current timeout value to use for the peer.
532    pub(crate) timeout: Arc<AtomicU64>,
533    /// The range info for the peer.
534    pub(crate) range_info: Option<BlockRangeInfo>,
535    /// Whether the connection negotiated `snap/2` and can serve [`DownloadRequest::GetSnap`].
536    pub(crate) supports_snap: bool,
537}
538
539/// Represents a connected peer
540#[derive(Debug)]
541struct Peer {
542    /// The state this peer currently resides in.
543    state: PeerState,
544    /// Best known hash that the peer has
545    best_hash: B256,
546    /// Tracks the best number of the peer.
547    best_number: u64,
548    /// Capabilities announced by the peer.
549    #[allow(dead_code)]
550    capabilities: Arc<Capabilities>,
551    /// Tracks the current timeout value we use for the peer.
552    timeout: Arc<AtomicU64>,
553    /// Tracks whether the peer has recently responded with a likely bad response.
554    ///
555    /// This is used to de-rank the peer if there are other peers available.
556    /// This exists because empty responses may not be penalized (e.g. when blocks near the tip are
557    /// downloaded), but we still want to avoid requesting from the same peer again if it has the
558    /// lowest timeout.
559    last_response_likely_bad: bool,
560    /// Tracks the range info for the peer.
561    range_info: Option<BlockRangeInfo>,
562    /// Whether the connection negotiated `snap/2` and can serve [`DownloadRequest::GetSnap`].
563    supports_snap: bool,
564}
565
566impl Peer {
567    fn timeout(&self) -> u64 {
568        self.timeout.load(Ordering::Relaxed)
569    }
570
571    /// Returns the earliest block number available from the peer.
572    fn earliest(&self) -> u64 {
573        self.range_info.as_ref().map_or(0, |info| info.earliest())
574    }
575
576    /// Returns true if the peer has the full history available.
577    fn has_full_history(&self) -> bool {
578        self.earliest() == 0
579    }
580
581    fn range(&self) -> Option<RangeInclusive<u64>> {
582        self.range_info.as_ref().map(|info| info.range())
583    }
584
585    /// Returns whether this peer can serve requests with the given hard requirements.
586    fn satisfies(&self, requirement: &BestPeerRequirements) -> bool {
587        match requirement {
588            BestPeerRequirements::EthVersion(ver) => self.capabilities.supports_eth_at_least(ver),
589            BestPeerRequirements::SupportsSnap => self.supports_snap,
590            BestPeerRequirements::None |
591            BestPeerRequirements::FullBlock |
592            BestPeerRequirements::FullBlockRange(_) => true,
593        }
594    }
595
596    /// Returns true if this peer has a better range than the other peer for serving the requested
597    /// range.
598    ///
599    /// A peer has a "better range" if:
600    /// 1. It can fully cover the requested range while the other cannot
601    /// 2. None can fully cover the range, but this peer has lower start value
602    /// 3. If a peer doesn't announce a range we assume it has full history, but check the other's
603    ///    range and treat that as better if it can cover the range
604    fn has_better_range(&self, other: &Self, range: &RangeInclusive<u64>) -> bool {
605        let self_range = self.range();
606        let other_range = other.range();
607
608        match (self_range, other_range) {
609            (Some(self_r), Some(other_r)) => {
610                // Check if each peer can fully cover the requested range
611                let self_covers = self_r.contains(range.start()) && self_r.contains(range.end());
612                let other_covers = other_r.contains(range.start()) && other_r.contains(range.end());
613
614                #[expect(clippy::match_same_arms)]
615                match (self_covers, other_covers) {
616                    (true, false) => true,  // Only self covers the range
617                    (false, true) => false, // Only other covers the range
618                    (true, true) => false,  // Both cover
619                    (false, false) => {
620                        // neither covers - prefer if peer has lower (better) start range
621                        self_r.start() < other_r.start()
622                    }
623                }
624            }
625            (Some(self_r), None) => {
626                // Self has range info, other doesn't (treated as full history with unknown latest)
627                // Self is better only if it covers the range
628                self_r.contains(range.start()) && self_r.contains(range.end())
629            }
630            (None, Some(other_r)) => {
631                // Self has no range info (full history), other has range info
632                // Self is better only if other doesn't cover the range
633                !(other_r.contains(range.start()) && other_r.contains(range.end()))
634            }
635            (None, None) => false, // Neither has range info - no one is better
636        }
637    }
638
639    /// Returns true if this peer is better than the other peer based on the given requirements.
640    fn is_better(&self, other: &Self, requirement: &BestPeerRequirements) -> bool {
641        match requirement {
642            BestPeerRequirements::FullBlockRange(range) => self.has_better_range(other, range),
643            BestPeerRequirements::FullBlock => self.has_full_history() && !other.has_full_history(),
644            // Version/capability-based filtering happens in `next_best_peer`, so by the time we
645            // get here both peers already satisfy the requirement.
646            BestPeerRequirements::None |
647            BestPeerRequirements::EthVersion(_) |
648            BestPeerRequirements::SupportsSnap => false,
649        }
650    }
651}
652
653/// Tracks the state of an individual peer
654#[derive(Debug)]
655enum PeerState {
656    /// Peer is currently not handling requests and is available.
657    Idle,
658    /// Peer is handling a `GetBlockHeaders` request.
659    GetBlockHeaders,
660    /// Peer is handling a `GetBlockBodies` request.
661    GetBlockBodies,
662    /// Peer is handling a `GetBlockAccessLists` request.
663    GetBlockAccessLists,
664    /// Peer is handling a `GetReceipts` request.
665    GetReceipts,
666    /// Peer is handling a `snap/2` request.
667    GetSnap,
668    /// Peer session is about to close
669    Closing,
670}
671
672// === impl PeerState ===
673
674impl PeerState {
675    /// Returns true if the peer is currently idle.
676    const fn is_idle(&self) -> bool {
677        matches!(self, Self::Idle)
678    }
679
680    /// Resets the state on a received response.
681    ///
682    /// If the state was already marked as `Closing` do nothing.
683    ///
684    /// Returns `true` if the peer is ready for another request.
685    const fn on_request_finished(&mut self) -> bool {
686        if !matches!(self, Self::Closing) {
687            *self = Self::Idle;
688            return true
689        }
690        false
691    }
692}
693
694/// A request that waits for a response from the network, so it can send it back through the
695/// response channel.
696#[derive(Debug)]
697struct Request<Req, Resp> {
698    /// The issued request object
699    // TODO: this can be attached to the response in error case
700    request: Req,
701    response: oneshot::Sender<Resp>,
702}
703
704/// Requests that can be sent to the Syncer from a [`FetchClient`]
705#[derive(Debug)]
706#[expect(clippy::enum_variant_names)]
707pub(crate) enum DownloadRequest<N: NetworkPrimitives> {
708    /// Download the requested headers and send response through channel
709    GetBlockHeaders {
710        request: HeadersRequest,
711        response: oneshot::Sender<PeerRequestResult<Vec<N::BlockHeader>>>,
712        priority: Priority,
713    },
714    /// Download the requested bodies and send response through channel
715    GetBlockBodies {
716        request: Vec<B256>,
717        response: oneshot::Sender<PeerRequestResult<Vec<N::BlockBody>>>,
718        priority: Priority,
719        range_hint: Option<RangeInclusive<u64>>,
720    },
721    /// Download the requested access lists and send response through channel
722    GetBlockAccessLists {
723        request: Vec<B256>,
724        response: oneshot::Sender<PeerRequestResult<BlockAccessLists>>,
725        priority: Priority,
726        requirement: BalRequirement,
727    },
728    /// Download receipts for the given block hashes and send response through channel
729    GetReceipts {
730        request: Vec<B256>,
731        response: oneshot::Sender<PeerRequestResult<ReceiptsResponse<N::Receipt>>>,
732        priority: Priority,
733    },
734    /// Send a `snap/2` request and send response through channel
735    GetSnap {
736        request: SnapProtocolMessage,
737        response: oneshot::Sender<PeerRequestResult<SnapResponse>>,
738        priority: Priority,
739    },
740}
741
742// === impl DownloadRequest ===
743
744impl<N: NetworkPrimitives> DownloadRequest<N> {
745    /// Returns the corresponding state for a peer that handles the request.
746    const fn peer_state(&self) -> PeerState {
747        match self {
748            Self::GetBlockHeaders { .. } => PeerState::GetBlockHeaders,
749            Self::GetBlockBodies { .. } => PeerState::GetBlockBodies,
750            Self::GetBlockAccessLists { .. } => PeerState::GetBlockAccessLists,
751            Self::GetReceipts { .. } => PeerState::GetReceipts,
752            Self::GetSnap { .. } => PeerState::GetSnap,
753        }
754    }
755
756    /// Returns the requested priority of this request
757    const fn get_priority(&self) -> &Priority {
758        match self {
759            Self::GetBlockHeaders { priority, .. } |
760            Self::GetBlockBodies { priority, .. } |
761            Self::GetBlockAccessLists { priority, .. } |
762            Self::GetReceipts { priority, .. } |
763            Self::GetSnap { priority, .. } => priority,
764        }
765    }
766
767    /// Returns `true` if this request is normal priority.
768    const fn is_normal_priority(&self) -> bool {
769        self.get_priority().is_normal()
770    }
771
772    /// Returns `true` if this is an optional BAL request.
773    const fn is_optional_bal(&self) -> bool {
774        matches!(self, Self::GetBlockAccessLists { requirement: BalRequirement::Optional, .. })
775    }
776
777    /// Returns `true` if this is a `snap/2` request.
778    const fn is_snap(&self) -> bool {
779        matches!(self, Self::GetSnap { .. })
780    }
781
782    /// Sends an error response to the waiting caller.
783    fn send_err_response(self, err: RequestError) {
784        let _ = match self {
785            Self::GetBlockHeaders { response, .. } => response.send(Err(err)).ok(),
786            Self::GetBlockBodies { response, .. } => response.send(Err(err)).ok(),
787            Self::GetBlockAccessLists { response, .. } => response.send(Err(err)).ok(),
788            Self::GetReceipts { response, .. } => response.send(Err(err)).ok(),
789            Self::GetSnap { response, .. } => response.send(Err(err)).ok(),
790        };
791    }
792
793    /// Returns the best peer requirements for this request.
794    fn best_peer_requirements(&self) -> BestPeerRequirements {
795        match self {
796            Self::GetBlockHeaders { .. } => BestPeerRequirements::None,
797            Self::GetBlockAccessLists { .. } => BestPeerRequirements::EthVersion(EthVersion::Eth71),
798            Self::GetBlockBodies { range_hint, .. } => {
799                if let Some(range) = range_hint {
800                    BestPeerRequirements::FullBlockRange(range.clone())
801                } else {
802                    BestPeerRequirements::FullBlock
803                }
804            }
805            Self::GetReceipts { .. } => BestPeerRequirements::FullBlock,
806            Self::GetSnap { .. } => BestPeerRequirements::SupportsSnap,
807        }
808    }
809}
810
811/// An action the syncer can emit.
812pub(crate) enum FetchAction {
813    /// Dispatch an eth request to the given peer.
814    BlockRequest {
815        /// The targeted recipient for the request
816        peer_id: PeerId,
817        /// The request to send
818        request: BlockRequest,
819    },
820}
821
822/// Outcome of a processed response.
823///
824/// Returned after processing a response.
825#[derive(Debug, PartialEq, Eq)]
826pub(crate) enum BlockResponseOutcome {
827    /// Continue with another request to the peer.
828    Request(PeerId, BlockRequest),
829    /// How to handle a bad response and the reputation change to apply, if any.
830    BadResponse(PeerId, ReputationChangeKind),
831}
832
833/// Additional requirements for how to rank peers during selection.
834enum BestPeerRequirements {
835    /// No additional requirements
836    None,
837    /// Peer must have this block range available.
838    FullBlockRange(RangeInclusive<u64>),
839    /// Peer must have full range.
840    FullBlock,
841    /// Peer must support at least this eth protocol version.
842    EthVersion(EthVersion),
843    /// Peer must have negotiated `snap/2`.
844    SupportsSnap,
845}
846
847#[cfg(test)]
848mod tests {
849    use super::*;
850    use crate::{peers::PeersManager, PeersConfig};
851    use alloy_consensus::Header;
852    use alloy_primitives::B512;
853    use reth_eth_wire::Capability;
854    use reth_eth_wire_types::snap::{AccountRangeMessage, GetAccountRangeMessage};
855    use std::future::poll_fn;
856
857    #[tokio::test(flavor = "multi_thread")]
858    async fn test_poll_fetcher() {
859        let manager = PeersManager::new(PeersConfig::default());
860        let mut fetcher =
861            StateFetcher::<EthNetworkPrimitives>::new(manager.handle(), Default::default());
862
863        poll_fn(move |cx| {
864            assert!(fetcher.poll(cx).is_pending());
865            let (tx, _rx) = oneshot::channel();
866            fetcher.queued_requests.push_back(DownloadRequest::GetBlockBodies {
867                request: vec![],
868                response: tx,
869                priority: Priority::default(),
870                range_hint: None,
871            });
872            assert!(fetcher.poll(cx).is_pending());
873
874            Poll::Ready(())
875        })
876        .await;
877    }
878
879    #[tokio::test]
880    async fn test_peer_rotation() {
881        let manager = PeersManager::new(PeersConfig::default());
882        let mut fetcher =
883            StateFetcher::<EthNetworkPrimitives>::new(manager.handle(), Default::default());
884        // Add a few random peers
885        let peer1 = B512::random();
886        let peer2 = B512::random();
887        let capabilities = Arc::new(Capabilities::from(vec![]));
888        fetcher.new_active_peer(NewPeerInfo {
889            peer_id: peer1,
890            best_hash: B256::random(),
891            best_number: 1,
892            capabilities: Arc::clone(&capabilities),
893            timeout: Arc::new(AtomicU64::new(1)),
894            range_info: None,
895            supports_snap: false,
896        });
897        fetcher.new_active_peer(NewPeerInfo {
898            peer_id: peer2,
899            best_hash: B256::random(),
900            best_number: 2,
901            capabilities: Arc::clone(&capabilities),
902            timeout: Arc::new(AtomicU64::new(1)),
903            range_info: None,
904            supports_snap: false,
905        });
906
907        let first_peer = fetcher.next_best_peer(BestPeerRequirements::None).unwrap();
908        assert!(first_peer == peer1 || first_peer == peer2);
909        // Pending disconnect for first_peer
910        fetcher.on_pending_disconnect(&first_peer);
911        // first_peer now isn't idle, so we should get other peer
912        let second_peer = fetcher.next_best_peer(BestPeerRequirements::None).unwrap();
913        assert!(first_peer == peer1 || first_peer == peer2);
914        assert_ne!(first_peer, second_peer);
915        // without idle peers, returns None
916        fetcher.on_pending_disconnect(&second_peer);
917        assert_eq!(fetcher.next_best_peer(BestPeerRequirements::None), None);
918    }
919
920    #[tokio::test]
921    async fn test_peer_prioritization() {
922        let manager = PeersManager::new(PeersConfig::default());
923        let mut fetcher =
924            StateFetcher::<EthNetworkPrimitives>::new(manager.handle(), Default::default());
925        // Add a few random peers
926        let peer1 = B512::random();
927        let peer2 = B512::random();
928        let peer3 = B512::random();
929
930        let peer2_timeout = Arc::new(AtomicU64::new(300));
931
932        let capabilities = Arc::new(Capabilities::from(vec![]));
933        fetcher.new_active_peer(NewPeerInfo {
934            peer_id: peer1,
935            best_hash: B256::random(),
936            best_number: 1,
937            capabilities: Arc::clone(&capabilities),
938            timeout: Arc::new(AtomicU64::new(30)),
939            range_info: None,
940            supports_snap: false,
941        });
942        fetcher.new_active_peer(NewPeerInfo {
943            peer_id: peer2,
944            best_hash: B256::random(),
945            best_number: 2,
946            capabilities: Arc::clone(&capabilities),
947            timeout: Arc::clone(&peer2_timeout),
948            range_info: None,
949            supports_snap: false,
950        });
951        fetcher.new_active_peer(NewPeerInfo {
952            peer_id: peer3,
953            best_hash: B256::random(),
954            best_number: 3,
955            capabilities: Arc::clone(&capabilities),
956            timeout: Arc::new(AtomicU64::new(50)),
957            range_info: None,
958            supports_snap: false,
959        });
960
961        // Must always get peer1 (lowest timeout)
962        assert_eq!(fetcher.next_best_peer(BestPeerRequirements::None), Some(peer1));
963        assert_eq!(fetcher.next_best_peer(BestPeerRequirements::None), Some(peer1));
964        // peer2's timeout changes below peer1's
965        peer2_timeout.store(10, Ordering::Relaxed);
966        // Then we get peer 2 always (now lowest)
967        assert_eq!(fetcher.next_best_peer(BestPeerRequirements::None), Some(peer2));
968        assert_eq!(fetcher.next_best_peer(BestPeerRequirements::None), Some(peer2));
969    }
970
971    #[tokio::test]
972    async fn test_on_block_headers_response() {
973        let manager = PeersManager::new(PeersConfig::default());
974        let mut fetcher =
975            StateFetcher::<EthNetworkPrimitives>::new(manager.handle(), Default::default());
976        let peer_id = B512::random();
977
978        assert_eq!(fetcher.on_block_headers_response(peer_id, Ok(vec![Header::default()])), None);
979
980        assert_eq!(
981            fetcher.on_block_headers_response(peer_id, Err(RequestError::Timeout)),
982            Some(BlockResponseOutcome::BadResponse(peer_id, ReputationChangeKind::Timeout))
983        );
984        assert_eq!(
985            fetcher.on_block_headers_response(peer_id, Err(RequestError::BadResponse)),
986            None
987        );
988        assert_eq!(
989            fetcher.on_block_headers_response(peer_id, Err(RequestError::ChannelClosed)),
990            None
991        );
992        assert_eq!(
993            fetcher.on_block_headers_response(peer_id, Err(RequestError::ConnectionDropped)),
994            None
995        );
996        assert_eq!(
997            fetcher.on_block_headers_response(peer_id, Err(RequestError::UnsupportedCapability)),
998            None
999        );
1000    }
1001
1002    #[tokio::test]
1003    async fn test_header_response_outcome() {
1004        let manager = PeersManager::new(PeersConfig::default());
1005        let mut fetcher =
1006            StateFetcher::<EthNetworkPrimitives>::new(manager.handle(), Default::default());
1007        let peer_id = B512::random();
1008
1009        let request_pair = || {
1010            let (tx, _rx) = oneshot::channel();
1011            let req = Request {
1012                request: HeadersRequest {
1013                    start: 0u64.into(),
1014                    limit: 1,
1015                    direction: Default::default(),
1016                },
1017                response: tx,
1018            };
1019            let header = Header { number: 0, ..Default::default() };
1020            (req, header)
1021        };
1022
1023        fetcher.new_active_peer(NewPeerInfo {
1024            peer_id,
1025            best_hash: Default::default(),
1026            best_number: Default::default(),
1027            capabilities: Arc::new(Capabilities::from(vec![])),
1028            timeout: Default::default(),
1029            range_info: None,
1030            supports_snap: false,
1031        });
1032
1033        let (req, header) = request_pair();
1034        fetcher.inflight_headers_requests.insert(peer_id, req);
1035
1036        let outcome = fetcher.on_block_headers_response(peer_id, Ok(vec![header]));
1037        assert!(outcome.is_none());
1038        assert!(fetcher.peers[&peer_id].state.is_idle());
1039
1040        let outcome =
1041            fetcher.on_block_headers_response(peer_id, Err(RequestError::Timeout)).unwrap();
1042
1043        assert!(EthResponseValidator::reputation_change_err(&Err::<Vec<Header>, _>(
1044            RequestError::Timeout
1045        ))
1046        .is_some());
1047
1048        match outcome {
1049            BlockResponseOutcome::BadResponse(peer, _) => {
1050                assert_eq!(peer, peer_id)
1051            }
1052            BlockResponseOutcome::Request(_, _) => {
1053                unreachable!()
1054            }
1055        };
1056
1057        assert!(fetcher.peers[&peer_id].state.is_idle());
1058    }
1059
1060    #[tokio::test]
1061    async fn test_partial_header_response_triggers_followup() {
1062        let (mut fetcher, peer_id) = fetcher_with_peer();
1063
1064        let (followup_tx, _followup_rx) = oneshot::channel();
1065        fetcher.queued_requests.push_back(DownloadRequest::GetBlockHeaders {
1066            request: HeadersRequest::falling(1u64.into(), 1),
1067            response: followup_tx,
1068            priority: Priority::High,
1069        });
1070
1071        let (response_tx, mut response_rx) = oneshot::channel();
1072        fetcher.inflight_headers_requests.insert(
1073            peer_id,
1074            Request { request: HeadersRequest::falling(2u64.into(), 2), response: response_tx },
1075        );
1076        fetcher.peers.get_mut(&peer_id).unwrap().state = PeerState::GetBlockHeaders;
1077
1078        let header = Header { number: 2, ..Default::default() };
1079        let outcome = fetcher.on_block_headers_response(peer_id, Ok(vec![header]));
1080
1081        let Some(BlockResponseOutcome::Request(
1082            dispatched_peer,
1083            BlockRequest::GetBlockHeaders(request),
1084        )) = outcome
1085        else {
1086            panic!("expected a header followup request")
1087        };
1088        assert_eq!(dispatched_peer, peer_id);
1089        assert_eq!(request.start_block, 1u64.into());
1090        assert_eq!(request.limit, 1);
1091        assert!(!fetcher.peers[&peer_id].last_response_likely_bad);
1092        assert!(response_rx.try_recv().is_ok());
1093    }
1094
1095    #[test]
1096    fn test_peer_is_better_none_requirement() {
1097        let peer1 = Peer {
1098            state: PeerState::Idle,
1099            best_hash: B256::random(),
1100            best_number: 100,
1101            capabilities: Arc::new(Capabilities::new(vec![])),
1102            timeout: Arc::new(AtomicU64::new(10)),
1103            last_response_likely_bad: false,
1104            range_info: Some(BlockRangeInfo::new(0, 100, B256::random())),
1105            supports_snap: false,
1106        };
1107
1108        let peer2 = Peer {
1109            state: PeerState::Idle,
1110            best_hash: B256::random(),
1111            best_number: 50,
1112            capabilities: Arc::new(Capabilities::new(vec![])),
1113            timeout: Arc::new(AtomicU64::new(20)),
1114            last_response_likely_bad: false,
1115            range_info: None,
1116            supports_snap: false,
1117        };
1118
1119        // With None requirement, is_better should always return false
1120        assert!(!peer1.is_better(&peer2, &BestPeerRequirements::None));
1121        assert!(!peer2.is_better(&peer1, &BestPeerRequirements::None));
1122    }
1123
1124    #[test]
1125    fn test_peer_is_better_full_block_requirement() {
1126        // Peer with full history (earliest = 0)
1127        let peer_full = Peer {
1128            state: PeerState::Idle,
1129            best_hash: B256::random(),
1130            best_number: 100,
1131            capabilities: Arc::new(Capabilities::new(vec![])),
1132            timeout: Arc::new(AtomicU64::new(10)),
1133            last_response_likely_bad: false,
1134            range_info: Some(BlockRangeInfo::new(0, 100, B256::random())),
1135            supports_snap: false,
1136        };
1137
1138        // Peer without full history (earliest = 50)
1139        let peer_partial = Peer {
1140            state: PeerState::Idle,
1141            best_hash: B256::random(),
1142            best_number: 100,
1143            capabilities: Arc::new(Capabilities::new(vec![])),
1144            timeout: Arc::new(AtomicU64::new(10)),
1145            last_response_likely_bad: false,
1146            range_info: Some(BlockRangeInfo::new(50, 100, B256::random())),
1147            supports_snap: false,
1148        };
1149
1150        // Peer without range info (treated as full history)
1151        let peer_no_range = Peer {
1152            state: PeerState::Idle,
1153            best_hash: B256::random(),
1154            best_number: 100,
1155            capabilities: Arc::new(Capabilities::new(vec![])),
1156            timeout: Arc::new(AtomicU64::new(10)),
1157            last_response_likely_bad: false,
1158            range_info: None,
1159            supports_snap: false,
1160        };
1161
1162        // Peer with full history is better than peer without
1163        assert!(peer_full.is_better(&peer_partial, &BestPeerRequirements::FullBlock));
1164        assert!(!peer_partial.is_better(&peer_full, &BestPeerRequirements::FullBlock));
1165
1166        // Peer without range info (full history) is better than partial
1167        assert!(peer_no_range.is_better(&peer_partial, &BestPeerRequirements::FullBlock));
1168        assert!(!peer_partial.is_better(&peer_no_range, &BestPeerRequirements::FullBlock));
1169
1170        // Both have full history - no improvement
1171        assert!(!peer_full.is_better(&peer_no_range, &BestPeerRequirements::FullBlock));
1172        assert!(!peer_no_range.is_better(&peer_full, &BestPeerRequirements::FullBlock));
1173    }
1174
1175    #[test]
1176    fn test_peer_is_better_full_block_range_requirement() {
1177        let range = RangeInclusive::new(40, 60);
1178
1179        // Peer that covers the requested range
1180        let peer_covers = Peer {
1181            state: PeerState::Idle,
1182            best_hash: B256::random(),
1183            best_number: 100,
1184            capabilities: Arc::new(Capabilities::new(vec![])),
1185            timeout: Arc::new(AtomicU64::new(10)),
1186            last_response_likely_bad: false,
1187            range_info: Some(BlockRangeInfo::new(0, 100, B256::random())),
1188            supports_snap: false,
1189        };
1190
1191        // Peer that doesn't cover the range (earliest too high)
1192        let peer_no_cover = Peer {
1193            state: PeerState::Idle,
1194            best_hash: B256::random(),
1195            best_number: 100,
1196            capabilities: Arc::new(Capabilities::new(vec![])),
1197            timeout: Arc::new(AtomicU64::new(10)),
1198            last_response_likely_bad: false,
1199            range_info: Some(BlockRangeInfo::new(70, 100, B256::random())),
1200            supports_snap: false,
1201        };
1202
1203        // Peer that covers the requested range is better than one that doesn't
1204        assert!(peer_covers
1205            .is_better(&peer_no_cover, &BestPeerRequirements::FullBlockRange(range.clone())));
1206        assert!(
1207            !peer_no_cover.is_better(&peer_covers, &BestPeerRequirements::FullBlockRange(range))
1208        );
1209    }
1210
1211    #[test]
1212    fn test_peer_is_better_both_cover_range() {
1213        let range = RangeInclusive::new(30, 50);
1214
1215        // Peer with full history that covers the range
1216        let peer_full = Peer {
1217            state: PeerState::Idle,
1218            best_hash: B256::random(),
1219            best_number: 100,
1220            capabilities: Arc::new(Capabilities::new(vec![])),
1221            timeout: Arc::new(AtomicU64::new(10)),
1222            last_response_likely_bad: false,
1223            range_info: Some(BlockRangeInfo::new(0, 50, B256::random())),
1224            supports_snap: false,
1225        };
1226
1227        // Peer without full history that also covers the range
1228        let peer_partial = Peer {
1229            state: PeerState::Idle,
1230            best_hash: B256::random(),
1231            best_number: 100,
1232            capabilities: Arc::new(Capabilities::new(vec![])),
1233            timeout: Arc::new(AtomicU64::new(10)),
1234            last_response_likely_bad: false,
1235            range_info: Some(BlockRangeInfo::new(30, 50, B256::random())),
1236            supports_snap: false,
1237        };
1238
1239        // When both cover the range, prefer none
1240        assert!(!peer_full
1241            .is_better(&peer_partial, &BestPeerRequirements::FullBlockRange(range.clone())));
1242        assert!(!peer_partial.is_better(&peer_full, &BestPeerRequirements::FullBlockRange(range)));
1243    }
1244
1245    #[test]
1246    fn test_peer_is_better_lower_start() {
1247        let range = RangeInclusive::new(30, 60);
1248
1249        // Peer with full history that covers the range
1250        let peer_full = Peer {
1251            state: PeerState::Idle,
1252            best_hash: B256::random(),
1253            best_number: 100,
1254            capabilities: Arc::new(Capabilities::new(vec![])),
1255            timeout: Arc::new(AtomicU64::new(10)),
1256            last_response_likely_bad: false,
1257            range_info: Some(BlockRangeInfo::new(0, 50, B256::random())),
1258            supports_snap: false,
1259        };
1260
1261        // Peer without full history that also covers the range
1262        let peer_partial = Peer {
1263            state: PeerState::Idle,
1264            best_hash: B256::random(),
1265            best_number: 100,
1266            capabilities: Arc::new(Capabilities::new(vec![])),
1267            timeout: Arc::new(AtomicU64::new(10)),
1268            last_response_likely_bad: false,
1269            range_info: Some(BlockRangeInfo::new(30, 50, B256::random())),
1270            supports_snap: false,
1271        };
1272
1273        // When both cover the range, prefer lower start value
1274        assert!(peer_full
1275            .is_better(&peer_partial, &BestPeerRequirements::FullBlockRange(range.clone())));
1276        assert!(!peer_partial.is_better(&peer_full, &BestPeerRequirements::FullBlockRange(range)));
1277    }
1278
1279    #[test]
1280    fn test_peer_is_better_neither_covers_range() {
1281        let range = RangeInclusive::new(40, 60);
1282
1283        // Peer with full history that doesn't cover the range (latest too low)
1284        let peer_full = Peer {
1285            state: PeerState::Idle,
1286            best_hash: B256::random(),
1287            best_number: 30,
1288            capabilities: Arc::new(Capabilities::new(vec![])),
1289            timeout: Arc::new(AtomicU64::new(10)),
1290            last_response_likely_bad: false,
1291            range_info: Some(BlockRangeInfo::new(0, 30, B256::random())),
1292            supports_snap: false,
1293        };
1294
1295        // Peer without full history that also doesn't cover the range
1296        let peer_partial = Peer {
1297            state: PeerState::Idle,
1298            best_hash: B256::random(),
1299            best_number: 30,
1300            capabilities: Arc::new(Capabilities::new(vec![])),
1301            timeout: Arc::new(AtomicU64::new(10)),
1302            last_response_likely_bad: false,
1303            range_info: Some(BlockRangeInfo::new(10, 30, B256::random())),
1304            supports_snap: false,
1305        };
1306
1307        // When neither covers the range, prefer full history
1308        assert!(peer_full
1309            .is_better(&peer_partial, &BestPeerRequirements::FullBlockRange(range.clone())));
1310        assert!(!peer_partial.is_better(&peer_full, &BestPeerRequirements::FullBlockRange(range)));
1311    }
1312
1313    #[test]
1314    fn test_peer_is_better_no_range_info() {
1315        let range = RangeInclusive::new(40, 60);
1316
1317        // Peer with range info
1318        let peer_with_range = Peer {
1319            state: PeerState::Idle,
1320            best_hash: B256::random(),
1321            best_number: 100,
1322            capabilities: Arc::new(Capabilities::new(vec![])),
1323            timeout: Arc::new(AtomicU64::new(10)),
1324            last_response_likely_bad: false,
1325            range_info: Some(BlockRangeInfo::new(30, 100, B256::random())),
1326            supports_snap: false,
1327        };
1328
1329        // Peer without range info
1330        let peer_no_range = Peer {
1331            state: PeerState::Idle,
1332            best_hash: B256::random(),
1333            best_number: 100,
1334            capabilities: Arc::new(Capabilities::new(vec![])),
1335            timeout: Arc::new(AtomicU64::new(10)),
1336            last_response_likely_bad: false,
1337            range_info: None,
1338            supports_snap: false,
1339        };
1340
1341        // Peer without range info is not better (we prefer peers with known ranges)
1342        assert!(!peer_no_range
1343            .is_better(&peer_with_range, &BestPeerRequirements::FullBlockRange(range.clone())));
1344
1345        // Peer with range info is better than peer without
1346        assert!(
1347            peer_with_range.is_better(&peer_no_range, &BestPeerRequirements::FullBlockRange(range))
1348        );
1349    }
1350
1351    #[test]
1352    fn test_peer_is_better_one_peer_no_range_covers() {
1353        let range = RangeInclusive::new(40, 60);
1354
1355        // Peer with range info that covers the requested range
1356        let peer_with_range_covers = Peer {
1357            state: PeerState::Idle,
1358            best_hash: B256::random(),
1359            best_number: 100,
1360            capabilities: Arc::new(Capabilities::new(vec![])),
1361            timeout: Arc::new(AtomicU64::new(10)),
1362            last_response_likely_bad: false,
1363            range_info: Some(BlockRangeInfo::new(30, 100, B256::random())),
1364            supports_snap: false,
1365        };
1366
1367        // Peer without range info (treated as full history with unknown latest)
1368        let peer_no_range = Peer {
1369            state: PeerState::Idle,
1370            best_hash: B256::random(),
1371            best_number: 100,
1372            capabilities: Arc::new(Capabilities::new(vec![])),
1373            timeout: Arc::new(AtomicU64::new(10)),
1374            last_response_likely_bad: false,
1375            range_info: None,
1376            supports_snap: false,
1377        };
1378
1379        // Peer with range that covers is better than peer without range info
1380        assert!(peer_with_range_covers
1381            .is_better(&peer_no_range, &BestPeerRequirements::FullBlockRange(range.clone())));
1382
1383        // Peer without range info is not better when other covers
1384        assert!(!peer_no_range
1385            .is_better(&peer_with_range_covers, &BestPeerRequirements::FullBlockRange(range)));
1386    }
1387
1388    #[test]
1389    fn test_peer_is_better_one_peer_no_range_doesnt_cover() {
1390        let range = RangeInclusive::new(40, 60);
1391
1392        // Peer with range info that does NOT cover the requested range (too high)
1393        let peer_with_range_no_cover = Peer {
1394            state: PeerState::Idle,
1395            best_hash: B256::random(),
1396            best_number: 100,
1397            capabilities: Arc::new(Capabilities::new(vec![])),
1398            timeout: Arc::new(AtomicU64::new(10)),
1399            last_response_likely_bad: false,
1400            range_info: Some(BlockRangeInfo::new(70, 100, B256::random())),
1401            supports_snap: false,
1402        };
1403
1404        // Peer without range info (treated as full history)
1405        let peer_no_range = Peer {
1406            state: PeerState::Idle,
1407            best_hash: B256::random(),
1408            best_number: 100,
1409            capabilities: Arc::new(Capabilities::new(vec![])),
1410            timeout: Arc::new(AtomicU64::new(10)),
1411            last_response_likely_bad: false,
1412            range_info: None,
1413            supports_snap: false,
1414        };
1415
1416        // Peer with range that doesn't cover is not better
1417        assert!(!peer_with_range_no_cover
1418            .is_better(&peer_no_range, &BestPeerRequirements::FullBlockRange(range.clone())));
1419
1420        // Peer without range info (full history) is better when other doesn't cover
1421        assert!(peer_no_range
1422            .is_better(&peer_with_range_no_cover, &BestPeerRequirements::FullBlockRange(range)));
1423    }
1424
1425    #[test]
1426    fn test_peer_is_better_edge_cases() {
1427        // Test exact range boundaries
1428        let range = RangeInclusive::new(50, 100);
1429
1430        // Peer that exactly covers the range
1431        let peer_exact = Peer {
1432            state: PeerState::Idle,
1433            best_hash: B256::random(),
1434            best_number: 100,
1435            capabilities: Arc::new(Capabilities::new(vec![])),
1436            timeout: Arc::new(AtomicU64::new(10)),
1437            last_response_likely_bad: false,
1438            range_info: Some(BlockRangeInfo::new(50, 100, B256::random())),
1439            supports_snap: false,
1440        };
1441
1442        // Peer that's one block short at the start
1443        let peer_short_start = Peer {
1444            state: PeerState::Idle,
1445            best_hash: B256::random(),
1446            best_number: 100,
1447            capabilities: Arc::new(Capabilities::new(vec![])),
1448            timeout: Arc::new(AtomicU64::new(10)),
1449            last_response_likely_bad: false,
1450            range_info: Some(BlockRangeInfo::new(51, 100, B256::random())),
1451            supports_snap: false,
1452        };
1453
1454        // Peer that's one block short at the end
1455        let peer_short_end = Peer {
1456            state: PeerState::Idle,
1457            best_hash: B256::random(),
1458            best_number: 100,
1459            capabilities: Arc::new(Capabilities::new(vec![])),
1460            timeout: Arc::new(AtomicU64::new(10)),
1461            last_response_likely_bad: false,
1462            range_info: Some(BlockRangeInfo::new(50, 99, B256::random())),
1463            supports_snap: false,
1464        };
1465
1466        // Exact coverage is better than short coverage
1467        assert!(peer_exact
1468            .is_better(&peer_short_start, &BestPeerRequirements::FullBlockRange(range.clone())));
1469        assert!(peer_exact
1470            .is_better(&peer_short_end, &BestPeerRequirements::FullBlockRange(range.clone())));
1471
1472        // Short coverage is not better than exact coverage
1473        assert!(!peer_short_start
1474            .is_better(&peer_exact, &BestPeerRequirements::FullBlockRange(range.clone())));
1475        assert!(
1476            !peer_short_end.is_better(&peer_exact, &BestPeerRequirements::FullBlockRange(range))
1477        );
1478    }
1479
1480    /// Creates a `StateFetcher` with a single idle peer and returns both.
1481    fn fetcher_with_peer() -> (StateFetcher<EthNetworkPrimitives>, PeerId) {
1482        let manager = PeersManager::new(PeersConfig::default());
1483        let mut fetcher =
1484            StateFetcher::<EthNetworkPrimitives>::new(manager.handle(), Default::default());
1485        let peer_id = B512::random();
1486
1487        fetcher.new_active_peer(NewPeerInfo {
1488            peer_id,
1489            best_hash: Default::default(),
1490            best_number: Default::default(),
1491            capabilities: Arc::new(Capabilities::from(vec![])),
1492            timeout: Default::default(),
1493            range_info: None,
1494            supports_snap: false,
1495        });
1496        (fetcher, peer_id)
1497    }
1498
1499    /// Inserts an inflight receipts request into the fetcher and returns the
1500    /// `oneshot::Receiver` that the final response will be sent through.
1501    fn insert_inflight_receipts(
1502        fetcher: &mut StateFetcher<EthNetworkPrimitives>,
1503        peer_id: PeerId,
1504    ) -> oneshot::Receiver<PeerRequestResult<ReceiptsResponse<reth_ethereum_primitives::Receipt>>>
1505    {
1506        let (tx, rx) = oneshot::channel();
1507        fetcher.inflight_receipts_requests.insert(peer_id, Request { request: (), response: tx });
1508        fetcher.peers.get_mut(&peer_id).unwrap().state = PeerState::GetReceipts;
1509        rx
1510    }
1511
1512    // ---- Receipts: basic dispatch ----
1513
1514    #[tokio::test]
1515    async fn test_poll_dispatches_receipts_to_peer() {
1516        let (mut fetcher, peer_id) = fetcher_with_peer();
1517
1518        poll_fn(move |cx| {
1519            let (tx, _rx) = oneshot::channel();
1520            fetcher.queued_requests.push_back(DownloadRequest::GetReceipts {
1521                request: vec![B256::ZERO],
1522                response: tx,
1523                priority: Priority::default(),
1524            });
1525
1526            let Poll::Ready(FetchAction::BlockRequest { peer_id: dispatched_peer, request }) =
1527                fetcher.poll(cx)
1528            else {
1529                panic!("expected Ready(BlockRequest)");
1530            };
1531            assert_eq!(dispatched_peer, peer_id);
1532            assert!(matches!(request, BlockRequest::GetReceipts(_)));
1533
1534            // Peer should now be in GetReceipts state
1535            assert!(matches!(fetcher.peers[&peer_id].state, PeerState::GetReceipts));
1536            // Inflight request should be tracked
1537            assert!(fetcher.inflight_receipts_requests.contains_key(&peer_id));
1538
1539            Poll::Ready(())
1540        })
1541        .await;
1542    }
1543
1544    // ---- Receipts: response handling ----
1545
1546    #[tokio::test]
1547    async fn test_receipts_complete_response_resolves_and_idles_peer() {
1548        let (mut fetcher, peer_id) = fetcher_with_peer();
1549
1550        let rx = insert_inflight_receipts(&mut fetcher, peer_id);
1551
1552        let resp = ReceiptsResponse::new(vec![vec![]]);
1553        let outcome = fetcher.on_receipts_response(peer_id, Ok(resp));
1554
1555        // No queued requests, so no followup
1556        assert!(outcome.is_none());
1557        // Peer back to idle
1558        assert!(fetcher.peers[&peer_id].state.is_idle());
1559        // Inflight cleaned up
1560        assert!(!fetcher.inflight_receipts_requests.contains_key(&peer_id));
1561
1562        // Caller receives the response
1563        let result = rx.await.unwrap().unwrap();
1564        assert_eq!(result.1.receipts.len(), 1);
1565    }
1566
1567    #[tokio::test]
1568    async fn test_receipts_empty_response_marks_peer_bad() {
1569        let (mut fetcher, peer_id) = fetcher_with_peer();
1570        let _rx = insert_inflight_receipts(&mut fetcher, peer_id);
1571
1572        let resp = ReceiptsResponse::new(vec![]);
1573        let _ = fetcher.on_receipts_response(peer_id, Ok(resp));
1574
1575        assert!(fetcher.peers[&peer_id].last_response_likely_bad);
1576    }
1577
1578    #[tokio::test]
1579    async fn test_receipts_error_forwards_and_marks_peer_bad() {
1580        let (mut fetcher, peer_id) = fetcher_with_peer();
1581        let rx = insert_inflight_receipts(&mut fetcher, peer_id);
1582
1583        let _ = fetcher.on_receipts_response(peer_id, Err(RequestError::Timeout));
1584
1585        assert!(fetcher.peers[&peer_id].last_response_likely_bad);
1586        // Error is forwarded to the caller
1587        let result = rx.await.unwrap();
1588        assert_eq!(result.unwrap_err(), RequestError::Timeout);
1589    }
1590
1591    #[tokio::test]
1592    async fn test_session_closed_cancels_inflight_receipts() {
1593        let (mut fetcher, peer_id) = fetcher_with_peer();
1594        let rx = insert_inflight_receipts(&mut fetcher, peer_id);
1595
1596        fetcher.on_session_closed(&peer_id);
1597
1598        assert!(!fetcher.peers.contains_key(&peer_id));
1599        assert!(!fetcher.inflight_receipts_requests.contains_key(&peer_id));
1600
1601        let result = rx.await.unwrap();
1602        assert_eq!(result.unwrap_err(), RequestError::ConnectionDropped);
1603    }
1604
1605    #[tokio::test]
1606    async fn test_receipts_response_triggers_followup() {
1607        let (mut fetcher, peer_id) = fetcher_with_peer();
1608
1609        // Queue a bodies request as a followup candidate
1610        let (followup_tx, _followup_rx) = oneshot::channel();
1611        fetcher.queued_requests.push_back(DownloadRequest::GetBlockBodies {
1612            request: vec![B256::random()],
1613            response: followup_tx,
1614            priority: Priority::default(),
1615            range_hint: None,
1616        });
1617
1618        let _rx = insert_inflight_receipts(&mut fetcher, peer_id);
1619
1620        let resp = ReceiptsResponse::new(vec![vec![]]);
1621        let outcome = fetcher.on_receipts_response(peer_id, Ok(resp));
1622
1623        assert!(matches!(outcome, Some(BlockResponseOutcome::Request(pid, _)) if pid == peer_id));
1624    }
1625
1626    #[tokio::test]
1627    async fn test_followup_skips_request_peer_cannot_serve() {
1628        let (mut fetcher, peer_id) = fetcher_with_peer();
1629
1630        let peer_71 = B512::random();
1631        let caps_71 = Arc::new(Capabilities::from(vec![Capability::new("eth".into(), 71)]));
1632        fetcher.new_active_peer(NewPeerInfo {
1633            peer_id: peer_71,
1634            best_hash: B256::random(),
1635            best_number: 100,
1636            capabilities: caps_71,
1637            timeout: Arc::new(AtomicU64::new(10)),
1638            range_info: None,
1639            supports_snap: false,
1640        });
1641        fetcher.peers.get_mut(&peer_71).expect("peer exists").state = PeerState::GetBlockHeaders;
1642
1643        let (followup_tx, _followup_rx) = oneshot::channel();
1644        fetcher.queued_requests.push_back(DownloadRequest::GetBlockAccessLists {
1645            request: vec![B256::random()],
1646            response: followup_tx,
1647            priority: Priority::Normal,
1648            requirement: BalRequirement::Optional,
1649        });
1650
1651        let _rx = insert_inflight_receipts(&mut fetcher, peer_id);
1652
1653        let resp = ReceiptsResponse::new(vec![vec![]]);
1654        assert!(fetcher.on_receipts_response(peer_id, Ok(resp)).is_none());
1655        assert!(fetcher.peers[&peer_id].state.is_idle());
1656        assert!(!fetcher.inflight_bals_requests.contains_key(&peer_id));
1657        assert!(matches!(
1658            fetcher.queued_requests.front(),
1659            Some(DownloadRequest::GetBlockAccessLists {
1660                requirement: BalRequirement::Optional,
1661                ..
1662            })
1663        ));
1664    }
1665
1666    #[tokio::test]
1667    async fn test_followup_uses_first_satisfiable_request() {
1668        let (mut fetcher, peer_id) = fetcher_with_peer();
1669
1670        let peer_71 = B512::random();
1671        let caps_71 = Arc::new(Capabilities::from(vec![Capability::new("eth".into(), 71)]));
1672        fetcher.new_active_peer(NewPeerInfo {
1673            peer_id: peer_71,
1674            best_hash: B256::random(),
1675            best_number: 100,
1676            capabilities: caps_71,
1677            timeout: Arc::new(AtomicU64::new(10)),
1678            range_info: None,
1679            supports_snap: false,
1680        });
1681        fetcher.peers.get_mut(&peer_71).expect("peer exists").state = PeerState::GetBlockHeaders;
1682
1683        let (bal_tx, _bal_rx) = oneshot::channel();
1684        fetcher.queued_requests.push_back(DownloadRequest::GetBlockAccessLists {
1685            request: vec![B256::random()],
1686            response: bal_tx,
1687            priority: Priority::Normal,
1688            requirement: BalRequirement::Optional,
1689        });
1690
1691        let (bodies_tx, _bodies_rx) = oneshot::channel();
1692        fetcher.queued_requests.push_back(DownloadRequest::GetBlockBodies {
1693            request: vec![B256::random()],
1694            response: bodies_tx,
1695            priority: Priority::Normal,
1696            range_hint: None,
1697        });
1698
1699        let _rx = insert_inflight_receipts(&mut fetcher, peer_id);
1700
1701        let resp = ReceiptsResponse::new(vec![vec![]]);
1702        let outcome = fetcher.on_receipts_response(peer_id, Ok(resp));
1703
1704        assert!(matches!(
1705            outcome,
1706            Some(BlockResponseOutcome::Request(pid, BlockRequest::GetBlockBodies(_))) if pid == peer_id
1707        ));
1708        assert!(fetcher.inflight_bodies_requests.contains_key(&peer_id));
1709        assert!(matches!(fetcher.peers[&peer_id].state, PeerState::GetBlockBodies));
1710        assert_eq!(fetcher.queued_requests.len(), 1);
1711        assert!(matches!(
1712            fetcher.queued_requests.front(),
1713            Some(DownloadRequest::GetBlockAccessLists {
1714                requirement: BalRequirement::Optional,
1715                ..
1716            })
1717        ));
1718    }
1719
1720    #[tokio::test]
1721    async fn test_prepare_block_request_creates_inflight_receipts() {
1722        let (mut fetcher, peer_id) = fetcher_with_peer();
1723        let hashes = vec![B256::with_last_byte(1), B256::with_last_byte(2)];
1724
1725        let (tx, _rx) = oneshot::channel();
1726        let req = DownloadRequest::GetReceipts {
1727            request: hashes.clone(),
1728            response: tx,
1729            priority: Priority::default(),
1730        };
1731
1732        let block_request = fetcher.prepare_block_request(peer_id, req);
1733
1734        // Returns a GetReceipts block request with the same hashes
1735        match block_request {
1736            BlockRequest::GetReceipts(ref get) => {
1737                assert_eq!(get.0, hashes);
1738            }
1739            other => panic!("expected GetReceipts, got {other:?}"),
1740        }
1741
1742        // Peer state transitions to GetReceipts
1743        assert!(matches!(fetcher.peers[&peer_id].state, PeerState::GetReceipts));
1744
1745        // Inflight request is tracked
1746        assert!(fetcher.inflight_receipts_requests.contains_key(&peer_id));
1747    }
1748    #[tokio::test]
1749    async fn test_next_best_peer_eth71_no_support() {
1750        let manager = PeersManager::new(PeersConfig::default());
1751        let mut fetcher =
1752            StateFetcher::<EthNetworkPrimitives>::new(manager.handle(), Default::default());
1753
1754        let peer = B512::random();
1755
1756        // Capabilities WITHOUT eth71
1757        let capabilities = Arc::new(Capabilities::new(vec![]));
1758
1759        fetcher.new_active_peer(NewPeerInfo {
1760            peer_id: peer,
1761            best_hash: B256::random(),
1762            best_number: 100,
1763            capabilities,
1764            timeout: Arc::new(AtomicU64::new(10)),
1765            range_info: None,
1766            supports_snap: false,
1767        });
1768
1769        // Should return None because peer doesn't support eth71
1770        assert_eq!(
1771            fetcher.next_best_peer(BestPeerRequirements::EthVersion(EthVersion::Eth71)),
1772            None
1773        );
1774    }
1775
1776    #[tokio::test]
1777    async fn test_next_best_peer_eth71_supported() {
1778        let manager = PeersManager::new(PeersConfig::default());
1779        let mut fetcher =
1780            StateFetcher::<EthNetworkPrimitives>::new(manager.handle(), Default::default());
1781
1782        let peer = B512::random();
1783
1784        // Build capability list that includes Eth71
1785        let capabilities = Arc::new(Capabilities::from(vec![Capability::new("eth".into(), 71)]));
1786
1787        fetcher.new_active_peer(NewPeerInfo {
1788            peer_id: peer,
1789            best_hash: B256::random(),
1790            best_number: 100,
1791            capabilities,
1792            timeout: Arc::new(AtomicU64::new(10)),
1793            range_info: None,
1794            supports_snap: false,
1795        });
1796
1797        assert_eq!(
1798            fetcher.next_best_peer(BestPeerRequirements::EthVersion(EthVersion::Eth71)),
1799            Some(peer)
1800        );
1801    }
1802
1803    #[tokio::test]
1804    async fn test_next_best_peer_eth71_filters_correctly() {
1805        let manager = PeersManager::new(PeersConfig::default());
1806        let mut fetcher =
1807            StateFetcher::<EthNetworkPrimitives>::new(manager.handle(), Default::default());
1808
1809        let peer_no_71 = B512::random();
1810        let peer_with_71 = B512::random();
1811
1812        // Peer without eth71
1813        let caps_old = Arc::new(Capabilities::new(vec![]));
1814
1815        // Peer with eth71
1816        let caps_71 = Arc::new(Capabilities::from(vec![Capability::new("eth".into(), 71)]));
1817
1818        fetcher.new_active_peer(NewPeerInfo {
1819            peer_id: peer_no_71,
1820            best_hash: B256::random(),
1821            best_number: 100,
1822            capabilities: caps_old,
1823            timeout: Arc::new(AtomicU64::new(5)),
1824            range_info: None,
1825            supports_snap: false,
1826        });
1827
1828        fetcher.new_active_peer(NewPeerInfo {
1829            peer_id: peer_with_71,
1830            best_hash: B256::random(),
1831            best_number: 100,
1832            capabilities: caps_71,
1833            timeout: Arc::new(AtomicU64::new(50)),
1834            range_info: None,
1835            supports_snap: false,
1836        });
1837
1838        // Even though peer_no_71 has lower timeout,
1839        // it must NOT be selected.
1840        assert_eq!(
1841            fetcher.next_best_peer(BestPeerRequirements::EthVersion(EthVersion::Eth71)),
1842            Some(peer_with_71)
1843        );
1844    }
1845
1846    #[tokio::test]
1847    async fn test_wakes_when_eth71_peer_connects() {
1848        use futures::task::noop_waker;
1849        use std::task::{Context, Poll};
1850
1851        let manager = PeersManager::new(PeersConfig::default());
1852        let mut fetcher =
1853            StateFetcher::<EthNetworkPrimitives>::new(manager.handle(), Default::default());
1854
1855        // Queue Eth71-required request
1856        let (tx, _rx) = oneshot::channel();
1857        fetcher.queued_requests.push_back(DownloadRequest::GetBlockAccessLists {
1858            request: vec![],
1859            response: tx,
1860            priority: Priority::Normal,
1861            requirement: BalRequirement::Mandatory,
1862        });
1863
1864        let waker = noop_waker();
1865        let mut cx = Context::from_waker(&waker);
1866
1867        // No peers -> must be Pending
1868        assert!(matches!(fetcher.poll(&mut cx), Poll::Pending));
1869
1870        // Add peer WITHOUT Eth71 support
1871        let peer_old = B512::random();
1872        let caps_old = Arc::new(Capabilities::new(vec![]));
1873
1874        fetcher.new_active_peer(NewPeerInfo {
1875            peer_id: peer_old,
1876            best_hash: B256::random(),
1877            best_number: 100,
1878            capabilities: caps_old,
1879            timeout: Arc::new(AtomicU64::new(10)),
1880            range_info: None,
1881            supports_snap: false,
1882        });
1883
1884        // Still Pending
1885        assert!(matches!(fetcher.poll(&mut cx), Poll::Pending));
1886
1887        // Add peer WITH Eth71 support
1888        let peer_71 = B512::random();
1889        let caps_71 = Arc::new(Capabilities::from(vec![Capability::new("eth".into(), 71)]));
1890
1891        fetcher.new_active_peer(NewPeerInfo {
1892            peer_id: peer_71,
1893            best_hash: B256::random(),
1894            best_number: 100,
1895            capabilities: caps_71,
1896            timeout: Arc::new(AtomicU64::new(10)),
1897            range_info: None,
1898            supports_snap: false,
1899        });
1900
1901        // Now we must get Ready(BlockRequest)
1902        if let Poll::Ready(FetchAction::BlockRequest { peer_id, .. }) = fetcher.poll(&mut cx) {
1903            assert_eq!(peer_id, peer_71);
1904        }
1905    }
1906
1907    #[tokio::test]
1908    async fn test_optional_bal_request_rejected_without_eth71_peer() {
1909        use futures::task::noop_waker;
1910        use std::task::{Context, Poll};
1911
1912        let manager = PeersManager::new(PeersConfig::default());
1913        let mut fetcher =
1914            StateFetcher::<EthNetworkPrimitives>::new(manager.handle(), Default::default());
1915
1916        let peer_old = B512::random();
1917        let caps_old = Arc::new(Capabilities::new(vec![]));
1918        fetcher.new_active_peer(NewPeerInfo {
1919            peer_id: peer_old,
1920            best_hash: B256::random(),
1921            best_number: 100,
1922            capabilities: caps_old,
1923            timeout: Arc::new(AtomicU64::new(10)),
1924            range_info: None,
1925            supports_snap: false,
1926        });
1927
1928        let (tx, rx) = oneshot::channel();
1929        fetcher
1930            .download_requests_tx
1931            .send(DownloadRequest::GetBlockAccessLists {
1932                request: vec![],
1933                response: tx,
1934                priority: Priority::Normal,
1935                requirement: BalRequirement::Optional,
1936            })
1937            .unwrap();
1938
1939        let waker = noop_waker();
1940        let mut cx = Context::from_waker(&waker);
1941
1942        assert!(matches!(fetcher.poll(&mut cx), Poll::Pending));
1943        assert!(fetcher.queued_requests.is_empty());
1944        assert_eq!(rx.await.unwrap().unwrap_err(), RequestError::UnsupportedCapability);
1945    }
1946
1947    #[tokio::test]
1948    async fn test_optional_bal_request_waits_for_busy_eth71_peer() {
1949        use futures::task::noop_waker;
1950        use std::task::{Context, Poll};
1951
1952        let manager = PeersManager::new(PeersConfig::default());
1953        let mut fetcher =
1954            StateFetcher::<EthNetworkPrimitives>::new(manager.handle(), Default::default());
1955
1956        let peer_71 = B512::random();
1957        let caps_71 = Arc::new(Capabilities::from(vec![Capability::new("eth".into(), 71)]));
1958        fetcher.new_active_peer(NewPeerInfo {
1959            peer_id: peer_71,
1960            best_hash: B256::random(),
1961            best_number: 100,
1962            capabilities: caps_71,
1963            timeout: Arc::new(AtomicU64::new(10)),
1964            range_info: None,
1965            supports_snap: false,
1966        });
1967        fetcher.peers.get_mut(&peer_71).expect("peer exists").state = PeerState::GetBlockHeaders;
1968
1969        let (tx, _rx) = oneshot::channel();
1970        fetcher
1971            .download_requests_tx
1972            .send(DownloadRequest::GetBlockAccessLists {
1973                request: vec![],
1974                response: tx,
1975                priority: Priority::Normal,
1976                requirement: BalRequirement::Optional,
1977            })
1978            .unwrap();
1979
1980        let waker = noop_waker();
1981        let mut cx = Context::from_waker(&waker);
1982
1983        assert!(matches!(fetcher.poll(&mut cx), Poll::Pending));
1984        assert_eq!(fetcher.queued_requests.len(), 1);
1985    }
1986
1987    #[tokio::test]
1988    async fn test_queued_optional_bal_request_rejected_after_eth71_disconnect() {
1989        use futures::task::noop_waker;
1990        use std::task::{Context, Poll};
1991
1992        let manager = PeersManager::new(PeersConfig::default());
1993        let mut fetcher =
1994            StateFetcher::<EthNetworkPrimitives>::new(manager.handle(), Default::default());
1995
1996        let peer_old = B512::random();
1997        let caps_old = Arc::new(Capabilities::new(vec![]));
1998        fetcher.new_active_peer(NewPeerInfo {
1999            peer_id: peer_old,
2000            best_hash: B256::random(),
2001            best_number: 100,
2002            capabilities: caps_old,
2003            timeout: Arc::new(AtomicU64::new(10)),
2004            range_info: None,
2005            supports_snap: false,
2006        });
2007
2008        let peer_71 = B512::random();
2009        let caps_71 = Arc::new(Capabilities::from(vec![Capability::new("eth".into(), 71)]));
2010        fetcher.new_active_peer(NewPeerInfo {
2011            peer_id: peer_71,
2012            best_hash: B256::random(),
2013            best_number: 100,
2014            capabilities: caps_71,
2015            timeout: Arc::new(AtomicU64::new(10)),
2016            range_info: None,
2017            supports_snap: false,
2018        });
2019        fetcher.peers.get_mut(&peer_71).expect("peer exists").state = PeerState::GetBlockHeaders;
2020
2021        let (tx, rx) = oneshot::channel();
2022        fetcher
2023            .download_requests_tx
2024            .send(DownloadRequest::GetBlockAccessLists {
2025                request: vec![],
2026                response: tx,
2027                priority: Priority::Normal,
2028                requirement: BalRequirement::Optional,
2029            })
2030            .unwrap();
2031
2032        let waker = noop_waker();
2033        let mut cx = Context::from_waker(&waker);
2034
2035        assert!(matches!(fetcher.poll(&mut cx), Poll::Pending));
2036        assert_eq!(fetcher.queued_requests.len(), 1);
2037
2038        fetcher.on_session_closed(&peer_71);
2039
2040        assert!(matches!(fetcher.poll(&mut cx), Poll::Pending));
2041        assert!(fetcher.queued_requests.is_empty());
2042        assert_eq!(rx.await.unwrap().unwrap_err(), RequestError::UnsupportedCapability);
2043    }
2044
2045    #[tokio::test]
2046    async fn test_next_best_peer_snap_no_support() {
2047        let manager = PeersManager::new(PeersConfig::default());
2048        let mut fetcher =
2049            StateFetcher::<EthNetworkPrimitives>::new(manager.handle(), Default::default());
2050
2051        let peer = B512::random();
2052        fetcher.new_active_peer(NewPeerInfo {
2053            peer_id: peer,
2054            best_hash: B256::random(),
2055            best_number: 100,
2056            capabilities: Arc::new(Capabilities::new(vec![])),
2057            timeout: Arc::new(AtomicU64::new(10)),
2058            range_info: None,
2059            supports_snap: false,
2060        });
2061
2062        assert_eq!(fetcher.next_best_peer(BestPeerRequirements::SupportsSnap), None);
2063    }
2064
2065    #[tokio::test]
2066    async fn test_next_best_peer_snap_supported() {
2067        let manager = PeersManager::new(PeersConfig::default());
2068        let mut fetcher =
2069            StateFetcher::<EthNetworkPrimitives>::new(manager.handle(), Default::default());
2070
2071        let peer = B512::random();
2072        fetcher.new_active_peer(NewPeerInfo {
2073            peer_id: peer,
2074            best_hash: B256::random(),
2075            best_number: 100,
2076            capabilities: Arc::new(Capabilities::new(vec![])),
2077            timeout: Arc::new(AtomicU64::new(10)),
2078            range_info: None,
2079            supports_snap: true,
2080        });
2081
2082        assert_eq!(fetcher.next_best_peer(BestPeerRequirements::SupportsSnap), Some(peer));
2083    }
2084
2085    #[tokio::test]
2086    async fn test_next_best_peer_snap_filters_correctly() {
2087        let manager = PeersManager::new(PeersConfig::default());
2088        let mut fetcher =
2089            StateFetcher::<EthNetworkPrimitives>::new(manager.handle(), Default::default());
2090
2091        let peer_no_snap = B512::random();
2092        let peer_with_snap = B512::random();
2093
2094        fetcher.new_active_peer(NewPeerInfo {
2095            peer_id: peer_no_snap,
2096            best_hash: B256::random(),
2097            best_number: 100,
2098            capabilities: Arc::new(Capabilities::new(vec![])),
2099            timeout: Arc::new(AtomicU64::new(5)),
2100            range_info: None,
2101            supports_snap: false,
2102        });
2103        fetcher.new_active_peer(NewPeerInfo {
2104            peer_id: peer_with_snap,
2105            best_hash: B256::random(),
2106            best_number: 100,
2107            capabilities: Arc::new(Capabilities::new(vec![])),
2108            timeout: Arc::new(AtomicU64::new(50)),
2109            range_info: None,
2110            supports_snap: true,
2111        });
2112
2113        // Even though peer_no_snap has a lower timeout, it must NOT be selected.
2114        assert_eq!(
2115            fetcher.next_best_peer(BestPeerRequirements::SupportsSnap),
2116            Some(peer_with_snap)
2117        );
2118    }
2119
2120    #[tokio::test]
2121    async fn test_snap_request_rejected_without_snap_peer() {
2122        use futures::task::noop_waker;
2123        use std::task::{Context, Poll};
2124
2125        let manager = PeersManager::new(PeersConfig::default());
2126        let mut fetcher =
2127            StateFetcher::<EthNetworkPrimitives>::new(manager.handle(), Default::default());
2128
2129        // Only an eth-only peer is connected.
2130        fetcher.new_active_peer(NewPeerInfo {
2131            peer_id: B512::random(),
2132            best_hash: B256::random(),
2133            best_number: 100,
2134            capabilities: Arc::new(Capabilities::new(vec![])),
2135            timeout: Arc::new(AtomicU64::new(10)),
2136            range_info: None,
2137            supports_snap: false,
2138        });
2139
2140        let (tx, rx) = oneshot::channel();
2141        fetcher
2142            .download_requests_tx
2143            .send(DownloadRequest::GetSnap {
2144                request: SnapProtocolMessage::GetAccountRange(GetAccountRangeMessage {
2145                    request_id: 0,
2146                    root_hash: B256::ZERO,
2147                    starting_hash: B256::ZERO,
2148                    limit_hash: B256::ZERO,
2149                    response_bytes: 0,
2150                }),
2151                response: tx,
2152                priority: Priority::Normal,
2153            })
2154            .unwrap();
2155
2156        let waker = noop_waker();
2157        let mut cx = Context::from_waker(&waker);
2158
2159        assert!(matches!(fetcher.poll(&mut cx), Poll::Pending));
2160        assert!(fetcher.queued_requests.is_empty());
2161        assert_eq!(rx.await.unwrap().unwrap_err(), RequestError::UnsupportedCapability);
2162    }
2163
2164    #[tokio::test]
2165    async fn test_snap_response_triggers_followup() {
2166        let manager = PeersManager::new(PeersConfig::default());
2167        let mut fetcher =
2168            StateFetcher::<EthNetworkPrimitives>::new(manager.handle(), Default::default());
2169
2170        let peer_id = B512::random();
2171        fetcher.new_active_peer(NewPeerInfo {
2172            peer_id,
2173            best_hash: B256::random(),
2174            best_number: 100,
2175            capabilities: Arc::new(Capabilities::new(vec![])),
2176            timeout: Arc::new(AtomicU64::new(10)),
2177            range_info: None,
2178            supports_snap: true,
2179        });
2180
2181        // Queue a followup snap request for the same peer.
2182        let (followup_tx, _followup_rx) = oneshot::channel();
2183        fetcher.queued_requests.push_back(DownloadRequest::GetSnap {
2184            request: SnapProtocolMessage::GetAccountRange(GetAccountRangeMessage {
2185                request_id: 0,
2186                root_hash: B256::ZERO,
2187                starting_hash: B256::ZERO,
2188                limit_hash: B256::ZERO,
2189                response_bytes: 0,
2190            }),
2191            response: followup_tx,
2192            priority: Priority::Normal,
2193        });
2194
2195        let (tx, mut rx) = oneshot::channel();
2196        fetcher.inflight_snap_requests.insert(peer_id, Request { request: (), response: tx });
2197        fetcher.peers.get_mut(&peer_id).unwrap().state = PeerState::GetSnap;
2198
2199        let resp = SnapResponse::AccountRange(AccountRangeMessage {
2200            request_id: 1,
2201            accounts: vec![],
2202            proof: vec![],
2203        });
2204        let outcome = fetcher.on_snap_response(peer_id, Ok(resp));
2205
2206        assert!(matches!(outcome, Some(BlockResponseOutcome::Request(pid, _)) if pid == peer_id));
2207        assert!(rx.try_recv().is_ok());
2208    }
2209
2210    #[tokio::test]
2211    async fn test_queued_snap_request_rejected_after_last_peer_disconnects() {
2212        use futures::task::noop_waker;
2213        use std::task::{Context, Poll};
2214
2215        let manager = PeersManager::new(PeersConfig::default());
2216        let mut fetcher =
2217            StateFetcher::<EthNetworkPrimitives>::new(manager.handle(), Default::default());
2218
2219        // The only connected peer supports snap but is busy, so the request gets queued.
2220        let peer = B512::random();
2221        fetcher.new_active_peer(NewPeerInfo {
2222            peer_id: peer,
2223            best_hash: B256::random(),
2224            best_number: 100,
2225            capabilities: Arc::new(Capabilities::new(vec![])),
2226            timeout: Arc::new(AtomicU64::new(10)),
2227            range_info: None,
2228            supports_snap: true,
2229        });
2230        fetcher.peers.get_mut(&peer).expect("peer exists").state = PeerState::GetBlockHeaders;
2231
2232        let (tx, rx) = oneshot::channel();
2233        fetcher
2234            .download_requests_tx
2235            .send(DownloadRequest::GetSnap {
2236                request: SnapProtocolMessage::GetAccountRange(GetAccountRangeMessage {
2237                    request_id: 0,
2238                    root_hash: B256::ZERO,
2239                    starting_hash: B256::ZERO,
2240                    limit_hash: B256::ZERO,
2241                    response_bytes: 0,
2242                }),
2243                response: tx,
2244                priority: Priority::Normal,
2245            })
2246            .unwrap();
2247
2248        let waker = noop_waker();
2249        let mut cx = Context::from_waker(&waker);
2250
2251        assert!(matches!(fetcher.poll(&mut cx), Poll::Pending));
2252        assert_eq!(fetcher.queued_requests.len(), 1);
2253
2254        // The only peer disconnects, leaving `self.peers` empty. The still-queued request must
2255        // resolve immediately instead of waiting for a peer that can never come back.
2256        fetcher.on_session_closed(&peer);
2257
2258        assert!(matches!(fetcher.poll(&mut cx), Poll::Pending));
2259        assert!(fetcher.queued_requests.is_empty());
2260        assert_eq!(rx.await.unwrap().unwrap_err(), RequestError::UnsupportedCapability);
2261    }
2262}