Skip to main content

macula_rust/
station_link.rs

1//! A client's link to one macula station, as macula_station_link and
2//! macula-go's stationlink are: a QUIC connection dialed to the station its
3//! target pins, one bidirectional control stream, the connection handshake on
4//! it (v5, or v4 once after a station refuses v5), then status statements
5//! both ways and every frame of the session.
6//!
7//! After HELLO the link sends its own status statement at every reissue of
8//! the node's statement issuer, and ends when the station's statement is five
9//! minutes past its expiry or the station's TLS binding reaches its
10//! not_after. On a v5 link no frame carries a neighbour signature: the
11//! session proofs authenticated the station once, and a neighbour-signed
12//! frame ends the link. On a v4 link, in pq_hybrid, every control frame is
13//! neighbour-signed with a sequence number per direction, from 0 after HELLO;
14//! a frame out of sequence ends the link. A liveness probe every 30 seconds
15//! ends the link after two misses in a row: on v5 a liveness_ping, on v4 a
16//! `_macula.ping` call.
17
18mod 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
65/// How long the handshake may take, as macula's 30 seconds.
66pub const HANDSHAKE_TIMEOUT: Duration = Duration::from_secs(30);
67
68/// How long Close waits, after its GOODBYE, for the station to close the
69/// connection before closing it itself.
70const CLOSE_LINGER: Duration = Duration::from_secs(1);
71
72/// How long past a statement's expiry the link keeps a station whose next
73/// statement has not arrived: macula's five minutes.
74const STATUS_GRACE_MS: i64 = 5 * 60 * 1000;
75
76/// Why a link or one of its operations failed.
77#[derive(Debug, Clone, PartialEq, Eq)]
78pub enum LinkError {
79    /// The configuration cannot be dialed: a key of another profile than the
80    /// target's, or admission limits macula would not start with.
81    InvalidConfig(String),
82    /// The dial failed.
83    Dial(String),
84    /// The handshake refused, or was refused.
85    Handshake(HandshakeError),
86    /// The handshake did not finish within [`HANDSHAKE_TIMEOUT`].
87    HandshakeTimeout,
88    /// The issuer had no CONNECT material in force.
89    Issuer(IssuerError),
90    /// The QUIC connection or a stream failed.
91    Io(String),
92    /// A frame longer than the cap in force.
93    FrameTooLarge(usize),
94    /// A frame that could not be built, or one received that does not
95    /// verify.
96    Frame(FrameError),
97    /// A record that could not be built, or one found that does not verify.
98    Record(RecordError),
99    /// The station's status statement lapsed past the grace.
100    StatusExpired,
101    /// The station's TLS binding reached its not_after.
102    BindingExpired,
103    /// The link was closed by its owner.
104    Closed,
105    /// The station ended the link with a GOODBYE, for this reason.
106    Goodbye(String),
107    /// The station missed two liveness probes in a row.
108    LivenessLost,
109    /// A station this process completed handshake v5 with that now answers
110    /// v4, or a v4 link to it that a v5 completion ended: refused, not
111    /// retried, until [`forget_v5_peer`].
112    V5DowngradeRefused,
113    /// No verified reply answered the call within its timeout.
114    CallTimeout,
115    /// A provider's ERROR for the call.
116    Provider {
117        responded_by: [u8; 32],
118        code: String,
119        detail: Option<String>,
120    },
121    /// A relay error the connected station reported for the call.
122    Relay { reported_by: [u8; 32], code: String },
123    /// A find_record the station answered not_found.
124    RecordNotFound,
125    /// A station reply to a DHT call in a shape that call never answers with.
126    UnexpectedReply(String),
127    /// An offer without exactly one handler, or an org procedure's without
128    /// the realm key.
129    InvalidOffer,
130    /// A procedure without a namespace, which nobody can authorize.
131    NoOrg,
132    /// A procedure already served on this link.
133    AlreadyServed,
134    /// A required-confidential procedure on a link whose `kem_advertise` is
135    /// off: it would name no key and refuse every clear call.
136    KemAdvertiseDisabled,
137    /// A served procedure withdrawn by its owner.
138    Stopped,
139    /// A stream ended by a STREAM_ERROR: the peer's, the station's relay
140    /// error (`relay`), or this side's own abort.
141    Stream {
142        code: String,
143        message: String,
144        relay: bool,
145    },
146    /// A stream that ended normally.
147    EndOfStream,
148    /// A send on a stream this side has ended.
149    StreamClosed,
150    /// A STREAM_OPEN over the 1 MiB a peer reads of one.
151    StreamOpenTooLarge(usize),
152    /// A call or an open that could not be kept confidential, so it was not
153    /// made, or failed rather than be taken in the clear.
154    Confidentiality(ConfidentialityError),
155    /// The provider's sealed_refused: it could not open the request. `named`
156    /// is the key it holds now, `None` when it holds none.
157    SealedRefused { named: Option<[u8; 8]> },
158    /// A clear answer to a sealed request that nothing clear may give: not a
159    /// relay error, not a refusal from the closed set, not sealed_refused.
160    /// Refused, never taken as the answer.
161    ClearAnswerToSealed,
162    /// A provider's sealed stream that has sealed all the frames GCM's
163    /// bound allows under random nonces: another would pass it.
164    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
235/// What a link is dialed with: the station to reach, the node's identity
236/// key, the statement issuer that holds its CONNECT key and statements, and
237/// the realm membership endorsement to present, empty for none. Links of one
238/// node share their publication seq, admission and dedup; a link given none
239/// makes its own.
240pub 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    /// This link's place in the admission: the station it dialed, host:port,
249    /// when `None`.
250    pub share: Option<String>,
251    /// This identity's KEM keys (macula 13, E2E design amendment A1): with
252    /// them the link opens sealed requests; without them it refuses them,
253    /// holding no key. Links of one identity share one.
254    pub keyring: Option<Arc<Keyring>>,
255    /// Names the keyring's current key in the advertisements of procedures
256    /// served confidentially ([`Confidentiality::Preferred`] or
257    /// [`Confidentiality::Required`]). Off by default: switch it on only once
258    /// every station runs macula 12.11 or later, which stores and routes a
259    /// keyed advertisement, and every caller runs 13. It needs a keyring.
260    pub kem_advertise: bool,
261}
262
263impl Config {
264    /// A link's configuration with none of the shared parts.
265    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/// A handshaked link to one station. Cloning it shares the link.
282#[derive(Clone)]
283pub struct Link {
284    inner: Arc<Inner>,
285}
286
287struct Inner {
288    /// Distinct for every link a process dials.
289    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    /// The handshake version the link completed: 4 or 5.
300    version: i64,
301    /// v5 liveness answers, for the probe.
302    pongs: tokio::sync::mpsc::Sender<[u8; frame::LIVENESS_NONCE_SIZE]>,
303    /// A liveness_pong is being written (macula-rust#9): at most one is in
304    /// flight, so a station's pings cannot pile up tasks on this link.
305    pong_in_flight: AtomicBool,
306    pongs_rx: tokio::sync::Mutex<tokio::sync::mpsc::Receiver<[u8; frame::LIVENESS_NONCE_SIZE]>>,
307    /// Orders a neighbour signature's seq with its write.
308    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    /// The owner is closing the link: however it then ends, it ends
324    /// [`LinkError::Closed`].
325    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    /// Dials `cfg.target`, runs the handshake as a client, and returns the
335    /// link once the station's HELLO accepts it. It dials with version 5
336    /// unless the station refused v5 in the last 10 minutes; a station never
337    /// seen on v5 that refuses v5 with `unsupported_version` is dialled once
338    /// more, on a new connection, with v4. A station seen on v5 that answers
339    /// v4 is [`LinkError::V5DowngradeRefused`]. Each handshake is bounded by
340    /// [`HANDSHAKE_TIMEOUT`].
341    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    /// The handshake version the link completed: 4 or 5.
381    pub fn handshake_version(&self) -> i64 {
382        self.inner.version
383    }
384
385    /// The node_id of the station the link reached.
386    pub fn station_node_id(&self) -> [u8; 32] {
387        self.inner.station.node_id
388    }
389
390    /// A number distinct for every link this process dials.
391    pub fn serial(&self) -> u64 {
392        self.inner.serial
393    }
394
395    /// The node_id this link connected as.
396    pub fn node_id(&self) -> [u8; 32] {
397        self.inner.self_id
398    }
399
400    /// The capability bits the station's HELLO announced.
401    pub fn station_capabilities(&self) -> u64 {
402        self.inner.station_capabilities
403    }
404
405    /// The profile the link runs.
406    pub fn profile(&self) -> Profile {
407        self.inner.profile
408    }
409
410    /// Why the link ended, or `None` while it runs.
411    pub fn error(&self) -> Option<LinkError> {
412        self.inner.lock().ended.clone()
413    }
414
415    /// Waits until the link has ended, and says why.
416    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    /// The frames received that nothing on this link handles, by what they
423    /// were.
424    pub fn unrouted(&self) -> HashMap<String, u64> {
425        self.inner.lock().unrouted.clone()
426    }
427
428    /// Sends a GOODBYE with `reason`, closes the control stream's sending
429    /// side, waits up to a second for the station to close the connection,
430    /// and ends the link. Closing an ended link does nothing.
431    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
450/// One QUIC connection and one handshake in `version`, within
451/// [`HANDSHAKE_TIMEOUT`]. A failed handshake closes its connection.
452async 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
468/// Runs the client's side of the handshake, in `version`, on a new control
469/// stream of `dialed`.
470async 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    /// A version-2 frame on the control stream: on v5 as it is, on v4
591    /// neighbour-signed with the next seq when the profile signs its type.
592    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    /// A frame on the control stream as it is: a data frame, which carries
612    /// its own end-to-end signature and no neighbour signature.
613    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    /// Ends the link once, with `err`: pending calls, subscriptions, served
620    /// procedures and streams end with it, and the connection closes.
621    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
659/// Sends each statement the issuer reissues as a STATUS, until the link
660/// ends. A STATUS carries no neighbour signature and takes no seq.
661async 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
685/// Reads the station's frames until the link ends.
686async 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
705/// One frame from the station: a STATUS renews the station's statement, a
706/// liveness frame is answered or handed to the probe, every other frame is
707/// opened (on v4 from its neighbour signature at the next seq), and a GOODBYE
708/// ends the link.
709fn 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
747/// A received frame as the link reads it: on v5 with no neighbour signature,
748/// on v4 opened from its neighbour signature at the next seq.
749fn 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
773/// A liveness frame as a link of `version` reads it: its kind and nonce, or
774/// `None` for any other frame. It exists only on v5, and there it is read as
775/// every v5 frame is, so one that carries a neighbour signature ends the link
776/// as malformed, as on v4 any liveness frame does.
777fn 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
791/// Answers the station's liveness_ping with a liveness_pong of the same
792/// nonce, and hands a liveness_pong to the probe. Neither goes further.
793fn liveness(
794    inner: &Arc<Inner>,
795    kind: Liveness,
796    nonce: [u8; frame::LIVENESS_NONCE_SIZE],
797) -> Result<(), LinkError> {
798    match kind {
799        // A ping while a pong is still being written is not answered again:
800        // the station's probe needs one answer per interval, not one per
801        // ping, and a burst of pings starts one task, not one each.
802        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
822/// Takes the one pong slot: true when no pong was in flight, and it is now.
823fn pong_slot(in_flight: &AtomicBool) -> bool {
824    !in_flight.swap(true, Ordering::AcqRel)
825}
826
827/// Ends the link when the station's statement lapses past the grace, or its
828/// TLS binding reaches its not_after.
829async 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
852/// Ends the link when its QUIC connection closes.
853async 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    /// macula-rust#9: a burst of pings takes one pong slot until its pong is
887    /// written.
888    #[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    /// macula-go 2294f25: liveness frames are read after the session-mode
898    /// check, so a neighbour-signed one ends a v5 link.
899    #[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}