Macula SDK — Streaming Protocol
View SourceThe raw wire primitives underneath macula_streamer / macula_stream_sink.
Audience: building something the supervised wrappers don't fit — custom retry logic, observability, an SDK for another language. Most applications want the Streaming Guide instead — it covers the same capability via
macula_streamer/macula_stream_sink, with an addressable pid, cancel, and mesh facts already wired in.
Opening a stream
Two ways to open one, mirroring unary RPC:
call_stream/5: resolves the procedure's provider through itsprocedure_advertisementand opens the stream at its serving station, naming the provider as the target, exactly ascall/5does for unary RPC. Good when you don't know or care which station serves it.call_stream_station/7(direct-dial): dials a specific station and opens the stream there in one hop, naming a provider's node_id as its target, exactly likecall_station/7for unary RPC. Use it after resolving a provider'sprocedure_advertisementandstation_endpointin the DHT (see the RPC Guide), so a stream reaches its provider the same way a unary call does.
Every STREAM_OPEN names its target. A provider serves only an open addressed to its own node_id; one for another node is dropped, and its stream closes with nothing written.
{ok, Stream} = macula:call_stream_station(Pool, StationUrl, ProviderNodeId, Realm,
Procedure, Args, #{}).Opts may set dial_timeout_ms (default 10_000) for the dial and handshake,
plus the same per-call TLS trust override as call_station/8: verify,
expected_node_id, pin_tls_cert (see the RPC Guide).
This is what macula_streamer/macula_stream_sink wrap —
an addressable pid you can monitor and cancel, streaming.*_v1 mesh facts
around each session. Reach for the raw calls below directly only if you're
building something the wrapper doesn't fit.
Consumer side
server_stream — read a feed
{ok, Stream} = macula:call_stream(Pool, Realm, <<"live.feed">>, Request, #{}),
loop(Stream).
loop(Stream) ->
case macula:recv(Stream) of
{chunk, Bin} -> handle(Bin), loop(Stream); %% raw bytes
{data, Term} -> handle(Term), loop(Stream); %% decoded (msgpack)
eof -> ok; %% source stopped
{error, R} -> {error, R} %% e.g. peer_down -> re-open
end.recv/1 blocks for the next chunk; recv/2 takes a timeout. eof means the
source closed the stream cleanly.
client_stream — push then await the result
{ok, Stream} = macula:call_stream(Pool, Realm, <<"bulk.ingest">>, Meta,
#{mode => client_stream}),
[ok = macula:send(Stream, Chunk) || Chunk <- Chunks],
ok = macula:close_send(Stream), %% signal "no more input"
{ok, Result} = macula:await_reply(Stream). %% the provider's final replysend/2 sends raw bytes; send/3 takes an encoding (raw | msgpack).
close_send/1 half-closes your direction; await_reply/1,2 returns the
provider's single final result.
Provider side
Advertise a streaming procedure with a mode and a fun(Stream, Args) handler.
The handler drives the stream with the same send / recv primitives, and ends
it with set_reply (a final result) or abort (an error).
The handler's process owns the stream, and the stream ends when that process ends: when the handler returns or crashes. A handler that lets another process keep using the stream hands it over first, and the stream then ends with that process:
ok = macula:advertise_stream(
Pool, Realm, <<"live.feed">>, server_stream,
fun(Stream, _Args) ->
Feeder = spawn(fun() -> feed(Stream) end),
ok = macula_stream:controlling_process(Stream, Feeder)
end),When a session ends (both sides closed, an abort, or the link lost), the
stream's owner gets {macula_stream, ended, Stream, How} once, with How
being closed, {error, {Code, Message}} or peer_down. An owner the stream
is handed to after that is told at once. macula_streamer takes its stream
over this way and stops on that message.
The provider verifies each STREAM_OPEN's signature against its caller before
the handler runs; a STREAM_OPEN that does not verify runs no handler.
advertise_stream/6 takes an auth policy in Opts, the same policies as
advertise/5. A STREAM_OPEN the policy refuses gets a STREAM_ERROR with code
unauthorized and runs no handler. A consumer presents its token with
call_stream/5's ucan_token opt. See the
Authorization Guide for the policies.
%% server_stream: push N chunks, then CLOSE — that is what produces `eof'
%% for a consumer looping on `recv'. `set_reply' is for client_stream /
%% bidi (see below); a pure push-only server_stream does not use it.
ok = macula:advertise_stream(
Pool, Realm, <<"live.feed">>, server_stream,
fun(Stream, _Args) ->
lists:foreach(fun(Frame) -> macula:send(Stream, Frame) end, frames()),
macula:close_stream(Stream)
end),
%% client_stream: drain the consumer's chunks, then reply
ok = macula:advertise_stream(
Pool, Realm, <<"bulk.ingest">>, client_stream,
fun(Stream, _Args) ->
N = drain(Stream, 0),
macula:set_reply(Stream, #{ingested => N})
end),
drain(Stream, N) ->
case macula:recv(Stream) of
{chunk, Bin} -> store(Bin), drain(Stream, N + 1);
eof -> N;
{error, _} -> N
end.
close_streamvs.set_reply— do not mix them forserver_stream.close_stream/1is what makes a consumer'srecvloop seeeof.set_reply/2only resolvesawait_reply/1,2; it does not close the stream. Aserver_streamhandler that callsset_replywithout also closing leaves a consumer'srecv-until-eofloop waiting forever — useclose_streamfor a pure push, and reserveset_reply+await_replyforclient_stream/bidi, where the consumer already knows to stop sending and ask for the result instead of draining chunks.
Abort with a BOLT#4-style code and message when something goes wrong:
macula:abort(Stream, <<"0F">>, <<"source unavailable">>).Over the mesh the message is text for people, at most 256 bytes of UTF-8;
a longer message, or one that is not valid UTF-8, travels empty. An error
reply from macula_stream:set_error/2 travels as the same STREAM_ERROR, with
code error and the reason as its message when the reason is a binary or an
atom.
A stream takes chunks only from the side its mode lets send: the server in
server_stream, the client in client_stream, and both in bidi. A caller's
chunk in a server_stream is refused as a malformed frame, charged to the
connection that carried it, and reaches no reader; the session goes on. A send
the mode does not allow returns {error, {send_not_allowed, Mode}} without
sending anything.
A provider serves at most 16 sessions per verified caller and 1000 on the
node at once (max_served_sessions_per_caller and max_served_sessions in
the macula application env), and one session per dedicated stream.
A provider serves a STREAM_OPEN once per caller and request id. A copy of an
admitted open gets a STREAM_ERROR with code request_copy, and an open the
node's request admission refuses gets the refusal's name: not_yet_valid or
expired for a deadline outside the window it accepts (up to 10 minutes ahead,
5 minutes past), request_id_reused, caller_quota, share_full,
admission_full or reply_not_kept. An open past a session cap gets
too_many_sessions, one its procedure's policy refuses gets unauthorized,
one for a procedure the provider does not advertise gets not_found, one in a
mode other than the one its procedure is advertised in gets mode_mismatch,
and a second STREAM_OPEN on a stream that already carries a session gets
refused. When the node's session counter or its request admission does not
answer, an open gets unavailable, and a later try may be served.
Each refusal is the provider's first frame under the open it refuses, signed and verifiable like any other. A refusal on a stream that carries no session also closes that stream, so open each session on a new stream.
A stream a peer opens must start with a STREAM_OPEN that verifies and names
this node as its target, and bring it within 10 seconds
(dedicated_stream_open_timeout_ms). A stream whose first frame is anything
else, whose STREAM_OPEN does not verify or names another node, or that stays
silent that long closes without a STREAM_ERROR.
A STREAM_OPEN is at most 1 MiB (max_stream_open_bytes). A longer first
frame closes its stream the same way, and call_stream refuses such an open
before sending anything, with {error, {open_too_large, Limit}}. Send bulk
data as chunks once the stream is open. Each node reads its own setting: a
provider whose max_stream_open_bytes is lower than a caller's closes a
larger open without an answer, and the caller waits out its deadline, so keep
the default unless both sides change it.
A stream holds at most 16 MiB of memory for chunks no reader has taken,
counting each chunk's bytes, a decoded term's heap size and the cell that
queues it. The streams a provider serves also share a budget for those
chunks: one caller's streams together at most 16 MiB
(max_served_inbox_bytes_per_caller), and all served streams on the node at
most 256 MiB (max_served_inbox_bytes). A chunk that would take a stream
past its bound, or a budget past its limit, ends the session with code
resource_exhausted. A reader that takes chunks as they come never meets
either.
The two codes tell a peer different things. stream_protocol_error says it
broke the protocol, and sending the same again fails the same way.
resource_exhausted says the receiving side had no room for what it sent,
which a well-behaved sender cannot see coming, and a later session may
succeed.
A queued chunk is a copy, but a chunk handed straight to a waiting reader can
still be part of the frame it arrived in. A handler that keeps chunks after
reading them, in a list, a table or its state, should copy what it keeps
(binary:copy/1), or it keeps whole frames alive.
Local (in-process) streams
macula:open_stream/3,4, macula:advertise_stream/2,3, and call_stream/2,3
drive streams inside one BEAM (no mesh), backed by macula_stream_local.
They are for unit tests and same-node dispatch. The pool forms
(call_stream/5, advertise_stream/5,6) are the ones that go over the mesh.
Reference
| Function | Role |
|---|---|
call_stream(Pool, Realm, Proc, Args, Opts) | raw consumer: open a stream on the pool's own link (Opts may set mode, and ucan_token for a gated procedure) |
call_stream_station(Pool, Station, Realm, Proc, Args, Opts) | raw consumer: direct-dial — dial Station and open the stream there in one hop |
advertise_stream(Pool, Realm, Proc, Mode, Handler) | raw provider: serve a streaming procedure |
advertise_stream(Pool, Realm, Proc, Mode, Handler, Opts) | raw provider: as /5, with an auth policy in Opts |
unadvertise_stream(Pool, Realm, Proc) | raw provider: stop serving it |
send(Stream, Bin) / send(Stream, Body, Enc) | send a chunk (Enc = raw | msgpack) |
recv(Stream) / recv(Stream, Timeout) | read the next {chunk,_} / {data,_} / eof |
close_send(Stream) | half-close your send direction |
await_reply(Stream) / /2 | consumer: get the provider's final result |
set_reply(Stream, Result) | provider: set the final result |
abort(Stream, Code, Message) | provider: end the stream with an error |
close_stream(Stream) | tear the stream down |
open_stream/3,4, advertise_stream/2,3 (2-arity family), call_stream/2,3 | local, in-process streams (no mesh) — unit tests and same-node dispatch |
See also
- STREAMING_GUIDE.md — the supervised wrappers most applications should use instead of these raw primitives.