Skip to main content

zakura_network/zakura/discovery/
service.rs

1//! Native discovery service (stream kind 4) on the Zakura transport.
2//!
3//! Discovery is a single long-lived ordered stream per peer. Each side runs a
4//! [`DiscoverySink`] (the reader, which imports peer records and answers
5//! `GetPeers`) and a [`DiscoverySource`] (the writer, which periodically gossips
6//! the local self-record and asks for more peers). The wire format is the
7//! [`DiscoveryMessage`] payload carried inside a generic transport [`Frame`]
8//! (`message_type = DISCOVERY_FRAME_MESSAGE_TYPE`, `flags = 0`), identical to the
9//! original native-discovery wire so peers interoperate.
10
11use std::{
12    sync::{
13        atomic::{AtomicBool, Ordering},
14        Arc,
15    },
16    time::Duration,
17};
18
19use iroh::NodeId;
20use tokio::sync::Notify;
21use tokio_util::sync::CancellationToken;
22
23use crate::zakura::{
24    handle_pipe_exit, spawn_supervised_peer_task, spawn_supervised_pipe, BlockSyncHandle,
25    CloseCause, Flow, Frame, FramedRecv, FramedSend, HeaderSyncEvent, HeaderSyncHandle,
26    OrderedSendError, Peer, PeerStreamSession, Pipe, Service, ServiceAdmissionDecision,
27    ServicePeerDirection, SinkReject, Stream, StreamMode, ZakuraConnId, ZakuraPeerId,
28    LOCAL_MAX_CONTROL_FRAME_BYTES, ZAKURA_CAP_DISCOVERY, ZAKURA_CAP_HEADER_SYNC,
29};
30
31#[cfg(test)]
32use super::pipe::decode_discovery_frame;
33use super::pipe::{discovery_pipe, DsEnv, DsLocal, DISCOVERY_FRAME_MESSAGE_TYPE};
34use super::protocol::{
35    BlockSyncServiceSummary, DiscoveryBookError, DiscoveryMessage, DiscoveryRecordError,
36    GetServices, HeaderSyncServiceSummary, ServiceSummaryEnvelope, Services, ZakuraDiscoveryHandle,
37    ZakuraNodeRecord, ZakuraServiceId, MAX_DISCOVERY_RECORDS_PER_RESPONSE,
38    ZAKURA_DISCOVERY_STREAM_VERSION, ZAKURA_STREAM_DISCOVERY,
39};
40
41/// Maximum time discovery waits for first-party exchange responses before releasing the session.
42const DISCOVERY_EXCHANGE_SETTLE_TIMEOUT: Duration = Duration::from_secs(2);
43
44const DISCOVERY_SERVICE_STREAMS: [Stream; 1] = [Stream {
45    kind: ZAKURA_STREAM_DISCOVERY,
46    version: ZAKURA_DISCOVERY_STREAM_VERSION,
47    // Advisory until the transport wires Stream::frame_cap end-to-end; the
48    // authoritative inbound cap is app_frame_cap_for_stream_kind.
49    frame_cap: LOCAL_MAX_CONTROL_FRAME_BYTES,
50    capability: ZAKURA_CAP_DISCOVERY,
51    mode: StreamMode::Ordered,
52}];
53
54/// Service-declared streams for native discovery.
55pub(crate) fn discovery_streams() -> &'static [Stream] {
56    &DISCOVERY_SERVICE_STREAMS
57}
58
59/// Cloneable typed sender for one native discovery ordered stream.
60#[derive(Clone, Debug)]
61pub struct DiscoveryPeerSession {
62    peer_id: ZakuraPeerId,
63    direction: ServicePeerDirection,
64    send: FramedSend,
65    cancel: CancellationToken,
66}
67
68impl DiscoveryPeerSession {
69    fn new(session: &PeerStreamSession, direction: ServicePeerDirection) -> Self {
70        Self {
71            peer_id: session.peer_id().clone(),
72            direction,
73            send: session.sender(),
74            cancel: session.cancel_token(),
75        }
76    }
77
78    /// Authenticated peer identity for this discovery stream.
79    pub fn peer_id(&self) -> &ZakuraPeerId {
80        &self.peer_id
81    }
82
83    /// Direction of the underlying Zakura connection.
84    pub fn direction(&self) -> ServicePeerDirection {
85        self.direction
86    }
87
88    /// Peer disconnect/local shutdown cancellation token.
89    pub fn cancel_token(&self) -> CancellationToken {
90        self.cancel.clone()
91    }
92
93    /// Send this node's signed self-record.
94    pub fn try_send_hello(&self, record: ZakuraNodeRecord) -> Result<(), OrderedSendError> {
95        self.try_send_message(DiscoveryMessage::Hello { record })
96    }
97
98    /// Ask this peer for more peer records.
99    pub fn try_send_get_peers(
100        &self,
101        limit: u16,
102        wanted_services: Vec<ZakuraServiceId>,
103        exclude_node_ids: Vec<NodeId>,
104    ) -> Result<(), OrderedSendError> {
105        self.try_send_message(DiscoveryMessage::GetPeers {
106            limit,
107            wanted_services,
108            exclude_node_ids,
109        })
110    }
111
112    /// Send peer records to this peer.
113    pub fn try_send_peers(&self, records: Vec<ZakuraNodeRecord>) -> Result<(), OrderedSendError> {
114        self.try_send_message(DiscoveryMessage::Peers { records })
115    }
116
117    /// Ask this peer for its own live service summaries.
118    pub fn try_send_get_services(
119        &self,
120        wanted_services: Vec<ZakuraServiceId>,
121    ) -> Result<(), OrderedSendError> {
122        self.try_send_message(DiscoveryMessage::GetServices(GetServices {
123            wanted_services,
124        }))
125    }
126
127    /// Send this node's first-party live service summaries.
128    pub fn try_send_services(&self, services: Services) -> Result<(), OrderedSendError> {
129        self.try_send_message(DiscoveryMessage::Services(services))
130    }
131
132    fn try_send_message(&self, message: DiscoveryMessage) -> Result<(), OrderedSendError> {
133        let payload = message
134            .encode()
135            .map_err(|error| OrderedSendError::Encode(Box::new(error)))?;
136        match self.send.try_send(Frame {
137            message_type: DISCOVERY_FRAME_MESSAGE_TYPE,
138            flags: 0,
139            payload,
140        }) {
141            Ok(()) => Ok(()),
142            Err(tokio::sync::mpsc::error::TrySendError::Full(_frame)) => {
143                Err(OrderedSendError::Full)
144            }
145            Err(tokio::sync::mpsc::error::TrySendError::Closed(_frame)) => {
146                Err(OrderedSendError::Closed)
147            }
148        }
149    }
150}
151
152/// Native discovery service backed by a [`ZakuraDiscoveryHandle`] runtime.
153#[derive(Clone, Debug)]
154pub struct DiscoveryService {
155    handle: ZakuraDiscoveryHandle,
156    header_sync: Option<HeaderSyncHandle>,
157    block_sync: Option<BlockSyncHandle>,
158}
159
160impl DiscoveryService {
161    /// Builds a discovery service driven by `handle`.
162    pub fn new(handle: ZakuraDiscoveryHandle) -> Self {
163        Self {
164            handle,
165            header_sync: None,
166            block_sync: None,
167        }
168    }
169
170    /// Builds a discovery service with header-sync and block-sync summary providers.
171    pub(crate) fn with_sync_services(
172        handle: ZakuraDiscoveryHandle,
173        header_sync: HeaderSyncHandle,
174        block_sync: Option<BlockSyncHandle>,
175    ) -> Self {
176        Self {
177            handle,
178            header_sync: Some(header_sync),
179            block_sync,
180        }
181    }
182
183    /// Returns the underlying discovery runtime handle.
184    pub fn handle(&self) -> &ZakuraDiscoveryHandle {
185        &self.handle
186    }
187}
188
189impl Service for DiscoveryService {
190    fn name(&self) -> &'static str {
191        "discovery"
192    }
193
194    fn streams(&self) -> &[Stream] {
195        discovery_streams()
196    }
197
198    fn wants_peer(
199        &self,
200        _peer: &ZakuraPeerId,
201        _negotiated: u64,
202        direction: ServicePeerDirection,
203    ) -> bool {
204        // Discovery escalation only checks this reactor's local room; live
205        // summaries are first-party advisory data imported by the runtime.
206        let snapshot = self.handle.peer_snapshot();
207        match direction {
208            ServicePeerDirection::Inbound => snapshot.inbound_slots_free > 0,
209            ServicePeerDirection::Outbound => snapshot.outbound_slots_free > 0,
210        }
211    }
212
213    fn add_peer(&self, mut peer: Peer) {
214        let Some((recv, send)) = peer.take_stream(ZAKURA_STREAM_DISCOVERY) else {
215            return;
216        };
217        let Some(peer_node_id) = node_id_from_peer_id(&peer.id) else {
218            // A peer id that is not a 32-byte node id cannot be a discovery
219            // author; drop the stream without registering an exchange.
220            return;
221        };
222        let session = PeerStreamSession::new(
223            peer.id.clone(),
224            ZAKURA_STREAM_DISCOVERY,
225            recv,
226            send,
227            peer.service_cancel_token(),
228        );
229        let discovery_session = DiscoveryPeerSession::new(&session, peer.direction);
230        let conn_id = peer.conn_id;
231        let service_cancel = discovery_session.cancel_token();
232        let connection_cancel = peer.cancel_token();
233        let close_cause = peer.close_cause();
234        let other_service_negotiated = has_other_negotiated_service(peer.negotiated);
235        let (_peer_id, _stream_kind, recv, _send, _session_cancel) = session.into_parts();
236
237        let handle = self.handle.clone();
238        let header_sync = self.header_sync.clone();
239        let block_sync = self.block_sync.clone();
240        // SR-1: a panic in the admission task (before it hands off to the
241        // exchange) must still disconnect this one peer and cancel its discovery
242        // session instead of leaving admitted state behind a half-live
243        // connection. Normal/parked exits cancel `service_cancel` inline below;
244        // `on_panic` covers the unwind path only.
245        let admit_peer_id = discovery_session.peer_id().clone();
246        let panic_service_cancel = service_cancel.clone();
247        let panic_connection_cancel = connection_cancel.clone();
248        let panic_close_cause = close_cause.clone();
249        spawn_supervised_peer_task(
250            admit_peer_id,
251            || {},
252            move || {
253                panic_close_cause.record("service_panic");
254                panic_service_cancel.cancel();
255                panic_connection_cancel.cancel();
256            },
257            async move {
258                let decision = handle
259                    .admit_peer(
260                        conn_id,
261                        discovery_session.peer_id().clone(),
262                        discovery_session.direction(),
263                    )
264                    .await;
265                if decision != ServiceAdmissionDecision::Admit {
266                    metrics::counter!("zakura.discovery.peer.parked").increment(1);
267                    tracing::info!(
268                        peer = ?discovery_session.peer_id(),
269                        direction = ?discovery_session.direction(),
270                        ?decision,
271                        "locally parking Zakura discovery service session"
272                    );
273                    service_cancel.cancel();
274                    return;
275                }
276
277                spawn_discovery_exchange(DiscoveryExchangeStart {
278                    handle,
279                    header_sync,
280                    block_sync,
281                    peer_node_id,
282                    discovery_session,
283                    conn_id,
284                    recv,
285                    service_cancel,
286                    connection_cancel,
287                    close_cause,
288                    other_service_negotiated,
289                });
290            },
291        );
292    }
293
294    fn remove_peer(&self, peer: &ZakuraPeerId, conn_id: ZakuraConnId) {
295        let handle = self.handle.clone();
296        let peer = peer.clone();
297        tokio::spawn(async move {
298            handle.remove_peer(&peer, conn_id).await;
299        });
300    }
301}
302
303struct DiscoveryExchangeStart {
304    handle: ZakuraDiscoveryHandle,
305    header_sync: Option<HeaderSyncHandle>,
306    block_sync: Option<BlockSyncHandle>,
307    peer_node_id: NodeId,
308    discovery_session: DiscoveryPeerSession,
309    conn_id: ZakuraConnId,
310    recv: FramedRecv,
311    service_cancel: CancellationToken,
312    connection_cancel: CancellationToken,
313    close_cause: CloseCause,
314    other_service_negotiated: bool,
315}
316
317fn spawn_discovery_exchange(start: DiscoveryExchangeStart) {
318    let DiscoveryExchangeStart {
319        handle,
320        header_sync,
321        block_sync,
322        peer_node_id,
323        discovery_session,
324        conn_id,
325        recv,
326        service_cancel,
327        connection_cancel,
328        close_cause,
329        other_service_negotiated,
330    } = start;
331    let peer_id = discovery_session.peer_id().clone();
332    let progress = Arc::new(DiscoveryExchangeProgress::default());
333    let source_header_sync = header_sync.clone();
334    let sink = DiscoverySink {
335        handle: handle.clone(),
336        header_sync,
337        block_sync,
338        peer_node_id,
339        session: discovery_session.clone(),
340        progress: progress.clone(),
341    };
342    let sink_service_cancel = service_cancel.clone();
343    let reject_connection_cancel = connection_cancel.clone();
344    let panic_connection_cancel = connection_cancel.clone();
345    let reject_close_cause = close_cause.clone();
346    let panic_close_cause = close_cause.clone();
347    let sink_peer_id = peer_id.clone();
348    // A protocol reject is fatal to the connection; normal/parked exits leave it
349    // for the source task to tear down once it knows no other service owns the
350    // peer (below). Panic teardown is in `on_panic`.
351    let pipe = async move {
352        let mut pipe = discovery_pipe(sink_peer_id);
353        handle_pipe_exit(
354            "discovery",
355            &reject_connection_cancel,
356            &reject_close_cause,
357            run_discovery_pipe(&mut pipe, recv, sink).await,
358        );
359    };
360    let on_panic = move || {
361        panic_close_cause.record("service_panic");
362        panic_connection_cancel.cancel();
363    };
364    // Let the returned handle drop to detach the supervised reader task; the
365    // `PipeTeardown` still runs on every exit path.
366    spawn_supervised_pipe(peer_id.clone(), sink_service_cancel, || {}, on_panic, pipe);
367
368    let source = DiscoverySource {
369        handle: handle.clone(),
370        session: discovery_session,
371        progress,
372    };
373    // SR-1: a panic in the source task skips its `service_cancel.cancel()`,
374    // `handle.remove_peer()`, and discovery-only connection cancellation,
375    // leaving admitted discovery state behind a half-live connection. On the
376    // unwind path, disconnect this one peer; the connection teardown then drives
377    // the async `remove_peer` through the registry. Normal exits run the inline
378    // cleanup below, so `on_panic` is the panic-only path.
379    let source_task_peer_id = peer_id.clone();
380    let panic_source_service_cancel = service_cancel.clone();
381    let panic_source_connection_cancel = connection_cancel.clone();
382    let panic_source_close_cause = close_cause.clone();
383    let source_close_cause = close_cause.clone();
384    spawn_supervised_peer_task(
385        source_task_peer_id,
386        || {},
387        move || {
388            panic_source_close_cause.record("service_panic");
389            panic_source_service_cancel.cancel();
390            panic_source_connection_cancel.cancel();
391        },
392        async move {
393            let exchanged = source.run().await;
394            if exchanged {
395                handle.mark_short_lived_exchange(&peer_node_id).await;
396            }
397            service_cancel.cancel();
398            handle.remove_peer(&peer_id, conn_id).await;
399            if exchanged
400                && !peer_has_other_service_owner(
401                    source_header_sync.as_ref(),
402                    peer_node_id,
403                    other_service_negotiated,
404                )
405            {
406                source_close_cause.record("discovery_exchange_complete");
407                connection_cancel.cancel();
408            }
409        },
410    );
411}
412
413/// Reader half of the discovery stream: imports peer records and answers queries.
414struct DiscoverySink {
415    handle: ZakuraDiscoveryHandle,
416    header_sync: Option<HeaderSyncHandle>,
417    block_sync: Option<BlockSyncHandle>,
418    peer_node_id: NodeId,
419    session: DiscoveryPeerSession,
420    progress: Arc<DiscoveryExchangeProgress>,
421}
422
423async fn run_discovery_pipe(
424    pipe: &mut Pipe<DsLocal, DsEnv>,
425    mut recv: FramedRecv,
426    sink: DiscoverySink,
427) -> Result<(), SinkReject> {
428    let cancel = sink.session.cancel_token();
429    loop {
430        let frame = tokio::select! {
431            biased;
432            _ = cancel.cancelled() => return Ok(()),
433            frame = recv.recv() => frame,
434        };
435        let Some(frame) = frame else {
436            return Ok(());
437        };
438
439        match pipe.run_one(frame) {
440            Flow::Continue(()) | Flow::Done => {}
441            Flow::Reject(reject) => return Err(reject),
442        }
443
444        let Some(message) = pipe.local_mut().take_decoded() else {
445            continue;
446        };
447        sink.handle_message(message).await?;
448    }
449}
450
451impl DiscoverySink {
452    async fn handle_message(&self, message: DiscoveryMessage) -> Result<(), SinkReject> {
453        match message {
454            DiscoveryMessage::Hello { record } => self.handle_hello(record).await,
455            DiscoveryMessage::GetPeers {
456                limit,
457                wanted_services,
458                exclude_node_ids,
459            } => {
460                let records = self
461                    .handle
462                    .sample_peers(usize::from(limit), &wanted_services, &exclude_node_ids)
463                    .await;
464                self.send_peers(records)
465            }
466            DiscoveryMessage::Peers { records } => {
467                self.handle
468                    .import_peer_records(records, Some(self.peer_node_id))
469                    .await;
470                self.progress.mark_peers();
471                Ok(())
472            }
473            DiscoveryMessage::GetServices(query) => {
474                let services = self.local_services_response(query).await?;
475                self.send_services(services)
476            }
477            DiscoveryMessage::Services(services) => self.handle_services(services).await,
478        }
479    }
480
481    async fn local_services_response(&self, query: GetServices) -> Result<Services, SinkReject> {
482        let mut summaries = Vec::new();
483
484        if service_wanted(&query.wanted_services, &ZakuraServiceId::header_sync()) {
485            if let Some(header_sync) = &self.header_sync {
486                let (best_height, best_hash) = header_sync.best_header_tip();
487                let summary = HeaderSyncServiceSummary::from_snapshot(
488                    best_height,
489                    best_hash,
490                    None,
491                    true,
492                    header_sync.peer_snapshot(),
493                );
494                summaries.push(
495                    ServiceSummaryEnvelope::header_sync(&summary).map_err(SinkReject::local)?,
496                );
497            }
498        }
499
500        if service_wanted(&query.wanted_services, &ZakuraServiceId::discovery()) {
501            let summary = self.handle.local_discovery_summary().await;
502            summaries.push(ServiceSummaryEnvelope::discovery(&summary).map_err(SinkReject::local)?);
503        }
504
505        if service_wanted(&query.wanted_services, &ZakuraServiceId::block_sync()) {
506            if let Some(block_sync) = &self.block_sync {
507                let summary = BlockSyncServiceSummary::from_status_and_snapshot(
508                    block_sync.local_status(),
509                    block_sync.peer_snapshot(),
510                );
511                summaries
512                    .push(ServiceSummaryEnvelope::block_sync(&summary).map_err(SinkReject::local)?);
513            }
514        }
515
516        Ok(self.handle.local_services_response(summaries))
517    }
518
519    async fn handle_services(&self, services: Services) -> Result<(), SinkReject> {
520        if services.node_id != self.peer_node_id {
521            return Err(SinkReject::protocol(
522                "Zakura discovery SERVICES authored by a different node id",
523            ));
524        }
525
526        let header_summaries =
527            decode_header_sync_summaries(&services).map_err(SinkReject::protocol)?;
528        self.handle
529            .import_connected_peer_services(services, self.peer_node_id)
530            .await
531            .map_err(SinkReject::protocol)?;
532        self.progress.mark_services();
533
534        if let Some(header_sync) = &self.header_sync {
535            for summary in header_summaries {
536                if let Err(error) = header_sync
537                    .send(HeaderSyncEvent::AdvisoryHeaderSummary {
538                        peer: self.session.peer_id().clone(),
539                        summary,
540                    })
541                    .await
542                {
543                    tracing::debug!(
544                        ?error,
545                        peer = ?self.session.peer_id(),
546                        "failed to queue first-party Zakura header-sync advisory summary"
547                    );
548                    break;
549                }
550            }
551        }
552
553        Ok(())
554    }
555
556    async fn handle_hello(&self, record: ZakuraNodeRecord) -> Result<(), SinkReject> {
557        if record.body.node_id != self.peer_node_id {
558            return Err(SinkReject::protocol(
559                "Zakura discovery hello authored by a different node id",
560            ));
561        }
562        match self
563            .handle
564            .import_connected_peer_record(record, self.peer_node_id)
565            .await
566        {
567            Ok(_) => Ok(()),
568            Err(error) if is_advisory_self_record_import_error(&error) => {
569                tracing::debug!(?error, "ignoring advisory discovery hello import error");
570                Ok(())
571            }
572            Err(error) => Err(SinkReject::protocol(error)),
573        }?;
574        self.progress.mark_hello();
575        Ok(())
576    }
577
578    fn send_peers(&self, records: Vec<ZakuraNodeRecord>) -> Result<(), SinkReject> {
579        match self.session.try_send_peers(records) {
580            Ok(()) | Err(OrderedSendError::Full) => Ok(()),
581            Err(OrderedSendError::Closed) => {
582                Err(SinkReject::local("Zakura discovery send channel closed"))
583            }
584            Err(OrderedSendError::Encode(error)) => Err(SinkReject::local(error)),
585        }
586    }
587
588    fn send_services(&self, services: Services) -> Result<(), SinkReject> {
589        match self.session.try_send_services(services) {
590            Ok(()) | Err(OrderedSendError::Full) => Ok(()),
591            Err(OrderedSendError::Closed) => {
592                Err(SinkReject::local("Zakura discovery send channel closed"))
593            }
594            Err(OrderedSendError::Encode(error)) => Err(SinkReject::local(error)),
595        }
596    }
597}
598
599fn service_wanted(wanted_services: &[ZakuraServiceId], service_id: &ZakuraServiceId) -> bool {
600    wanted_services.is_empty() || wanted_services.iter().any(|wanted| wanted == service_id)
601}
602
603fn decode_header_sync_summaries(
604    services: &Services,
605) -> Result<Vec<HeaderSyncServiceSummary>, crate::BoxError> {
606    let mut summaries = Vec::new();
607    for envelope in &services.summaries {
608        if let Some(summary) = envelope.decode_header_sync()? {
609            summaries.push(summary);
610        }
611    }
612    Ok(summaries)
613}
614
615/// Writer half of the discovery stream: periodic self-record gossip + peer asks.
616struct DiscoverySource {
617    handle: ZakuraDiscoveryHandle,
618    session: DiscoveryPeerSession,
619    progress: Arc<DiscoveryExchangeProgress>,
620}
621
622impl DiscoverySource {
623    async fn run(self) -> bool {
624        if self.exchange().await.is_err() {
625            return false;
626        }
627        let cancel = self.session.cancel_token();
628        tokio::select! {
629            biased;
630            _ = cancel.cancelled() => {}
631            _ = self.progress.wait_complete() => {}
632            _ = tokio::time::sleep(DISCOVERY_EXCHANGE_SETTLE_TIMEOUT) => {}
633        }
634        true
635    }
636
637    /// Gossips the current self-record and asks the peer for records and services.
638    ///
639    /// Returns `Err(())` once the stream's send side is gone, so the caller
640    /// stops the periodic loop.
641    async fn exchange(&self) -> Result<(), ()> {
642        let record = (*self.handle.current_self_record()).clone();
643        self.handle_send_result(self.session.try_send_hello(record))?;
644
645        let limit = self
646            .handle
647            .peer_sample_limit()
648            .await
649            .min(MAX_DISCOVERY_RECORDS_PER_RESPONSE);
650        // `peer_sample_limit` is bounded by MAX_DISCOVERY_RECORDS_PER_RESPONSE
651        // (<= u16::MAX), so the cast cannot truncate.
652        let exclude_node_ids = self.handle.peer_sample_exclusions().await;
653        self.handle_send_result(self.session.try_send_get_peers(
654            limit as u16,
655            Vec::new(),
656            exclude_node_ids,
657        ))?;
658
659        self.handle_send_result(self.session.try_send_get_services(Vec::new()))
660    }
661
662    fn handle_send_result(&self, result: Result<(), OrderedSendError>) -> Result<(), ()> {
663        match result {
664            Ok(()) | Err(OrderedSendError::Full) => Ok(()),
665            Err(OrderedSendError::Closed) => Err(()),
666            Err(OrderedSendError::Encode(error)) => {
667                tracing::debug!(
668                    ?error,
669                    peer = ?self.session.peer_id(),
670                    "failed to encode Zakura discovery message"
671                );
672                Ok(())
673            }
674        }
675    }
676}
677
678#[derive(Default)]
679struct DiscoveryExchangeProgress {
680    hello: AtomicBool,
681    peers: AtomicBool,
682    services: AtomicBool,
683    notify: Notify,
684}
685
686impl DiscoveryExchangeProgress {
687    fn mark_hello(&self) {
688        self.hello.store(true, Ordering::Relaxed);
689        self.notify.notify_waiters();
690    }
691
692    fn mark_peers(&self) {
693        self.peers.store(true, Ordering::Relaxed);
694        self.notify.notify_waiters();
695    }
696
697    fn mark_services(&self) {
698        self.services.store(true, Ordering::Relaxed);
699        self.notify.notify_waiters();
700    }
701
702    fn complete(&self) -> bool {
703        self.hello.load(Ordering::Relaxed)
704            && self.peers.load(Ordering::Relaxed)
705            && self.services.load(Ordering::Relaxed)
706    }
707
708    async fn wait_complete(&self) {
709        while !self.complete() {
710            self.notify.notified().await;
711        }
712    }
713}
714
715fn peer_has_other_service_owner(
716    header_sync: Option<&HeaderSyncHandle>,
717    peer_node_id: NodeId,
718    other_service_negotiated: bool,
719) -> bool {
720    if other_service_negotiated {
721        return true;
722    }
723
724    header_sync.is_some_and(|header_sync| {
725        header_sync
726            .candidate_state()
727            .admitted_node_ids
728            .contains(&peer_node_id)
729    })
730}
731
732fn has_other_negotiated_service(negotiated: u64) -> bool {
733    negotiated & !(ZAKURA_CAP_DISCOVERY | ZAKURA_CAP_HEADER_SYNC) != 0
734}
735
736/// Returns the iroh node id encoded by a discovery peer id, if it is a 32-byte
737/// node id.
738fn node_id_from_peer_id(peer_id: &ZakuraPeerId) -> Option<NodeId> {
739    let bytes: [u8; 32] = peer_id.as_bytes().try_into().ok()?;
740    NodeId::from_bytes(&bytes).ok()
741}
742
743/// A peer-hello import error that should be logged and ignored rather than
744/// closing the live connection. These mean the peer's record is not locally
745/// dialable or has drifted out of the freshness window, neither of which is the
746/// connected peer's fault.
747fn is_advisory_self_record_import_error(error: &DiscoveryBookError) -> bool {
748    matches!(
749        error,
750        DiscoveryBookError::NoUsableDirectAddress
751            | DiscoveryBookError::NonDialableDirectAddress { .. }
752            | DiscoveryBookError::Record(DiscoveryRecordError::Expired)
753            | DiscoveryBookError::Record(DiscoveryRecordError::FarFutureExpiry)
754    )
755}
756
757#[cfg(test)]
758mod tests {
759    use std::{
760        collections::HashMap,
761        net::{IpAddr, Ipv4Addr, SocketAddr},
762        time::{Duration, SystemTime, UNIX_EPOCH},
763    };
764
765    use iroh::SecretKey;
766    use tokio::{sync::watch, task::JoinHandle};
767
768    use super::*;
769    use crate::zakura::discovery::protocol::{
770        DiscoveryServiceSummary, ZakuraLiveServiceSummary, ZakuraNodeRecordBody,
771        SUMMARY_TAG_HEADER_SYNC_RETIRED,
772    };
773    use crate::zakura::{
774        framed_channel, spawn_block_sync_reactor, spawn_header_sync_reactor, BlockSyncFrontiers,
775        BlockSyncStartup, HeaderSyncAction, HeaderSyncFrontiers, HeaderSyncMessage,
776        HeaderSyncPeerSession, HeaderSyncStartup, HeaderSyncStatus, ServicePeerLimits,
777        ZakuraBlockSyncConfig, ZakuraDiscoveryConfig, ZakuraDiscoveryLocalConfig,
778        ZakuraHandshakeConfig, ZakuraHeaderSyncConfig, LOCAL_MAX_MESSAGE_BYTES,
779        MAX_BS_RESPONSE_BYTES, ZAKURA_CAP_BLOCK_SYNC, ZAKURA_CAP_DISCOVERY, ZAKURA_CAP_HEADER_SYNC,
780    };
781    use zakura_chain::{block, parameters::Network};
782
783    #[test]
784    fn discovery_and_header_sync_do_not_count_as_another_service_owner() {
785        assert!(!has_other_negotiated_service(ZAKURA_CAP_DISCOVERY));
786        assert!(!has_other_negotiated_service(
787            ZAKURA_CAP_DISCOVERY | ZAKURA_CAP_HEADER_SYNC
788        ));
789        assert!(has_other_negotiated_service(
790            ZAKURA_CAP_DISCOVERY | ZAKURA_CAP_HEADER_SYNC | ZAKURA_CAP_BLOCK_SYNC
791        ));
792    }
793
794    struct HeaderAdvisoryFixture {
795        discovery_handle: ZakuraDiscoveryHandle,
796        header_sync: HeaderSyncHandle,
797        header_actions: tokio::sync::mpsc::Receiver<HeaderSyncAction>,
798        header_task: JoinHandle<()>,
799        peer_node_id: NodeId,
800        peer_id: ZakuraPeerId,
801        peer_send: FramedSend,
802        _peer_recv: FramedRecv,
803    }
804
805    impl Drop for HeaderAdvisoryFixture {
806        fn drop(&mut self) {
807            self.header_task.abort();
808        }
809    }
810
811    fn current_test_unix_secs() -> u64 {
812        SystemTime::now()
813            .duration_since(UNIX_EPOCH)
814            .expect("system clock is after Unix epoch")
815            .as_secs()
816    }
817
818    fn header_summary(best_height: block::Height) -> HeaderSyncServiceSummary {
819        HeaderSyncServiceSummary {
820            best_height,
821            best_hash: block::Hash([7; 32]),
822            finalized_height: None,
823            serving_headers: true,
824            inbound_slots_free: 1,
825            inbound_slots_max: 1,
826            outbound_slots_free: 1,
827            outbound_slots_max: 1,
828        }
829    }
830
831    fn spawn_test_header_sync() -> Result<
832        (
833            HeaderSyncHandle,
834            tokio::sync::mpsc::Receiver<HeaderSyncAction>,
835            JoinHandle<()>,
836        ),
837        crate::BoxError,
838    > {
839        let network = Network::new_regtest(Default::default());
840        let anchor = (block::Height(0), network.genesis_hash());
841        let mut startup = HeaderSyncStartup::new(
842            network,
843            anchor,
844            HeaderSyncFrontiers {
845                finalized_height: anchor.0,
846                verified_block_tip: anchor.0,
847                verified_block_hash: anchor.1,
848            },
849            Some(anchor),
850            ZakuraHeaderSyncConfig::default(),
851            LOCAL_MAX_MESSAGE_BYTES,
852        );
853        startup.range_state_actions_enabled = true;
854        spawn_header_sync_reactor(startup).map_err(Into::into)
855    }
856
857    fn signed_header_sync_record(
858        secret_key: &SecretKey,
859        handshake: &ZakuraHandshakeConfig,
860    ) -> Result<ZakuraNodeRecord, crate::BoxError> {
861        let body = ZakuraNodeRecordBody {
862            node_id: secret_key.public(),
863            direct_addrs: vec![SocketAddr::new(
864                IpAddr::V4(Ipv4Addr::new(192, 0, 2, 44)),
865                8233,
866            )],
867            services: vec![ZakuraServiceId::header_sync()],
868            zakura_protocol_min: handshake.zakura_protocol_min,
869            zakura_protocol_max: handshake.zakura_protocol_max,
870            network_id: handshake.network_id,
871            chain_id: handshake.chain_id,
872            sequence: 1,
873            expires_at_unix_secs: current_test_unix_secs().saturating_add(60),
874        };
875        Ok(ZakuraNodeRecord::sign(body, secret_key)?)
876    }
877
878    fn spawn_header_advisory_fixture(
879        peer_seed: u8,
880    ) -> Result<HeaderAdvisoryFixture, crate::BoxError> {
881        let (connected_tx, connected_rx) = watch::channel(Vec::new());
882        let handshake = ZakuraHandshakeConfig::for_network(&Network::Mainnet);
883        let local_secret = SecretKey::from_bytes(&[31u8; 32]);
884        let discovery_handle = ZakuraDiscoveryHandle::new(
885            ZakuraDiscoveryLocalConfig {
886                secret_key: local_secret,
887                direct_addrs: Vec::new(),
888                services: vec![ZakuraServiceId::discovery()],
889                zakura_protocol_min: handshake.zakura_protocol_min,
890                zakura_protocol_max: handshake.zakura_protocol_max,
891                network_id: handshake.network_id,
892                chain_id: handshake.chain_id,
893                last_authored_sequence: None,
894            },
895            ZakuraDiscoveryConfig::default(),
896            connected_rx,
897        )?;
898        let (header_sync, header_actions, header_task) = spawn_test_header_sync()?;
899        let service = DiscoveryService::with_sync_services(
900            discovery_handle.clone(),
901            header_sync.clone(),
902            None,
903        );
904        let peer_node_id = SecretKey::from_bytes(&[peer_seed; 32]).public();
905        let peer_id = ZakuraPeerId::new(peer_node_id.as_bytes().to_vec())?;
906        connected_tx.send_replace(vec![peer_id.clone()]);
907
908        let (peer_send, service_recv) = framed_channel(8);
909        let (service_send, peer_recv) = framed_channel(8);
910        let streams = HashMap::from([(ZAKURA_STREAM_DISCOVERY, (service_recv, service_send))]);
911
912        service.add_peer(Peer::new(
913            peer_id.clone(),
914            None,
915            ZAKURA_CAP_DISCOVERY,
916            streams,
917            CancellationToken::new(),
918        ));
919
920        Ok(HeaderAdvisoryFixture {
921            discovery_handle,
922            header_sync,
923            header_actions,
924            header_task,
925            peer_node_id,
926            peer_id,
927            peer_send,
928            _peer_recv: peer_recv,
929        })
930    }
931
932    async fn send_discovery_message(
933        fixture: &HeaderAdvisoryFixture,
934        message: DiscoveryMessage,
935    ) -> Result<(), crate::BoxError> {
936        fixture
937            .peer_send
938            .send(Frame {
939                message_type: DISCOVERY_FRAME_MESSAGE_TYPE,
940                flags: 0,
941                payload: message.encode()?,
942            })
943            .await?;
944        Ok(())
945    }
946
947    fn discovery_frame(message: DiscoveryMessage) -> Result<Frame, crate::BoxError> {
948        Ok(Frame {
949            message_type: DISCOVERY_FRAME_MESSAGE_TYPE,
950            flags: 0,
951            payload: message.encode()?,
952        })
953    }
954
955    fn signed_discovery_record(
956        secret_key: &SecretKey,
957        handshake: &ZakuraHandshakeConfig,
958    ) -> Result<ZakuraNodeRecord, crate::BoxError> {
959        let body = ZakuraNodeRecordBody {
960            node_id: secret_key.public(),
961            direct_addrs: vec![SocketAddr::new(
962                IpAddr::V4(Ipv4Addr::new(192, 0, 2, 45)),
963                8233,
964            )],
965            services: vec![ZakuraServiceId::discovery()],
966            zakura_protocol_min: handshake.zakura_protocol_min,
967            zakura_protocol_max: handshake.zakura_protocol_max,
968            network_id: handshake.network_id,
969            chain_id: handshake.chain_id,
970            sequence: 1,
971            expires_at_unix_secs: current_test_unix_secs().saturating_add(60),
972        };
973        Ok(ZakuraNodeRecord::sign(body, secret_key)?)
974    }
975
976    async fn complete_peer_side_discovery_exchange(
977        peer_send: &FramedSend,
978        peer_recv: &mut FramedRecv,
979        peer_secret: &SecretKey,
980        handshake: &ZakuraHandshakeConfig,
981    ) -> Result<(), crate::BoxError> {
982        let mut saw_hello = false;
983        let mut saw_get_peers = false;
984        let mut saw_get_services = false;
985        while !(saw_hello && saw_get_peers && saw_get_services) {
986            let frame = tokio::time::timeout(Duration::from_secs(2), peer_recv.recv())
987                .await?
988                .expect("discovery source sends exchange frames");
989            match decode_discovery_frame(&frame)? {
990                DiscoveryMessage::Hello { .. } => saw_hello = true,
991                DiscoveryMessage::GetPeers { .. } => saw_get_peers = true,
992                DiscoveryMessage::GetServices(_) => saw_get_services = true,
993                DiscoveryMessage::Peers { .. } | DiscoveryMessage::Services(_) => {}
994            }
995        }
996
997        peer_send
998            .send(discovery_frame(DiscoveryMessage::Hello {
999                record: signed_discovery_record(peer_secret, handshake)?,
1000            })?)
1001            .await?;
1002        peer_send
1003            .send(discovery_frame(DiscoveryMessage::Peers {
1004                records: Vec::new(),
1005            })?)
1006            .await?;
1007        let summary = DiscoveryServiceSummary {
1008            peer_exchange_slots_free: 1,
1009            max_records_per_response: 1,
1010            expected_disconnect_after_exchange: true,
1011        };
1012        peer_send
1013            .send(discovery_frame(DiscoveryMessage::Services(Services {
1014                node_id: peer_secret.public(),
1015                expires_at_unix_secs: u64::MAX,
1016                summaries: vec![ServiceSummaryEnvelope::discovery(&summary)?],
1017            }))?)
1018            .await?;
1019
1020        Ok(())
1021    }
1022
1023    async fn wait_for_discovery_inbound_peers(handle: &ZakuraDiscoveryHandle, expected: usize) {
1024        tokio::time::timeout(Duration::from_secs(2), async {
1025            loop {
1026                if handle.peer_snapshot().inbound_peers == expected {
1027                    return;
1028                }
1029                tokio::time::sleep(Duration::from_millis(10)).await;
1030            }
1031        })
1032        .await
1033        .expect("discovery peer snapshot reaches expected inbound count");
1034    }
1035
1036    async fn advisory_backoff_after_empty_headers(
1037        fixture: &mut HeaderAdvisoryFixture,
1038    ) -> Result<bool, crate::BoxError> {
1039        let (send, _recv) = framed_channel(32);
1040        let session = HeaderSyncPeerSession::from_parts_with_direction(
1041            fixture.peer_id.clone(),
1042            ServicePeerDirection::Inbound,
1043            send,
1044            CancellationToken::new(),
1045        );
1046        fixture
1047            .header_sync
1048            .send(HeaderSyncEvent::PeerConnected(session))
1049            .await?;
1050        fixture
1051            .header_sync
1052            .send(HeaderSyncEvent::WireMessage {
1053                peer: fixture.peer_id.clone(),
1054                msg: HeaderSyncMessage::Status(HeaderSyncStatus {
1055                    tip_height: block::Height(1),
1056                    tip_hash: block::Hash([9; 32]),
1057                    anchor_height: block::Height(0),
1058                    max_headers_per_response: 1,
1059                    max_inflight_requests: 1,
1060                }),
1061            })
1062            .await?;
1063
1064        let request_id = tokio::time::timeout(Duration::from_secs(2), async {
1065            loop {
1066                if let Some(HeaderSyncAction::SendMessage {
1067                    peer,
1068                    request_id,
1069                    msg: HeaderSyncMessage::GetHeaders { .. },
1070                }) = fixture.header_actions.recv().await
1071                {
1072                    if peer == fixture.peer_id {
1073                        return request_id
1074                            .expect("an outbound GetHeaders always carries a request ID");
1075                    }
1076                }
1077            }
1078        })
1079        .await
1080        .expect("header sync schedules a request before empty response");
1081
1082        fixture
1083            .header_sync
1084            .send(HeaderSyncEvent::WireHeaders {
1085                peer: fixture.peer_id.clone(),
1086                session_id: 0,
1087                request_id,
1088                headers: Vec::new(),
1089                body_sizes: Vec::new(),
1090                tree_aux_roots: Vec::new(),
1091            })
1092            .await?;
1093        tokio::time::sleep(Duration::from_millis(20)).await;
1094
1095        Ok(fixture
1096            .header_sync
1097            .candidate_state()
1098            .backed_off_node_ids
1099            .contains(&fixture.peer_node_id))
1100    }
1101
1102    #[tokio::test]
1103    async fn get_services_returns_local_first_party_discovery_summary(
1104    ) -> Result<(), crate::BoxError> {
1105        let (_connected_tx, connected_rx) = watch::channel(Vec::new());
1106        let handshake = ZakuraHandshakeConfig::for_network(&Network::Mainnet);
1107        let local_secret = SecretKey::from_bytes(&[21u8; 32]);
1108        let handle = ZakuraDiscoveryHandle::new(
1109            ZakuraDiscoveryLocalConfig {
1110                secret_key: local_secret.clone(),
1111                direct_addrs: Vec::new(),
1112                services: vec![ZakuraServiceId::discovery()],
1113                zakura_protocol_min: handshake.zakura_protocol_min,
1114                zakura_protocol_max: handshake.zakura_protocol_max,
1115                network_id: handshake.network_id,
1116                chain_id: handshake.chain_id,
1117                last_authored_sequence: None,
1118            },
1119            ZakuraDiscoveryConfig {
1120                peer_limits: ServicePeerLimits {
1121                    max_inbound_peers: 4,
1122                    ..ServicePeerLimits::default()
1123                },
1124                ..ZakuraDiscoveryConfig::default()
1125            },
1126            connected_rx,
1127        )?;
1128        let service = DiscoveryService::new(handle.clone());
1129        let peer_node_id = SecretKey::from_bytes(&[22u8; 32]).public();
1130        let peer_id = ZakuraPeerId::new(peer_node_id.as_bytes().to_vec())?;
1131        let (peer_send, service_recv) = framed_channel(8);
1132        let (service_send, mut peer_recv) = framed_channel(8);
1133        let streams = HashMap::from([(ZAKURA_STREAM_DISCOVERY, (service_recv, service_send))]);
1134
1135        service.add_peer(Peer::new(
1136            peer_id,
1137            None,
1138            ZAKURA_CAP_DISCOVERY,
1139            streams,
1140            CancellationToken::new(),
1141        ));
1142
1143        peer_send
1144            .send(Frame {
1145                message_type: DISCOVERY_FRAME_MESSAGE_TYPE,
1146                flags: 0,
1147                payload: DiscoveryMessage::GetServices(GetServices {
1148                    wanted_services: vec![ZakuraServiceId::discovery()],
1149                })
1150                .encode()?,
1151            })
1152            .await?;
1153
1154        let services = tokio::time::timeout(Duration::from_secs(2), async {
1155            loop {
1156                let frame = peer_recv.recv().await.expect("discovery stream stays open");
1157                let message = decode_discovery_frame(&frame).expect("outbound frame decodes");
1158                if let DiscoveryMessage::Services(services) = message {
1159                    return services;
1160                }
1161            }
1162        })
1163        .await
1164        .expect("service response is sent");
1165
1166        assert_eq!(services.node_id, local_secret.public());
1167        assert_eq!(services.summaries.len(), 1);
1168        assert_eq!(
1169            services.summaries[0].service_id,
1170            ZakuraServiceId::discovery()
1171        );
1172        let summary = services.summaries[0]
1173            .decode_discovery()?
1174            .expect("discovery summary tag decodes");
1175        assert_eq!(summary.peer_exchange_slots_free, 3);
1176        assert!(summary.expected_disconnect_after_exchange);
1177        assert_eq!(
1178            summary.max_records_per_response,
1179            u16::try_from(MAX_DISCOVERY_RECORDS_PER_RESPONSE)
1180                .expect("record response cap fits in u16")
1181        );
1182
1183        Ok(())
1184    }
1185
1186    #[tokio::test]
1187    async fn get_services_returns_local_first_party_block_sync_summary(
1188    ) -> Result<(), crate::BoxError> {
1189        let (_connected_tx, connected_rx) = watch::channel(Vec::new());
1190        let handshake = ZakuraHandshakeConfig::for_network(&Network::Mainnet);
1191        let local_secret = SecretKey::from_bytes(&[24u8; 32]);
1192        let discovery_handle = ZakuraDiscoveryHandle::new(
1193            ZakuraDiscoveryLocalConfig {
1194                secret_key: local_secret.clone(),
1195                direct_addrs: Vec::new(),
1196                services: vec![ZakuraServiceId::discovery(), ZakuraServiceId::block_sync()],
1197                zakura_protocol_min: handshake.zakura_protocol_min,
1198                zakura_protocol_max: handshake.zakura_protocol_max,
1199                network_id: handshake.network_id,
1200                chain_id: handshake.chain_id,
1201                last_authored_sequence: None,
1202            },
1203            ZakuraDiscoveryConfig::default(),
1204            connected_rx,
1205        )?;
1206        let (header_sync, _header_actions, header_task) = spawn_test_header_sync()?;
1207        let (tip_tx, tip_rx) = watch::channel((block::Height(5), block::Hash([5; 32])));
1208        drop(tip_tx);
1209        let (block_sync, _block_actions, block_task) =
1210            spawn_block_sync_reactor(BlockSyncStartup::new(
1211                BlockSyncFrontiers {
1212                    finalized_height: block::Height(0),
1213                    verified_block_tip: block::Height(5),
1214                    verified_block_hash: block::Hash([5; 32]),
1215                },
1216                (block::Height(5), block::Hash([5; 32])),
1217                tip_rx,
1218                ZakuraBlockSyncConfig::default(),
1219            ));
1220        let service = DiscoveryService::with_sync_services(
1221            discovery_handle,
1222            header_sync,
1223            Some(block_sync.clone()),
1224        );
1225        let peer_node_id = SecretKey::from_bytes(&[25u8; 32]).public();
1226        let peer_id = ZakuraPeerId::new(peer_node_id.as_bytes().to_vec())?;
1227        let (peer_send, service_recv) = framed_channel(8);
1228        let (service_send, mut peer_recv) = framed_channel(8);
1229        let streams = HashMap::from([(ZAKURA_STREAM_DISCOVERY, (service_recv, service_send))]);
1230
1231        service.add_peer(Peer::new(
1232            peer_id,
1233            None,
1234            ZAKURA_CAP_DISCOVERY | ZAKURA_CAP_BLOCK_SYNC,
1235            streams,
1236            CancellationToken::new(),
1237        ));
1238
1239        peer_send
1240            .send(Frame {
1241                message_type: DISCOVERY_FRAME_MESSAGE_TYPE,
1242                flags: 0,
1243                payload: DiscoveryMessage::GetServices(GetServices {
1244                    wanted_services: vec![ZakuraServiceId::block_sync()],
1245                })
1246                .encode()?,
1247            })
1248            .await?;
1249
1250        let services = tokio::time::timeout(Duration::from_secs(2), async {
1251            loop {
1252                let frame = peer_recv.recv().await.expect("discovery stream stays open");
1253                let message = decode_discovery_frame(&frame).expect("outbound frame decodes");
1254                if let DiscoveryMessage::Services(services) = message {
1255                    return services;
1256                }
1257            }
1258        })
1259        .await
1260        .expect("service response is sent");
1261
1262        assert_eq!(services.node_id, local_secret.public());
1263        assert_eq!(services.summaries.len(), 1);
1264        assert_eq!(
1265            services.summaries[0].service_id,
1266            ZakuraServiceId::block_sync()
1267        );
1268        let summary = services.summaries[0]
1269            .decode_block_sync()?
1270            .expect("block summary tag decodes");
1271        assert_eq!(summary.servable_low, block::Height(0));
1272        assert_eq!(summary.servable_high, block::Height(5));
1273        assert_eq!(summary.tip_hash, block::Hash([5; 32]));
1274        assert_eq!(
1275            usize::from(summary.free_slots),
1276            block_sync.peer_snapshot().inbound_slots_free
1277        );
1278        assert_eq!(
1279            summary.max_blocks_per_response,
1280            ZakuraBlockSyncConfig::default().advertised_max_blocks_per_response()
1281        );
1282        assert_eq!(summary.max_response_bytes, MAX_BS_RESPONSE_BYTES);
1283
1284        header_task.abort();
1285        block_task.abort();
1286        Ok(())
1287    }
1288
1289    #[tokio::test]
1290    async fn inbound_services_updates_first_party_live_summary_cache() -> Result<(), crate::BoxError>
1291    {
1292        let (connected_tx, connected_rx) = watch::channel(Vec::new());
1293        let handshake = ZakuraHandshakeConfig::for_network(&Network::Mainnet);
1294        let local_secret = SecretKey::from_bytes(&[23u8; 32]);
1295        let handle = ZakuraDiscoveryHandle::new(
1296            ZakuraDiscoveryLocalConfig {
1297                secret_key: local_secret,
1298                direct_addrs: Vec::new(),
1299                services: vec![ZakuraServiceId::discovery()],
1300                zakura_protocol_min: handshake.zakura_protocol_min,
1301                zakura_protocol_max: handshake.zakura_protocol_max,
1302                network_id: handshake.network_id,
1303                chain_id: handshake.chain_id,
1304                last_authored_sequence: None,
1305            },
1306            ZakuraDiscoveryConfig::default(),
1307            connected_rx,
1308        )?;
1309        let service = DiscoveryService::new(handle.clone());
1310        let peer_node_id = SecretKey::from_bytes(&[24u8; 32]).public();
1311        let peer_id = ZakuraPeerId::new(peer_node_id.as_bytes().to_vec())?;
1312        connected_tx.send_replace(vec![peer_id.clone()]);
1313
1314        let (peer_send, service_recv) = framed_channel(8);
1315        let (service_send, _peer_recv) = framed_channel(8);
1316        let streams = HashMap::from([(ZAKURA_STREAM_DISCOVERY, (service_recv, service_send))]);
1317
1318        service.add_peer(Peer::new(
1319            peer_id,
1320            None,
1321            ZAKURA_CAP_DISCOVERY,
1322            streams,
1323            CancellationToken::new(),
1324        ));
1325
1326        let summary = DiscoveryServiceSummary {
1327            peer_exchange_slots_free: 7,
1328            max_records_per_response: 11,
1329            expected_disconnect_after_exchange: false,
1330        };
1331        peer_send
1332            .send(Frame {
1333                message_type: DISCOVERY_FRAME_MESSAGE_TYPE,
1334                flags: 0,
1335                payload: DiscoveryMessage::Services(Services {
1336                    node_id: peer_node_id,
1337                    expires_at_unix_secs: u64::MAX,
1338                    summaries: vec![ServiceSummaryEnvelope::discovery(&summary)?],
1339                })
1340                .encode()?,
1341            })
1342            .await?;
1343
1344        let cached = tokio::time::timeout(Duration::from_secs(2), async {
1345            loop {
1346                if let Some(cached) = handle.live_service_summaries(peer_node_id).await {
1347                    if !cached.is_empty() {
1348                        return cached;
1349                    }
1350                }
1351                tokio::time::sleep(Duration::from_millis(10)).await;
1352            }
1353        })
1354        .await
1355        .expect("inbound SERVICES is imported");
1356
1357        assert_eq!(cached.len(), 1);
1358        assert_eq!(
1359            cached[0].summary,
1360            ZakuraLiveServiceSummary::Discovery(summary)
1361        );
1362
1363        Ok(())
1364    }
1365
1366    #[tokio::test]
1367    async fn first_party_header_services_emit_header_sync_advisory() -> Result<(), crate::BoxError>
1368    {
1369        let mut fixture = spawn_header_advisory_fixture(25)?;
1370        let summary = header_summary(block::Height(10));
1371
1372        send_discovery_message(
1373            &fixture,
1374            DiscoveryMessage::Services(Services {
1375                node_id: fixture.peer_node_id,
1376                expires_at_unix_secs: u64::MAX,
1377                summaries: vec![ServiceSummaryEnvelope::header_sync(&summary)?],
1378            }),
1379        )
1380        .await?;
1381
1382        tokio::time::timeout(Duration::from_secs(2), async {
1383            loop {
1384                if let Some(cached) = fixture
1385                    .discovery_handle
1386                    .live_service_summaries(fixture.peer_node_id)
1387                    .await
1388                {
1389                    if cached.iter().any(|cached_summary| {
1390                        cached_summary.summary == ZakuraLiveServiceSummary::HeaderSync(summary)
1391                    }) {
1392                        return;
1393                    }
1394                }
1395                tokio::time::sleep(Duration::from_millis(10)).await;
1396            }
1397        })
1398        .await
1399        .expect("first-party header summary is cached");
1400
1401        assert!(
1402            advisory_backoff_after_empty_headers(&mut fixture).await?,
1403            "first-party header SERVICES should emit a header-sync advisory event"
1404        );
1405
1406        Ok(())
1407    }
1408
1409    #[tokio::test]
1410    async fn first_party_retired_header_services_emit_header_sync_advisory(
1411    ) -> Result<(), crate::BoxError> {
1412        let mut fixture = spawn_header_advisory_fixture(35)?;
1413        let summary = header_summary(block::Height(10));
1414        let mut legacy_envelope = ServiceSummaryEnvelope::header_sync(&summary)?;
1415        legacy_envelope.service_id = ZakuraServiceId::header_sync_retired();
1416        legacy_envelope.summary_tag = SUMMARY_TAG_HEADER_SYNC_RETIRED;
1417
1418        send_discovery_message(
1419            &fixture,
1420            DiscoveryMessage::Services(Services {
1421                node_id: fixture.peer_node_id,
1422                expires_at_unix_secs: u64::MAX,
1423                summaries: vec![legacy_envelope],
1424            }),
1425        )
1426        .await?;
1427
1428        tokio::time::timeout(Duration::from_secs(2), async {
1429            loop {
1430                if let Some(cached) = fixture
1431                    .discovery_handle
1432                    .live_service_summaries(fixture.peer_node_id)
1433                    .await
1434                {
1435                    if cached.iter().any(|cached_summary| {
1436                        cached_summary.service_id == ZakuraServiceId::header_sync_retired()
1437                            && cached_summary.summary
1438                                == ZakuraLiveServiceSummary::HeaderSync(summary)
1439                    }) {
1440                        return;
1441                    }
1442                }
1443                tokio::time::sleep(Duration::from_millis(10)).await;
1444            }
1445        })
1446        .await
1447        .expect("retired first-party header summary is cached");
1448
1449        assert!(
1450            advisory_backoff_after_empty_headers(&mut fixture).await?,
1451            "retired first-party header SERVICES should emit a header-sync advisory event"
1452        );
1453
1454        Ok(())
1455    }
1456
1457    #[tokio::test]
1458    async fn mismatched_services_node_id_does_not_emit_header_sync_advisory(
1459    ) -> Result<(), crate::BoxError> {
1460        let mut fixture = spawn_header_advisory_fixture(26)?;
1461        let claimed_node_id = SecretKey::from_bytes(&[27u8; 32]).public();
1462        let summary = header_summary(block::Height(10));
1463
1464        send_discovery_message(
1465            &fixture,
1466            DiscoveryMessage::Services(Services {
1467                node_id: claimed_node_id,
1468                expires_at_unix_secs: u64::MAX,
1469                summaries: vec![ServiceSummaryEnvelope::header_sync(&summary)?],
1470            }),
1471        )
1472        .await?;
1473        tokio::time::sleep(Duration::from_millis(20)).await;
1474
1475        assert_eq!(
1476            fixture
1477                .discovery_handle
1478                .live_service_summaries(fixture.peer_node_id)
1479                .await,
1480            None
1481        );
1482        assert_eq!(
1483            fixture
1484                .discovery_handle
1485                .live_service_summaries(claimed_node_id)
1486                .await,
1487            None
1488        );
1489        assert!(
1490            !advisory_backoff_after_empty_headers(&mut fixture).await?,
1491            "mismatched SERVICES node id must not emit a header-sync advisory event"
1492        );
1493
1494        Ok(())
1495    }
1496
1497    #[tokio::test]
1498    async fn peers_response_does_not_emit_header_sync_advisory() -> Result<(), crate::BoxError> {
1499        let mut fixture = spawn_header_advisory_fixture(28)?;
1500        let handshake = ZakuraHandshakeConfig::for_network(&Network::Mainnet);
1501        let record_secret = SecretKey::from_bytes(&[29u8; 32]);
1502        let record = signed_header_sync_record(&record_secret, &handshake)?;
1503
1504        send_discovery_message(
1505            &fixture,
1506            DiscoveryMessage::Peers {
1507                records: vec![record.clone()],
1508            },
1509        )
1510        .await?;
1511        tokio::time::sleep(Duration::from_millis(20)).await;
1512
1513        assert_eq!(
1514            fixture
1515                .discovery_handle
1516                .live_service_summaries(record.body.node_id)
1517                .await,
1518            None
1519        );
1520        assert!(
1521            !advisory_backoff_after_empty_headers(&mut fixture).await?,
1522            "PEERS/gossiped records must not emit live header-sync advisory events"
1523        );
1524
1525        Ok(())
1526    }
1527
1528    #[tokio::test]
1529    async fn discovery_only_short_lived_exchange_closes_connection_and_backs_off(
1530    ) -> Result<(), crate::BoxError> {
1531        let (connected_tx, connected_rx) = watch::channel(Vec::new());
1532        let handshake = ZakuraHandshakeConfig::for_network(&Network::Mainnet);
1533        let local_secret = SecretKey::from_bytes(&[40u8; 32]);
1534        let handle = ZakuraDiscoveryHandle::new(
1535            ZakuraDiscoveryLocalConfig {
1536                secret_key: local_secret,
1537                direct_addrs: Vec::new(),
1538                services: vec![ZakuraServiceId::discovery()],
1539                zakura_protocol_min: handshake.zakura_protocol_min,
1540                zakura_protocol_max: handshake.zakura_protocol_max,
1541                network_id: handshake.network_id,
1542                chain_id: handshake.chain_id,
1543                last_authored_sequence: None,
1544            },
1545            ZakuraDiscoveryConfig::default(),
1546            connected_rx,
1547        )?;
1548        let service = DiscoveryService::new(handle.clone());
1549        let peer_secret = SecretKey::from_bytes(&[41u8; 32]);
1550        let peer_node_id = peer_secret.public();
1551        let peer_id = ZakuraPeerId::new(peer_node_id.as_bytes().to_vec())?;
1552        connected_tx.send_replace(vec![peer_id.clone()]);
1553
1554        let connection_cancel = CancellationToken::new();
1555        let (peer_send, service_recv) = framed_channel(16);
1556        let (service_send, mut peer_recv) = framed_channel(16);
1557        let streams = HashMap::from([(ZAKURA_STREAM_DISCOVERY, (service_recv, service_send))]);
1558
1559        service.add_peer(Peer::new(
1560            peer_id,
1561            None,
1562            ZAKURA_CAP_DISCOVERY,
1563            streams,
1564            connection_cancel.clone(),
1565        ));
1566
1567        wait_for_discovery_inbound_peers(&handle, 1).await;
1568        complete_peer_side_discovery_exchange(&peer_send, &mut peer_recv, &peer_secret, &handshake)
1569            .await?;
1570        tokio::time::timeout(Duration::from_secs(2), connection_cancel.cancelled())
1571            .await
1572            .expect("discovery-only exchange closes the shared connection");
1573        wait_for_discovery_inbound_peers(&handle, 0).await;
1574
1575        connected_tx.send_replace(Vec::new());
1576        assert!(handle
1577            .dial_candidates(&[ZakuraServiceId::discovery()], &[])
1578            .await
1579            .is_empty());
1580
1581        Ok(())
1582    }
1583
1584    #[tokio::test]
1585    async fn discovery_short_lived_exchange_keeps_header_sync_connection(
1586    ) -> Result<(), crate::BoxError> {
1587        let (connected_tx, connected_rx) = watch::channel(Vec::new());
1588        let handshake = ZakuraHandshakeConfig::for_network(&Network::Mainnet);
1589        let local_secret = SecretKey::from_bytes(&[42u8; 32]);
1590        let discovery_handle = ZakuraDiscoveryHandle::new(
1591            ZakuraDiscoveryLocalConfig {
1592                secret_key: local_secret,
1593                direct_addrs: Vec::new(),
1594                services: vec![ZakuraServiceId::discovery()],
1595                zakura_protocol_min: handshake.zakura_protocol_min,
1596                zakura_protocol_max: handshake.zakura_protocol_max,
1597                network_id: handshake.network_id,
1598                chain_id: handshake.chain_id,
1599                last_authored_sequence: None,
1600            },
1601            ZakuraDiscoveryConfig::default(),
1602            connected_rx,
1603        )?;
1604        let (header_sync, _header_actions, header_task) = spawn_test_header_sync()?;
1605        let service = DiscoveryService::with_sync_services(
1606            discovery_handle.clone(),
1607            header_sync.clone(),
1608            None,
1609        );
1610        let peer_secret = SecretKey::from_bytes(&[43u8; 32]);
1611        let peer_node_id = peer_secret.public();
1612        let peer_id = ZakuraPeerId::new(peer_node_id.as_bytes().to_vec())?;
1613        connected_tx.send_replace(vec![peer_id.clone()]);
1614
1615        let (header_send, _header_recv) = framed_channel(8);
1616        let header_session = HeaderSyncPeerSession::from_parts_with_direction(
1617            peer_id.clone(),
1618            ServicePeerDirection::Inbound,
1619            header_send,
1620            CancellationToken::new(),
1621        );
1622        header_sync
1623            .send(HeaderSyncEvent::PeerConnected(header_session))
1624            .await?;
1625        tokio::time::timeout(Duration::from_secs(2), async {
1626            loop {
1627                if header_sync
1628                    .candidate_state()
1629                    .admitted_node_ids
1630                    .contains(&peer_node_id)
1631                {
1632                    return;
1633                }
1634                tokio::time::sleep(Duration::from_millis(10)).await;
1635            }
1636        })
1637        .await
1638        .expect("header sync admits the peer");
1639
1640        let connection_cancel = CancellationToken::new();
1641        let (peer_send, service_recv) = framed_channel(16);
1642        let (service_send, mut peer_recv) = framed_channel(16);
1643        let streams = HashMap::from([(ZAKURA_STREAM_DISCOVERY, (service_recv, service_send))]);
1644        service.add_peer(Peer::new(
1645            peer_id,
1646            None,
1647            ZAKURA_CAP_DISCOVERY | ZAKURA_CAP_HEADER_SYNC,
1648            streams,
1649            connection_cancel.clone(),
1650        ));
1651
1652        wait_for_discovery_inbound_peers(&discovery_handle, 1).await;
1653        complete_peer_side_discovery_exchange(&peer_send, &mut peer_recv, &peer_secret, &handshake)
1654            .await?;
1655        wait_for_discovery_inbound_peers(&discovery_handle, 0).await;
1656        assert_eq!(header_sync.peer_snapshot().inbound_peers, 1);
1657        assert!(
1658            tokio::time::timeout(Duration::from_millis(100), connection_cancel.cancelled())
1659                .await
1660                .is_err(),
1661            "discovery releases only its own session while header sync owns the connection"
1662        );
1663
1664        header_task.abort();
1665        Ok(())
1666    }
1667}