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::{
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
67/// How long the handshake may take, as macula's 30 seconds.
68pub const HANDSHAKE_TIMEOUT: Duration = Duration::from_secs(30);
69
70/// How long Close waits, after its GOODBYE, for the station to close the
71/// connection before closing it itself.
72const CLOSE_LINGER: Duration = Duration::from_secs(1);
73
74/// How long past a statement's expiry the link keeps a station whose next
75/// statement has not arrived: macula's five minutes.
76const STATUS_GRACE_MS: i64 = 5 * 60 * 1000;
77
78/// Why a link or one of its operations failed.
79#[derive(Debug, Clone, PartialEq, Eq)]
80pub enum LinkError {
81    /// The configuration cannot be dialed: a key of another profile than the
82    /// target's, or admission limits macula would not start with.
83    InvalidConfig(String),
84    /// The dial failed.
85    Dial(String),
86    /// The handshake refused, or was refused.
87    Handshake(HandshakeError),
88    /// The handshake did not finish within [`HANDSHAKE_TIMEOUT`].
89    HandshakeTimeout,
90    /// The issuer had no CONNECT material in force.
91    Issuer(IssuerError),
92    /// The QUIC connection or a stream failed.
93    Io(String),
94    /// A frame longer than the cap in force.
95    FrameTooLarge(usize),
96    /// A frame that could not be built, or one received that does not
97    /// verify.
98    Frame(FrameError),
99    /// A record that could not be built, or one found that does not verify.
100    Record(RecordError),
101    /// The station's status statement lapsed past the grace.
102    StatusExpired,
103    /// The station's TLS binding reached its not_after.
104    BindingExpired,
105    /// The link was closed by its owner.
106    Closed,
107    /// The station ended the link with a GOODBYE, for this reason.
108    Goodbye(String),
109    /// The station missed two liveness probes in a row.
110    LivenessLost,
111    /// A station this process completed handshake v5 with that now answers
112    /// v4, or a v4 link to it that a v5 completion ended: refused, not
113    /// retried, until [`forget_v5_peer`].
114    V5DowngradeRefused,
115    /// No verified reply answered the call within its timeout.
116    CallTimeout,
117    /// A provider's ERROR for the call.
118    Provider {
119        responded_by: [u8; 32],
120        code: String,
121        detail: Option<String>,
122    },
123    /// A relay error the connected station reported for the call.
124    Relay { reported_by: [u8; 32], code: String },
125    /// A find_record the station answered not_found.
126    RecordNotFound,
127    /// A station reply to a DHT call in a shape that call never answers with.
128    UnexpectedReply(String),
129    /// An offer without exactly one handler, or an org procedure's without
130    /// the realm key.
131    InvalidOffer,
132    /// A procedure without a namespace, which nobody can authorize.
133    NoOrg,
134    /// A procedure already served on this link.
135    AlreadyServed,
136    /// A required-confidential procedure on a link whose `kem_advertise` is
137    /// off: it would name no key and refuse every clear call.
138    KemAdvertiseDisabled,
139    /// A served procedure withdrawn by its owner.
140    Stopped,
141    /// A stream ended by a STREAM_ERROR: the peer's, the station's relay
142    /// error (`relay`), or this side's own abort.
143    Stream {
144        code: String,
145        message: String,
146        relay: bool,
147    },
148    /// A stream that ended normally.
149    EndOfStream,
150    /// A send on a stream this side has ended.
151    StreamClosed,
152    /// A STREAM_OPEN over the 1 MiB a peer reads of one.
153    StreamOpenTooLarge(usize),
154    /// A call or an open that could not be kept confidential, so it was not
155    /// made, or failed rather than be taken in the clear.
156    Confidentiality(ConfidentialityError),
157    /// The provider's sealed_refused: it could not open the request. `named`
158    /// is the key it holds now, `None` when it holds none.
159    SealedRefused { named: Option<[u8; 8]> },
160    /// A clear answer to a sealed request that nothing clear may give: not a
161    /// relay error, not a refusal from the closed set, not sealed_refused.
162    /// Refused, never taken as the answer.
163    ClearAnswerToSealed,
164    /// A provider's sealed stream that has sealed all the frames GCM's
165    /// bound allows under random nonces: another would pass it.
166    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
237/// What a link is dialed with: the station to reach, the node's identity
238/// key, the statement issuer that holds its CONNECT key and statements, and
239/// the realm membership endorsement to present, empty for none. Links of one
240/// node share their publication seq, admission and dedup; a link given none
241/// makes its own.
242pub 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    /// This link's place in the admission: the station it dialed, host:port,
251    /// when `None`.
252    pub share: Option<String>,
253    /// This identity's KEM keys (macula 13, E2E design amendment A1): with
254    /// them the link opens sealed requests; without them it refuses them,
255    /// holding no key. Links of one identity share one.
256    pub keyring: Option<Arc<Keyring>>,
257    /// Names the keyring's current key in the advertisements of procedures
258    /// served confidentially ([`Confidentiality::Preferred`] or
259    /// [`Confidentiality::Required`]). Off by default: switch it on only once
260    /// every station runs macula 12.11 or later, which stores and routes a
261    /// keyed advertisement, and every caller runs 13. It needs a keyring.
262    pub kem_advertise: bool,
263}
264
265impl Config {
266    /// A link's configuration with none of the shared parts.
267    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/// A handshaked link to one station. Cloning it shares the link.
284#[derive(Clone)]
285pub struct Link {
286    inner: Arc<Inner>,
287}
288
289struct Inner {
290    /// Distinct for every link a process dials.
291    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    /// The handshake version the link completed: 4 or 5.
302    version: i64,
303    /// v5 liveness answers, for the probe.
304    pongs: tokio::sync::mpsc::Sender<[u8; frame::LIVENESS_NONCE_SIZE]>,
305    /// A liveness_pong is being written (macula-rust#9): at most one is in
306    /// flight, so a station's pings cannot pile up tasks on this link.
307    pong_in_flight: AtomicBool,
308    pongs_rx: tokio::sync::Mutex<tokio::sync::mpsc::Receiver<[u8; frame::LIVENESS_NONCE_SIZE]>>,
309    /// Orders a neighbour signature's seq with its write.
310    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    /// The owner is closing the link: however it then ends, it ends
326    /// [`LinkError::Closed`].
327    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    /// Dials `cfg.target`, runs the handshake as a client, and returns the
337    /// link once the station's HELLO accepts it. It dials with version 5
338    /// unless the station refused v5 in the last 10 minutes; a station never
339    /// seen on v5 that refuses v5 with `unsupported_version` is dialled once
340    /// more, on a new connection, with v4. A station seen on v5 that answers
341    /// v4 is [`LinkError::V5DowngradeRefused`]. Each handshake is bounded by
342    /// [`HANDSHAKE_TIMEOUT`].
343    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    /// The handshake version the link completed: 4 or 5.
385    pub fn handshake_version(&self) -> i64 {
386        self.inner.version
387    }
388
389    /// The node_id of the station the link reached.
390    pub fn station_node_id(&self) -> [u8; 32] {
391        self.inner.station.node_id
392    }
393
394    /// A number distinct for every link this process dials.
395    pub fn serial(&self) -> u64 {
396        self.inner.serial
397    }
398
399    /// The node_id this link connected as.
400    pub fn node_id(&self) -> [u8; 32] {
401        self.inner.self_id
402    }
403
404    /// The capability bits the station's HELLO announced.
405    pub fn station_capabilities(&self) -> u64 {
406        self.inner.station_capabilities
407    }
408
409    /// The profile the link runs.
410    pub fn profile(&self) -> Profile {
411        self.inner.profile
412    }
413
414    /// Why the link ended, or `None` while it runs.
415    pub fn error(&self) -> Option<LinkError> {
416        self.inner.lock().ended.clone()
417    }
418
419    /// Waits until the link has ended, and says why.
420    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    /// The frames received that nothing on this link handles, by what they
427    /// were.
428    pub fn unrouted(&self) -> HashMap<String, u64> {
429        self.inner.lock().unrouted.clone()
430    }
431
432    /// Sends a GOODBYE with `reason`, closes the control stream's sending
433    /// side, waits up to a second for the station to close the connection,
434    /// and ends the link. Closing an ended link does nothing.
435    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
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 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
518/// The TLS exporter of `connection`, as the v5 handshake asks for keying
519/// material.
520fn 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
532/// What the client brings to the handshake in `version`: the configured
533/// identity and target, the leaf it received, and the issuer's CONNECT
534/// material, at the time now.
535fn 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
558/// The state of a link whose handshake completed with `station` on the
559/// control stream `control` of `dialed`, numbered with the next serial.
560fn 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
616/// Starts the tasks a linked link runs: its statements, its control
617/// stream's reader, the streams it accepts, its liveness probe, and the
618/// watches on its expiries and its connection.
619fn 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    /// Marks the link closing, so it ends as closed: false when it has
644    /// already ended.
645    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    /// A version-2 frame on the control stream: on v5 as it is, on v4
655    /// neighbour-signed with the next seq when the profile signs its type.
656    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    /// A frame on the control stream as it is: a data frame, which carries
676    /// its own end-to-end signature and no neighbour signature.
677    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    /// Ends the link once, with `err`: pending calls, subscriptions, served
684    /// procedures and streams end with it, and the connection closes.
685    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    /// Marks the link ended once, with `err`, or as closed when its owner was
715    /// closing it: the error it ended with, or `None` when it had already
716    /// ended.
717    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
727/// Sends each statement the issuer reissues as a STATUS, until the link
728/// ends. A STATUS carries no neighbour signature and takes no seq.
729async 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
753/// Reads the station's frames until the link ends.
754async 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
773/// One frame from the station: a STATUS renews the station's statement, a
774/// liveness frame is answered or handed to the probe, every other frame is
775/// opened (on v4 from its neighbour signature at the next seq), and a GOODBYE
776/// ends the link.
777fn 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
815/// A received frame as the link reads it: on v5 with no neighbour signature,
816/// on v4 opened from its neighbour signature at the next seq.
817fn 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
841/// A liveness frame as a link of `version` reads it: its kind and nonce, or
842/// `None` for any other frame. It exists only on v5, and there it is read as
843/// every v5 frame is, so one that carries a neighbour signature ends the link
844/// as malformed, as on v4 any liveness frame does.
845fn 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
859/// Answers the station's liveness_ping with a liveness_pong of the same
860/// nonce, and hands a liveness_pong to the probe. Neither goes further.
861fn liveness(
862    inner: &Arc<Inner>,
863    kind: Liveness,
864    nonce: [u8; frame::LIVENESS_NONCE_SIZE],
865) -> Result<(), LinkError> {
866    match kind {
867        // A ping while a pong is still being written is not answered again:
868        // the station's probe needs one answer per interval, not one per
869        // ping, and a burst of pings starts one task, not one each.
870        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
881/// Writes the liveness_pong of `nonce`, gives back the pong slot, and ends
882/// the link when the write failed.
883async 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
893/// Takes the one pong slot: true when no pong was in flight, and it is now.
894fn pong_slot(in_flight: &AtomicBool) -> bool {
895    !in_flight.swap(true, Ordering::AcqRel)
896}
897
898/// Ends the link when the station's statement lapses past the grace, or its
899/// TLS binding reaches its not_after.
900async 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
923/// Ends the link when its QUIC connection closes.
924async 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    /// macula-rust#9: a burst of pings takes one pong slot until its pong is
958    /// written.
959    #[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    /// macula-go 2294f25: liveness frames are read after the session-mode
969    /// check, so a neighbour-signed one ends a v5 link.
970    #[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}