Packages

macula

0.20.23
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_rpc_system macula_rpc_async.erl
Raw

src/macula_rpc_system/macula_rpc_async.erl

%%%-------------------------------------------------------------------
%%% @doc
%%% Async RPC Module (NATS-style Request/Reply)
%%%
%%% Handles asynchronous RPC operations with callback-based responses:
%%% - Callback management (fun callbacks and pid callbacks)
%%% - Request ID generation and tracking
%%% - Request message building for P2P delivery
%%% - Reply processing and callback invocation
%%% - Timeout handling for async requests
%%%
%%% This module provides stateless helper functions used by macula_rpc_handler.
%%% The actual state (pending_requests map) remains in the handler.
%%%
%%% Extracted from macula_rpc_handler.erl (Dec 2025) to improve testability
%%% and separation of concerns.
%%% @end
%%%-------------------------------------------------------------------
-module(macula_rpc_async).
-include_lib("kernel/include/logger.hrl").
%% API - Callback management
-export([
get_callback/2,
invoke_callback/3
]).
%% API - Request building
-export([
build_request_message/5,
get_local_endpoint/0
]).
%% API - Reply processing
-export([
extract_result/1,
calculate_rtt/1
]).
%% Types
-export_type([callback/0]).
-type callback() :: {fun_cb, fun((term()) -> any())} | {pid_cb, pid()}.
%%%===================================================================
%%% Callback Management
%%%===================================================================
%% @doc Extract callback from options, defaulting to pid callback.
%% If opts contains a callback function, use it. Otherwise send to caller pid.
-spec get_callback(map(), pid()) -> callback().
get_callback(Opts, CallerPid) ->
case maps:get(callback, Opts, undefined) of
undefined -> {pid_cb, CallerPid};
Fun when is_function(Fun) -> {fun_cb, Fun}
end.
%% @doc Invoke async callback with result.
%% For function callbacks, spawns a process to avoid blocking.
%% For pid callbacks, sends a message.
-spec invoke_callback(callback(), binary(), term()) -> ok.
invoke_callback({fun_cb, Fun}, _RequestId, Result) ->
%% Spawn to avoid blocking the gen_server
spawn(fun() -> invoke_fun_callback_safe(Fun, Result) end),
ok;
invoke_callback({pid_cb, Pid}, RequestId, Result) ->
Pid ! {rpc_reply, RequestId, Result},
ok.
%% @private Invoke function callback safely
invoke_fun_callback_safe(Fun, Result) ->
handle_fun_callback_result(catch Fun(Result)).
%% @private Handle function callback result
handle_fun_callback_result({'EXIT', {Error, Stacktrace}}) ->
?LOG_ERROR("Async callback crashed: error:~p~n~p", [Error, Stacktrace]);
handle_fun_callback_result({'EXIT', Error}) ->
?LOG_ERROR("Async callback crashed: ~p", [Error]);
handle_fun_callback_result(_) ->
ok.
%%%===================================================================
%%% Request Building
%%%===================================================================
%% @doc Build an RPC_REQUEST message for NATS-style async RPC.
%% Includes from_endpoint so receiver can route reply back directly.
-spec build_request_message(binary(), binary(), binary(), binary(), binary()) -> map().
build_request_message(RequestId, Procedure, EncodedArgs, FromNodeId, Realm) ->
LocalEndpoint = get_local_endpoint(),
#{
type => <<"rpc_request">>,
request_id => RequestId,
procedure => Procedure,
args => EncodedArgs,
from_node => FromNodeId,
from_endpoint => LocalEndpoint,
realm => Realm,
timestamp => erlang:system_time(millisecond)
}.
%% @doc Get local endpoint from environment variables.
%% Used to include sender's endpoint in RPC requests so receivers can route replies back.
%% Format: "hostname:port" (e.g., "fc01:4433" in Docker)
-spec get_local_endpoint() -> binary().
get_local_endpoint() ->
Hostname = get_hostname_for_endpoint(),
Port = get_port_for_endpoint(),
iolist_to_binary([Hostname, <<":">>, Port]).
%%%===================================================================
%%% Reply Processing
%%%===================================================================
%% @doc Extract result from RPC reply message.
%% Returns {ok, DecodedValue} or {error, ErrorReason}.
-spec extract_result(map()) -> {ok, term()} | {error, term()}.
extract_result(Msg) ->
case maps:get(<<"result">>, Msg, maps:get(result, Msg, undefined)) of
undefined ->
Error = maps:get(<<"error">>, Msg, maps:get(error, Msg, <<"Unknown error">>)),
{error, Error};
Value ->
%% Try to decode JSON result
DecodedValue = safe_decode_json(Value),
{ok, DecodedValue}
end.
%% @private Safely decode JSON, returning original on failure
safe_decode_json(Value) ->
handle_decode_result(catch macula_utils:decode_json(Value), Value).
%% @private Handle decode result
handle_decode_result({'EXIT', _}, Original) ->
Original;
handle_decode_result(Decoded, _Original) ->
Decoded.
%% @doc Calculate RTT from sent_at timestamp.
-spec calculate_rtt(integer()) -> non_neg_integer().
calculate_rtt(SentAt) ->
erlang:system_time(millisecond) - SentAt.
%%%===================================================================
%%% Internal Functions
%%%===================================================================
%% @private Get hostname from environment variables.
-spec get_hostname_for_endpoint() -> binary().
get_hostname_for_endpoint() ->
case os:getenv("MACULA_HOSTNAME") of
false ->
case os:getenv("HOSTNAME") of
false -> <<"localhost">>;
Hostname -> list_to_binary(Hostname)
end;
Hostname -> list_to_binary(Hostname)
end.
%% @private Get QUIC port from environment or default.
-spec get_port_for_endpoint() -> binary().
get_port_for_endpoint() ->
case os:getenv("QUIC_PORT") of
false -> <<"4433">>; %% Default QUIC port for Docker setup
Port -> list_to_binary(Port)
end.