Skip to main content

macula_rust/station_link/
stream.rs

1//! Streaming RPC, as macula 12's link does it. Each session has a QUIC stream
2//! of its own: the caller opens it with a signed STREAM_OPEN naming its mode,
3//! the station opens one of its own to the provider and relays between them.
4//! After the open, each side sends frames signed by its own key (the
5//! provider's under MACULA-PQ-STREAM-V1, the caller's under
6//! MACULA-PQ-CALLER-STREAM-V1), each numbered from 0 on its side and bound to
7//! the open's request hash.
8//!
9//! A stream is released, both its QUIC directions finished, on every path:
10//! when it ends normally, when either side aborts or refuses it, when its
11//! inbox is over its bound, when its link ends, and when an open or an
12//! accepted stream fails before a session exists. The bounds are macula
13//! 12.3.0's: an open of at most 1 MiB, read within 10 seconds of a stream
14//! being opened to the provider, and at most 16 MiB of a stream's frames
15//! received and not yet read.
16//!
17//! A caller's sealed stream a pool opened reopens itself once when the
18//! provider refuses its open sealed_refused before this side has sent
19//! anything, after the provider's KEM key rotated: sealed to the key the
20//! refusal names, on a new session, behind the same handle; never in the
21//! clear (macula's stream reseal, E2E design Amendment A1, macula-rust#21).
22
23use std::collections::VecDeque;
24use std::future::Future;
25use std::pin::Pin;
26use std::sync::{Arc, Mutex, MutexGuard, Weak};
27use std::time::Duration;
28
29use tokio::sync::{watch, Notify};
30
31use crate::cbor::{self, Value};
32use crate::frame::{
33    self, RequestSpec, StreamEncoding, StreamFields, StreamMode, StreamRole, StreamState,
34    VerifiedRequest,
35};
36use crate::seal::KEY_ID_SIZE;
37
38use super::admission::{Admission, SessionPlace, Verdict};
39use super::authorize::authorize;
40use super::confidential::{
41    clear_allowed, opened_request, refused_key, sealed_request, stated, unsealed, Seal, StreamSeal,
42    CODE_SEALED_REFUSED, CODE_SEALED_REQUIRED,
43};
44use super::framing::{read_frame, FrameWriter, MAX_FRAME_BYTES};
45use super::serve::{
46    bounded_detail, without_caller, BoxFuture, Offer, StreamOffer, CODE_REQUEST_COPY,
47};
48use super::{frame_type_of, now_ms, Inner, Link, LinkError};
49
50const STREAM_OPEN_BYTES: usize = 1024 * 1024;
51const STREAM_OPEN_WAIT: Duration = Duration::from_secs(10);
52const STREAM_INBOX: usize = 16 * 1024 * 1024;
53
54/// How far ahead a STREAM_OPEN's deadline lies when its [`StreamCall`] names
55/// none, as macula's default.
56pub const DEFAULT_STREAM_DEADLINE: Duration = Duration::from_secs(30);
57
58/// The refusal codes of a STREAM_OPEN, besides the admission's own, as
59/// macula's refuse_open sends them, and the code of a failed handler.
60const CODE_STREAM_NOT_FOUND: &str = "not_found";
61const CODE_MODE_MISMATCH: &str = "mode_mismatch";
62const CODE_TOO_MANY_SESSIONS: &str = "too_many_sessions";
63const CODE_STREAM_HANDLER_ERROR: &str = "error";
64
65/// Serves one streaming session. When it returns `Ok` and has not ended the
66/// stream, the stream is closed on both sides; an `Err` or a panic aborts it
67/// with code `error` and the error's text, as macula aborts a stream whose
68/// handler failed. A handler still running when its stream ends is dropped.
69pub type StreamHandler = Arc<dyn Fn(Stream) -> BoxFuture<Result<(), String>> + Send + Sync>;
70
71/// A [`StreamHandler`] from an async closure.
72pub fn stream_handler<F, Fut>(f: F) -> StreamHandler
73where
74    F: Fn(Stream) -> Fut + Send + Sync + 'static,
75    Fut: Future<Output = Result<(), String>> + Send + 'static,
76{
77    Arc::new(move |s| Box::pin(f(s)))
78}
79
80/// A streaming session to open: the realm and procedure, the provider it
81/// targets, the mode, the open's payload, how far ahead its deadline lies
82/// ([`DEFAULT_STREAM_DEADLINE`] when zero), a UCAN and its proofs for a gated
83/// procedure, and how it is kept, which an open must state as a call does
84/// (see [`super::Call`]). Its default mode is server_stream.
85#[derive(Debug, Clone, PartialEq)]
86pub struct StreamCall {
87    pub realm: [u8; 32],
88    pub procedure: String,
89    pub target: [u8; 32],
90    pub mode: StreamMode,
91    pub payload: Value,
92    pub deadline: Duration,
93    pub token: Option<Vec<u8>>,
94    pub proofs: Vec<Vec<u8>>,
95    pub seal: Option<Seal>,
96}
97
98impl Default for StreamCall {
99    fn default() -> Self {
100        StreamCall {
101            realm: [0; 32],
102            procedure: String::new(),
103            target: [0; 32],
104            mode: StreamMode::ServerStream,
105            payload: Value::Map(Vec::new()),
106            deadline: Duration::ZERO,
107            token: None,
108            proofs: Vec::new(),
109            seal: None,
110        }
111    }
112}
113
114/// One frame the peer sent, verified: a chunk, the peer's end (role `Send`
115/// ends its sending only, `Both` the stream), or the provider's terminal
116/// value.
117#[derive(Debug, Clone, PartialEq)]
118pub enum StreamEvent {
119    Data {
120        encoding: StreamEncoding,
121        body: Value,
122    },
123    End {
124        role: StreamRole,
125    },
126    Reply {
127        payload: Value,
128    },
129}
130
131/// One streaming session, on either side. Cloning it shares the session.
132#[derive(Clone)]
133pub struct Stream {
134    pub(super) inner: Arc<StreamInner>,
135}
136
137/// What a caller's sealed stream runs, once, when its open is refused
138/// sealed_refused before it has sent anything: given the key id the refusal
139/// names (`None` for none), the stream reopened sealed to that key, or why
140/// there is none. It never opens in the clear.
141pub(crate) type Reseal = Box<
142    dyn FnOnce(
143            Option<[u8; KEY_ID_SIZE]>,
144        ) -> Pin<Box<dyn Future<Output = Result<Stream, LinkError>> + Send>>
145        + Send,
146>;
147
148/// What a served stream's inbox holds is charged to its caller's budget in
149/// the node's admission, and the session holds its place there.
150struct Budget {
151    admission: Arc<Admission>,
152    caller: [u8; 32],
153    place: Option<SessionPlace>,
154}
155
156pub(super) struct StreamInner {
157    link: Arc<Inner>,
158    writer: FrameWriter,
159    pub(super) open: VerifiedRequest,
160    pub(super) caller: bool,
161    /// This side's keys when the stream is sealed.
162    pub(super) sealing: Option<StreamSeal>,
163    /// Orders a frame's seq with its write.
164    send_seq: tokio::sync::Mutex<u64>,
165    state: Mutex<StreamSide>,
166    budget: Mutex<Option<Budget>>,
167    notify: Notify,
168    done_tx: watch::Sender<bool>,
169    /// A caller's sealed stream's one reseal, until it is spent.
170    reseal: Mutex<Option<Reseal>>,
171}
172
173#[derive(Default)]
174pub(super) struct StreamSide {
175    /// The reseal is running: this session reads and writes no more.
176    reopening: bool,
177    /// The session the stream reopened on, which carries it from now on.
178    successor: Option<Arc<StreamInner>>,
179    /// This side sent its last frame, or will send no more.
180    sent_end: bool,
181    /// The peer sent its last frame.
182    peer_ended: bool,
183    inbox: VecDeque<(StreamEvent, usize)>,
184    held: usize,
185    ended: bool,
186    err: Option<LinkError>,
187    /// A caller's seal report has settled (see [`Stream::report`]).
188    pub(super) settled: bool,
189}
190
191impl Stream {
192    /// The stream's verified STREAM_OPEN: its caller, procedure, mode and
193    /// payload. On a provider's side of a sealed stream the payload is the
194    /// open's opened plaintext, and on a provider's side a map payload has
195    /// no text "caller" key: the caller is `caller`, as verified.
196    pub fn request(&self) -> &VerifiedRequest {
197        &self.inner.open
198    }
199
200    /// Whether the stream is sealed end to end: every chunk, reply and
201    /// error after the open travels sealed.
202    pub fn sealed(&self) -> bool {
203        self.inner.sealing.is_some()
204    }
205
206    /// Sends a raw chunk.
207    pub async fn send(&self, body: &[u8]) -> Result<(), LinkError> {
208        let at = |seq| StreamFields::Data {
209            seq,
210            encoding: StreamEncoding::Raw,
211            body: Value::Bytes(body.to_vec()),
212        };
213        self.sent(at, false).await.1
214    }
215
216    /// Sends a structured chunk.
217    pub async fn send_value(&self, v: Value) -> Result<(), LinkError> {
218        let at = |seq| StreamFields::Data {
219            seq,
220            encoding: StreamEncoding::Msgpack,
221            body: v.clone(),
222        };
223        self.sent(at, false).await.1
224    }
225
226    /// Ends this side's sending; the peer may still send.
227    pub async fn close_send(&self) -> Result<(), LinkError> {
228        let at = |seq| StreamFields::End {
229            seq,
230            role: StreamRole::Send,
231        };
232        self.sent(at, true).await.1
233    }
234
235    /// Ends the stream on both sides.
236    pub async fn close(&self) -> Result<(), LinkError> {
237        let at = |seq| StreamFields::End {
238            seq,
239            role: StreamRole::Both,
240        };
241        let (inner, sent) = self.sent(at, true).await;
242        StreamInner::end(&inner, None);
243        sent
244    }
245
246    /// Sends the provider's terminal value and ends the stream.
247    pub async fn reply(&self, payload: Value) -> Result<(), LinkError> {
248        let at = |seq| StreamFields::Reply {
249            seq,
250            payload: payload.clone(),
251        };
252        let (inner, sent) = self.sent(at, true).await;
253        StreamInner::end(&inner, None);
254        sent
255    }
256
257    /// Ends the stream with a STREAM_ERROR of `code` and `message`.
258    pub async fn abort(&self, code: &str, message: &str) -> Result<(), LinkError> {
259        let at = |seq| StreamFields::Error {
260            seq,
261            code: code.to_string(),
262            message: message.to_string(),
263        };
264        let (inner, sent) = self.sent(at, true).await;
265        StreamInner::end(&inner, Some(aborted(code, message)));
266        sent
267    }
268
269    /// The next frame the peer sent. After the stream ends, once every event
270    /// before it is read, it returns [`LinkError::EndOfStream`] for a normal
271    /// end and the error that ended it otherwise.
272    pub async fn recv(&self) -> Result<StreamEvent, LinkError> {
273        loop {
274            let inner = self.live().await;
275            let notified = inner.notify.notified();
276            match inner.next_event() {
277                Some(outcome) => return outcome,
278                None => notified.await,
279            }
280        }
281    }
282
283    /// Waits until the stream has ended and been released, and says why:
284    /// `None` for a normal end.
285    pub async fn done(&self) -> Option<LinkError> {
286        let mut inner = self.live().await;
287        while !inner.released().await {
288            inner = self.live().await;
289        }
290        let err = inner.side().err.clone();
291        err
292    }
293
294    /// The fields `at` builds, sent on the session that carries the stream
295    /// now, and that session.
296    async fn sent(
297        &self,
298        at: impl Fn(u64) -> StreamFields,
299        last: bool,
300    ) -> (Arc<StreamInner>, Result<(), LinkError>) {
301        let mut inner = self.live().await;
302        let mut sent = inner.send_here(&at, last).await;
303        while sent.is_none() {
304            inner = self.live().await;
305            sent = inner.send_here(&at, last).await;
306        }
307        (inner, sent.unwrap_or(Err(LinkError::StreamClosed)))
308    }
309
310    /// The session that carries the stream: the one it opened on, or the
311    /// one it reopened on after a reseal, waited for while the reseal runs.
312    async fn live(&self) -> Arc<StreamInner> {
313        let mut at = self.current();
314        while at.side().reopening {
315            at.reopen_awaited().await;
316            at = self.current();
317        }
318        at
319    }
320
321    /// The session that carries the stream now, without waiting.
322    pub(super) fn current(&self) -> Arc<StreamInner> {
323        let mut at = self.inner.clone();
324        loop {
325            let next = at.side().successor.clone();
326            match next {
327                Some(next) => at = next,
328                None => return at,
329            }
330        }
331    }
332}
333
334/// The error a stream this side aborted ends with.
335fn aborted(code: &str, message: &str) -> LinkError {
336    LinkError::Stream {
337        code: code.to_string(),
338        message: message.to_string(),
339        relay: false,
340    }
341}
342
343impl StreamInner {
344    fn new(
345        link: Arc<Inner>,
346        send: quinn::SendStream,
347        open: VerifiedRequest,
348        caller: bool,
349        sealing: Option<StreamSeal>,
350    ) -> Arc<StreamInner> {
351        Arc::new(StreamInner {
352            link,
353            writer: FrameWriter::new(send),
354            open,
355            caller,
356            sealing,
357            send_seq: tokio::sync::Mutex::new(0),
358            state: Mutex::new(StreamSide::default()),
359            budget: Mutex::new(None),
360            notify: Notify::new(),
361            done_tx: watch::channel(false).0,
362            reseal: Mutex::new(None),
363        })
364    }
365
366    pub(super) fn side(&self) -> MutexGuard<'_, StreamSide> {
367        self.state.lock().unwrap_or_else(|p| p.into_inner())
368    }
369
370    fn budget(&self) -> MutexGuard<'_, Option<Budget>> {
371        self.budget.lock().unwrap_or_else(|p| p.into_inner())
372    }
373
374    /// The inbox's next event, its room given back; once the stream has
375    /// ended and nothing is queued, why it ended ([`LinkError::EndOfStream`]
376    /// for a normal end); `None` while there is nothing yet.
377    fn next_event(&self) -> Option<Result<StreamEvent, LinkError>> {
378        let mut side = self.side();
379        if let Some((event, size)) = side.inbox.pop_front() {
380            side.held -= size;
381            drop(side);
382            self.release_inbox(size);
383            return Some(Ok(event));
384        }
385        if side.ended {
386            return Some(Err(side.err.clone().unwrap_or(LinkError::EndOfStream)));
387        }
388        None
389    }
390
391    /// Signs the fields `at(seq)` builds at this side's next seq and writes
392    /// them; `last` marks this side's last frame, after which its QUIC
393    /// direction is finished. On a sealed stream the seq is spent before
394    /// anything is sealed under it: a frame that then fails to go out ends
395    /// this side's sending, so nothing is ever sealed twice under one
396    /// (key, seq).
397    async fn send(
398        self: &Arc<Self>,
399        at: impl FnOnce(u64) -> StreamFields,
400        last: bool,
401    ) -> Result<(), LinkError> {
402        let seq = self.send_seq.lock().await;
403        self.send_at(seq, at, last).await
404    }
405
406    /// [`StreamInner::send`], unless a reseal has taken the stream off this
407    /// session: then `None`, nothing sent, and the frame goes on the
408    /// session it reopens on.
409    async fn send_here(
410        self: &Arc<Self>,
411        at: &impl Fn(u64) -> StreamFields,
412        last: bool,
413    ) -> Option<Result<(), LinkError>> {
414        let seq = self.send_seq.lock().await;
415        if self.side().reopening || self.side().successor.is_some() {
416            return None;
417        }
418        Some(self.send_at(seq, at, last).await)
419    }
420
421    async fn send_at(
422        self: &Arc<Self>,
423        mut seq: tokio::sync::MutexGuard<'_, u64>,
424        at: impl FnOnce(u64) -> StreamFields,
425        last: bool,
426    ) -> Result<(), LinkError> {
427        if self.side().sent_end {
428            return Err(LinkError::StreamClosed);
429        }
430        let fields = at(*seq);
431        let plain = match &self.sealing {
432            Some(sealing) => sealing.plain_of(&fields)?,
433            None => None,
434        };
435        let spent = plain.is_some();
436        let fields = match (&self.sealing, plain) {
437            (Some(sealing), Some(plain)) => {
438                *seq += 1;
439                let sealed = sealing.sealed(fields, plain);
440                sealed.inspect_err(|_| self.side().sent_end = true)?
441            }
442            _ => fields,
443        };
444        let encoded = self
445            .signed(&fields)
446            .inspect_err(|_| self.unsent_after_spending(spent))?;
447        if let Err(e) = self.writer.write(&encoded, MAX_FRAME_BYTES).await {
448            self.side().sent_end = true;
449            return Err(e);
450        }
451        if !spent {
452            *seq += 1;
453        }
454        if last {
455            self.sent_last().await;
456        }
457        Ok(())
458    }
459
460    /// Ends this side's sending when a frame that spent its seq did not go
461    /// out, so nothing is ever sealed twice under one (key, seq).
462    fn unsent_after_spending(&self, spent: bool) {
463        if spent {
464            self.side().sent_end = true;
465        }
466    }
467
468    /// After this side's last frame: its QUIC direction finished, and the
469    /// stream ended when the peer had ended too.
470    async fn sent_last(self: &Arc<Self>) {
471        let peer_ended = {
472            let mut side = self.side();
473            side.sent_end = true;
474            side.peer_ended
475        };
476        self.writer.finish().await;
477        if peer_ended {
478            StreamInner::end(self, None);
479        }
480    }
481
482    /// `fields` signed by this side's key and encoded.
483    fn signed(&self, fields: &StreamFields) -> Result<Vec<u8>, LinkError> {
484        let signed = if self.caller {
485            frame::sign_caller_stream(fields, &self.open, &self.link.key)?
486        } else {
487            frame::sign_provider_stream(fields, &self.open, &self.link.key)?
488        };
489        cbor::encode(&signed)
490            .map_err(|e| LinkError::Frame(frame::FrameError::Payload(e.to_string())))
491    }
492
493    async fn abort(self: &Arc<Self>, code: &str, message: &str) -> Result<(), LinkError> {
494        let sent = self
495            .send(
496                |seq| StreamFields::Error {
497                    seq,
498                    code: code.to_string(),
499                    message: message.to_string(),
500                },
501                true,
502            )
503            .await;
504        StreamInner::end(self, Some(aborted(code, message)));
505        sent
506    }
507
508    /// Waits until this session has ended and been released: true when the
509    /// stream ended with it, false when it went on on a reopened session.
510    async fn released(&self) -> bool {
511        let mut done = self.done_tx.subscribe();
512        let _ = done.wait_for(|ended| *ended).await;
513        self.side().successor.is_none()
514    }
515
516    /// Returns once this session's reseal has ended, either way.
517    async fn reopen_awaited(&self) {
518        let notified = self.notify.notified();
519        // Checked again once registered, so the reseal's end is not missed
520        // between the two.
521        if self.side().reopening {
522            notified.await;
523        }
524    }
525
526    fn reseal_slot(&self) -> MutexGuard<'_, Option<Reseal>> {
527        self.reseal.lock().unwrap_or_else(|p| p.into_inner())
528    }
529
530    /// The provider refused the open sealed_refused, naming `detail`'s key:
531    /// a caller's sealed stream holding its reseal, and that has sent
532    /// nothing, starts it and reads no more on this session (true).
533    /// Checked and marked under the send lock, so no frame goes out on this
534    /// session after it. Otherwise the refusal ends the stream (false).
535    async fn reopened(self: &Arc<Self>, detail: &str) -> bool {
536        let Some(reseal) = self.reseal_slot().take() else {
537            return false;
538        };
539        let seq = self.send_seq.lock().await;
540        if *seq != 0 || self.side().sent_end {
541            return false;
542        }
543        self.side().reopening = true;
544        drop(seq);
545        self.notify.notify_waiters();
546        let s = self.clone();
547        let named = refused_key(Some(detail));
548        tokio::spawn(async move {
549            let outcome = reseal(named).await;
550            s.adopt(outcome);
551        });
552        true
553    }
554
555    /// Carries the stream on the session its reseal reopened, this one
556    /// released; or ends it with why there is none.
557    fn adopt(self: &Arc<Self>, outcome: Result<Stream, LinkError>) {
558        let reopened = {
559            let mut side = self.side();
560            side.reopening = false;
561            outcome.map(|next| side.successor = Some(next.inner))
562        };
563        match reopened {
564            Ok(()) => StreamInner::end(self, None),
565            Err(e) => self.peer_finished(Some(e)),
566        }
567        // Whatever waits on the reseal goes on, even when this session had
568        // ended before (its link ended meanwhile), when end wakes nobody.
569        self.notify.notify_waiters();
570    }
571
572    /// Queues `event` for recv, refusing it when it would take the inbox, or
573    /// the node's budget for served streams, past its bound.
574    fn deliver(&self, event: StreamEvent, size: usize) -> bool {
575        if !self.queue(event, size) {
576            return false;
577        }
578        self.notify.notify_one();
579        true
580    }
581
582    /// Puts `event` in the inbox, unless it would take the inbox, or the
583    /// node's budget for served streams, past its bound.
584    fn queue(&self, event: StreamEvent, size: usize) -> bool {
585        let mut side = self.side();
586        if side.held + size > STREAM_INBOX {
587            return false;
588        }
589        if !self.charge_inbox(size) {
590            return false;
591        }
592        side.inbox.push_back((event, size));
593        side.held += size;
594        true
595    }
596
597    /// Charges `size` to a served stream's caller's inbox budget, and
598    /// whether it fits; a stream with no budget always fits.
599    fn charge_inbox(&self, size: usize) -> bool {
600        let budget = self.budget();
601        let Some(budget) = &*budget else {
602            return true;
603        };
604        budget.admission.charge_inbox(budget.caller, size)
605    }
606
607    fn release_inbox(&self, size: usize) {
608        if let Some(budget) = &*self.budget() {
609            budget.admission.release_inbox(budget.caller, size);
610        }
611    }
612
613    /// Ends the stream after the peer's last frame: this side sends no more,
614    /// and `err` is why it ended, `None` for a normal end.
615    fn peer_finished(self: &Arc<Self>, err: Option<LinkError>) {
616        self.side().peer_ended = true;
617        StreamInner::end(self, err);
618    }
619
620    /// Aborts the stream from this side for a fault it found in what the
621    /// peer sent, or an inbox over its bound, telling the peer when it still
622    /// can.
623    async fn fail(self: &Arc<Self>, code: &str, cause: Option<String>) {
624        let message = cause
625            .as_deref()
626            .map(bounded_detail)
627            .unwrap_or("")
628            .to_string();
629        let _ = self.abort(code, &message).await;
630    }
631
632    /// Releases the stream once: its sending side finished after this side's
633    /// last frame and reset otherwise, its reader stopped, what its inbox
634    /// held and its session's place given back.
635    pub(super) fn end(this: &Arc<StreamInner>, err: Option<LinkError>) {
636        let Some(graceful) = this.mark_ended(err) else {
637            return;
638        };
639        if !graceful {
640            StreamInner::reset_sending(this);
641        }
642        if let Some(budget) = this.budget().take() {
643            let held = std::mem::take(&mut this.side().held);
644            budget.admission.release_inbox(budget.caller, held);
645            drop(budget.place);
646        }
647        this.link
648            .lock()
649            .streams
650            .retain(|w| w.strong_count() > 0 && !std::ptr::eq(w.as_ptr(), Arc::as_ptr(this)));
651        let _ = this.done_tx.send_replace(true);
652        this.notify.notify_waiters();
653        this.notify.notify_one();
654    }
655}
656
657impl StreamInner {
658    /// Marks the stream ended, keeping `err` unless an error is kept
659    /// already, and this side's sending with it; whether this side had sent
660    /// its last frame, or `None` when the stream had ended before.
661    fn mark_ended(&self, err: Option<LinkError>) -> Option<bool> {
662        let mut side = self.side();
663        if side.ended {
664            return None;
665        }
666        side.ended = true;
667        if side.err.is_none() {
668            side.err = err;
669        }
670        let graceful = side.sent_end;
671        side.sent_end = true;
672        Some(graceful)
673    }
674
675    /// Resets this side's sending direction, on the runtime when there is
676    /// one.
677    fn reset_sending(this: &Arc<StreamInner>) {
678        let released = this.clone();
679        let Ok(runtime) = tokio::runtime::Handle::try_current() else {
680            return;
681        };
682        runtime.spawn(async move { released.writer.reset().await });
683    }
684}
685
686/// Keeps `s` among the link's streams, ended with the link; false once the
687/// link has ended.
688fn hold_stream(inner: &Inner, s: &Arc<StreamInner>) -> bool {
689    let mut state = inner.lock();
690    if state.ended.is_some() {
691        return false;
692    }
693    state.streams.push(Arc::downgrade(s));
694    true
695}
696
697/// Releases a QUIC stream no session holds, in both directions.
698fn abandon(mut send: quinn::SendStream, mut recv: quinn::RecvStream) {
699    let _ = send.reset(0u32.into());
700    let _ = recv.stop(0u32.into());
701}
702
703impl Link {
704    /// Opens a streaming session: a QUIC stream of its own, on which it
705    /// writes the signed STREAM_OPEN. A stream it opens but cannot write the
706    /// open on is released before the error returns.
707    pub async fn open_stream(&self, c: StreamCall) -> Result<Stream, LinkError> {
708        self.open_stream_resealing(c, None).await
709    }
710
711    /// [`Link::open_stream`], the stream holding `reseal` for a
712    /// sealed_refused of its open (see the module's docs).
713    pub(crate) async fn open_stream_resealing(
714        &self,
715        c: StreamCall,
716        reseal: Option<Reseal>,
717    ) -> Result<Stream, LinkError> {
718        let inner = &self.inner;
719        stated(&c.target, &inner.station.node_id, &c.seal)?;
720        let deadline = if c.deadline.is_zero() {
721            DEFAULT_STREAM_DEADLINE
722        } else {
723            c.deadline
724        };
725        let mut request_id = [0u8; 16];
726        aws_lc_rs::rand::fill(&mut request_id)
727            .map_err(|_| LinkError::Io("no randomness".into()))?;
728        let deadline = (now_ms() + deadline.as_millis() as i64) as u64;
729        let (sealed, sealing) = sealed_open(inner, &c, request_id, deadline)?;
730        let signed = frame::sign_stream_open(
731            &RequestSpec {
732                request_id,
733                realm: c.realm,
734                procedure: c.procedure,
735                target: c.target,
736                deadline,
737                payload: c.payload,
738                sealed,
739                mode: Some(c.mode),
740                token: c.token,
741                proofs: c.proofs,
742                source_route: None,
743                retry_budget: None,
744            },
745            &inner.key,
746        )?;
747        let encoded = cbor::encode(&signed)
748            .map_err(|e| LinkError::Frame(frame::FrameError::Payload(e.to_string())))?;
749        if encoded.len() > STREAM_OPEN_BYTES {
750            return Err(LinkError::StreamOpenTooLarge(encoded.len()));
751        }
752        let open = frame::verify_request(&signed, inner.profile)?;
753        let state = frame::open_stream(&open)?;
754        let (send, recv) = inner
755            .connection
756            .open_bi()
757            .await
758            .map_err(|e| LinkError::Io(format!("open a stream: {e}")))?;
759        let s = StreamInner::new(inner.clone(), send, open, true, sealing);
760        let held = hold_stream(inner, &s);
761        let written = match held {
762            true => s.writer.write(&encoded, STREAM_OPEN_BYTES).await,
763            false => Err(inner.lock().ended.clone().unwrap_or(LinkError::Closed)),
764        };
765        if let Err(e) = written {
766            StreamInner::end(&s, Some(e.clone()));
767            let mut recv = recv;
768            let _ = recv.stop(0u32.into());
769            return Err(e);
770        }
771        *s.reseal_slot() = reseal;
772        tokio::spawn(read(s.clone(), recv, state));
773        Ok(Stream { inner: s })
774    }
775}
776
777/// The open's payload sealed to the key its seal names, with the caller's
778/// keys for the stream; nothing when the stream goes clear.
779fn sealed_open(
780    inner: &Inner,
781    c: &StreamCall,
782    request_id: [u8; 16],
783    deadline: u64,
784) -> Result<(Option<frame::Sealed>, Option<StreamSeal>), LinkError> {
785    match &c.seal {
786        Some(Seal::To(key)) => {
787            let (sealed, s) = sealed_request(
788                inner.profile,
789                key,
790                crate::seal::FRAME_STREAM_OPEN,
791                c.realm,
792                &c.procedure,
793                inner.self_id,
794                c.target,
795                request_id,
796                deadline,
797                &c.payload,
798            )?;
799            Ok((Some(sealed), Some(StreamSeal::caller(&s))))
800        }
801        _ => Ok((None, None)),
802    }
803}
804
805/// Takes each stream the station opens to this link, until the link ends.
806pub(super) async fn accept_streams(link: Weak<Inner>) {
807    let Some(connection) = link.upgrade().map(|l| l.connection.clone()) else {
808        return;
809    };
810    while let Ok((send, recv)) = connection.accept_bi().await {
811        tokio::spawn(incoming(link.clone(), send, recv));
812    }
813}
814
815/// Reads a stream's first frame within 10 seconds and starts the session it
816/// opens, or refuses it. A stream that fails before a session exists is
817/// released: one that does not deliver a STREAM_OPEN in time, whose first
818/// frame is not one, that does not verify or targets another node is dropped
819/// without a word, and one the provider refuses is told why at seq 0.
820async fn incoming(link: Weak<Inner>, send: quinn::SendStream, mut recv: quinn::RecvStream) {
821    let Some(inner) = link.upgrade() else { return };
822    let Some((open, state)) = read_open(&inner, &mut recv).await else {
823        abandon(send, recv);
824        return;
825    };
826    let offer = inner
827        .lock()
828        .served
829        .get(&(open.realm, open.procedure.clone()))
830        .map(|s| s.offer.clone());
831    let refuse_clear = |code: &str, message: &str, send: quinn::SendStream, recv| {
832        let s = StreamInner::new(inner.clone(), send, open.clone(), false, None);
833        let (code, message) = (code.to_string(), message.to_string());
834        async move { refuse(&s, &code, &message, recv).await }
835    };
836    // Before the open is opened, so a caller over its admission or at its
837    // session cap costs no decapsulation: in the clear, from the closed set.
838    let place = match admit_stream(&inner, &open) {
839        Ok(place) => place,
840        Err(code) => return refuse_clear(code, "", send, recv).await,
841    };
842    let (session_open, sealing) = match session_open(&inner, &open, offer.as_ref()) {
843        Ok(opened) => opened,
844        Err((code, message)) => return refuse_clear(code, &message, send, recv).await,
845    };
846    // From here a refusal of a sealed open goes sealed.
847    let s = StreamInner::new(inner.clone(), send, session_open, false, sealing);
848    let Some((offer, policy)) = offer.and_then(|o| Some((o.stream?, o.policy))) else {
849        return refuse(&s, CODE_STREAM_NOT_FOUND, "", recv).await;
850    };
851    if let Some(code) = authorize(inner.profile, policy.as_ref(), &open) {
852        return refuse(&s, code, "", recv).await;
853    }
854    if Some(offer.mode) != open.mode {
855        return refuse(&s, CODE_MODE_MISMATCH, "", recv).await;
856    }
857    *s.budget() = Some(Budget {
858        admission: inner.admission.clone(),
859        caller: open.caller,
860        place: Some(place),
861    });
862    if !hold_stream(&inner, &s) {
863        StreamInner::end(&s, Some(LinkError::Closed));
864        let _ = recv.stop(0u32.into());
865        return;
866    }
867    tokio::spawn(read(s.clone(), recv, state));
868    tokio::spawn(serve(s, offer));
869}
870
871/// Reads a stream's first frame within 10 seconds: its verified STREAM_OPEN
872/// for this node and the verifier state it starts, or `None`, counted, when
873/// the stream does not deliver one.
874async fn read_open(
875    inner: &Inner,
876    recv: &mut quinn::RecvStream,
877) -> Option<(VerifiedRequest, StreamState)> {
878    let payload =
879        match tokio::time::timeout(STREAM_OPEN_WAIT, read_frame(recv, STREAM_OPEN_BYTES)).await {
880            Ok(Ok(payload)) => payload,
881            _ => {
882                inner.count("stream_open_unread");
883                return None;
884            }
885        };
886    let v = match cbor::decode(&payload) {
887        Ok(v) if frame_type_of(&v) == "stream_open" => v,
888        _ => {
889            inner.count("stream_open_malformed");
890            return None;
891        }
892    };
893    let Ok(open) = frame::verify_request(&v, inner.profile) else {
894        inner.count("stream_open_unverified");
895        return None;
896    };
897    if open.target != inner.self_id {
898        inner.count("stream_for_another_node");
899        return None;
900    }
901    let state = frame::open_stream(&open).ok()?;
902    Some((open, state))
903}
904
905/// The open as its session sees it, opened when it came sealed, with the
906/// provider's keys for the stream; or the code and message to refuse it with
907/// in the clear: a sealed open that does not open, or a clear one to a
908/// procedure past its keyless window. A map payload loses a text "caller"
909/// key (macula-rust#13), clear or opened.
910fn session_open(
911    inner: &Inner,
912    open: &VerifiedRequest,
913    offer: Option<&Offer>,
914) -> Result<(VerifiedRequest, Option<StreamSeal>), (&'static str, String)> {
915    match &open.sealed {
916        Some(_) => opened_request(inner.keyring.as_deref(), open)
917            .map(|(payload, sealed)| {
918                (
919                    VerifiedRequest {
920                        payload: without_caller(payload),
921                        ..open.clone()
922                    },
923                    Some(StreamSeal::provider(&sealed)),
924                )
925            })
926            .map_err(|detail| (CODE_SEALED_REFUSED, detail)),
927        None if offer
928            .is_some_and(|o| !clear_allowed(o.confidential, inner.keyed_since(o), now_ms())) =>
929        {
930            let message = "this procedure takes sealed opens only";
931            Err((CODE_SEALED_REQUIRED, message.to_string()))
932        }
933        None => Ok((
934            VerifiedRequest {
935                payload: without_caller(open.payload.clone()),
936                ..open.clone()
937            },
938            None,
939        )),
940    }
941}
942
943/// Admits an open as macula's link does, before anything of it is opened:
944/// one run per request, the deadline window and its bounds, then a place
945/// among the caller's sessions. The place, or the code to refuse with.
946fn admit_stream(inner: &Inner, open: &VerifiedRequest) -> Result<SessionPlace, &'static str> {
947    match inner.admission.admit(open, &inner.share, now_ms()) {
948        Verdict::Refused(code) => return Err(code),
949        Verdict::Copy(_) => return Err(CODE_REQUEST_COPY),
950        Verdict::New => {}
951    }
952    inner
953        .admission
954        .open_session(open.caller)
955        .ok_or(CODE_TOO_MANY_SESSIONS)
956}
957
958/// Answers an open with a STREAM_ERROR of `code` and `message` at seq 0 and
959/// releases the stream.
960async fn refuse(s: &Arc<StreamInner>, code: &str, message: &str, mut recv: quinn::RecvStream) {
961    s.link.count(&format!("stream_refused_{code}"));
962    let _ = s.abort(code, message).await;
963    let _ = recv.stop(0u32.into());
964}
965
966/// Runs the handler for the session, and ends the stream as the handler
967/// leaves it: closed when it returns `Ok` without ending it, aborted with its
968/// error or panic. A handler still running when the stream ends is dropped.
969async fn serve(s: Arc<StreamInner>, offer: StreamOffer) {
970    let stream = Stream { inner: s.clone() };
971    let mut running = tokio::spawn((offer.handler)(stream.clone()));
972    let mut done = s.done_tx.subscribe();
973    let outcome = tokio::select! {
974        outcome = &mut running => outcome,
975        _ = done.wait_for(|ended| *ended) => {
976            running.abort();
977            return;
978        }
979    };
980    match outcome {
981        Ok(Ok(())) => {
982            let _ = stream.close().await;
983        }
984        Ok(Err(e)) => {
985            let _ = stream
986                .abort(CODE_STREAM_HANDLER_ERROR, bounded_detail(&e))
987                .await;
988        }
989        Err(panicked) => {
990            let _ = stream
991                .abort(
992                    CODE_STREAM_HANDLER_ERROR,
993                    bounded_detail(&panicked.to_string()),
994                )
995                .await;
996        }
997    }
998}
999
1000/// Verifies the peer's frames until the stream ends. It is the stream's one
1001/// reader, the only holder of its verifier state; when it returns, its
1002/// receiving side is dropped, which stops it.
1003async fn read(s: Arc<StreamInner>, mut recv: quinn::RecvStream, mut state: StreamState) {
1004    let mut done = s.done_tx.subscribe();
1005    loop {
1006        let payload = tokio::select! {
1007            _ = done.wait_for(|ended| *ended) => return,
1008            payload = read_frame(&mut recv, MAX_FRAME_BYTES) => payload,
1009        };
1010        let payload = match payload {
1011            Ok(payload) => payload,
1012            Err(e) => return read_ended(&s, e),
1013        };
1014        match received(&s, &payload, &state).await {
1015            Some(next) => state = next,
1016            None => return,
1017        }
1018    }
1019}
1020
1021/// Ends a stream whose peer direction finished: after the peer's last frame
1022/// that is expected, and before it the stream was lost.
1023fn read_ended(s: &Arc<StreamInner>, e: LinkError) {
1024    if s.side().peer_ended {
1025        return;
1026    }
1027    let err = s.link.lock().ended.clone().unwrap_or(e);
1028    StreamInner::end(s, Some(err));
1029}
1030
1031/// Handles one frame from the peer; the next verifier state, or `None` when
1032/// reading stops.
1033async fn received(
1034    s: &Arc<StreamInner>,
1035    payload: &[u8],
1036    state: &StreamState,
1037) -> Option<StreamState> {
1038    let v = match cbor::decode(payload) {
1039        Ok(v) => v,
1040        Err(e) => {
1041            s.fail("malformed_frame", Some(e.to_string())).await;
1042            return None;
1043        }
1044    };
1045    if s.caller && v.get("relay_error").is_some() {
1046        relay_failed(s, &v).await;
1047        return None;
1048    }
1049    let verified = if s.caller {
1050        frame::verify_provider_stream(&v, state, s.link.profile)
1051    } else {
1052        frame::verify_caller_stream(&v, state, s.link.profile)
1053    };
1054    let (verified, next) = match verified {
1055        Ok(verified) => verified,
1056        Err(e) => {
1057            s.fail("malformed_frame", Some(e.to_string())).await;
1058            return None;
1059        }
1060    };
1061    let size = payload.len();
1062    let fields = match unsealed(verified, s.sealing.as_ref()) {
1063        Ok(fields) => fields,
1064        Err(e) => {
1065            StreamInner::end(s, Some(e));
1066            return None;
1067        }
1068    };
1069    // Before the frame is delivered, so a recv that returns it sees the
1070    // report settled.
1071    if s.caller && settles(&fields, s.sealing.is_some()) {
1072        s.side().settled = true;
1073    }
1074    taken(s, fields, size, next).await
1075}
1076
1077/// Ends a caller's stream on the station's relay error once it verifies, or
1078/// fails the stream on one that does not.
1079async fn relay_failed(s: &Arc<StreamInner>, v: &Value) {
1080    match frame::verify_relay_error(v, &s.open, s.link.profile, &s.link.station.node_id) {
1081        Ok(relayed) => s.peer_finished(Some(LinkError::Stream {
1082            code: relayed.code,
1083            message: String::new(),
1084            relay: true,
1085        })),
1086        Err(e) => s.fail("malformed_frame", Some(e.to_string())).await,
1087    }
1088}
1089
1090/// Acts on one verified, unsealed frame of `size` bytes from the peer; the
1091/// next verifier state, or `None` when reading stops.
1092async fn taken(
1093    s: &Arc<StreamInner>,
1094    fields: StreamFields,
1095    size: usize,
1096    next: StreamState,
1097) -> Option<StreamState> {
1098    match fields {
1099        StreamFields::Error {
1100            seq: 0,
1101            ref code,
1102            ref message,
1103        } if s.caller && code == CODE_SEALED_REFUSED && s.reopened(message).await => None,
1104        StreamFields::Error { code, message, .. } => {
1105            s.peer_finished(Some(LinkError::Stream {
1106                code,
1107                message,
1108                relay: false,
1109            }));
1110            None
1111        }
1112        StreamFields::Reply { payload, .. } => {
1113            s.deliver(StreamEvent::Reply { payload }, size);
1114            s.peer_finished(None);
1115            None
1116        }
1117        StreamFields::End { role, .. } => {
1118            s.deliver(StreamEvent::End { role }, size);
1119            if role == StreamRole::Both {
1120                s.peer_finished(None);
1121                return None;
1122            }
1123            let mine = {
1124                let mut side = s.side();
1125                side.peer_ended = true;
1126                side.sent_end
1127            };
1128            if mine {
1129                StreamInner::end(s, None);
1130            }
1131            None
1132        }
1133        StreamFields::Data { encoding, body, .. } => {
1134            if !s.deliver(StreamEvent::Data { encoding, body }, size) {
1135                s.fail("resource_exhausted", None).await;
1136                return None;
1137            }
1138            Some(next)
1139        }
1140        // unsealed opens every sealed frame or refuses it.
1141        StreamFields::SealedData { .. }
1142        | StreamFields::SealedError { .. }
1143        | StreamFields::SealedReply { .. } => {
1144            StreamInner::end(s, Some(LinkError::ClearAnswerToSealed));
1145            None
1146        }
1147    }
1148}
1149
1150/// Whether a provider's frame settles its caller's seal report: on a sealed
1151/// stream a data or reply frame, which [`unsealed`] returns only once it
1152/// opened under the stream's key; on a clear stream a data, reply or end
1153/// frame. A sealed stream's end travels clear and settles nothing, and no
1154/// error settles a stream.
1155fn settles(fields: &StreamFields, sealed: bool) -> bool {
1156    match fields {
1157        StreamFields::Data { .. } | StreamFields::Reply { .. } => true,
1158        StreamFields::End { .. } => !sealed,
1159        _ => false,
1160    }
1161}