1use 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
41const 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 frame_cap: LOCAL_MAX_CONTROL_FRAME_BYTES,
50 capability: ZAKURA_CAP_DISCOVERY,
51 mode: StreamMode::Ordered,
52}];
53
54pub(crate) fn discovery_streams() -> &'static [Stream] {
56 &DISCOVERY_SERVICE_STREAMS
57}
58
59#[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 pub fn peer_id(&self) -> &ZakuraPeerId {
80 &self.peer_id
81 }
82
83 pub fn direction(&self) -> ServicePeerDirection {
85 self.direction
86 }
87
88 pub fn cancel_token(&self) -> CancellationToken {
90 self.cancel.clone()
91 }
92
93 pub fn try_send_hello(&self, record: ZakuraNodeRecord) -> Result<(), OrderedSendError> {
95 self.try_send_message(DiscoveryMessage::Hello { record })
96 }
97
98 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 pub fn try_send_peers(&self, records: Vec<ZakuraNodeRecord>) -> Result<(), OrderedSendError> {
114 self.try_send_message(DiscoveryMessage::Peers { records })
115 }
116
117 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 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#[derive(Clone, Debug)]
154pub struct DiscoveryService {
155 handle: ZakuraDiscoveryHandle,
156 header_sync: Option<HeaderSyncHandle>,
157 block_sync: Option<BlockSyncHandle>,
158}
159
160impl DiscoveryService {
161 pub fn new(handle: ZakuraDiscoveryHandle) -> Self {
163 Self {
164 handle,
165 header_sync: None,
166 block_sync: None,
167 }
168 }
169
170 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 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 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 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 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 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 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 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
413struct 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
615struct 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 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 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
736fn 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
743fn 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}