Packages

macula

0.22.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
macula src macula_gateway_system macula_gateway_rpc_router.erl
Raw

src/macula_gateway_system/macula_gateway_rpc_router.erl

%%%-------------------------------------------------------------------
%%% @doc
%%% RPC Router Module - handles routed RPC messages (CALL/REPLY).
%%%
%%% Responsibilities:
%%% - Process routed CALL messages delivered locally
%%% - Process routed REPLY messages delivered locally
%%% - Send REPLY back via routing path
%%% - Forward rpc_route messages to next hop
%%% - Coordinate between RPC handler, mesh, and routing modules
%%%
%%% Pattern: Stateless delegation module
%%% - No GenServer (no state to manage)
%%% - Pure functions coordinating between modules
%%% - Consistent error handling ({ok, Result} | {error, Reason})
%%%
%%% Extracted from macula_gateway.erl (Phase 11)
%%% @end
%%%-------------------------------------------------------------------
-module(macula_gateway_rpc_router).
-include_lib("kernel/include/logger.hrl").
%% API
-export([
handle_routed_call/5,
handle_routed_reply/4,
send_reply_via_routing/4,
forward_rpc_route/3
]).
%%%===================================================================
%%% API
%%%===================================================================
%% @doc Handle routed CALL message delivered locally.
%% Looks up RPC handler, invokes it, sends reply via routing path.
%% Authorization check added in v0.17.0.
-spec handle_routed_call(map(), map(), binary(), pid(), pid()) ->
ok | {error, term()}.
handle_routed_call(CallMsg, RpcRouteMsg, NodeId, RpcPid, MeshPid) ->
?LOG_DEBUG("Processing routed CALL message"),
%% Extract source node ID from rpc_route envelope
#{<<"source_node_id">> := SourceNodeId} = RpcRouteMsg,
%% Extract call data
#{<<"procedure">> := Procedure,
<<"args">> := ArgsJson,
<<"call_id">> := CallId} = CallMsg,
%% Authorization check (v0.17.0+)
CallerDID = extract_caller_did(RpcRouteMsg, CallMsg),
UcanToken = maps:get(<<"ucan_token">>, CallMsg, undefined),
case macula_authorization:check_rpc_call(CallerDID, Procedure, UcanToken, #{}) of
{ok, authorized} ->
%% Authorized - look up handler
case macula_gateway_rpc:get_handler(RpcPid, Procedure) of
not_found ->
?LOG_WARNING("No handler for procedure: ~p", [Procedure]),
ErrorReply = #{
call_id => CallId,
error => #{
code => <<"no_handler">>,
message => <<"No handler registered for ", Procedure/binary>>
}
},
send_reply_via_routing(ErrorReply, SourceNodeId, NodeId, MeshPid);
{ok, Handler} ->
invoke_handler_and_reply(Handler, ArgsJson, CallId, SourceNodeId, NodeId, MeshPid)
end;
{error, Reason} ->
?LOG_WARNING("RPC call to ~s denied: ~p (caller: ~s)", [Procedure, Reason, CallerDID]),
ErrorReply = #{
call_id => CallId,
error => #{
code => <<"unauthorized">>,
message => iolist_to_binary(io_lib:format("~p", [Reason]))
}
},
send_reply_via_routing(ErrorReply, SourceNodeId, NodeId, MeshPid)
end.
%% @doc Handle routed REPLY message delivered locally.
%% Routes to connection via gproc (local node) or to client stream (remote client).
-spec handle_routed_reply(map(), map(), binary(), map()) ->
ok | {error, term()}.
handle_routed_reply(ReplyMsg, RpcRouteMsg, NodeId, ClientStreams) ->
?LOG_DEBUG("Processing routed REPLY message"),
%% Extract destination node ID
#{<<"destination_node_id">> := DestinationNodeId} = RpcRouteMsg,
%% Check if this is for local node or remote client
case DestinationNodeId of
NodeId ->
%% Local node - send to connection via gproc (send full rpc_route message)
deliver_reply_to_local_connection(RpcRouteMsg);
_RemoteNodeId ->
%% Remote client - send to client stream
deliver_reply_to_remote_client(ReplyMsg, RpcRouteMsg, ClientStreams)
end.
%% @doc Send REPLY back via routing path.
%% Wraps reply in rpc_route envelope and routes to destination.
%% Crashes on routing failures - exposes mesh/routing issues immediately.
-spec send_reply_via_routing(map(), binary(), binary(), pid()) -> ok.
send_reply_via_routing(ReplyMsg, DestNodeId, NodeId, MeshPid) ->
?LOG_DEBUG("Sending reply via routing to ~p", [binary:encode_hex(DestNodeId)]),
do_send_reply_via_routing(whereis(macula_routing_server), ReplyMsg, DestNodeId, NodeId, MeshPid).
%% @doc Forward rpc_route message to next hop.
%% Uses async (fire-and-forget) pattern to avoid blocking.
%% Graceful error handling - logs errors but doesn't crash gateway.
-spec forward_rpc_route(map(), map(), pid()) -> ok.
forward_rpc_route(NextHopNodeInfo, RpcRouteMsg, MeshPid) ->
?LOG_DEBUG("Forwarding rpc_route to next hop (async)"),
%% Extract next hop info (routing_bucket:node_info uses atom keys, not binary)
#{node_id := NextHopNodeId,
address := Address} = NextHopNodeInfo,
%% Encode message
EncodedMsg = macula_protocol_encoder:encode(rpc_route, RpcRouteMsg),
%% Send asynchronously - does NOT block
%% Connection creation and sending happens in a spawned process
macula_gateway_mesh:send_async(MeshPid, NextHopNodeId, Address, EncodedMsg),
?LOG_DEBUG("Queued rpc_route for async send to ~s",
[binary:encode_hex(NextHopNodeId)]),
ok.
%%%===================================================================
%%% Internal functions
%%%===================================================================
%% @private
%% @doc Invoke RPC handler and send reply.
invoke_handler_and_reply(Handler, ArgsJson, CallId, SourceNodeId, NodeId, MeshPid) ->
ExecutionResult = catch execute_rpc_handler(Handler, ArgsJson),
handle_rpc_execution_result(ExecutionResult, CallId, SourceNodeId, NodeId, MeshPid).
%% @private Execute RPC handler with decoded args
execute_rpc_handler(Handler, ArgsJson) ->
Args = json:decode(ArgsJson),
Handler(Args).
%% @private Handle RPC execution result
handle_rpc_execution_result({'EXIT', {Reason, _Stack}}, CallId, SourceNodeId, NodeId, MeshPid) ->
?LOG_ERROR("Handler error: ~p", [Reason]),
ErrorReply = #{
call_id => CallId,
error => #{
code => <<"handler_error">>,
message => format_error(error, Reason)
}
},
send_reply_via_routing(ErrorReply, SourceNodeId, NodeId, MeshPid);
handle_rpc_execution_result({'EXIT', Reason}, CallId, SourceNodeId, NodeId, MeshPid) ->
?LOG_ERROR("Handler error: ~p", [Reason]),
ErrorReply = #{
call_id => CallId,
error => #{
code => <<"handler_error">>,
message => format_error(error, Reason)
}
},
send_reply_via_routing(ErrorReply, SourceNodeId, NodeId, MeshPid);
handle_rpc_execution_result(Result, CallId, SourceNodeId, NodeId, MeshPid) ->
Reply = #{
call_id => CallId,
result => encode_json(Result)
},
?LOG_DEBUG("Handler invoked successfully"),
send_reply_via_routing(Reply, SourceNodeId, NodeId, MeshPid).
%% @private
%% @doc Deliver reply to local connection via gproc.
%% Sends the full rpc_route message to all local connections.
deliver_reply_to_local_connection(RpcRouteMsg) ->
?LOG_DEBUG("Delivering reply to local connections"),
Pids = gproc:lookup_pids({p, l, macula_connection}),
do_deliver_to_local_connections(Pids, RpcRouteMsg).
%% @private
%% @doc Deliver reply to remote client via stream.
%% Sends the full rpc_route message (not just the reply).
deliver_reply_to_remote_client(_ReplyMsg, RpcRouteMsg, ClientStreams) ->
ConnectionId = maps:get(<<"connection_id">>, RpcRouteMsg, undefined),
do_deliver_to_remote_client(ConnectionId, RpcRouteMsg, ClientStreams).
%% @private
%% @doc Encode result to JSON binary.
encode_json(Term) when is_binary(Term) ->
Term;
encode_json(Term) ->
macula_utils:encode_json(Term).
%% @private
%% @doc Format error for reply message.
format_error(Class, Reason) ->
iolist_to_binary(io_lib:format("~p: ~p", [Class, Reason])).
%%%===================================================================
%%% Routing reply helpers
%%%===================================================================
%% @private Routing server not available
do_send_reply_via_routing(undefined, _ReplyMsg, _DestNodeId, _NodeId, _MeshPid) ->
?LOG_ERROR("Routing server not available"),
error(routing_server_not_available);
%% @private Route the reply
do_send_reply_via_routing(RoutingServerPid, ReplyMsg, DestNodeId, NodeId, MeshPid) ->
RpcRouteMsg = macula_rpc_routing:wrap_reply(NodeId, DestNodeId, ReplyMsg, 10),
handle_route_result(macula_rpc_routing:route_or_deliver(NodeId, RpcRouteMsg, RoutingServerPid), MeshPid).
handle_route_result({deliver, _, _}, _MeshPid) ->
?LOG_WARNING("Reply is for local node (unexpected)"),
ok;
handle_route_result({forward, NextHopNodeInfo, UpdatedRpcRouteMsg}, MeshPid) ->
forward_rpc_route(NextHopNodeInfo, UpdatedRpcRouteMsg, MeshPid);
handle_route_result({error, Reason}, _MeshPid) ->
?LOG_ERROR("Routing error: ~p", [Reason]),
error({routing_failed, Reason}).
%%%===================================================================
%%% Local connection delivery helpers
%%%===================================================================
%% @private No local connections
do_deliver_to_local_connections([], _RpcRouteMsg) ->
?LOG_WARNING("No local connection process found"),
{error, connection_not_found};
%% @private Deliver to all local connections
do_deliver_to_local_connections(Pids, RpcRouteMsg) ->
?LOG_DEBUG("Found ~p local connection process(es)", [length(Pids)]),
lists:foreach(fun(Pid) -> gen_server:cast(Pid, {rpc_route_reply, RpcRouteMsg}) end, Pids),
ok.
%%%===================================================================
%%% Remote client delivery helpers
%%%===================================================================
%% @private No connection_id in message
do_deliver_to_remote_client(undefined, _RpcRouteMsg, _ClientStreams) ->
?LOG_WARNING("No connection_id in rpc_route"),
{error, no_connection_id};
%% @private Look up client stream
do_deliver_to_remote_client(ConnectionId, RpcRouteMsg, ClientStreams) ->
StreamPid = maps:get(ConnectionId, ClientStreams, undefined),
send_to_client_stream(StreamPid, RpcRouteMsg, ConnectionId).
send_to_client_stream(undefined, _RpcRouteMsg, ConnectionId) ->
?LOG_WARNING("Client stream not found for connection ~p", [ConnectionId]),
{error, client_not_found};
send_to_client_stream(StreamPid, RpcRouteMsg, _ConnectionId) ->
RpcRouteBinary = macula_protocol_encoder:encode(rpc_route, RpcRouteMsg),
macula_quic:send(StreamPid, RpcRouteBinary),
?LOG_DEBUG("Delivered reply to remote client"),
ok.
%%%===================================================================
%%% Authorization Helpers (v0.17.0+)
%%%===================================================================
%% @private
%% @doc Extract caller DID from rpc_route envelope or call message.
%% Priority: caller_did from CallMsg > from RpcRouteMsg > derive from source_node
-spec extract_caller_did(map(), map()) -> binary().
extract_caller_did(RpcRouteMsg, CallMsg) ->
case maps:get(<<"caller_did">>, CallMsg, undefined) of
undefined ->
case maps:get(<<"caller_did">>, RpcRouteMsg, undefined) of
undefined ->
%% Fallback: derive from source node ID (temporary until TLS cert extraction)
SourceNodeId = maps:get(<<"source_node_id">>, RpcRouteMsg, <<"unknown">>),
<<"did:macula:", SourceNodeId/binary>>;
CallerDID ->
CallerDID
end;
CallerDID ->
CallerDID
end.