Packages
macula
4.5.0
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.
%%%
%%% 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.
%%%
%%% Cross-node streaming travels through `macula_station_link' (V2
%%% pool); the public SDK surface in `macula.erl' is identical
%%% between the LOCAL and pool paths — only the entry point arity
%%% differs (`call_stream/3' vs `call_stream/5').
%%% @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).