Packages

macula

0.20.13
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_gateway_system macula_gateway_rpc.erl
Raw

src/macula_gateway_system/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).
-include_lib("kernel/include/logger.hrl").
%% API
-export([
start_link/1,
stop/1,
register_handler/3,
unregister_handler/2,
call/4,
invoke_handler/3,
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.
%% Handler can be either a PID or a function.
-spec register_handler(pid(), binary(), pid() | fun()) -> 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).
%% @doc Invoke a handler directly for local calls.
%% If handler is a function, invokes it directly.
%% If handler is a PID, sends rpc_call message and waits for response.
-spec invoke_handler(pid(), binary(), map()) -> {ok, term()} | {error, term()}.
invoke_handler(Pid, Procedure, Args) ->
gen_server:call(Pid, {invoke_handler, Procedure, Args}, 5000).
%%%===================================================================
%%% gen_server callbacks
%%%===================================================================
init(Opts) ->
?LOG_INFO("Initializing RPC handler"),
State = #state{
opts = Opts,
registrations = #{},
monitors = #{}
},
?LOG_INFO("RPC handler initialized"),
{ok, State}.
handle_call({register_handler, Procedure, Handler}, _From, State)
when is_binary(Procedure), (is_pid(Handler) orelse is_function(Handler)) ->
ExistingHandler = maps:get(Procedure, State#state.registrations, undefined),
NewMonitors = cleanup_existing_handler(ExistingHandler, Procedure, State#state.monitors),
{Registrations, Monitors} = register_new_handler(Handler, Procedure, State#state.registrations, NewMonitors),
NewState = State#state{registrations = Registrations, monitors = Monitors},
{reply, ok, NewState};
handle_call({unregister_handler, Procedure}, _From, State) when is_binary(Procedure) ->
Handler = maps:get(Procedure, State#state.registrations, undefined),
NewMonitors = cleanup_existing_handler(Handler, Procedure, State#state.monitors),
NewRegistrations = maps:remove(Procedure, State#state.registrations),
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({invoke_handler, Procedure, Args}, _From, State) when is_binary(Procedure) ->
case maps:get(Procedure, State#state.registrations, undefined) of
undefined ->
{reply, {error, no_handler}, State};
Handler when is_function(Handler) ->
%% Direct function call
{reply, invoke_function_handler(Handler, Args), State};
Handler when is_pid(Handler) ->
%% Synchronous call to handler process
{reply, invoke_pid_handler(Handler, Procedure, Args), 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) ->
NewState = handle_monitor_down(maps:get(MonitorRef, State#state.monitors, undefined), MonitorRef, State),
{noreply, NewState};
handle_info(_Info, State) ->
{noreply, State}.
terminate(_Reason, _State) ->
ok.
%%%===================================================================
%%% Internal functions
%%%===================================================================
%% @private Invoke function handler safely
invoke_function_handler(Handler, Args) ->
handle_function_invoke_result(catch Handler(Args)).
%% @private Handle function invoke result
handle_function_invoke_result({'EXIT', Reason}) ->
{error, Reason};
handle_function_invoke_result(Result) ->
{ok, Result}.
%% @private Invoke pid handler with timeout handling
invoke_pid_handler(Handler, Procedure, Args) ->
handle_pid_invoke_result(catch gen_server:call(Handler, {invoke, Procedure, Args}, 5000)).
%% @private Handle pid invoke result
handle_pid_invoke_result({'EXIT', {timeout, _}}) ->
{error, timeout};
handle_pid_invoke_result({'EXIT', Reason}) ->
{error, Reason};
handle_pid_invoke_result(Result) ->
Result.
%% @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.
%% First argument is ignored (kept for API compatibility).
-spec find_monitor_ref(term(), 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.
%%%===================================================================
%%% Handler registration helpers
%%%===================================================================
%% @private No existing handler to clean up
cleanup_existing_handler(undefined, _Procedure, Monitors) ->
Monitors;
%% @private Clean up existing handler's monitor
cleanup_existing_handler(_ExistingHandler, Procedure, Monitors) ->
OldMonRef = find_monitor_ref(undefined, Procedure, Monitors),
demonitor_and_remove(OldMonRef, Monitors).
demonitor_and_remove(undefined, Monitors) ->
Monitors;
demonitor_and_remove(Ref, Monitors) ->
erlang:demonitor(Ref, [flush]),
maps:remove(Ref, Monitors).
%% @private Register PID handler with monitoring
register_new_handler(Handler, Procedure, Registrations, Monitors) when is_pid(Handler) ->
MonitorRef = erlang:monitor(process, Handler),
{maps:put(Procedure, Handler, Registrations),
maps:put(MonitorRef, Procedure, Monitors)};
%% @private Register function handler (no monitoring)
register_new_handler(Handler, Procedure, Registrations, Monitors) ->
{maps:put(Procedure, Handler, Registrations), Monitors}.
%%%===================================================================
%%% Monitor down helpers
%%%===================================================================
%% @private Unknown monitor - ignore
handle_monitor_down(undefined, _MonitorRef, State) ->
State;
%% @private Known monitor - clean up registration
handle_monitor_down(Procedure, MonitorRef, State) ->
State#state{
registrations = maps:remove(Procedure, State#state.registrations),
monitors = maps:remove(MonitorRef, State#state.monitors)
}.