1mod 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#[derive(Debug)]
50pub struct StateFetcher<N: NetworkPrimitives = EthNetworkPrimitives> {
51 inflight_headers_requests: HashMap<PeerId, InflightHeadersRequest<N::BlockHeader>>,
53 inflight_bodies_requests: HashMap<PeerId, InflightBodiesRequest<N::BlockBody>>,
55 inflight_bals_requests: HashMap<PeerId, InflightBlockAccessListsRequest>,
57 inflight_receipts_requests: HashMap<PeerId, InflightReceiptsRequest<N::Receipt>>,
59 inflight_snap_requests: HashMap<PeerId, InflightSnapRequest>,
61 peers: HashMap<PeerId, Peer>,
63 peers_handle: PeersHandle,
65 num_active_peers: Arc<AtomicUsize>,
67 queued_requests: VecDeque<DownloadRequest<N>>,
69 download_requests_rx: UnboundedReceiverStream<DownloadRequest<N>>,
71 download_requests_tx: UnboundedSender<DownloadRequest<N>>,
73}
74
75impl<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 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 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 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 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 fn next_best_peer(&self, requirement: BestPeerRequirements) -> Option<PeerId> {
172 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 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 if maybe_better.1.is_better(best_peer.1, &requirement) {
189 best_peer = maybe_better;
190 continue
191 }
192
193 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 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 fn poll_action(&mut self) -> PollAction {
214 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 if self.should_fail_fast(&request) {
224 request.send_err_response(RequestError::UnsupportedCapability);
225 } else {
226 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 pub(crate) fn poll(&mut self, cx: &mut Context<'_>) -> Poll<FetchAction> {
240 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 match self.download_requests_rx.poll_next_unpin(cx) {
251 Poll::Ready(Some(request)) => {
252 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 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 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 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 fn prepare_block_request(&mut self, peer_id: PeerId, req: DownloadRequest<N>) -> BlockRequest {
306 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 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 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 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 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 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 peer.last_response_likely_bad = is_likely_bad_response;
401
402 if peer.state.on_request_finished() && !is_error && !is_likely_bad_response {
405 return self.followup_request(peer_id)
406 }
407 }
408
409 maybe_reputation_change
412 .map(|reputation_change| BlockResponseOutcome::BadResponse(peer_id, reputation_change))
413 }
414
415 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 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 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 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 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 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
513enum PollAction {
515 Ready(FetchAction),
516 NoRequests,
517 NoPeersAvailable,
518}
519
520#[derive(Debug)]
522pub(crate) struct NewPeerInfo {
523 pub(crate) peer_id: PeerId,
525 pub(crate) best_hash: B256,
527 pub(crate) best_number: u64,
529 pub(crate) capabilities: Arc<Capabilities>,
531 pub(crate) timeout: Arc<AtomicU64>,
533 pub(crate) range_info: Option<BlockRangeInfo>,
535 pub(crate) supports_snap: bool,
537}
538
539#[derive(Debug)]
541struct Peer {
542 state: PeerState,
544 best_hash: B256,
546 best_number: u64,
548 #[allow(dead_code)]
550 capabilities: Arc<Capabilities>,
551 timeout: Arc<AtomicU64>,
553 last_response_likely_bad: bool,
560 range_info: Option<BlockRangeInfo>,
562 supports_snap: bool,
564}
565
566impl Peer {
567 fn timeout(&self) -> u64 {
568 self.timeout.load(Ordering::Relaxed)
569 }
570
571 fn earliest(&self) -> u64 {
573 self.range_info.as_ref().map_or(0, |info| info.earliest())
574 }
575
576 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 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 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 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, (false, true) => false, (true, true) => false, (false, false) => {
620 self_r.start() < other_r.start()
622 }
623 }
624 }
625 (Some(self_r), None) => {
626 self_r.contains(range.start()) && self_r.contains(range.end())
629 }
630 (None, Some(other_r)) => {
631 !(other_r.contains(range.start()) && other_r.contains(range.end()))
634 }
635 (None, None) => false, }
637 }
638
639 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 BestPeerRequirements::None |
647 BestPeerRequirements::EthVersion(_) |
648 BestPeerRequirements::SupportsSnap => false,
649 }
650 }
651}
652
653#[derive(Debug)]
655enum PeerState {
656 Idle,
658 GetBlockHeaders,
660 GetBlockBodies,
662 GetBlockAccessLists,
664 GetReceipts,
666 GetSnap,
668 Closing,
670}
671
672impl PeerState {
675 const fn is_idle(&self) -> bool {
677 matches!(self, Self::Idle)
678 }
679
680 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#[derive(Debug)]
697struct Request<Req, Resp> {
698 request: Req,
701 response: oneshot::Sender<Resp>,
702}
703
704#[derive(Debug)]
706#[expect(clippy::enum_variant_names)]
707pub(crate) enum DownloadRequest<N: NetworkPrimitives> {
708 GetBlockHeaders {
710 request: HeadersRequest,
711 response: oneshot::Sender<PeerRequestResult<Vec<N::BlockHeader>>>,
712 priority: Priority,
713 },
714 GetBlockBodies {
716 request: Vec<B256>,
717 response: oneshot::Sender<PeerRequestResult<Vec<N::BlockBody>>>,
718 priority: Priority,
719 range_hint: Option<RangeInclusive<u64>>,
720 },
721 GetBlockAccessLists {
723 request: Vec<B256>,
724 response: oneshot::Sender<PeerRequestResult<BlockAccessLists>>,
725 priority: Priority,
726 requirement: BalRequirement,
727 },
728 GetReceipts {
730 request: Vec<B256>,
731 response: oneshot::Sender<PeerRequestResult<ReceiptsResponse<N::Receipt>>>,
732 priority: Priority,
733 },
734 GetSnap {
736 request: SnapProtocolMessage,
737 response: oneshot::Sender<PeerRequestResult<SnapResponse>>,
738 priority: Priority,
739 },
740}
741
742impl<N: NetworkPrimitives> DownloadRequest<N> {
745 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 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 const fn is_normal_priority(&self) -> bool {
769 self.get_priority().is_normal()
770 }
771
772 const fn is_optional_bal(&self) -> bool {
774 matches!(self, Self::GetBlockAccessLists { requirement: BalRequirement::Optional, .. })
775 }
776
777 const fn is_snap(&self) -> bool {
779 matches!(self, Self::GetSnap { .. })
780 }
781
782 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 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
811pub(crate) enum FetchAction {
813 BlockRequest {
815 peer_id: PeerId,
817 request: BlockRequest,
819 },
820}
821
822#[derive(Debug, PartialEq, Eq)]
826pub(crate) enum BlockResponseOutcome {
827 Request(PeerId, BlockRequest),
829 BadResponse(PeerId, ReputationChangeKind),
831}
832
833enum BestPeerRequirements {
835 None,
837 FullBlockRange(RangeInclusive<u64>),
839 FullBlock,
841 EthVersion(EthVersion),
843 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 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 fetcher.on_pending_disconnect(&first_peer);
911 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 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 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 assert_eq!(fetcher.next_best_peer(BestPeerRequirements::None), Some(peer1));
963 assert_eq!(fetcher.next_best_peer(BestPeerRequirements::None), Some(peer1));
964 peer2_timeout.store(10, Ordering::Relaxed);
966 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 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 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 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 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 assert!(peer_full.is_better(&peer_partial, &BestPeerRequirements::FullBlock));
1164 assert!(!peer_partial.is_better(&peer_full, &BestPeerRequirements::FullBlock));
1165
1166 assert!(peer_no_range.is_better(&peer_partial, &BestPeerRequirements::FullBlock));
1168 assert!(!peer_partial.is_better(&peer_no_range, &BestPeerRequirements::FullBlock));
1169
1170 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 assert!(!peer_no_range
1343 .is_better(&peer_with_range, &BestPeerRequirements::FullBlockRange(range.clone())));
1344
1345 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 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 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 assert!(peer_with_range_covers
1381 .is_better(&peer_no_range, &BestPeerRequirements::FullBlockRange(range.clone())));
1382
1383 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 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 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 assert!(!peer_with_range_no_cover
1418 .is_better(&peer_no_range, &BestPeerRequirements::FullBlockRange(range.clone())));
1419
1420 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 let range = RangeInclusive::new(50, 100);
1429
1430 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 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 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 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 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 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 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 #[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 assert!(matches!(fetcher.peers[&peer_id].state, PeerState::GetReceipts));
1536 assert!(fetcher.inflight_receipts_requests.contains_key(&peer_id));
1538
1539 Poll::Ready(())
1540 })
1541 .await;
1542 }
1543
1544 #[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 assert!(outcome.is_none());
1557 assert!(fetcher.peers[&peer_id].state.is_idle());
1559 assert!(!fetcher.inflight_receipts_requests.contains_key(&peer_id));
1561
1562 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 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 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 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 assert!(matches!(fetcher.peers[&peer_id].state, PeerState::GetReceipts));
1744
1745 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 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 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 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 let caps_old = Arc::new(Capabilities::new(vec![]));
1814
1815 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 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 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 assert!(matches!(fetcher.poll(&mut cx), Poll::Pending));
1869
1870 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 assert!(matches!(fetcher.poll(&mut cx), Poll::Pending));
1886
1887 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 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 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 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 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 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 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}