macula_streamer behaviour (macula v9.8.1)
View SourceBehaviour for supervised streaming RPC providers.
advertise_stream/5 on the raw SDK takes a bare handler fun invoked as Handler(StreamPid, Args) in a transient process spawned per inbound STREAM_OPEN (see the internal macula_station_link advertise_stream/5 — "this link spawns a server-side macula_stream and dispatches Handler(StreamPid, Args) in a transient process"). This is the provider-side counterpart to macula_stream_sink: each inbound stream starts one supervised macula_streamer child (under a simple_one_for_one factory this module owns), threading state through Module:init/1 and Module:handle_open/3, and publishing streaming.started_v1 / streaming.completed_v1 mesh facts around the stream's lifetime.
Sending is push-based and driven from outside the callback: once Module:handle_open/3 has done whatever registration it needs (e.g. stashing self() in a registry keyed by some connection id), any process holding this streamer's pid can call send/2,3 / close/1 on it. This module does not prescribe the discovery mechanism.
This is the general-purpose RPC streaming feature (call_stream/5, advertise_stream/5, e.g. a logs.tail_v1-style procedure) — unrelated to content sharing's own chunked-transfer protocol; see macula_feeder / macula_download for that.
Direct-dial
advertise/5,6 registers the handler with the pool's advertise- gossip mechanism only — nothing published lets a caller on another station find this procedure without a route having propagated between the two stations first. advertise_direct/6,7 does that AND publishes a signed procedure_advertisement DHT record naming this pool's currently-connected station as the server — the exact same record type and publish function macula_response:advertise_direct/6,7 uses for plain RPC (a procedure_advertisement does not distinguish RPC from streaming), so a caller using macula_stream_sink:start_link_direct/5,6 can resolve and dial here directly, in one hop, regardless of whether the two stations have a routing edge between them.
Example
-module(log_tailer_provider).
-behaviour(macula_streamer).
-export([init/1, handle_open/3]).
init(Registry) -> {ok, Registry}.
handle_open(#{topic := Topic}, Registry, State) ->
Registry ! {tailer_ready, Topic, self()},
{ok, State}. {ok, _Sup} = macula_streamer:advertise(Pool, Realm,
<<"logs.tail_v1">>, log_tailer_provider, self()).
%% elsewhere, once the provider has announced its pid:
ok = macula_streamer:send(TailerPid, <<"a log line\n">>).
Summary
Functions
Advertise Procedure on Pool/Realm. Starts a private factory supervisor for per-stream provider children and registers a dispatch handler with macula:advertise_stream/5. Returns the supervisor pid so the caller can supervise it (or ignore it).
As advertise/5. Opts may include announce (default true) and mode (default server_stream).
As advertise/5, and additionally publishes a signed procedure_advertisement DHT record naming this pool's connected station as the server, so macula_stream_sink:start_link_direct/5,6 can resolve and dial here directly. Identity signs it — reuse the same one across re-advertises so each one doesn't mint a fresh advertiser identity.
As advertise_direct/6, with Opts forwarded to macula_direct_dial:publish_advertisement/5 — e.g. cert_chain => ChainPem (Slice 7c Direction B, managed realms only).
Close the send side of the stream.
Send a chunk out on the stream this streamer owns.
As send/2, with an explicit encoding.
Stop advertising. Does not stop the factory supervisor returned by advertise/5,6 — callers that want to tear it down should exit(Sup, shutdown) themselves.
Callbacks
Functions
-spec advertise(macula:pool(), macula:realm(), macula:procedure(), module(), term()) -> {ok, pid()} | {error, term()}.
Advertise Procedure on Pool/Realm. Starts a private factory supervisor for per-stream provider children and registers a dispatch handler with macula:advertise_stream/5. Returns the supervisor pid so the caller can supervise it (or ignore it).
-spec advertise(macula:pool(), macula:realm(), macula:procedure(), module(), term(), map()) -> {ok, pid()} | {error, term()}.
As advertise/5. Opts may include announce (default true) and mode (default server_stream).
-spec advertise_direct(macula:pool(), macula:realm(), macula:procedure(), module(), term(), macula_identity:key_pair()) -> {ok, pid()} | {error, term()}.
As advertise/5, and additionally publishes a signed procedure_advertisement DHT record naming this pool's connected station as the server, so macula_stream_sink:start_link_direct/5,6 can resolve and dial here directly. Identity signs it — reuse the same one across re-advertises so each one doesn't mint a fresh advertiser identity.
The DHT publish is best-effort: if it fails, the handler is still advertised and reachable via the ordinary pooled path — direct-dial callers just won't be able to resolve it until a later publish succeeds.
-spec advertise_direct(macula:pool(), macula:realm(), macula:procedure(), module(), term(), macula_identity:key_pair(), map()) -> {ok, pid()} | {error, term()}.
As advertise_direct/6, with Opts forwarded to macula_direct_dial:publish_advertisement/5 — e.g. cert_chain => ChainPem (Slice 7c Direction B, managed realms only).
-spec close(pid()) -> ok.
Close the send side of the stream.
Send a chunk out on the stream this streamer owns.
-spec send(pid(), binary() | term(), macula_stream:encoding()) -> ok | {error, term()}.
As send/2, with an explicit encoding.
-spec unadvertise(macula:pool(), macula:realm(), macula:procedure()) -> ok.
Stop advertising. Does not stop the factory supervisor returned by advertise/5,6 — callers that want to tear it down should exit(Sup, shutdown) themselves.