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::{IssuerError, StatementIssuer, StatementSubscription};
61use crate::transport::{self, DialError, Target};
62
63use framing::{read_frame, FrameWriter, HANDSHAKE_FRAME_BYTES, MAX_FRAME_BYTES};
64
65pub const HANDSHAKE_TIMEOUT: Duration = Duration::from_secs(30);
67
68const CLOSE_LINGER: Duration = Duration::from_secs(1);
71
72const STATUS_GRACE_MS: i64 = 5 * 60 * 1000;
75
76#[derive(Debug, Clone, PartialEq, Eq)]
78pub enum LinkError {
79 InvalidConfig(String),
82 Dial(String),
84 Handshake(HandshakeError),
86 HandshakeTimeout,
88 Issuer(IssuerError),
90 Io(String),
92 FrameTooLarge(usize),
94 Frame(FrameError),
97 Record(RecordError),
99 StatusExpired,
101 BindingExpired,
103 Closed,
105 Goodbye(String),
107 LivenessLost,
109 V5DowngradeRefused,
113 CallTimeout,
115 Provider {
117 responded_by: [u8; 32],
118 code: String,
119 detail: Option<String>,
120 },
121 Relay { reported_by: [u8; 32], code: String },
123 RecordNotFound,
125 UnexpectedReply(String),
127 InvalidOffer,
130 NoOrg,
132 AlreadyServed,
134 KemAdvertiseDisabled,
137 Stopped,
139 Stream {
142 code: String,
143 message: String,
144 relay: bool,
145 },
146 EndOfStream,
148 StreamClosed,
150 StreamOpenTooLarge(usize),
152 Confidentiality(ConfidentialityError),
155 SealedRefused { named: Option<[u8; 8]> },
158 ClearAnswerToSealed,
162 SealedFramesExhausted,
165}
166
167impl fmt::Display for LinkError {
168 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
169 match self {
170 LinkError::Provider {
171 code,
172 detail: Some(d),
173 ..
174 } => write!(f, "the provider answered {code}: {d}"),
175 LinkError::Provider { code, .. } => write!(f, "the provider answered {code}"),
176 LinkError::Relay { code, .. } => {
177 write!(f, "the station could not relay the call: {code}")
178 }
179 LinkError::Stream { code, message, .. } if !message.is_empty() => {
180 write!(f, "stream error {code}: {message}")
181 }
182 LinkError::Stream { code, .. } => write!(f, "stream error {code}"),
183 LinkError::Handshake(e) => write!(f, "handshake: {e}"),
184 LinkError::Frame(e) => write!(f, "frame: {e}"),
185 LinkError::Record(e) => write!(f, "record: {e}"),
186 LinkError::Issuer(e) => write!(f, "{e}"),
187 LinkError::Goodbye(reason) => write!(f, "the station said goodbye: {reason}"),
188 LinkError::Confidentiality(e) => write!(f, "{e}"),
189 LinkError::SealedRefused { named: Some(id) } => write!(
190 f,
191 "the provider could not open the request; it holds key {}",
192 id.iter().map(|b| format!("{b:02x}")).collect::<String>()
193 ),
194 LinkError::SealedRefused { named: None } => {
195 f.write_str("the provider opens no sealed payload")
196 }
197 LinkError::ClearAnswerToSealed => f.write_str("a clear answer to a sealed request"),
198 LinkError::KemAdvertiseDisabled => {
199 f.write_str("a required-confidential procedure needs kem_advertise on")
200 }
201 LinkError::SealedFramesExhausted => {
202 f.write_str("this sealed stream has sealed all the frames it may")
203 }
204 other => write!(f, "{other:?}"),
205 }
206 }
207}
208
209impl std::error::Error for LinkError {}
210
211impl From<FrameError> for LinkError {
212 fn from(e: FrameError) -> Self {
213 LinkError::Frame(e)
214 }
215}
216
217impl From<RecordError> for LinkError {
218 fn from(e: RecordError) -> Self {
219 LinkError::Record(e)
220 }
221}
222
223impl From<HandshakeError> for LinkError {
224 fn from(e: HandshakeError) -> Self {
225 LinkError::Handshake(e)
226 }
227}
228
229impl From<DialError> for LinkError {
230 fn from(e: DialError) -> Self {
231 LinkError::Dial(e.to_string())
232 }
233}
234
235pub struct Config {
241 pub target: Target,
242 pub identity: Arc<NodeKey>,
243 pub issuer: StatementIssuer,
244 pub member_endorsement: Vec<u8>,
245 pub publication_seq: Option<Arc<PublicationSeq>>,
246 pub admission: Option<Arc<Admission>>,
247 pub dedup: Option<Arc<EventDedup>>,
248 pub share: Option<String>,
251 pub keyring: Option<Arc<Keyring>>,
255 pub kem_advertise: bool,
261}
262
263impl Config {
264 pub fn new(target: Target, identity: Arc<NodeKey>, issuer: StatementIssuer) -> Config {
266 Config {
267 target,
268 identity,
269 issuer,
270 member_endorsement: Vec::new(),
271 publication_seq: None,
272 admission: None,
273 dedup: None,
274 share: None,
275 keyring: None,
276 kem_advertise: false,
277 }
278 }
279}
280
281#[derive(Clone)]
283pub struct Link {
284 inner: Arc<Inner>,
285}
286
287struct Inner {
288 serial: u64,
290 connection: quinn::Connection,
291 _endpoint: quinn::Endpoint,
292 control: FrameWriter,
293 profile: Profile,
294 key: Arc<NodeKey>,
295 self_id: [u8; 32],
296 station: Station,
297 station_capabilities: u64,
298 connection_hash: [u8; 48],
299 version: i64,
301 pongs: tokio::sync::mpsc::Sender<[u8; frame::LIVENESS_NONCE_SIZE]>,
303 pong_in_flight: AtomicBool,
306 pongs_rx: tokio::sync::Mutex<tokio::sync::mpsc::Receiver<[u8; frame::LIVENESS_NONCE_SIZE]>>,
307 send_seq: tokio::sync::Mutex<u64>,
309 status_deadline: AtomicI64,
310 publication_seq: Arc<PublicationSeq>,
311 admission: Arc<Admission>,
312 dedup: Arc<EventDedup>,
313 share: String,
314 keyring: Option<Arc<Keyring>>,
315 kem_advertise: bool,
316 state: Mutex<State>,
317 done_tx: watch::Sender<bool>,
318 done_rx: watch::Receiver<bool>,
319}
320
321struct State {
322 ended: Option<LinkError>,
323 closing: bool,
326 unrouted: HashMap<String, u64>,
327 pending: HashMap<[u8; 16], call::Pending>,
328 subs: HashMap<([u8; 32], String), Vec<pubsub::SubscriberSlot>>,
329 served: HashMap<([u8; 32], String), serve::ServedEntry>,
330 streams: Vec<std::sync::Weak<stream::StreamInner>>,
331}
332
333impl Link {
334 pub async fn dial(cfg: Config) -> Result<Link, LinkError> {
342 if cfg.identity.profile() != cfg.target.profile {
343 return Err(LinkError::InvalidConfig(
344 "the identity key is of another profile than the target's".into(),
345 ));
346 }
347 if let Some(admission) = &cfg.admission {
348 admission.limits().validate()?;
349 }
350 if cfg
351 .keyring
352 .as_ref()
353 .is_some_and(|k| k.profile() != cfg.target.profile)
354 {
355 return Err(LinkError::InvalidConfig(
356 "the keyring is of another profile than the target's".into(),
357 ));
358 }
359 if cfg.kem_advertise && cfg.keyring.is_none() {
360 return Err(LinkError::InvalidConfig(
361 "kem_advertise names a key: it needs a keyring".into(),
362 ));
363 }
364 let node_id = cfg.target.expected_node_id;
365 let version = versions::dial_version(&node_id, std::time::Instant::now());
366 let linked = dial_once(&cfg, version).await;
367 match linked {
368 Err(LinkError::Handshake(HandshakeError::Refused(RefusalCode::UnsupportedVersion)))
369 if version == VERSION_5 =>
370 {
371 if !versions::unsupported_version(&node_id, std::time::Instant::now()) {
372 return Err(LinkError::V5DowngradeRefused);
373 }
374 dial_once(&cfg, VERSION).await
375 }
376 other => other,
377 }
378 }
379
380 pub fn handshake_version(&self) -> i64 {
382 self.inner.version
383 }
384
385 pub fn station_node_id(&self) -> [u8; 32] {
387 self.inner.station.node_id
388 }
389
390 pub fn serial(&self) -> u64 {
392 self.inner.serial
393 }
394
395 pub fn node_id(&self) -> [u8; 32] {
397 self.inner.self_id
398 }
399
400 pub fn station_capabilities(&self) -> u64 {
402 self.inner.station_capabilities
403 }
404
405 pub fn profile(&self) -> Profile {
407 self.inner.profile
408 }
409
410 pub fn error(&self) -> Option<LinkError> {
412 self.inner.lock().ended.clone()
413 }
414
415 pub async fn done(&self) -> LinkError {
417 let mut done = self.inner.done_rx.clone();
418 let _ = done.wait_for(|ended| *ended).await;
419 self.error().unwrap_or(LinkError::Closed)
420 }
421
422 pub fn unrouted(&self) -> HashMap<String, u64> {
425 self.inner.lock().unrouted.clone()
426 }
427
428 pub async fn close(&self, reason: &str) -> Result<(), LinkError> {
432 {
433 let mut state = self.inner.lock();
434 if state.ended.is_some() {
435 return Ok(());
436 }
437 state.closing = true;
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 exporter_connection = dialed.connection.clone();
476 let export = move |label: &str, context: &[u8], length: usize| {
477 let mut out = vec![0; length];
478 exporter_connection
479 .export_keying_material(&mut out, label.as_bytes(), context)
480 .ok()?;
481 Some(out)
482 };
483 let export: &Exporter = &export;
484 let (send, mut recv) = dialed
485 .connection
486 .open_bi()
487 .await
488 .map_err(|e| LinkError::Io(format!("open the control stream: {e}")))?;
489 let control = FrameWriter::new(send);
490 control
491 .write(&handshake::opener(), HANDSHAKE_FRAME_BYTES)
492 .await?;
493 let challenge = read_frame(&mut recv, HANDSHAKE_FRAME_BYTES).await?;
494 let material = cfg.issuer.connect_material().map_err(LinkError::Issuer)?;
495 let (connect, station) = handshake::answer_challenge(
496 &challenge,
497 &ClientSession {
498 profile: cfg.target.profile,
499 expected_node_id: cfg.target.expected_node_id,
500 leaf: &dialed.leaf,
501 identity_key: cfg.identity.public_key(),
502 connect_key: &material.key,
503 connect_binding: &material.binding,
504 connect_status: &material.status,
505 capabilities: 0,
506 now_ms: now_ms(),
507 member_endorsement: cfg.member_endorsement.clone(),
508 version,
509 export: Some(export),
510 },
511 )?;
512 control.write(&connect, HANDSHAKE_FRAME_BYTES).await?;
513 let hello = read_frame(&mut recv, HANDSHAKE_FRAME_BYTES).await?;
514 let capabilities = handshake::read_hello(&hello, &station)?;
515 let self_id = cfg
516 .identity
517 .node_id()
518 .map_err(|e| LinkError::InvalidConfig(e.to_string()))?;
519 let statements = cfg
520 .issuer
521 .subscribe(&material.binding)
522 .map_err(LinkError::Issuer)?;
523 let (done_tx, done_rx) = watch::channel(false);
524 let share = cfg
525 .share
526 .clone()
527 .unwrap_or_else(|| format!("{}:{}", cfg.target.host, cfg.target.port));
528 let (pongs, pongs_rx) = tokio::sync::mpsc::channel(1);
529 static SERIALS: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(1);
530 let inner = Arc::new(Inner {
531 serial: SERIALS.fetch_add(1, Ordering::Relaxed),
532 connection: dialed.connection,
533 _endpoint: dialed.endpoint,
534 control,
535 profile: cfg.target.profile,
536 key: cfg.identity.clone(),
537 self_id,
538 status_deadline: AtomicI64::new(station.status_expires_at + STATUS_GRACE_MS),
539 station_capabilities: capabilities,
540 connection_hash: Sha384::digest(&challenge).into(),
541 version: station.version,
542 pongs,
543 pong_in_flight: AtomicBool::new(false),
544 pongs_rx: tokio::sync::Mutex::new(pongs_rx),
545 station,
546 send_seq: tokio::sync::Mutex::new(0),
547 publication_seq: cfg.publication_seq.clone().unwrap_or_default(),
548 admission: cfg
549 .admission
550 .clone()
551 .unwrap_or_else(|| Arc::new(Admission::new(AdmissionLimits::default()))),
552 dedup: cfg.dedup.clone().unwrap_or_default(),
553 share,
554 keyring: cfg.keyring.clone(),
555 kem_advertise: cfg.kem_advertise,
556 state: Mutex::new(State {
557 ended: None,
558 closing: false,
559 unrouted: HashMap::new(),
560 pending: HashMap::new(),
561 subs: HashMap::new(),
562 served: HashMap::new(),
563 streams: Vec::new(),
564 }),
565 done_tx,
566 done_rx,
567 });
568 let station_node_id = inner.station.node_id;
569 versions::completed(&station_node_id, &inner)?;
570 tokio::spawn(send_statements(Arc::downgrade(&inner), statements));
571 tokio::spawn(read_control(inner.clone(), recv));
572 tokio::spawn(stream::accept_streams(Arc::downgrade(&inner)));
573 tokio::spawn(call::probe(Arc::downgrade(&inner)));
574 tokio::spawn(watch_expiries(Arc::downgrade(&inner)));
575 tokio::spawn(watch_connection(Arc::downgrade(&inner)));
576 Ok(Link { inner })
577}
578
579impl Inner {
580 fn lock(&self) -> MutexGuard<'_, State> {
581 self.state
582 .lock()
583 .unwrap_or_else(|poisoned| poisoned.into_inner())
584 }
585
586 fn count(&self, what: &str) {
587 *self.lock().unrouted.entry(what.to_string()).or_default() += 1;
588 }
589
590 async fn send_control(&self, v: &Value) -> Result<(), LinkError> {
593 if self.version == VERSION_5 {
594 return self.write_control(v).await;
595 }
596 let mut seq = self.send_seq.lock().await;
597 let signed = frame::sign_neighbour(
598 v,
599 &self.key,
600 &NeighbourLink {
601 connection: self.connection_hash,
602 seq: *seq,
603 },
604 )?;
605 if frame::neighbour_signed(self.profile, &frame_type_of(v)) {
606 *seq += 1;
607 }
608 self.write_control(&signed).await
609 }
610
611 async fn write_control(&self, v: &Value) -> Result<(), LinkError> {
614 let encoded =
615 cbor::encode(v).map_err(|e| LinkError::Frame(FrameError::Payload(e.to_string())))?;
616 self.control.write(&encoded, MAX_FRAME_BYTES).await
617 }
618
619 fn end(&self, err: LinkError) {
622 let (err, pending, subs, served, streams) = {
623 let mut state = self.lock();
624 if state.ended.is_some() {
625 return;
626 }
627 let err = if state.closing {
628 LinkError::Closed
629 } else {
630 err
631 };
632 state.ended = Some(err.clone());
633 (
634 err,
635 std::mem::take(&mut state.pending),
636 std::mem::take(&mut state.subs),
637 std::mem::take(&mut state.served),
638 std::mem::take(&mut state.streams),
639 )
640 };
641 for (_, p) in pending {
642 let _ = p.outcome.send(Err(err.clone()));
643 }
644 drop(subs);
645 for s in served.into_values() {
646 s.end(err.clone());
647 }
648 for s in streams.into_iter().filter_map(|w| w.upgrade()) {
649 stream::StreamInner::end(&s, Some(err.clone()));
650 }
651 self.connection.close(0u32.into(), b"link ended");
652 if self.version == VERSION {
653 versions::v4_ended(&self.station.node_id, self.serial);
654 }
655 let _ = self.done_tx.send_replace(true);
656 }
657}
658
659async fn send_statements(link: std::sync::Weak<Inner>, mut statements: StatementSubscription) {
662 loop {
663 let Some(done) = link.upgrade().map(|l| l.done_rx.clone()) else {
664 return;
665 };
666 let mut done = done;
667 let statement = tokio::select! {
668 _ = done.wait_for(|ended| *ended) => return,
669 statement = statements.recv() => statement,
670 };
671 let (Some(statement), Some(inner)) = (statement, link.upgrade()) else {
672 return;
673 };
674 if let Err(e) = inner
675 .control
676 .write(&handshake::status_frame(&statement), MAX_FRAME_BYTES)
677 .await
678 {
679 inner.end(e);
680 return;
681 }
682 }
683}
684
685async fn read_control(inner: Arc<Inner>, mut recv: quinn::RecvStream) {
687 let mut recv_seq = 0u64;
688 let mut done = inner.done_rx.clone();
689 loop {
690 let payload = tokio::select! {
691 _ = done.wait_for(|ended| *ended) => return,
692 payload = read_frame(&mut recv, MAX_FRAME_BYTES) => payload,
693 };
694 let outcome = match payload {
695 Ok(payload) => received(&inner, &payload, &mut recv_seq),
696 Err(e) => Err(e),
697 };
698 if let Err(e) = outcome {
699 inner.end(e);
700 return;
701 }
702 }
703}
704
705fn received(inner: &Arc<Inner>, payload: &[u8], recv_seq: &mut u64) -> Result<(), LinkError> {
710 let v = cbor::decode(payload).map_err(|_| LinkError::Frame(FrameError::Malformed))?;
711 let frame_type = frame_type_of(&v);
712 if frame_type == "status" {
713 let expires_at = handshake::read_status(
714 payload,
715 &Peer {
716 profile: inner.profile,
717 identity_key: inner.station.identity_key.clone(),
718 binding: inner.station.tls_binding.clone(),
719 now_ms: now_ms(),
720 },
721 )?;
722 inner
723 .status_deadline
724 .store(expires_at + STATUS_GRACE_MS, Ordering::SeqCst);
725 return Ok(());
726 }
727 if let Some((kind, nonce)) = liveness_frame(inner.version, &v)? {
728 return liveness(inner, kind, nonce);
729 }
730 let opened = opened(inner, &v, &frame_type, recv_seq)?;
731 match frame_type.as_str() {
732 "event" => pubsub::evented(inner, &opened),
733 "result" | "error" => call::replied(inner, &opened),
734 "call" => serve::called(inner, &opened),
735 "goodbye" => {
736 let reason = match opened.get("reason") {
737 Some(Value::Text(r)) => r.clone(),
738 _ => String::new(),
739 };
740 return Err(LinkError::Goodbye(reason));
741 }
742 _ => inner.count(&frame_type),
743 }
744 Ok(())
745}
746
747fn opened(
750 inner: &Inner,
751 v: &Value,
752 frame_type: &str,
753 recv_seq: &mut u64,
754) -> Result<Value, LinkError> {
755 if inner.version == VERSION_5 {
756 return Ok(frame::verify_session_frame(v)?);
757 }
758 let opened = frame::verify_neighbour(
759 v,
760 &NeighbourPeer {
761 profile: inner.profile,
762 peer_key: inner.station.identity_key.clone(),
763 connection: inner.connection_hash,
764 seq: *recv_seq,
765 },
766 )?;
767 if frame::neighbour_signed(inner.profile, frame_type) {
768 *recv_seq += 1;
769 }
770 Ok(opened)
771}
772
773fn liveness_frame(
778 version: i64,
779 v: &Value,
780) -> Result<Option<(Liveness, [u8; frame::LIVENESS_NONCE_SIZE])>, LinkError> {
781 let Some(found) = frame::liveness_nonce(v) else {
782 return Ok(None);
783 };
784 if version != VERSION_5 {
785 return Err(LinkError::Frame(FrameError::Malformed));
786 }
787 frame::verify_session_frame(v)?;
788 Ok(Some(found))
789}
790
791fn liveness(
794 inner: &Arc<Inner>,
795 kind: Liveness,
796 nonce: [u8; frame::LIVENESS_NONCE_SIZE],
797) -> Result<(), LinkError> {
798 match kind {
799 Liveness::Ping if !pong_slot(&inner.pong_in_flight) => {}
803 Liveness::Ping => {
804 let inner = inner.clone();
805 tokio::spawn(async move {
806 let sent = inner
807 .send_control(&frame::liveness_pong_frame(&nonce))
808 .await;
809 inner.pong_in_flight.store(false, Ordering::Release);
810 if let Err(e) = sent {
811 inner.end(e);
812 }
813 });
814 }
815 Liveness::Pong => {
816 let _ = inner.pongs.try_send(nonce);
817 }
818 }
819 Ok(())
820}
821
822fn pong_slot(in_flight: &AtomicBool) -> bool {
824 !in_flight.swap(true, Ordering::AcqRel)
825}
826
827async fn watch_expiries(link: std::sync::Weak<Inner>) {
830 loop {
831 let Some(inner) = link.upgrade() else { return };
832 let now = now_ms();
833 if now >= inner.station.binding_not_after {
834 inner.end(LinkError::BindingExpired);
835 return;
836 }
837 let status_deadline = inner.status_deadline.load(Ordering::SeqCst);
838 if now >= status_deadline {
839 inner.end(LinkError::StatusExpired);
840 return;
841 }
842 let wait = (status_deadline.min(inner.station.binding_not_after) - now).clamp(1, 60_000);
843 let mut done = inner.done_rx.clone();
844 drop(inner);
845 tokio::select! {
846 _ = done.wait_for(|ended| *ended) => return,
847 _ = tokio::time::sleep(Duration::from_millis(wait as u64)) => {}
848 }
849 }
850}
851
852async fn watch_connection(link: std::sync::Weak<Inner>) {
854 let Some(connection) = link.upgrade().map(|l| l.connection.clone()) else {
855 return;
856 };
857 let cause = connection.closed().await;
858 if let Some(inner) = link.upgrade() {
859 inner.end(LinkError::Io(cause.to_string()));
860 }
861}
862
863fn frame_type_of(v: &Value) -> String {
864 match v.get("frame_type") {
865 Some(Value::Text(t)) => t.clone(),
866 _ => String::new(),
867 }
868}
869
870fn now_ms() -> i64 {
871 crate::uuid_v7::now_ms() as i64
872}
873
874#[cfg(test)]
875mod tests {
876 use super::*;
877
878 fn with_neighbour(v: Value) -> Value {
879 let Value::Map(mut pairs) = v else {
880 panic!("a frame is a map");
881 };
882 pairs.push((Value::text("neighbour"), Value::Map(Vec::new())));
883 Value::Map(pairs)
884 }
885
886 #[test]
889 fn a_burst_of_pings_starts_one_pong_at_a_time() {
890 let in_flight = AtomicBool::new(false);
891 let taken = (0..1000).filter(|_| pong_slot(&in_flight)).count();
892 assert_eq!(taken, 1);
893 in_flight.store(false, Ordering::Release);
894 assert!(pong_slot(&in_flight), "free again once the pong is written");
895 }
896
897 #[test]
900 fn a_liveness_frame_is_read_as_every_v5_frame_and_exists_only_on_v5() {
901 let nonce = [9; frame::LIVENESS_NONCE_SIZE];
902 let ping = frame::liveness_ping_frame(&nonce);
903 let pong = frame::liveness_pong_frame(&nonce);
904 assert_eq!(
905 liveness_frame(VERSION_5, &ping),
906 Ok(Some((Liveness::Ping, nonce)))
907 );
908 assert_eq!(
909 liveness_frame(VERSION_5, &pong),
910 Ok(Some((Liveness::Pong, nonce)))
911 );
912 for signed in [with_neighbour(ping.clone()), with_neighbour(pong.clone())] {
913 assert_eq!(
914 liveness_frame(VERSION_5, &signed),
915 Err(LinkError::Frame(FrameError::Malformed))
916 );
917 }
918 assert_eq!(
919 liveness_frame(VERSION, &ping),
920 Err(LinkError::Frame(FrameError::Malformed))
921 );
922 let goodbye = frame::goodbye_frame("bye", None).unwrap();
923 assert_eq!(liveness_frame(VERSION_5, &goodbye), Ok(None));
924 }
925}