1mod admission;
19mod call;
20mod confidential;
21mod dht;
22mod framing;
23mod pubsub;
24mod report;
25mod serve;
26mod stream;
27mod versions;
28
29pub use admission::{Admission, AdmissionLimits};
30pub use call::{Call, DEFAULT_CALL_TIMEOUT, MAX_CALL_TIMEOUT};
31pub use confidential::{
32 is_clear_refusal, Confidentiality, ConfidentialityError, ConfidentialityReason, Seal,
33};
34pub use pubsub::{Event, EventDedup, Publication, PublicationSeq, SignedPublication, Subscription};
35pub use report::{Report, ReportError};
36pub use serve::{handler, BoxFuture, Handler, Offer, Request, Served, StreamOffer};
37pub use stream::{
38 stream_handler, Stream, StreamCall, StreamEvent, StreamHandler, DEFAULT_STREAM_DEADLINE,
39};
40pub use versions::{forget_v5_peer, handshake_counters};
41
42use std::collections::HashMap;
43use std::fmt;
44use std::sync::atomic::{AtomicBool, AtomicI64, Ordering};
45use std::sync::{Arc, Mutex, MutexGuard};
46use std::time::Duration;
47
48use sha2::{Digest, Sha384};
49use tokio::sync::watch;
50
51use crate::cbor::{self, Value};
52use crate::frame::{self, FrameError, Liveness, NeighbourLink, NeighbourPeer};
53use crate::handshake::{
54 self, ClientSession, Exporter, HandshakeError, Peer, RefusalCode, Station, VERSION, VERSION_5,
55};
56use crate::node_key::NodeKey;
57use crate::profile::Profile;
58use crate::record::RecordError;
59use crate::seal::Keyring;
60use crate::statement_issuer::{
61 ConnectMaterial, IssuerError, StatementIssuer, StatementSubscription,
62};
63use crate::transport::{self, DialError, Target};
64
65use framing::{read_frame, FrameWriter, HANDSHAKE_FRAME_BYTES, MAX_FRAME_BYTES};
66
67pub const HANDSHAKE_TIMEOUT: Duration = Duration::from_secs(30);
69
70const CLOSE_LINGER: Duration = Duration::from_secs(1);
73
74const STATUS_GRACE_MS: i64 = 5 * 60 * 1000;
77
78#[derive(Debug, Clone, PartialEq, Eq)]
80pub enum LinkError {
81 InvalidConfig(String),
84 Dial(String),
86 Handshake(HandshakeError),
88 HandshakeTimeout,
90 Issuer(IssuerError),
92 Io(String),
94 FrameTooLarge(usize),
96 Frame(FrameError),
99 Record(RecordError),
101 StatusExpired,
103 BindingExpired,
105 Closed,
107 Goodbye(String),
109 LivenessLost,
111 V5DowngradeRefused,
115 CallTimeout,
117 Provider {
119 responded_by: [u8; 32],
120 code: String,
121 detail: Option<String>,
122 },
123 Relay { reported_by: [u8; 32], code: String },
125 RecordNotFound,
127 UnexpectedReply(String),
129 InvalidOffer,
132 NoOrg,
134 AlreadyServed,
136 KemAdvertiseDisabled,
139 Stopped,
141 Stream {
144 code: String,
145 message: String,
146 relay: bool,
147 },
148 EndOfStream,
150 StreamClosed,
152 StreamOpenTooLarge(usize),
154 Confidentiality(ConfidentialityError),
157 SealedRefused { named: Option<[u8; 8]> },
160 ClearAnswerToSealed,
164 SealedFramesExhausted,
167}
168
169impl fmt::Display for LinkError {
170 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
171 match self {
172 LinkError::Provider {
173 code,
174 detail: Some(d),
175 ..
176 } => write!(f, "the provider answered {code}: {d}"),
177 LinkError::Provider { code, .. } => write!(f, "the provider answered {code}"),
178 LinkError::Relay { code, .. } => {
179 write!(f, "the station could not relay the call: {code}")
180 }
181 LinkError::Stream { code, message, .. } if !message.is_empty() => {
182 write!(f, "stream error {code}: {message}")
183 }
184 LinkError::Stream { code, .. } => write!(f, "stream error {code}"),
185 LinkError::Handshake(e) => write!(f, "handshake: {e}"),
186 LinkError::Frame(e) => write!(f, "frame: {e}"),
187 LinkError::Record(e) => write!(f, "record: {e}"),
188 LinkError::Issuer(e) => write!(f, "{e}"),
189 LinkError::Goodbye(reason) => write!(f, "the station said goodbye: {reason}"),
190 LinkError::Confidentiality(e) => write!(f, "{e}"),
191 LinkError::SealedRefused { named: Some(id) } => write!(
192 f,
193 "the provider could not open the request; it holds key {}",
194 id.iter().map(|b| format!("{b:02x}")).collect::<String>()
195 ),
196 LinkError::SealedRefused { named: None } => {
197 f.write_str("the provider opens no sealed payload")
198 }
199 LinkError::ClearAnswerToSealed => f.write_str("a clear answer to a sealed request"),
200 LinkError::KemAdvertiseDisabled => {
201 f.write_str("a required-confidential procedure needs kem_advertise on")
202 }
203 LinkError::SealedFramesExhausted => {
204 f.write_str("this sealed stream has sealed all the frames it may")
205 }
206 other => write!(f, "{other:?}"),
207 }
208 }
209}
210
211impl std::error::Error for LinkError {}
212
213impl From<FrameError> for LinkError {
214 fn from(e: FrameError) -> Self {
215 LinkError::Frame(e)
216 }
217}
218
219impl From<RecordError> for LinkError {
220 fn from(e: RecordError) -> Self {
221 LinkError::Record(e)
222 }
223}
224
225impl From<HandshakeError> for LinkError {
226 fn from(e: HandshakeError) -> Self {
227 LinkError::Handshake(e)
228 }
229}
230
231impl From<DialError> for LinkError {
232 fn from(e: DialError) -> Self {
233 LinkError::Dial(e.to_string())
234 }
235}
236
237pub struct Config {
243 pub target: Target,
244 pub identity: Arc<NodeKey>,
245 pub issuer: StatementIssuer,
246 pub member_endorsement: Vec<u8>,
247 pub publication_seq: Option<Arc<PublicationSeq>>,
248 pub admission: Option<Arc<Admission>>,
249 pub dedup: Option<Arc<EventDedup>>,
250 pub share: Option<String>,
253 pub keyring: Option<Arc<Keyring>>,
257 pub kem_advertise: bool,
263}
264
265impl Config {
266 pub fn new(target: Target, identity: Arc<NodeKey>, issuer: StatementIssuer) -> Config {
268 Config {
269 target,
270 identity,
271 issuer,
272 member_endorsement: Vec::new(),
273 publication_seq: None,
274 admission: None,
275 dedup: None,
276 share: None,
277 keyring: None,
278 kem_advertise: false,
279 }
280 }
281}
282
283#[derive(Clone)]
285pub struct Link {
286 inner: Arc<Inner>,
287}
288
289struct Inner {
290 serial: u64,
292 connection: quinn::Connection,
293 _endpoint: quinn::Endpoint,
294 control: FrameWriter,
295 profile: Profile,
296 key: Arc<NodeKey>,
297 self_id: [u8; 32],
298 station: Station,
299 station_capabilities: u64,
300 connection_hash: [u8; 48],
301 version: i64,
303 pongs: tokio::sync::mpsc::Sender<[u8; frame::LIVENESS_NONCE_SIZE]>,
305 pong_in_flight: AtomicBool,
308 pongs_rx: tokio::sync::Mutex<tokio::sync::mpsc::Receiver<[u8; frame::LIVENESS_NONCE_SIZE]>>,
309 send_seq: tokio::sync::Mutex<u64>,
311 status_deadline: AtomicI64,
312 publication_seq: Arc<PublicationSeq>,
313 admission: Arc<Admission>,
314 dedup: Arc<EventDedup>,
315 share: String,
316 keyring: Option<Arc<Keyring>>,
317 kem_advertise: bool,
318 state: Mutex<State>,
319 done_tx: watch::Sender<bool>,
320 done_rx: watch::Receiver<bool>,
321}
322
323struct State {
324 ended: Option<LinkError>,
325 closing: bool,
328 unrouted: HashMap<String, u64>,
329 pending: HashMap<[u8; 16], call::Pending>,
330 subs: HashMap<([u8; 32], String), Vec<pubsub::SubscriberSlot>>,
331 served: HashMap<([u8; 32], String), serve::ServedEntry>,
332 streams: Vec<std::sync::Weak<stream::StreamInner>>,
333}
334
335impl Link {
336 pub async fn dial(cfg: Config) -> Result<Link, LinkError> {
344 if cfg.identity.profile() != cfg.target.profile {
345 return Err(LinkError::InvalidConfig(
346 "the identity key is of another profile than the target's".into(),
347 ));
348 }
349 if let Some(admission) = &cfg.admission {
350 admission.limits().validate()?;
351 }
352 if cfg
353 .keyring
354 .as_ref()
355 .is_some_and(|k| k.profile() != cfg.target.profile)
356 {
357 return Err(LinkError::InvalidConfig(
358 "the keyring is of another profile than the target's".into(),
359 ));
360 }
361 if cfg.kem_advertise && cfg.keyring.is_none() {
362 return Err(LinkError::InvalidConfig(
363 "kem_advertise names a key: it needs a keyring".into(),
364 ));
365 }
366 let node_id = cfg.target.expected_node_id;
367 let version = versions::dial_version(&node_id, std::time::Instant::now());
368 let linked = dial_once(&cfg, version).await;
369 let v5_refused = matches!(
370 linked,
371 Err(LinkError::Handshake(HandshakeError::Refused(
372 RefusalCode::UnsupportedVersion
373 )))
374 ) && version == VERSION_5;
375 if !v5_refused {
376 return linked;
377 }
378 if !versions::unsupported_version(&node_id, std::time::Instant::now()) {
379 return Err(LinkError::V5DowngradeRefused);
380 }
381 dial_once(&cfg, VERSION).await
382 }
383
384 pub fn handshake_version(&self) -> i64 {
386 self.inner.version
387 }
388
389 pub fn station_node_id(&self) -> [u8; 32] {
391 self.inner.station.node_id
392 }
393
394 pub fn serial(&self) -> u64 {
396 self.inner.serial
397 }
398
399 pub fn node_id(&self) -> [u8; 32] {
401 self.inner.self_id
402 }
403
404 pub fn station_capabilities(&self) -> u64 {
406 self.inner.station_capabilities
407 }
408
409 pub fn profile(&self) -> Profile {
411 self.inner.profile
412 }
413
414 pub fn error(&self) -> Option<LinkError> {
416 self.inner.lock().ended.clone()
417 }
418
419 pub async fn done(&self) -> LinkError {
421 let mut done = self.inner.done_rx.clone();
422 let _ = done.wait_for(|ended| *ended).await;
423 self.error().unwrap_or(LinkError::Closed)
424 }
425
426 pub fn unrouted(&self) -> HashMap<String, u64> {
429 self.inner.lock().unrouted.clone()
430 }
431
432 pub async fn close(&self, reason: &str) -> Result<(), LinkError> {
436 if !self.inner.mark_closing() {
437 return Ok(());
438 }
439 let goodbye = frame::goodbye_frame(reason, None)?;
440 let sent = self.inner.send_control(&goodbye).await;
441 if sent.is_ok() {
442 self.inner.control.finish().await;
443 let _ = tokio::time::timeout(CLOSE_LINGER, self.inner.connection.closed()).await;
444 }
445 self.inner.end(LinkError::Closed);
446 sent
447 }
448}
449
450async fn dial_once(cfg: &Config, version: i64) -> Result<Link, LinkError> {
453 tokio::time::timeout(HANDSHAKE_TIMEOUT, async {
454 let dialed = transport::dial_target(&cfg.target).await?;
455 let connection = dialed.connection.clone();
456 let linked = handshaken(cfg, dialed, version).await;
457 if let Err(e) = &linked {
458 versions::count_refusal(e);
459 connection.close(0u32.into(), b"handshake failed");
460 }
461 linked
462 })
463 .await
464 .map_err(|_| LinkError::HandshakeTimeout)
465 .and_then(|linked| linked)
466}
467
468async fn handshaken(
471 cfg: &Config,
472 dialed: transport::Dialed,
473 version: i64,
474) -> Result<Link, LinkError> {
475 let export = keying_exporter(dialed.connection.clone());
476 let export: &Exporter = &export;
477 let (send, mut recv) = dialed
478 .connection
479 .open_bi()
480 .await
481 .map_err(|e| LinkError::Io(format!("open the control stream: {e}")))?;
482 let control = FrameWriter::new(send);
483 control
484 .write(&handshake::opener(), HANDSHAKE_FRAME_BYTES)
485 .await?;
486 let challenge = read_frame(&mut recv, HANDSHAKE_FRAME_BYTES).await?;
487 let material = cfg.issuer.connect_material().map_err(LinkError::Issuer)?;
488 let (connect, station) = handshake::answer_challenge(
489 &challenge,
490 &client_session(cfg, &dialed.leaf, &material, version, export),
491 )?;
492 control.write(&connect, HANDSHAKE_FRAME_BYTES).await?;
493 let hello = read_frame(&mut recv, HANDSHAKE_FRAME_BYTES).await?;
494 let capabilities = handshake::read_hello(&hello, &station)?;
495 let self_id = cfg
496 .identity
497 .node_id()
498 .map_err(|e| LinkError::InvalidConfig(e.to_string()))?;
499 let statements = cfg
500 .issuer
501 .subscribe(&material.binding)
502 .map_err(LinkError::Issuer)?;
503 let inner = new_inner(
504 cfg,
505 dialed,
506 control,
507 station,
508 capabilities,
509 &challenge,
510 self_id,
511 );
512 let station_node_id = inner.station.node_id;
513 versions::completed(&station_node_id, &inner)?;
514 spawn_link_tasks(&inner, statements, recv);
515 Ok(Link { inner })
516}
517
518fn keying_exporter(
521 connection: quinn::Connection,
522) -> impl Fn(&str, &[u8], usize) -> Option<Vec<u8>> + Send + Sync {
523 move |label: &str, context: &[u8], length: usize| {
524 let mut out = vec![0; length];
525 connection
526 .export_keying_material(&mut out, label.as_bytes(), context)
527 .ok()?;
528 Some(out)
529 }
530}
531
532fn client_session<'a>(
536 cfg: &Config,
537 leaf: &'a [u8],
538 material: &'a ConnectMaterial,
539 version: i64,
540 export: &'a Exporter,
541) -> ClientSession<'a> {
542 ClientSession {
543 profile: cfg.target.profile,
544 expected_node_id: cfg.target.expected_node_id,
545 leaf,
546 identity_key: cfg.identity.public_key(),
547 connect_key: &material.key,
548 connect_binding: &material.binding,
549 connect_status: &material.status,
550 capabilities: 0,
551 now_ms: now_ms(),
552 member_endorsement: cfg.member_endorsement.clone(),
553 version,
554 export: Some(export),
555 }
556}
557
558fn new_inner(
561 cfg: &Config,
562 dialed: transport::Dialed,
563 control: FrameWriter,
564 station: Station,
565 capabilities: u64,
566 challenge: &[u8],
567 self_id: [u8; 32],
568) -> Arc<Inner> {
569 let (done_tx, done_rx) = watch::channel(false);
570 let share = cfg
571 .share
572 .clone()
573 .unwrap_or_else(|| format!("{}:{}", cfg.target.host, cfg.target.port));
574 let (pongs, pongs_rx) = tokio::sync::mpsc::channel(1);
575 static SERIALS: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(1);
576 Arc::new(Inner {
577 serial: SERIALS.fetch_add(1, Ordering::Relaxed),
578 connection: dialed.connection,
579 _endpoint: dialed.endpoint,
580 control,
581 profile: cfg.target.profile,
582 key: cfg.identity.clone(),
583 self_id,
584 status_deadline: AtomicI64::new(station.status_expires_at + STATUS_GRACE_MS),
585 station_capabilities: capabilities,
586 connection_hash: Sha384::digest(challenge).into(),
587 version: station.version,
588 pongs,
589 pong_in_flight: AtomicBool::new(false),
590 pongs_rx: tokio::sync::Mutex::new(pongs_rx),
591 station,
592 send_seq: tokio::sync::Mutex::new(0),
593 publication_seq: cfg.publication_seq.clone().unwrap_or_default(),
594 admission: cfg
595 .admission
596 .clone()
597 .unwrap_or_else(|| Arc::new(Admission::new(AdmissionLimits::default()))),
598 dedup: cfg.dedup.clone().unwrap_or_default(),
599 share,
600 keyring: cfg.keyring.clone(),
601 kem_advertise: cfg.kem_advertise,
602 state: Mutex::new(State {
603 ended: None,
604 closing: false,
605 unrouted: HashMap::new(),
606 pending: HashMap::new(),
607 subs: HashMap::new(),
608 served: HashMap::new(),
609 streams: Vec::new(),
610 }),
611 done_tx,
612 done_rx,
613 })
614}
615
616fn spawn_link_tasks(
620 inner: &Arc<Inner>,
621 statements: StatementSubscription,
622 recv: quinn::RecvStream,
623) {
624 tokio::spawn(send_statements(Arc::downgrade(inner), statements));
625 tokio::spawn(read_control(inner.clone(), recv));
626 tokio::spawn(stream::accept_streams(Arc::downgrade(inner)));
627 tokio::spawn(call::probe(Arc::downgrade(inner)));
628 tokio::spawn(watch_expiries(Arc::downgrade(inner)));
629 tokio::spawn(watch_connection(Arc::downgrade(inner)));
630}
631
632impl Inner {
633 fn lock(&self) -> MutexGuard<'_, State> {
634 self.state
635 .lock()
636 .unwrap_or_else(|poisoned| poisoned.into_inner())
637 }
638
639 fn count(&self, what: &str) {
640 *self.lock().unrouted.entry(what.to_string()).or_default() += 1;
641 }
642
643 fn mark_closing(&self) -> bool {
646 let mut state = self.lock();
647 if state.ended.is_some() {
648 return false;
649 }
650 state.closing = true;
651 true
652 }
653
654 async fn send_control(&self, v: &Value) -> Result<(), LinkError> {
657 if self.version == VERSION_5 {
658 return self.write_control(v).await;
659 }
660 let mut seq = self.send_seq.lock().await;
661 let signed = frame::sign_neighbour(
662 v,
663 &self.key,
664 &NeighbourLink {
665 connection: self.connection_hash,
666 seq: *seq,
667 },
668 )?;
669 if frame::neighbour_signed(self.profile, &frame_type_of(v)) {
670 *seq += 1;
671 }
672 self.write_control(&signed).await
673 }
674
675 async fn write_control(&self, v: &Value) -> Result<(), LinkError> {
678 let encoded =
679 cbor::encode(v).map_err(|e| LinkError::Frame(FrameError::Payload(e.to_string())))?;
680 self.control.write(&encoded, MAX_FRAME_BYTES).await
681 }
682
683 fn end(&self, err: LinkError) {
686 let mut state = self.lock();
687 let Some(err) = state.mark_ended(err) else {
688 return;
689 };
690 let pending = std::mem::take(&mut state.pending);
691 let subs = std::mem::take(&mut state.subs);
692 let served = std::mem::take(&mut state.served);
693 let streams = std::mem::take(&mut state.streams);
694 drop(state);
695 for (_, p) in pending {
696 let _ = p.outcome.send(Err(err.clone()));
697 }
698 drop(subs);
699 for s in served.into_values() {
700 s.end(err.clone());
701 }
702 for s in streams.into_iter().filter_map(|w| w.upgrade()) {
703 stream::StreamInner::end(&s, Some(err.clone()));
704 }
705 self.connection.close(0u32.into(), b"link ended");
706 if self.version == VERSION {
707 versions::v4_ended(&self.station.node_id, self.serial);
708 }
709 let _ = self.done_tx.send_replace(true);
710 }
711}
712
713impl State {
714 fn mark_ended(&mut self, err: LinkError) -> Option<LinkError> {
718 if self.ended.is_some() {
719 return None;
720 }
721 let err = if self.closing { LinkError::Closed } else { err };
722 self.ended = Some(err.clone());
723 Some(err)
724 }
725}
726
727async fn send_statements(link: std::sync::Weak<Inner>, mut statements: StatementSubscription) {
730 loop {
731 let Some(done) = link.upgrade().map(|l| l.done_rx.clone()) else {
732 return;
733 };
734 let mut done = done;
735 let statement = tokio::select! {
736 _ = done.wait_for(|ended| *ended) => return,
737 statement = statements.recv() => statement,
738 };
739 let (Some(statement), Some(inner)) = (statement, link.upgrade()) else {
740 return;
741 };
742 if let Err(e) = inner
743 .control
744 .write(&handshake::status_frame(&statement), MAX_FRAME_BYTES)
745 .await
746 {
747 inner.end(e);
748 return;
749 }
750 }
751}
752
753async fn read_control(inner: Arc<Inner>, mut recv: quinn::RecvStream) {
755 let mut recv_seq = 0u64;
756 let mut done = inner.done_rx.clone();
757 loop {
758 let payload = tokio::select! {
759 _ = done.wait_for(|ended| *ended) => return,
760 payload = read_frame(&mut recv, MAX_FRAME_BYTES) => payload,
761 };
762 let outcome = match payload {
763 Ok(payload) => received(&inner, &payload, &mut recv_seq),
764 Err(e) => Err(e),
765 };
766 if let Err(e) = outcome {
767 inner.end(e);
768 return;
769 }
770 }
771}
772
773fn received(inner: &Arc<Inner>, payload: &[u8], recv_seq: &mut u64) -> Result<(), LinkError> {
778 let v = cbor::decode(payload).map_err(|_| LinkError::Frame(FrameError::Malformed))?;
779 let frame_type = frame_type_of(&v);
780 if frame_type == "status" {
781 let expires_at = handshake::read_status(
782 payload,
783 &Peer {
784 profile: inner.profile,
785 identity_key: inner.station.identity_key.clone(),
786 binding: inner.station.tls_binding.clone(),
787 now_ms: now_ms(),
788 },
789 )?;
790 inner
791 .status_deadline
792 .store(expires_at + STATUS_GRACE_MS, Ordering::SeqCst);
793 return Ok(());
794 }
795 if let Some((kind, nonce)) = liveness_frame(inner.version, &v)? {
796 return liveness(inner, kind, nonce);
797 }
798 let opened = opened(inner, &v, &frame_type, recv_seq)?;
799 match frame_type.as_str() {
800 "event" => pubsub::evented(inner, &opened),
801 "result" | "error" => call::replied(inner, &opened),
802 "call" => serve::called(inner, &opened),
803 "goodbye" => {
804 let reason = match opened.get("reason") {
805 Some(Value::Text(r)) => r.clone(),
806 _ => String::new(),
807 };
808 return Err(LinkError::Goodbye(reason));
809 }
810 _ => inner.count(&frame_type),
811 }
812 Ok(())
813}
814
815fn opened(
818 inner: &Inner,
819 v: &Value,
820 frame_type: &str,
821 recv_seq: &mut u64,
822) -> Result<Value, LinkError> {
823 if inner.version == VERSION_5 {
824 return Ok(frame::verify_session_frame(v)?);
825 }
826 let opened = frame::verify_neighbour(
827 v,
828 &NeighbourPeer {
829 profile: inner.profile,
830 peer_key: inner.station.identity_key.clone(),
831 connection: inner.connection_hash,
832 seq: *recv_seq,
833 },
834 )?;
835 if frame::neighbour_signed(inner.profile, frame_type) {
836 *recv_seq += 1;
837 }
838 Ok(opened)
839}
840
841fn liveness_frame(
846 version: i64,
847 v: &Value,
848) -> Result<Option<(Liveness, [u8; frame::LIVENESS_NONCE_SIZE])>, LinkError> {
849 let Some(found) = frame::liveness_nonce(v) else {
850 return Ok(None);
851 };
852 if version != VERSION_5 {
853 return Err(LinkError::Frame(FrameError::Malformed));
854 }
855 frame::verify_session_frame(v)?;
856 Ok(Some(found))
857}
858
859fn liveness(
862 inner: &Arc<Inner>,
863 kind: Liveness,
864 nonce: [u8; frame::LIVENESS_NONCE_SIZE],
865) -> Result<(), LinkError> {
866 match kind {
867 Liveness::Ping if !pong_slot(&inner.pong_in_flight) => {}
871 Liveness::Ping => {
872 tokio::spawn(send_pong(inner.clone(), nonce));
873 }
874 Liveness::Pong => {
875 let _ = inner.pongs.try_send(nonce);
876 }
877 }
878 Ok(())
879}
880
881async fn send_pong(inner: Arc<Inner>, nonce: [u8; frame::LIVENESS_NONCE_SIZE]) {
884 let sent = inner
885 .send_control(&frame::liveness_pong_frame(&nonce))
886 .await;
887 inner.pong_in_flight.store(false, Ordering::Release);
888 if let Err(e) = sent {
889 inner.end(e);
890 }
891}
892
893fn pong_slot(in_flight: &AtomicBool) -> bool {
895 !in_flight.swap(true, Ordering::AcqRel)
896}
897
898async fn watch_expiries(link: std::sync::Weak<Inner>) {
901 loop {
902 let Some(inner) = link.upgrade() else { return };
903 let now = now_ms();
904 if now >= inner.station.binding_not_after {
905 inner.end(LinkError::BindingExpired);
906 return;
907 }
908 let status_deadline = inner.status_deadline.load(Ordering::SeqCst);
909 if now >= status_deadline {
910 inner.end(LinkError::StatusExpired);
911 return;
912 }
913 let wait = (status_deadline.min(inner.station.binding_not_after) - now).clamp(1, 60_000);
914 let mut done = inner.done_rx.clone();
915 drop(inner);
916 tokio::select! {
917 _ = done.wait_for(|ended| *ended) => return,
918 _ = tokio::time::sleep(Duration::from_millis(wait as u64)) => {}
919 }
920 }
921}
922
923async fn watch_connection(link: std::sync::Weak<Inner>) {
925 let Some(connection) = link.upgrade().map(|l| l.connection.clone()) else {
926 return;
927 };
928 let cause = connection.closed().await;
929 if let Some(inner) = link.upgrade() {
930 inner.end(LinkError::Io(cause.to_string()));
931 }
932}
933
934fn frame_type_of(v: &Value) -> String {
935 match v.get("frame_type") {
936 Some(Value::Text(t)) => t.clone(),
937 _ => String::new(),
938 }
939}
940
941fn now_ms() -> i64 {
942 crate::uuid_v7::now_ms() as i64
943}
944
945#[cfg(test)]
946mod tests {
947 use super::*;
948
949 fn with_neighbour(v: Value) -> Value {
950 let Value::Map(mut pairs) = v else {
951 panic!("a frame is a map");
952 };
953 pairs.push((Value::text("neighbour"), Value::Map(Vec::new())));
954 Value::Map(pairs)
955 }
956
957 #[test]
960 fn a_burst_of_pings_starts_one_pong_at_a_time() {
961 let in_flight = AtomicBool::new(false);
962 let taken = (0..1000).filter(|_| pong_slot(&in_flight)).count();
963 assert_eq!(taken, 1);
964 in_flight.store(false, Ordering::Release);
965 assert!(pong_slot(&in_flight), "free again once the pong is written");
966 }
967
968 #[test]
971 fn a_liveness_frame_is_read_as_every_v5_frame_and_exists_only_on_v5() {
972 let nonce = [9; frame::LIVENESS_NONCE_SIZE];
973 let ping = frame::liveness_ping_frame(&nonce);
974 let pong = frame::liveness_pong_frame(&nonce);
975 assert_eq!(
976 liveness_frame(VERSION_5, &ping),
977 Ok(Some((Liveness::Ping, nonce)))
978 );
979 assert_eq!(
980 liveness_frame(VERSION_5, &pong),
981 Ok(Some((Liveness::Pong, nonce)))
982 );
983 for signed in [with_neighbour(ping.clone()), with_neighbour(pong.clone())] {
984 assert_eq!(
985 liveness_frame(VERSION_5, &signed),
986 Err(LinkError::Frame(FrameError::Malformed))
987 );
988 }
989 assert_eq!(
990 liveness_frame(VERSION, &ping),
991 Err(LinkError::Frame(FrameError::Malformed))
992 );
993 let goodbye = frame::goodbye_frame("bye", None).unwrap();
994 assert_eq!(liveness_frame(VERSION_5, &goodbye), Ok(None));
995 }
996}