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