Skip to main content

Module stream

Module stream 

Source
Expand description

General-purpose streaming RPC, caller/consumer role (§13.1 of plans/PLAN_WIRE_PROTOCOL.md), ported from macula_stream_sink.erl. Like content transfer (src/content.rs), this is not a separate wire mechanism: it runs the frame types built in src/frame.rs §13 over a dedicated QUIC stream, opened via Session::open_dedicated_stream rather than the control stream.

Both roles are built. Caller/consumer (§13.1) opens a stream and is the natural fit for pulling/pushing against a procedure that already exists somewhere. Provider (§13.2) advertises a procedure (§6.9, Session::advertise) and answers inbound STREAM_OPENs the station routes back — Session::accept_dedicated_stream accepts the fresh dedicated stream the station opens toward us, StreamHandle::accept reads and parses the STREAM_OPEN that’s always its first frame. Both roles end up holding the same StreamHandle afterward — a stream’s wire vocabulary (STREAM_DATA/END/ERROR/REPLY) is symmetric regardless of which side opened it, so send_data/recv/close_send/abort all mean the same thing either way. StreamHandle::send_reply is the one provider-only addition: sending the terminal STREAM_REPLY a client_stream/bidi caller’s own await_reply is waiting on.

Caller/consumer usage, matching the reference’s own pattern:

  1. StreamHandle::open sends STREAM_OPEN and returns a handle once the frame is on the wire — there’s no open-time acknowledgement to wait for; the provider starts reacting to it directly.
  2. Drive a receive loop with StreamHandle::recv until StreamItem::Eof or an error.
  3. For client_stream/bidi modes wanting a result: StreamHandle::send_data each chunk in order, StreamHandle::close_send when done, then StreamHandle::await_reply.
  4. Non-normal termination must call StreamHandle::abort, not just drop the handle — the peer’s only signal to tell a cancellation/failure apart from a dropped connection (plans/PLAN_WIRE_PROTOCOL.md §13.1, point 4).

Provider usage:

  1. Session::advertise once per procedure this session will answer.
  2. Loop on StreamHandle::accept, which blocks for the next inbound STREAM_OPEN and hands back a ready-to-use handle plus the parsed frame::StreamOpenInfo (check its procedure — a single connection’s dedicated streams aren’t partitioned by which procedure they’re for, so a session that’s advertised more than one needs this to route).
  3. Drive it exactly like the caller side, from the opposite chair: server_stream mode pushes with send_data/close_send; client_stream mode drains with recv and finishes with send_reply.

Structs§

StreamHandle

Enums§

AcceptError
OpenError
RecvStreamError
StreamItem
One item StreamHandle::recv hands back: a chunk, or a clean end-of-stream.