Skip to main content

zakura_network/zakura/header_sync/
service.rs

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    // Advisory until the transport wires Stream::frame_cap end-to-end; the
24    // authoritative inbound cap is app_frame_cap_for_stream_kind. The cast is
25    // safe because both terms are small protocol constants checked against the
26    // local message cap in header_sync::wire.
27    frame_cap: (MAX_HS_MESSAGE_BYTES + FRAME_HEADER_BYTES) as u32,
28    capability: ZAKURA_CAP_HEADER_SYNC,
29    mode: StreamMode::Ordered,
30}];
31
32/// Service-declared streams for native header sync.
33pub(crate) fn header_sync_streams() -> &'static [Stream] {
34    &HEADER_SYNC_SERVICE_STREAMS
35}
36
37/// Cloneable typed header-sync sender and peer-local response expectations.
38#[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    /// Authenticated peer identity for this header-sync session.
137    pub fn peer_id(&self) -> &ZakuraPeerId {
138        &self.peer_id
139    }
140
141    /// Unique ordered-stream generation that owns this session.
142    pub fn session_id(&self) -> u64 {
143        self.session_id
144    }
145
146    /// Direction of the underlying Zakura connection.
147    pub fn direction(&self) -> ServicePeerDirection {
148        self.direction
149    }
150
151    /// Peer disconnect/local shutdown cancellation token.
152    pub fn cancel_token(&self) -> CancellationToken {
153        self.inner.cancel_token.clone()
154    }
155
156    /// Current free slots in this peer's bounded outbound stream queue.
157    pub fn outbound_capacity(&self) -> usize {
158        self.inner.send.capacity()
159    }
160
161    /// Total slots in this peer's bounded outbound stream queue.
162    pub fn outbound_max_capacity(&self) -> usize {
163        self.inner.send.max_capacity()
164    }
165
166    /// Retire a request ID so any late response is dropped without scoring.
167    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    /// Send a typed status advertisement.
201    pub fn try_send_status(&self, status: HeaderSyncStatus) -> Result<(), OrderedSendError> {
202        self.try_send_message(HeaderSyncMessage::Status(status), None)
203    }
204
205    /// Prepare a correlated header request without making its frame visible to the peer.
206    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    /// Send a typed header range response with one advisory body-size hint and
234    /// tree-aux root payload per header.
235    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    /// Send a typed full tip block announcement.
253    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    /// Wait for outbound capacity and publish the prepared frame.
286    ///
287    /// Dropping this future before publication synchronously cancels the pipe's
288    /// response reservation.
289    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/// Commands from shared scheduling state into one peer-owned header-sync pipe.
343#[derive(Debug)]
344pub(super) enum HeaderSyncPeerCommand {
345    /// Reserve an expected `Headers` response before `GetHeaders` is queued.
346    Reserve(ExpectedHeadersResponse),
347    /// Roll back an expectation when `GetHeaders` could not be queued.
348    Cancel(ExpectedHeadersResponse),
349    /// Retire an expected `Headers` response after timeout or cancellation.
350    Retire(HeaderSyncRequestId),
351}
352
353/// Pump actor actions that can be satisfied at the transport/service seam.
354pub(crate) async fn drive_header_sync_actions(
355    mut actions: mpsc::Receiver<HeaderSyncAction>,
356    handle: HeaderSyncHandle,
357    // Retained so the disconnect capability stays wired into the driver, even
358    // though peer scoring no longer drives disconnects (misbehavior is record-only).
359    _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                // Record-only: peer scoring no longer drives disconnects.
378                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/// Native versioned header-sync service.
429#[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        // Escalation is a local-room check. First-party summary usefulness is
467        // advisory and is applied by header-sync candidate selection upstream.
468        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        // The sink loop parks on the service token (a child of the connection
491        // token) exactly as the old `HeaderSyncSink::run` select did. The
492        // connection token is cancelled only on a protocol reject below, never on
493        // a normal/parked exit — parking one service must not tear down the
494        // shared connection that other services (discovery, block-sync) ride on.
495        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        // Phase 2 keeps request/response correlation in `HsLocal`: after the
538        // session queues an outbound `GetHeaders`, the peer-owned pipe records
539        // the expected `Headers` response in plain local state.
540        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        // The pipe future reproduces the old sink's connection handling: a
549        // protocol reject (the only way `run_peer` returns `Err`, since
550        // `run_inbound` maps a closed-queue `Local` to a benign continue)
551        // cancels the *connection*, matching the old
552        // `connection_cancel_token.cancel()` on `SinkReject::Protocol`. A normal
553        // or parked exit leaves the connection alone.
554        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        // The supervised teardown runs on every exit path — normal return,
567        // protocol reject, or panic. It cancels this peer's *service* token
568        // (idempotent; already cancelled on a park/protocol exit) and sends
569        // `PeerDisconnected`. Sending it from teardown is the latent-bug fix: the
570        // old sink only sent `PeerDisconnected` on the normal return path, so a
571        // panicking task leaked the peer's reactor state.
572        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        // Reuse the single supervised launcher; let the returned handle drop to
602        // detach the task (the `PipeTeardown` still runs on every exit path).
603        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        // The test/recorder path has no peer session, so a `Headers` response
640        // with no outstanding request is rejected as `UnsolicitedHeaders`. A
641        // `Local` reject (closed reactor queue) is surfaced to the registry
642        // exactly as the old `deliver_header_sync_frame` returned it.
643        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
650/// Service-level oversize cap for the header-sync guard.
651///
652/// Matches the decode stage's `MAX_HS_MESSAGE_BYTES` threshold so the guard
653/// rejects nothing the decode stage would have admitted; the transport already
654/// caps frames at this payload size before they reach the service, so this is a
655/// defense-in-depth bound that never changes which events fire.
656fn header_sync_guard_max_bytes() -> u32 {
657    // `MAX_HS_MESSAGE_BYTES` is a 2 MiB protocol constant that fits in `u32`;
658    // the `const` assertion in `wire.rs` keeps it below the local message cap.
659    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/// Testkit/no-reactor mode records header-sync inbound frames without running header sync.
664#[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}