Packages
macula
3.10.1
7.0.0
6.0.0
5.2.2
5.2.1
5.2.0
5.1.0
5.0.0
4.8.0
4.7.1
4.7.0
4.6.0
4.5.0
4.4.10
4.4.9
4.4.8
4.4.7
4.4.6
4.4.5
4.4.4
4.4.3
4.4.2
4.4.1
4.4.0
4.3.1
4.3.0
4.2.9
4.2.8
4.2.7
4.2.6
4.2.5
4.2.4
4.2.3
4.2.2
4.2.1
4.2.0
4.1.1
4.1.0
4.0.0
3.16.0
3.15.3
3.15.2
3.15.1
3.14.0
3.13.0
3.12.1
3.12.0
3.11.1
3.11.0
3.10.3
3.10.2
3.10.1
3.9.0
3.8.0
3.7.0
3.5.0
3.4.0
3.3.0
3.2.0
3.1.0
3.0.0
2.1.1
2.1.0
2.0.0
1.5.2
1.5.1
1.4.30
1.4.29
1.4.28
1.4.27
1.4.26
1.4.25
1.4.24
1.4.23
1.4.22
1.4.21
1.4.20
1.4.19
1.4.18
1.4.17
1.4.16
1.4.15
1.4.14
1.4.13
1.4.11
1.4.10
1.4.9
1.4.8
1.4.7
1.4.6
1.4.5
1.4.4
1.4.3
1.4.2
1.4.1
1.4.0
1.3.1
1.3.0
1.2.0
1.1.0
1.0.10
1.0.9
1.0.8
1.0.7
1.0.6
1.0.5
1.0.4
1.0.3
1.0.2
1.0.1
1.0.0
0.48.6
0.48.5
0.48.4
0.48.3
0.48.2
0.48.1
0.48.0
0.47.1
0.47.0
0.46.3
0.46.1
0.46.0
0.45.3
0.45.2
0.45.1
0.45.0
0.44.2
0.44.1
0.44.0
0.43.3
0.43.2
0.43.1
0.43.0
0.42.9
0.42.8
0.42.7
0.42.6
0.42.5
0.42.4
0.42.3
0.42.2
0.42.1
0.42.0
0.41.1
0.41.0
0.40.1
0.40.0
0.39.9
0.39.8
0.39.7
0.39.6
0.39.5
0.39.4
0.39.3
0.39.2
0.39.1
0.39.0
0.38.8
0.38.7
0.38.6
0.38.5
0.38.4
0.38.3
0.38.2
0.38.1
0.38.0
0.37.7
0.37.6
0.37.5
0.37.4
0.37.3
0.37.2
0.37.1
0.37.0
0.36.6
0.36.5
0.36.4
0.36.3
0.36.2
0.36.1
0.36.0
0.35.4
0.35.3
0.35.2
0.35.1
0.35.0
0.34.1
0.34.0
0.33.1
0.33.0
0.32.5
0.32.4
0.32.3
0.32.2
0.32.1
0.32.0
0.31.9
0.31.8
0.31.7
0.31.6
0.31.5
0.31.4
0.31.3
0.31.2
0.31.1
0.31.0
0.30.10
0.30.9
0.30.8
0.30.7
0.30.6
0.30.5
0.30.4
0.30.3
0.30.2
0.30.1
0.30.0
0.29.0
0.28.3
0.28.2
0.28.1
0.28.0
0.27.1
0.27.0
0.26.1
0.26.0
0.25.6
0.25.5
0.25.4
0.25.3
0.25.2
0.25.1
0.25.0
0.24.6
0.24.5
0.24.4
0.24.3
0.24.2
0.24.1
0.24.0
0.23.3
0.23.2
0.23.1
0.23.0
0.22.12
0.22.11
0.22.10
0.22.9
0.22.8
0.22.7
0.22.6
0.22.5
0.22.4
0.22.3
0.22.2
0.22.1
0.22.0
0.21.7
0.21.6
0.21.5
0.21.4
0.21.2
0.21.1
0.21.0
0.20.25
0.20.24
0.20.23
0.20.22
0.20.21
0.20.20
0.20.19
0.20.18
0.20.17
0.20.16
0.20.15
0.20.14
0.20.13
0.20.12
0.20.11
0.20.10
0.20.9
0.20.8
0.20.7
0.20.6
0.20.5
0.20.3
0.20.2
0.20.1
0.20.0
0.19.2
0.19.1
0.19.0
0.18.1
0.18.0
0.17.4
0.17.3
0.17.2
0.17.1
0.17.0
0.16.6
0.16.5
0.16.4
0.16.3
0.16.2
0.16.1
0.16.0
0.15.1
0.15.0
0.14.3
0.14.2
0.14.1
0.14.0
0.12.6
0.12.5
0.12.3
0.11.3
0.10.2
0.10.1
0.10.0
0.9.2
0.9.1
0.9.0
0.8.25
0.8.24
0.8.23
0.8.22
0.8.21
0.8.20
0.8.19
0.8.18
0.8.17
0.8.16
0.8.15
0.8.14
0.8.13
0.8.12
0.8.11
0.8.10
0.8.9
0.8.8
0.8.7
0.8.6
0.8.5
0.8.4
0.8.3
0.8.2
0.8.1
0.8.0
0.7.30
0.7.29
0.7.28
0.7.27
0.7.26
0.7.25
0.7.24
0.7.23
0.7.22
0.7.21
0.7.20
0.7.19
0.7.18
0.7.17
0.7.16
0.7.15
0.7.14
0.7.13
0.7.12
0.7.11
0.7.10
0.7.9
0.7.8
0.7.7
0.7.6
0.7.5
0.7.4
0.7.3
0.7.2
0.7.1
0.7.0
0.6.7
0.6.6
0.6.5
0.6.4
0.6.3
0.6.2
0.6.1
0.6.0
0.5.0
0.4.4
0.4.3
0.4.2
0.4.1
0.4.0
0.3.4
0.3.3
0.3.2
0.3.1
Macula HTTP/3 Mesh SDK — connect, subscribe, publish, call, advertise
Current section
Files
Jump to
Current section
Files
src/macula_stream_local.erl
%%%-------------------------------------------------------------------
%%% @doc Local registry + dispatcher for streaming RPC.
%%%
%%% Phase 1 of PLAN_MACULA_STREAMING.md ships LOCAL streaming only —
%%% the client-side and server-side macula_stream processes both live
%%% in the same BEAM and are paired with macula_stream:pair/2. This
%%% module is the registry that lets call_stream find a locally-
%%% advertised handler for a given procedure name.
%%%
%%% Phase 2 will add a parallel path through macula_mesh_client that
%%% bridges streams to QUIC. The public SDK surface in macula.erl
%%% stays the same; macula_stream_local becomes a fast in-process
%%% short-circuit for procedures advertised on the same node.
%%% @end
%%%-------------------------------------------------------------------
-module(macula_stream_local).
-behaviour(gen_server).
-export([
start_link/0,
advertise/2,
advertise/3,
unadvertise/1,
call_stream/3,
open_stream/3,
list_advertised/0
]).
-export([
init/1,
handle_call/3,
handle_cast/2,
handle_info/2,
terminate/2
]).
-define(SERVER, ?MODULE).
-type handler() :: fun((Stream :: pid(), Args :: term()) -> any()).
-type mode() :: macula_stream:mode().
-record(state, {
%% procedure (binary) -> {Mode, Handler}
handlers = #{} :: #{binary() => {mode(), handler()}}
}).
%%%===================================================================
%%% API
%%%===================================================================
-spec start_link() -> {ok, pid()} | {error, term()}.
start_link() ->
gen_server:start_link({local, ?SERVER}, ?MODULE, [], []).
%% @doc Advertise a streaming procedure (default mode: server_stream).
-spec advertise(binary(), handler()) -> ok | {error, term()}.
advertise(Procedure, Handler) ->
advertise(Procedure, server_stream, Handler).
-spec advertise(binary(), mode(), handler()) -> ok | {error, term()}.
advertise(Procedure, Mode, Handler)
when is_binary(Procedure), is_function(Handler, 2) ->
case Mode of
server_stream -> ok;
client_stream -> ok;
bidi -> ok
end,
gen_server:call(?SERVER, {advertise, Procedure, Mode, Handler}).
-spec unadvertise(binary()) -> ok.
unadvertise(Procedure) when is_binary(Procedure) ->
gen_server:call(?SERVER, {unadvertise, Procedure}).
%% @doc Open a server-stream call. Returns the client-side stream pid.
%% The caller drains chunks with macula_stream:recv/2.
-spec call_stream(binary(), term(), map()) ->
{ok, pid()} | {error, term()}.
call_stream(Procedure, Args, Opts) ->
open_kind(Procedure, Args, Opts, server_stream).
%% @doc Open a client-stream or bidi call. Returns the client-side
%% stream pid; caller writes with macula_stream:send/2,3 and reads
%% the terminal value with macula_stream:await_reply/1,2.
-spec open_stream(binary(), term(), map()) ->
{ok, pid()} | {error, term()}.
open_stream(Procedure, Args, Opts) ->
Mode = maps:get(mode, Opts, bidi),
open_kind(Procedure, Args, Opts, Mode).
-spec list_advertised() -> [{binary(), mode()}].
list_advertised() ->
gen_server:call(?SERVER, list_advertised).
%%%===================================================================
%%% gen_server callbacks
%%%===================================================================
init([]) ->
{ok, #state{}}.
handle_call({advertise, Procedure, Mode, Handler}, _From, State) ->
NewHandlers = maps:put(Procedure, {Mode, Handler}, State#state.handlers),
{reply, ok, State#state{handlers = NewHandlers}};
handle_call({unadvertise, Procedure}, _From, State) ->
{reply, ok, State#state{handlers = maps:remove(Procedure, State#state.handlers)}};
handle_call({lookup, Procedure}, _From, State) ->
case maps:find(Procedure, State#state.handlers) of
{ok, Entry} -> {reply, {ok, Entry}, State};
error -> {reply, {error, not_advertised}, State}
end;
handle_call(list_advertised, _From, State) ->
L = [{P, M} || {P, {M, _H}} <- maps:to_list(State#state.handlers)],
{reply, L, State};
handle_call(_Msg, _From, State) ->
{reply, {error, unknown}, State}.
handle_cast(_Msg, State) -> {noreply, State}.
handle_info(_Msg, State) -> {noreply, State}.
terminate(_Reason, _State) -> ok.
%%%===================================================================
%%% Internal
%%%===================================================================
open_kind(Procedure, Args, Opts, RequestedMode) ->
case gen_server:call(?SERVER, {lookup, Procedure}) of
{ok, {AdvertisedMode, Handler}} ->
ensure_compatible(RequestedMode, AdvertisedMode),
spawn_pair(Procedure, AdvertisedMode, Handler, Args, Opts);
{error, not_advertised} = Err ->
Err
end.
ensure_compatible(_Requested, _Advertised) ->
%% Mode mismatch handling left lenient in Phase 1: we trust the
%% advertise side. Phase 2 protocol will negotiate / reject up
%% front in STREAM_OPEN.
ok.
spawn_pair(Procedure, Mode, Handler, Args, Opts) ->
StreamId = stream_id(),
Caller = maps:get(owner, Opts, self()),
{ok, ClientPid} = macula_stream:start_link(#{
id => StreamId,
role => client,
mode => Mode,
owner => Caller
}),
HandlerHost = self_host_pid(),
{ok, ServerPid} = macula_stream:start_link(#{
id => StreamId,
role => server,
mode => Mode,
owner => HandlerHost
}),
ok = macula_stream:pair(ClientPid, ServerPid),
%% Run the handler in a dedicated process so a crashing handler
%% doesn't take down the caller, and so the handler can block on
%% recv/send without affecting the client.
_HandlerPid = spawn_link_handler(Handler, ServerPid, Args, Procedure),
{ok, ClientPid}.
%% Owner of the server-side stream is a no-op host process; it just
%% keeps the stream alive while the handler runs in a sibling process.
%% Using self() would cause the registry gen_server to exit if the
%% client linked to it; spawn a dedicated host instead.
self_host_pid() ->
spawn(fun() ->
receive stop -> ok end
end).
spawn_link_handler(Handler, Stream, Args, Procedure) ->
spawn(fun() ->
try Handler(Stream, Args)
catch
Class:Reason:Stack ->
ErrCode = atom_to_binary(Class, utf8),
ErrMsg = list_to_binary(io_lib:format(
"handler ~s crashed: ~p:~p~n~p",
[Procedure, Class, Reason, Stack])),
_ = macula_stream:abort(Stream, ErrCode, ErrMsg),
ok
end
end).
%% @private 16-byte unique stream id.
stream_id() ->
crypto:strong_rand_bytes(16).