Packages
macula
0.6.4
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
Current section
Files
src/macula_gateway_rpc.erl
%%%-------------------------------------------------------------------
%%% @doc
%%% RPC Handler GenServer - manages RPC handler registration and call routing.
%%%
%%% Responsibilities:
%%% - Register/unregister RPC handlers for procedures
%%% - Route RPC calls to registered handlers
%%% - Handle call/response matching
%%% - Monitor handler processes for automatic cleanup
%%%
%%% Extracted from macula_gateway.erl (Phase 4)
%%% @end
%%%-------------------------------------------------------------------
-module(macula_gateway_rpc).
-behaviour(gen_server).
%% API
-export([
start_link/1,
stop/1,
register_handler/3,
unregister_handler/2,
call/4,
get_handler/2,
list_handlers/1
]).
%% gen_server callbacks
-export([init/1, handle_call/3, handle_cast/2, handle_info/2, terminate/2]).
-record(state, {
opts :: map(),
registrations :: #{binary() => pid()}, % procedure => handler_pid
monitors :: #{reference() => binary()} % monitor_ref => procedure
}).
%%%===================================================================
%%% API
%%%===================================================================
%% @doc Start the RPC handler with options.
-spec start_link(map()) -> {ok, pid()} | {error, term()}.
start_link(Opts) ->
gen_server:start_link(?MODULE, Opts, []).
%% @doc Stop the RPC handler.
-spec stop(pid()) -> ok.
stop(Pid) ->
gen_server:stop(Pid).
%% @doc Register a handler for an RPC procedure.
-spec register_handler(pid(), binary(), pid()) -> ok.
register_handler(Pid, Procedure, Handler) ->
gen_server:call(Pid, {register_handler, Procedure, Handler}).
%% @doc Unregister a handler for an RPC procedure.
-spec unregister_handler(pid(), binary()) -> ok.
unregister_handler(Pid, Procedure) ->
gen_server:call(Pid, {unregister_handler, Procedure}).
%% @doc Make an RPC call to a procedure.
-spec call(pid(), binary(), map(), map()) -> {ok, term()} | {error, term()}.
call(Pid, Procedure, Args, Opts) ->
gen_server:call(Pid, {call, Procedure, Args, Opts}, get_timeout(Opts)).
%% @doc Get the handler for a procedure.
-spec get_handler(pid(), binary()) -> {ok, pid()} | not_found.
get_handler(Pid, Procedure) ->
gen_server:call(Pid, {get_handler, Procedure}).
%% @doc List all registered handlers.
-spec list_handlers(pid()) -> {ok, [{binary(), pid()}]}.
list_handlers(Pid) ->
gen_server:call(Pid, list_handlers).
%%%===================================================================
%%% gen_server callbacks
%%%===================================================================
init(Opts) ->
State = #state{
opts = Opts,
registrations = #{},
monitors = #{}
},
{ok, State}.
handle_call({register_handler, Procedure, Handler}, _From, State)
when is_binary(Procedure), is_pid(Handler) ->
%% Check if procedure already has a handler
NewMonitors = case maps:get(Procedure, State#state.registrations, undefined) of
undefined ->
State#state.monitors;
ExistingHandler ->
%% Demonitor old handler
OldMonRef = find_monitor_ref(ExistingHandler, Procedure, State#state.monitors),
case OldMonRef of
undefined -> ok;
Ref -> erlang:demonitor(Ref, [flush])
end,
case OldMonRef of
undefined -> State#state.monitors;
_ -> maps:remove(OldMonRef, State#state.monitors)
end
end,
%% Register new handler
MonitorRef = erlang:monitor(process, Handler),
Registrations = maps:put(Procedure, Handler, State#state.registrations),
Monitors = maps:put(MonitorRef, Procedure, NewMonitors),
NewState = State#state{
registrations = Registrations,
monitors = Monitors
},
{reply, ok, NewState};
handle_call({unregister_handler, Procedure}, _From, State) when is_binary(Procedure) ->
%% Remove registration
NewRegistrations = maps:remove(Procedure, State#state.registrations),
%% Find and remove monitor
Handler = maps:get(Procedure, State#state.registrations, undefined),
NewMonitors = case Handler of
undefined ->
State#state.monitors;
_ ->
MonRef = find_monitor_ref(Handler, Procedure, State#state.monitors),
case MonRef of
undefined -> State#state.monitors;
Ref ->
erlang:demonitor(Ref, [flush]),
maps:remove(Ref, State#state.monitors)
end
end,
NewState = State#state{
registrations = NewRegistrations,
monitors = NewMonitors
},
{reply, ok, NewState};
handle_call({call, Procedure, Args, _Opts}, From, State) when is_binary(Procedure) ->
case maps:get(Procedure, State#state.registrations, undefined) of
undefined ->
{reply, {error, no_handler}, State};
Handler ->
%% Send call to handler
Handler ! {rpc_call, Procedure, Args, From},
%% Don't reply yet - handler will reply to From directly
{noreply, State}
end;
handle_call({get_handler, Procedure}, _From, State) when is_binary(Procedure) ->
Result = case maps:get(Procedure, State#state.registrations, undefined) of
undefined -> not_found;
Handler -> {ok, Handler}
end,
{reply, Result, State};
handle_call(list_handlers, _From, State) ->
Handlers = maps:to_list(State#state.registrations),
{reply, {ok, Handlers}, State};
handle_call(_Request, _From, State) ->
{reply, {error, unknown_request}, State}.
handle_cast(_Msg, State) ->
{noreply, State}.
%% @doc Handle handler process death - automatic cleanup.
handle_info({'DOWN', MonitorRef, process, _HandlerPid, _Reason}, State) ->
case maps:get(MonitorRef, State#state.monitors, undefined) of
undefined ->
%% Unknown monitor (shouldn't happen)
{noreply, State};
Procedure ->
%% Remove registration and monitor
NewRegistrations = maps:remove(Procedure, State#state.registrations),
NewMonitors = maps:remove(MonitorRef, State#state.monitors),
NewState = State#state{
registrations = NewRegistrations,
monitors = NewMonitors
},
{noreply, NewState}
end;
handle_info(_Info, State) ->
{noreply, State}.
terminate(_Reason, _State) ->
ok.
%%%===================================================================
%%% Internal functions
%%%===================================================================
%% @doc Get timeout from options (default 5000ms).
-spec get_timeout(map()) -> non_neg_integer().
get_timeout(Opts) ->
maps:get(timeout, Opts, 5000).
%% @doc Find monitor reference for a handler and procedure.
-spec find_monitor_ref(pid(), binary(), #{reference() => binary()}) -> reference() | undefined.
find_monitor_ref(_Handler, Procedure, Monitors) ->
%% This is a bit inefficient but works for now
%% We need to find the monitor ref that corresponds to this handler
%% Since we don't store handler pids in monitors map, we need to search
MonitorList = maps:to_list(Monitors),
extract_monitor_ref(lists:filter(fun({_Ref, Proc}) -> Proc =:= Procedure end, MonitorList)).
%% @doc Extract monitor ref from filter result.
extract_monitor_ref([{Ref, _Proc}]) -> Ref;
extract_monitor_ref(_) -> undefined.