Packages

macula

3.15.2
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 v1 macula_stream_local.erl
Raw

src/v1/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_v1: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_v1: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_v1: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_v1:send/2,3 and reads
%% the terminal value with macula_stream_v1: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_v1:start_link(#{
id => StreamId,
role => client,
mode => Mode,
owner => Caller
}),
HandlerHost = self_host_pid(),
{ok, ServerPid} = macula_stream_v1:start_link(#{
id => StreamId,
role => server,
mode => Mode,
owner => HandlerHost
}),
ok = macula_stream_v1: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_v1:abort(Stream, ErrCode, ErrMsg),
ok
end
end).
%% @private 16-byte unique stream id.
stream_id() ->
crypto:strong_rand_bytes(16).