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}