1use std::{
2 collections::HashMap,
3 sync::{
4 atomic::{AtomicU64, Ordering},
5 Arc, Mutex as StdMutex,
6 },
7};
8
9use tokio::{sync::mpsc, task};
10use tokio_util::sync::CancellationToken;
11
12use super::{events::*, pipe::*, wire::*, *};
13use crate::zakura::{
14 handle_pipe_exit, spawn_supervised_pipe, BoxRunFuture, Flow, Frame, FramedRecv, FramedSend,
15 OrderedSendError, Peer, PeerStreamSession, Pipe, Service, ServicePeerDirection, SessionGuard,
16 Sink, SinkReject, Stream, StreamMode, ZakuraConnId, ZakuraPeerId, ZakuraSupervisorHandle,
17 ZAKURA_CAP_HEADER_SYNC,
18};
19
20const HEADER_SYNC_SERVICE_STREAMS: [Stream; 1] = [Stream {
21 kind: ZAKURA_STREAM_HEADER_SYNC,
22 version: ZAKURA_HEADER_SYNC_STREAM_VERSION,
23 frame_cap: (MAX_HS_MESSAGE_BYTES + FRAME_HEADER_BYTES) as u32,
28 capability: ZAKURA_CAP_HEADER_SYNC,
29 mode: StreamMode::Ordered,
30}];
31
32pub(crate) fn header_sync_streams() -> &'static [Stream] {
34 &HEADER_SYNC_SERVICE_STREAMS
35}
36
37#[derive(Clone, Debug)]
39pub struct HeaderSyncPeerSession {
40 peer_id: ZakuraPeerId,
41 session_id: u64,
42 direction: ServicePeerDirection,
43 inner: Arc<HeaderSyncPeerSessionInner>,
44}
45
46#[derive(Debug)]
47struct HeaderSyncPeerSessionInner {
48 send: FramedSend,
49 cancel_token: CancellationToken,
50 commands: Option<mpsc::UnboundedSender<HeaderSyncPeerCommand>>,
51 next_request_id: AtomicU64,
52}
53
54impl HeaderSyncPeerSession {
55 fn new_with_commands(
56 session: &PeerStreamSession,
57 direction: ServicePeerDirection,
58 commands: mpsc::UnboundedSender<HeaderSyncPeerCommand>,
59 session_id: u64,
60 ) -> Self {
61 Self::from_parts_with_direction_and_commands(
62 session.peer_id().clone(),
63 session_id,
64 direction,
65 session.sender(),
66 session.cancel_token(),
67 Some(commands),
68 )
69 }
70
71 #[cfg(test)]
72 pub(crate) fn from_parts(
73 peer_id: ZakuraPeerId,
74 send: FramedSend,
75 cancel_token: CancellationToken,
76 ) -> Self {
77 Self::from_parts_with_direction(peer_id, ServicePeerDirection::Inbound, send, cancel_token)
78 }
79
80 #[cfg(test)]
81 pub(crate) fn from_parts_with_direction(
82 peer_id: ZakuraPeerId,
83 direction: ServicePeerDirection,
84 send: FramedSend,
85 cancel_token: CancellationToken,
86 ) -> Self {
87 Self::from_parts_with_direction_and_commands(
88 peer_id,
89 0,
90 direction,
91 send,
92 cancel_token,
93 None,
94 )
95 }
96
97 #[cfg(test)]
98 pub(crate) fn from_parts_with_direction_and_session_id(
99 peer_id: ZakuraPeerId,
100 direction: ServicePeerDirection,
101 send: FramedSend,
102 cancel_token: CancellationToken,
103 session_id: u64,
104 ) -> Self {
105 Self::from_parts_with_direction_and_commands(
106 peer_id,
107 session_id,
108 direction,
109 send,
110 cancel_token,
111 None,
112 )
113 }
114
115 fn from_parts_with_direction_and_commands(
116 peer_id: ZakuraPeerId,
117 session_id: u64,
118 direction: ServicePeerDirection,
119 send: FramedSend,
120 cancel_token: CancellationToken,
121 commands: Option<mpsc::UnboundedSender<HeaderSyncPeerCommand>>,
122 ) -> Self {
123 Self {
124 peer_id,
125 session_id,
126 direction,
127 inner: Arc::new(HeaderSyncPeerSessionInner {
128 send,
129 cancel_token,
130 commands,
131 next_request_id: AtomicU64::new(1),
132 }),
133 }
134 }
135
136 pub fn peer_id(&self) -> &ZakuraPeerId {
138 &self.peer_id
139 }
140
141 pub fn session_id(&self) -> u64 {
143 self.session_id
144 }
145
146 pub fn direction(&self) -> ServicePeerDirection {
148 self.direction
149 }
150
151 pub fn cancel_token(&self) -> CancellationToken {
153 self.inner.cancel_token.clone()
154 }
155
156 pub fn outbound_capacity(&self) -> usize {
158 self.inner.send.capacity()
159 }
160
161 pub fn outbound_max_capacity(&self) -> usize {
163 self.inner.send.max_capacity()
164 }
165
166 pub fn retire_expected_headers(
168 &self,
169 request_id: HeaderSyncRequestId,
170 ) -> Result<(), OrderedSendError> {
171 let Some(commands) = &self.inner.commands else {
172 return Ok(());
173 };
174 commands
175 .send(HeaderSyncPeerCommand::Retire(request_id))
176 .map_err(|_| OrderedSendError::Closed)
177 }
178
179 fn next_request_id(&self) -> Result<HeaderSyncRequestId, OrderedSendError> {
180 let mut id = self.inner.next_request_id.load(Ordering::Relaxed);
181 loop {
182 let next_id = id.checked_add(1).ok_or_else(|| {
183 OrderedSendError::Encode("header-sync request ID counter exhausted".into())
184 })?;
185 match self.inner.next_request_id.compare_exchange_weak(
186 id,
187 next_id,
188 Ordering::Relaxed,
189 Ordering::Relaxed,
190 ) {
191 Ok(_) => break,
192 Err(current_id) => id = current_id,
193 }
194 }
195 HeaderSyncRequestId::new(id).ok_or_else(|| {
196 OrderedSendError::Encode("header-sync request ID counter exhausted".into())
197 })
198 }
199
200 pub fn try_send_status(&self, status: HeaderSyncStatus) -> Result<(), OrderedSendError> {
202 self.try_send_message(HeaderSyncMessage::Status(status), None)
203 }
204
205 pub(super) fn prepare_get_headers(
207 &self,
208 start_height: block::Height,
209 count: u32,
210 want_tree_aux_roots: bool,
211 ) -> Result<PreparedGetHeaders, OrderedSendError> {
212 let request_id = self.next_request_id()?;
213 let expected =
214 ExpectedHeadersResponse::new(request_id, start_height, count, want_tree_aux_roots)
215 .map_err(|error| OrderedSendError::Encode(Box::new(error)))?;
216 let frame = HeaderSyncMessage::GetHeaders {
217 start_height,
218 count,
219 want_tree_aux_roots,
220 }
221 .encode_frame(Some(request_id))
222 .map_err(|error| OrderedSendError::Encode(Box::new(error)))?;
223
224 let reservation = ExpectedHeadersReservation::new(self.inner.commands.clone(), expected)?;
225 Ok(PreparedGetHeaders {
226 request_id,
227 frame,
228 send: self.inner.send.clone(),
229 reservation,
230 })
231 }
232
233 pub fn try_send_headers_with_sizes_and_roots(
236 &self,
237 request_id: HeaderSyncRequestId,
238 headers: Vec<Arc<block::Header>>,
239 body_sizes: Vec<u32>,
240 tree_aux_roots: Vec<BlockCommitmentRoots>,
241 ) -> Result<(), OrderedSendError> {
242 self.try_send_message(
243 HeaderSyncMessage::Headers {
244 headers,
245 body_sizes,
246 tree_aux_roots,
247 },
248 Some(request_id),
249 )
250 }
251
252 pub fn try_send_new_block(&self, block: Arc<block::Block>) -> Result<(), OrderedSendError> {
254 self.try_send_message(HeaderSyncMessage::NewBlock(block), None)
255 }
256
257 fn try_send_message(
258 &self,
259 msg: HeaderSyncMessage,
260 request_id: Option<HeaderSyncRequestId>,
261 ) -> Result<(), OrderedSendError> {
262 let frame = msg
263 .encode_frame(request_id)
264 .map_err(|error| OrderedSendError::Encode(Box::new(error)))?;
265 match self.inner.send.try_send(frame) {
266 Ok(()) => Ok(()),
267 Err(mpsc::error::TrySendError::Full(_frame)) => Err(OrderedSendError::Full),
268 Err(mpsc::error::TrySendError::Closed(_frame)) => Err(OrderedSendError::Closed),
269 }
270 }
271}
272
273pub(super) struct PreparedGetHeaders {
274 request_id: HeaderSyncRequestId,
275 frame: Frame,
276 send: FramedSend,
277 reservation: ExpectedHeadersReservation,
278}
279
280impl PreparedGetHeaders {
281 pub(super) fn request_id(&self) -> HeaderSyncRequestId {
282 self.request_id
283 }
284
285 pub(super) async fn send(self) -> Result<HeaderSyncRequestId, OrderedSendError> {
290 let Self {
291 request_id,
292 frame,
293 send,
294 mut reservation,
295 } = self;
296 send.send(frame)
297 .await
298 .map_err(|_| OrderedSendError::Closed)?;
299 reservation.disarm();
300 Ok(request_id)
301 }
302}
303
304struct ExpectedHeadersReservation {
305 commands: Option<mpsc::UnboundedSender<HeaderSyncPeerCommand>>,
306 expected: ExpectedHeadersResponse,
307 armed: bool,
308}
309
310impl ExpectedHeadersReservation {
311 fn new(
312 commands: Option<mpsc::UnboundedSender<HeaderSyncPeerCommand>>,
313 expected: ExpectedHeadersResponse,
314 ) -> Result<Self, OrderedSendError> {
315 if let Some(commands) = &commands {
316 commands
317 .send(HeaderSyncPeerCommand::Reserve(expected))
318 .map_err(|_| OrderedSendError::Closed)?;
319 }
320 Ok(Self {
321 commands,
322 expected,
323 armed: true,
324 })
325 }
326
327 fn disarm(&mut self) {
328 self.armed = false;
329 }
330}
331
332impl Drop for ExpectedHeadersReservation {
333 fn drop(&mut self) {
334 if self.armed {
335 if let Some(commands) = &self.commands {
336 let _ = commands.send(HeaderSyncPeerCommand::Cancel(self.expected));
337 }
338 }
339 }
340}
341
342#[derive(Debug)]
344pub(super) enum HeaderSyncPeerCommand {
345 Reserve(ExpectedHeadersResponse),
347 Cancel(ExpectedHeadersResponse),
349 Retire(HeaderSyncRequestId),
351}
352
353pub(crate) async fn drive_header_sync_actions(
355 mut actions: mpsc::Receiver<HeaderSyncAction>,
356 handle: HeaderSyncHandle,
357 _supervisor: ZakuraSupervisorHandle,
360 shutdown: CancellationToken,
361) {
362 loop {
363 let action = tokio::select! {
364 _ = shutdown.cancelled() => return,
365 action = actions.recv() => {
366 let Some(action) = action else {
367 return;
368 };
369 action
370 }
371 };
372
373 match action {
374 #[cfg(test)]
375 HeaderSyncAction::SendMessage { .. } | HeaderSyncAction::ForwardNewBlock { .. } => {}
376 HeaderSyncAction::Misbehavior { peer, reason } => {
377 tracing::debug!(?peer, ?reason, "recorded Zakura header-sync peer violation");
379 }
380 HeaderSyncAction::NewBlockReceived { peer, hash, .. } => {
381 tracing::debug!(
382 ?peer,
383 ?hash,
384 "Zakura header-sync NewBlock body arrived before block-acceptance hook is wired"
385 );
386 }
387 HeaderSyncAction::QueryHeadersByHeightRange {
388 peer,
389 session_id,
390 request_id,
391 start,
392 count,
393 ..
394 } => {
395 let _ = handle
396 .send(HeaderSyncEvent::HeaderRangeResponseFinished {
397 peer,
398 session_id,
399 request_id,
400 start_height: start,
401 requested_count: count,
402 returned_count: 0,
403 })
404 .await;
405 }
406 HeaderSyncAction::CommitHeaderRange {
407 peer,
408 start_height,
409 headers,
410 ..
411 } => {
412 tracing::debug!(
413 ?peer,
414 ?start_height,
415 count = headers.len(),
416 "suppressing Zakura header range commit until state driver is wired"
417 );
418 }
419 HeaderSyncAction::QueryBestHeaderTip
420 | HeaderSyncAction::QueryMissingBlockBodies { .. }
421 | HeaderSyncAction::BodyGaps { .. }
422 | HeaderSyncAction::HeaderAdvanced { .. }
423 | HeaderSyncAction::HeaderReanchored { .. } => {}
424 }
425 }
426}
427
428#[derive(Debug)]
430pub(crate) struct HeaderSyncService {
431 header_sync: HeaderSyncHandle,
432 peers: Arc<StdMutex<HashMap<ZakuraPeerId, HeaderSyncPeerRecord>>>,
433}
434
435#[derive(Debug)]
436struct HeaderSyncPeerRecord {
437 conn_id: ZakuraConnId,
438 session_id: u64,
439 cancel_token: CancellationToken,
440}
441
442impl HeaderSyncService {
443 pub(crate) fn new(header_sync: HeaderSyncHandle) -> Self {
444 Self {
445 header_sync,
446 peers: Arc::new(StdMutex::new(HashMap::new())),
447 }
448 }
449}
450
451impl Service for HeaderSyncService {
452 fn name(&self) -> &'static str {
453 "header-sync"
454 }
455
456 fn streams(&self) -> &[Stream] {
457 header_sync_streams()
458 }
459
460 fn wants_peer(
461 &self,
462 _peer: &ZakuraPeerId,
463 _negotiated: u64,
464 direction: ServicePeerDirection,
465 ) -> bool {
466 let snapshot = self.header_sync.peer_snapshot();
469 match direction {
470 ServicePeerDirection::Inbound => snapshot.inbound_slots_free > 0,
471 ServicePeerDirection::Outbound => snapshot.outbound_slots_free > 0,
472 }
473 }
474
475 fn add_peer(&self, mut peer: Peer) {
476 let Some((session_id, recv, send)) =
477 peer.take_stream_with_session_id(ZAKURA_STREAM_HEADER_SYNC)
478 else {
479 return;
480 };
481
482 let peer_id = peer.id.clone();
483 let session = PeerStreamSession::new(
484 peer_id.clone(),
485 ZAKURA_STREAM_HEADER_SYNC,
486 recv,
487 send,
488 peer.service_cancel_token(),
489 );
490 let service_cancel_token = session.cancel_token();
496 let connection_cancel_token = peer.cancel_token();
497 let close_cause = peer.close_cause();
498 let conn_id = peer.conn_id;
499 let (commands_tx, commands_rx) = mpsc::unbounded_channel();
500 let header_sync_session = HeaderSyncPeerSession::new_with_commands(
501 &session,
502 peer.direction,
503 commands_tx,
504 session_id,
505 );
506
507 {
508 let mut peers = self
509 .peers
510 .lock()
511 .expect("header-sync peer map mutex is never poisoned");
512 if peers
513 .get(&peer_id)
514 .is_some_and(|record| record.conn_id > conn_id)
515 {
516 service_cancel_token.cancel();
517 return;
518 }
519 if let Some(old_record) = peers.insert(
520 peer_id.clone(),
521 HeaderSyncPeerRecord {
522 conn_id,
523 session_id,
524 cancel_token: header_sync_session.cancel_token(),
525 },
526 ) {
527 old_record.cancel_token.cancel();
528 }
529 }
530
531 let _ = self
532 .header_sync
533 .send_lifecycle(HeaderSyncEvent::PeerConnected(header_sync_session.clone()));
534
535 let (_session_peer, _stream_kind, recv, _send, _session_cancel) = session.into_parts();
536
537 let pipe = Pipe::new(
541 peer_id.clone(),
542 HsLocal::new(commands_rx, DEFAULT_HS_INBOUND_NEW_BLOCK_MIN_INTERVAL),
543 HsEnv::new_with_session_id(self.header_sync.clone(), session_id),
544 SessionGuard::oversize_only(header_sync_guard_max_bytes()),
545 run_inbound,
546 &PIPE_SHAPE,
547 );
548 let pipe_cancel_token = service_cancel_token.clone();
555 let protocol_connection_cancel_token = connection_cancel_token.clone();
556 let protocol_close_cause = close_cause.clone();
557 let pipe = async move {
558 handle_pipe_exit(
559 "header-sync",
560 &protocol_connection_cancel_token,
561 &protocol_close_cause,
562 run_peer(pipe, recv, pipe_cancel_token).await,
563 );
564 };
565
566 let teardown_handle = self.header_sync.clone();
573 let teardown_peers = self.peers.clone();
574 let teardown_peer = peer_id.clone();
575 let on_teardown = move || {
576 let should_notify = {
577 let mut peers = teardown_peers
578 .lock()
579 .expect("header-sync peer map mutex is never poisoned");
580 if peers.get(&teardown_peer).is_some_and(|record| {
581 record.conn_id == conn_id && record.session_id == session_id
582 }) {
583 peers.remove(&teardown_peer);
584 true
585 } else {
586 false
587 }
588 };
589 if should_notify {
590 let _ = teardown_handle
591 .send_lifecycle(HeaderSyncEvent::PeerDisconnected(teardown_peer));
592 }
593 };
594 let panic_connection_cancel_token = connection_cancel_token.clone();
595 let panic_close_cause = close_cause.clone();
596 let on_panic = move || {
597 panic_close_cause.record("service_panic");
598 panic_connection_cancel_token.cancel();
599 };
600
601 spawn_supervised_pipe(peer_id, service_cancel_token, on_teardown, on_panic, pipe);
604 }
605
606 fn remove_peer(&self, peer: &ZakuraPeerId, conn_id: ZakuraConnId) {
607 let removed = {
608 let mut peers = self
609 .peers
610 .lock()
611 .expect("header-sync peer map mutex is never poisoned");
612 if peers
613 .get(peer)
614 .is_some_and(|record| record.conn_id == conn_id)
615 {
616 peers.remove(peer)
617 } else {
618 None
619 }
620 };
621 if let Some(record) = removed {
622 record.cancel_token.cancel();
623 let _ = self
624 .header_sync
625 .send_lifecycle(HeaderSyncEvent::PeerDisconnected(peer.clone()));
626 }
627 }
628
629 fn deliver_frame(
630 &self,
631 peer_id: ZakuraPeerId,
632 stream_kind: u16,
633 frame: Frame,
634 ) -> Result<(), SinkReject> {
635 if stream_kind != ZAKURA_STREAM_HEADER_SYNC {
636 return Ok(());
637 }
638
639 match deliver(&self.header_sync, 0, None, peer_id, frame) {
644 Flow::Continue(()) | Flow::Done => Ok(()),
645 Flow::Reject(reject) => Err(reject),
646 }
647 }
648}
649
650fn header_sync_guard_max_bytes() -> u32 {
657 u32::try_from(MAX_HS_MESSAGE_BYTES)
660 .expect("MAX_HS_MESSAGE_BYTES is a 2 MiB constant that fits in u32")
661}
662
663#[derive(Debug)]
665pub(crate) struct HeaderSyncPassthroughService {
666 inner: Arc<dyn Service>,
667}
668
669impl HeaderSyncPassthroughService {
670 pub(crate) fn new(inner: Arc<dyn Service>) -> Self {
671 Self { inner }
672 }
673}
674
675impl Service for HeaderSyncPassthroughService {
676 fn name(&self) -> &'static str {
677 "header-sync-passthrough"
678 }
679
680 fn streams(&self) -> &[Stream] {
681 header_sync_streams()
682 }
683
684 fn wants_peer(
685 &self,
686 peer: &ZakuraPeerId,
687 negotiated: u64,
688 direction: ServicePeerDirection,
689 ) -> bool {
690 self.inner.wants_peer(peer, negotiated, direction)
691 }
692
693 fn add_peer(&self, mut peer: Peer) {
694 let Some((recv, _send)) = peer.take_stream(ZAKURA_STREAM_HEADER_SYNC) else {
695 return;
696 };
697
698 let inner = self.inner.clone();
699 let peer_id = peer.id.clone();
700 let cancel_token = peer.cancel_token();
701
702 task::spawn(async move {
703 let sink = Box::new(HeaderSyncPassthroughSink {
704 peer_id: peer_id.clone(),
705 inner,
706 cancel_token: cancel_token.clone(),
707 });
708
709 match sink.run(recv).await {
710 Ok(()) => {}
711 Err(SinkReject::Protocol(error)) => {
712 tracing::debug!(
713 ?error,
714 ?peer_id,
715 "header-sync passthrough rejected protocol-invalid frame"
716 );
717 cancel_token.cancel();
718 }
719 Err(SinkReject::Local(error)) => {
720 tracing::debug!(
721 ?error,
722 ?peer_id,
723 "header-sync passthrough could not deliver frame locally"
724 );
725 }
726 }
727 });
728 }
729
730 fn remove_peer(&self, _peer: &ZakuraPeerId, _conn_id: ZakuraConnId) {}
731
732 fn deliver_frame(
733 &self,
734 peer_id: ZakuraPeerId,
735 stream_kind: u16,
736 frame: Frame,
737 ) -> Result<(), SinkReject> {
738 self.inner.deliver_frame(peer_id, stream_kind, frame)
739 }
740}
741
742#[derive(Debug)]
743struct HeaderSyncPassthroughSink {
744 peer_id: ZakuraPeerId,
745 inner: Arc<dyn Service>,
746 cancel_token: CancellationToken,
747}
748
749impl Sink for HeaderSyncPassthroughSink {
750 fn run(self: Box<Self>, mut recv: FramedRecv) -> BoxRunFuture<'static, Result<(), SinkReject>> {
751 Box::pin(async move {
752 loop {
753 let frame = tokio::select! {
754 _ = self.cancel_token.cancelled() => return Ok(()),
755 frame = recv.recv() => {
756 let Some(frame) = frame else {
757 return Ok(());
758 };
759 frame
760 }
761 };
762
763 match self.inner.deliver_frame(
764 self.peer_id.clone(),
765 ZAKURA_STREAM_HEADER_SYNC,
766 frame,
767 ) {
768 Ok(()) => {}
769 Err(SinkReject::Protocol(error)) => return Err(SinkReject::Protocol(error)),
770 Err(SinkReject::Local(error)) => {
771 tracing::debug!(
772 ?error,
773 peer_id = ?self.peer_id,
774 "header-sync passthrough could not deliver frame locally"
775 );
776 }
777 }
778 }
779 })
780 }
781}
782
783#[cfg(test)]
784mod request_id_tests {
785 use super::*;
786 use crate::zakura::header_sync::{
787 requester::HeaderRequesterCommand,
788 state::{RangePriority, RangeRequest},
789 };
790
791 #[test]
792 fn request_id_exhaustion_remains_fail_closed() {
793 let (send, _recv) = crate::zakura::framed_channel(1);
794 let peer_id = ZakuraPeerId::new(vec![1; 32]).expect("test peer id is valid");
795 let session = HeaderSyncPeerSession::from_parts_with_direction(
796 peer_id,
797 ServicePeerDirection::Outbound,
798 send,
799 CancellationToken::new(),
800 );
801 session
802 .inner
803 .next_request_id
804 .store(u64::MAX, Ordering::Relaxed);
805
806 assert!(session.next_request_id().is_err());
807 assert!(session.next_request_id().is_err());
808 }
809
810 #[test]
811 fn requester_queue_failure_cancels_prepared_expectation() {
812 let (send, _recv) = crate::zakura::framed_channel(1);
813 let (commands_tx, mut commands_rx) = mpsc::unbounded_channel();
814 let peer_id = ZakuraPeerId::new(vec![5; 32]).expect("test peer id is valid");
815 let session = HeaderSyncPeerSession::from_parts_with_direction_and_commands(
816 peer_id,
817 1,
818 ServicePeerDirection::Outbound,
819 send,
820 CancellationToken::new(),
821 Some(commands_tx),
822 );
823 let range = RangeRequest {
824 start_height: block::Height(1),
825 count: 1,
826 anchor_hash: None,
827 finalized: false,
828 want_tree_aux_roots: true,
829 priority: RangePriority::Forward,
830 };
831 let prepared = session
832 .prepare_get_headers(range.start_height, range.count, range.want_tree_aux_roots)
833 .expect("valid test request is prepared");
834 let reserved = match commands_rx.try_recv().expect("reservation is published") {
835 HeaderSyncPeerCommand::Reserve(expected) => expected,
836 command => panic!("expected reservation command, got {command:?}"),
837 };
838 let (requester_tx, requester_rx) = mpsc::channel(1);
839 drop(requester_rx);
840
841 let rejected = match requester_tx.try_send(HeaderRequesterCommand { range, prepared }) {
842 Err(mpsc::error::TrySendError::Closed(command)) => command,
843 _ => panic!("closed requester queue rejects the prepared command"),
844 };
845 drop(rejected);
846
847 let cancelled = match commands_rx.try_recv().expect("cancellation is published") {
848 HeaderSyncPeerCommand::Cancel(expected) => expected,
849 command => panic!("expected cancellation command, got {command:?}"),
850 };
851 assert_eq!(cancelled, reserved);
852 }
853}