Packages

macula

11.3.1
11.4.0 11.3.1 11.3.0 11.2.0 11.1.0 11.0.0 10.25.0 10.24.0 10.23.0 10.22.0 10.21.0 10.20.3 10.20.2 10.20.0 10.19.2 10.19.1 10.19.0 10.18.0 10.17.0 10.16.0 10.15.0 10.14.5 10.14.4 10.14.2 10.14.1 10.14.0 10.13.2 10.13.1 10.11.0 10.10.2 10.10.1 10.10.0 10.9.1 10.9.0 10.8.0 10.7.0 10.5.8 10.5.7 10.5.6 10.5.5 10.5.4 10.5.3 10.5.2 10.5.1 10.5.0 10.4.0 10.2.0 10.1.1 10.1.0 10.0.2 10.0.1 10.0.0 9.13.8 9.8.2 9.8.1 9.8.0 9.5.0 9.4.0 9.3.1 9.3.0 9.2.0 9.1.1 9.1.0 9.0.0 8.7.0 8.6.0 8.5.0 8.4.1 8.4.0 8.3.0 8.2.0 8.1.0 8.0.2 8.0.1 8.0.0 7.1.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
macula src macula_stream_local.erl
Raw

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
]).
-ifdef(TEST).
%% The process a local call's handler runs in, spawned before its stream
%% exists: exported for macula_stream_tests.erl.
-export([spawn_handler/3]).
-endif.
-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
}),
%% The handler runs in a dedicated process that owns the server-side
%% stream: a crashing handler doesn't take down the caller, the handler
%% can block on recv/send without affecting the client, and the stream
%% ends when the handler does, unless the handler hands it over first
%% (`macula_stream:controlling_process/2').
HandlerPid = spawn_handler(Handler, Args, Procedure),
{ok, ServerPid} = macula_stream:start_link(#{
id => StreamId,
role => server,
mode => Mode,
owner => HandlerPid
}),
ok = macula_stream:pair(ClientPid, ServerPid),
HandlerPid ! {serve, ServerPid},
{ok, ClientPid}.
%% The handler process runs the handler once its stream is paired, and ends
%% without running it when the process that opened the call ends first. Once
%% the handler runs, the opener's end is no concern of it, so no notice of it
%% is left in the handler's mailbox.
spawn_handler(Handler, Args, Procedure) ->
Opener = self(),
spawn(fun() -> serve_when_paired(erlang:monitor(process, Opener), Handler, Args, Procedure) end).
serve_when_paired(OpenerRef, Handler, Args, Procedure) ->
receive
{serve, Stream} ->
true = erlang:demonitor(OpenerRef, [flush]),
run_handler(Handler, Stream, Args, Procedure);
{'DOWN', OpenerRef, process, _Opener, _Reason} ->
ok
end.
%% A handler crash aborts the stream with the crash class as the code
%% and the reason's name as the message, and none of the crash's
%% terms; the crash goes to the node's log.
run_handler(Handler, Stream, Args, Procedure) ->
try Handler(Stream, Args)
catch
Class:Reason:Stack ->
logger:warning(
"[macula_stream_local] handler ~ts crashed: ~ts",
[Procedure, macula_reason_name:logged("~p:~p~n stack=~p",
[Class, Reason, Stack])]),
_ = macula_stream:abort(Stream, atom_to_binary(Class, utf8),
macula_reason_name:text(Reason)),
ok
end.
%% @private 16-byte unique stream id.
stream_id() ->
crypto:strong_rand_bytes(16).