Skip to main content

macula_rust/
stream.rs

1//! General-purpose streaming RPC, caller/consumer role (§13.1 of
2//! `plans/PLAN_WIRE_PROTOCOL.md`), ported from `macula_stream_sink.erl`.
3//! Like content transfer (`src/content.rs`), this is not a separate wire
4//! mechanism: it runs the frame types built in `src/frame.rs` §13 over a
5//! dedicated QUIC stream, opened via
6//! [`Session::open_dedicated_stream`](crate::connection::Session::open_dedicated_stream)
7//! rather than the control stream.
8//!
9//! **Both roles are built.** Caller/consumer (§13.1) opens a stream and
10//! is the natural fit for pulling/pushing against a procedure that
11//! already exists somewhere. Provider (§13.2) advertises a procedure
12//! (§6.9, [`Session::advertise`](crate::connection::Session::advertise))
13//! and answers inbound STREAM_OPENs the station routes back —
14//! [`Session::accept_dedicated_stream`](crate::connection::Session::accept_dedicated_stream)
15//! accepts the fresh dedicated stream the station opens toward us,
16//! [`StreamHandle::accept`] reads and parses the STREAM_OPEN that's
17//! always its first frame. Both roles end up holding the same
18//! [`StreamHandle`] afterward — a stream's wire vocabulary
19//! (STREAM_DATA/END/ERROR/REPLY) is symmetric regardless of which side
20//! opened it, so `send_data`/`recv`/`close_send`/`abort` all mean the
21//! same thing either way. [`StreamHandle::send_reply`] is the one
22//! provider-only addition: sending the terminal STREAM_REPLY a
23//! `client_stream`/`bidi` caller's own
24//! [`await_reply`](StreamHandle::await_reply) is waiting on.
25//!
26//! Caller/consumer usage, matching the reference's own pattern:
27//! 1. [`StreamHandle::open`] sends STREAM_OPEN and returns a handle once
28//!    the frame is on the wire — there's no open-time acknowledgement to
29//!    wait for; the provider starts reacting to it directly.
30//! 2. Drive a receive loop with [`StreamHandle::recv`] until
31//!    [`StreamItem::Eof`] or an error.
32//! 3. For `client_stream`/`bidi` modes wanting a result:
33//!    [`StreamHandle::send_data`] each chunk in order,
34//!    [`StreamHandle::close_send`] when done, then
35//!    [`StreamHandle::await_reply`].
36//! 4. **Non-normal termination must call [`StreamHandle::abort`], not
37//!    just drop the handle** — the peer's only signal to tell a
38//!    cancellation/failure apart from a dropped connection
39//!    (`plans/PLAN_WIRE_PROTOCOL.md` §13.1, point 4).
40//!
41//! Provider usage:
42//! 1. [`Session::advertise`](crate::connection::Session::advertise) once
43//!    per procedure this session will answer.
44//! 2. Loop on [`StreamHandle::accept`], which blocks for the next
45//!    inbound STREAM_OPEN and hands back a ready-to-use handle plus the
46//!    parsed [`frame::StreamOpenInfo`] (check its `procedure` — a single
47//!    connection's dedicated streams aren't partitioned by which
48//!    procedure they're for, so a session that's advertised more than
49//!    one needs this to route).
50//! 3. Drive it exactly like the caller side, from the opposite chair:
51//!    `server_stream` mode pushes with `send_data`/`close_send`;
52//!    `client_stream` mode drains with `recv` and finishes with
53//!    `send_reply`.
54
55use std::time::Duration;
56
57use crate::cbor::Value;
58use crate::connection::{FrameStream, RecvFrameError, SendFrameError, Session};
59use crate::frame::{self, StreamEncoding, StreamMode, StreamRole};
60use crate::identity::KeyPair;
61
62pub struct StreamHandle {
63    stream: FrameStream,
64    pub stream_id: [u8; 16],
65    pub mode: StreamMode,
66    seq_out: u64,
67}
68
69#[derive(Debug)]
70pub enum OpenError {
71    OpenStream(quinn::ConnectionError),
72    Send(SendFrameError),
73}
74
75impl std::fmt::Display for OpenError {
76    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
77        match self {
78            OpenError::OpenStream(e) => write!(f, "opening a dedicated stream: {e}"),
79            OpenError::Send(e) => write!(f, "sending stream_open: {e}"),
80        }
81    }
82}
83
84impl std::error::Error for OpenError {}
85
86#[derive(Debug)]
87pub enum AcceptError {
88    AcceptStream(quinn::ConnectionError),
89    Timeout,
90    Recv(RecvFrameError),
91    Parse(frame::ParseStreamOpenError),
92}
93
94impl std::fmt::Display for AcceptError {
95    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
96        match self {
97            AcceptError::AcceptStream(e) => write!(f, "accepting a dedicated stream: {e}"),
98            AcceptError::Timeout => write!(f, "no inbound stream within the given timeout"),
99            AcceptError::Recv(e) => write!(f, "reading the stream's first frame: {e}"),
100            AcceptError::Parse(e) => write!(f, "expected a stream_open frame: {e}"),
101        }
102    }
103}
104
105impl std::error::Error for AcceptError {}
106
107/// One item [`StreamHandle::recv`] hands back: a chunk, or a clean
108/// end-of-stream.
109#[derive(Debug, Clone)]
110pub enum StreamItem {
111    Data {
112        seq: u64,
113        encoding: StreamEncoding,
114        body: Value,
115    },
116    Eof,
117}
118
119#[derive(Debug)]
120pub enum RecvStreamError {
121    Recv(RecvFrameError),
122    Parse(frame::ParseStreamEventError),
123    /// The peer sent an explicit STREAM_ERROR abort.
124    PeerAborted {
125        code: String,
126        message: String,
127    },
128    /// A frame for a *different* stream_id arrived on this stream —
129    /// never expected on a dedicated stream with a well-behaved peer,
130    /// surfaced distinctly rather than silently accepted.
131    StreamIdMismatch,
132    /// A frame arrived that isn't valid in the context this call is
133    /// waiting in — e.g. [`StreamHandle::recv`] got a STREAM_REPLY
134    /// (only [`StreamHandle::await_reply`] expects one), or
135    /// `await_reply` got a STREAM_DATA/STREAM_END before any reply.
136    UnexpectedFrame,
137}
138
139impl std::fmt::Display for RecvStreamError {
140    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
141        match self {
142            RecvStreamError::Recv(e) => write!(f, "{e}"),
143            RecvStreamError::Parse(e) => write!(f, "{e}"),
144            RecvStreamError::PeerAborted { code, message } => {
145                write!(f, "peer aborted the stream: {code} ({message})")
146            }
147            RecvStreamError::StreamIdMismatch => {
148                write!(f, "received a frame for a different stream_id")
149            }
150            RecvStreamError::UnexpectedFrame => {
151                write!(f, "received a frame not valid in this context")
152            }
153        }
154    }
155}
156
157impl std::error::Error for RecvStreamError {}
158
159impl StreamHandle {
160    /// Open a dedicated stream on `session`'s connection and send a
161    /// signed STREAM_OPEN. Fire-and-forget at the wire level — no reply
162    /// is expected here; drive [`recv`](Self::recv) (for
163    /// `server_stream`/`bidi`) or [`send_data`](Self::send_data) (for
164    /// `client_stream`/`bidi`) next, depending on `mode`.
165    pub async fn open(
166        session: &mut Session,
167        procedure: &str,
168        realm: [u8; 32],
169        mode: StreamMode,
170        args: Value,
171        deadline_ms: i128,
172        identity: &KeyPair,
173    ) -> Result<Self, OpenError> {
174        let mut stream = session
175            .open_dedicated_stream()
176            .await
177            .map_err(OpenError::OpenStream)?;
178        let stream_id: [u8; 16] = rand::random();
179        let spec = frame::StreamOpenSpec::new(
180            stream_id,
181            procedure,
182            realm,
183            mode,
184            args,
185            deadline_ms,
186            identity.node_id(),
187        );
188        let signed = frame::sign(frame::stream_open(&spec), identity);
189        stream.send_frame(signed).await.map_err(OpenError::Send)?;
190        Ok(Self {
191            stream,
192            stream_id,
193            mode,
194            seq_out: 0,
195        })
196    }
197
198    /// Provider role: block for the next inbound STREAM_OPEN on
199    /// `session`'s connection, bounded by `timeout`. Only ever succeeds
200    /// after [`Session::advertise`](crate::connection::Session::advertise)
201    /// has registered at least one procedure — otherwise the station has
202    /// nothing to route here. Returns the ready-to-use handle alongside
203    /// the parsed [`frame::StreamOpenInfo`] (check its `procedure` if
204    /// this session advertised more than one).
205    pub async fn accept(
206        session: &mut Session,
207        timeout: Duration,
208    ) -> Result<(Self, frame::StreamOpenInfo), AcceptError> {
209        let mut stream = tokio::time::timeout(timeout, session.accept_dedicated_stream())
210            .await
211            .map_err(|_| AcceptError::Timeout)?
212            .map_err(AcceptError::AcceptStream)?;
213        let first = stream.recv_frame().await.map_err(AcceptError::Recv)?;
214        let open = frame::parse_stream_open(&first).map_err(AcceptError::Parse)?;
215        let handle = Self {
216            stream,
217            stream_id: open.stream_id,
218            mode: open.mode,
219            seq_out: 0,
220        };
221        Ok((handle, open))
222    }
223
224    /// Provider role: send the terminal STREAM_REPLY a `client_stream`/
225    /// `bidi` caller's own [`await_reply`](Self::await_reply) is waiting
226    /// on, once this side has fully consumed and verified whatever the
227    /// caller streamed.
228    pub async fn send_reply(
229        &mut self,
230        payload: Value,
231        identity: &KeyPair,
232    ) -> Result<(), SendFrameError> {
233        let spec = frame::StreamReplySpec::new(self.stream_id, payload, identity.node_id());
234        let signed = frame::sign(frame::stream_reply(&spec), identity);
235        self.stream.send_frame(signed).await
236    }
237
238    /// Send one chunk. `seq` is tracked internally, starting at 0 and
239    /// incrementing per call — matches the reference's `seq_out` counter
240    /// (a sanity/debugging signal, not used for reordering: frames
241    /// arrive in order on a single QUIC stream by construction).
242    pub async fn send_data(
243        &mut self,
244        encoding: StreamEncoding,
245        body: Value,
246        identity: &KeyPair,
247    ) -> Result<(), SendFrameError> {
248        let spec = frame::StreamDataSpec::new(
249            self.stream_id,
250            self.seq_out,
251            encoding,
252            body,
253            Some(identity.public_bytes()),
254        );
255        self.seq_out += 1;
256        let signed = frame::sign(frame::stream_data(&spec), identity);
257        self.stream.send_frame(signed).await
258    }
259
260    /// Half-close: signal this side is done sending. For
261    /// `client_stream`/`bidi` modes, follow with
262    /// [`await_reply`](Self::await_reply).
263    pub async fn close_send(&mut self, identity: &KeyPair) -> Result<(), SendFrameError> {
264        let spec = frame::StreamEndSpec::new(
265            self.stream_id,
266            StreamRole::Send,
267            Some(identity.public_bytes()),
268        );
269        let signed = frame::sign(frame::stream_end(&spec), identity);
270        self.stream.send_frame(signed).await
271    }
272
273    /// Receive the next chunk or end-of-stream, bounded by `timeout`.
274    pub async fn recv(&mut self, timeout: Duration) -> Result<StreamItem, RecvStreamError> {
275        let value = self
276            .stream
277            .recv_frame_timeout(timeout)
278            .await
279            .map_err(RecvStreamError::Recv)?;
280        match frame::parse_stream_event(&value).map_err(RecvStreamError::Parse)? {
281            frame::StreamEvent::Data {
282                stream_id,
283                seq,
284                encoding,
285                body,
286            } => {
287                self.check_stream_id(stream_id)?;
288                Ok(StreamItem::Data {
289                    seq,
290                    encoding,
291                    body,
292                })
293            }
294            frame::StreamEvent::End { stream_id, role: _ } => {
295                self.check_stream_id(stream_id)?;
296                Ok(StreamItem::Eof)
297            }
298            frame::StreamEvent::Error {
299                stream_id,
300                code,
301                message,
302            } => {
303                self.check_stream_id(stream_id)?;
304                Err(RecvStreamError::PeerAborted { code, message })
305            }
306            frame::StreamEvent::Reply { .. } => Err(RecvStreamError::UnexpectedFrame),
307        }
308    }
309
310    /// Block for the provider's terminal STREAM_REPLY (`client_stream`/
311    /// `bidi` modes only) — call after [`close_send`](Self::close_send).
312    pub async fn await_reply(
313        &mut self,
314        timeout: Duration,
315    ) -> Result<(Value, [u8; 32]), RecvStreamError> {
316        let value = self
317            .stream
318            .recv_frame_timeout(timeout)
319            .await
320            .map_err(RecvStreamError::Recv)?;
321        match frame::parse_stream_event(&value).map_err(RecvStreamError::Parse)? {
322            frame::StreamEvent::Reply {
323                stream_id,
324                payload,
325                responded_by,
326            } => {
327                self.check_stream_id(stream_id)?;
328                Ok((payload, responded_by))
329            }
330            frame::StreamEvent::Error {
331                stream_id,
332                code,
333                message,
334            } => {
335                self.check_stream_id(stream_id)?;
336                Err(RecvStreamError::PeerAborted { code, message })
337            }
338            frame::StreamEvent::Data { .. } | frame::StreamEvent::End { .. } => {
339                Err(RecvStreamError::UnexpectedFrame)
340            }
341        }
342    }
343
344    fn check_stream_id(&self, stream_id: [u8; 16]) -> Result<(), RecvStreamError> {
345        if stream_id == self.stream_id {
346            Ok(())
347        } else {
348            Err(RecvStreamError::StreamIdMismatch)
349        }
350    }
351
352    /// Non-normal termination: explicitly tell the peer this stream is
353    /// aborting, per §13.1 point 4 — the only signal the peer gets to
354    /// distinguish a cancellation/failure from a dropped connection.
355    /// Best-effort, like [`Session::close`](crate::connection::Session::close)'s
356    /// GOODBYE — consumes `self` so the handle can't be used again after
357    /// aborting.
358    pub async fn abort(
359        mut self,
360        code: impl Into<String>,
361        message: impl Into<String>,
362        identity: &KeyPair,
363    ) {
364        let spec = frame::StreamErrorSpec::new(
365            self.stream_id,
366            code,
367            message,
368            Some(identity.public_bytes()),
369        );
370        let signed = frame::sign(frame::stream_error(&spec), identity);
371        let _ = self.stream.send_frame(signed).await;
372    }
373}