Packages

macula

14.1.0
14.2.1 14.2.0 14.1.0 14.0.0 13.6.0 13.5.0 13.4.0 13.3.0 13.2.2 13.2.1 13.2.0 13.1.0 13.0.1 13.0.0 12.12.0 12.11.1 12.11.0 12.10.0 12.9.1 12.9.0 12.8.0 12.7.0 12.6.0 12.5.1 12.5.0 12.4.0 12.3.0 12.2.1 12.2.0 12.1.0 12.0.0 11.5.0 11.4.0 11.3.1 11.3.0 11.2.0 11.1.0 11.0.0 10.25.0 10.24.0 10.23.0 10.22.0 10.21.0 10.20.3 10.20.2 10.20.0 10.19.2 10.19.1 10.19.0 10.18.0 10.17.0 10.16.0 10.15.0 10.14.5 10.14.4 10.14.2 10.14.1 10.14.0 10.13.2 10.13.1 10.11.0 10.10.2 10.10.1 10.10.0 10.9.1 10.9.0 10.8.0 10.7.0 10.5.8 10.5.7 10.5.6 10.5.5 10.5.4 10.5.3 10.5.2 10.5.1 10.5.0 10.4.0 10.2.0 10.1.1 10.1.0 10.0.2 10.0.1 10.0.0 9.13.8 9.8.2 9.8.1 9.8.0 9.5.0 9.4.0 9.3.1 9.3.0 9.2.0 9.1.1 9.1.0 9.0.0 8.7.0 8.6.0 8.5.0 8.4.1 8.4.0 8.3.0 8.2.0 8.1.0 8.0.2 8.0.1 8.0.0 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_response.erl
Raw

src/macula_response.erl

%%%-------------------------------------------------------------------
%%% @doc Behaviour for supervised RPC responses.
%%%
%%% `advertise/5' on the raw SDK takes a bare handler fun invoked in a
%%% transient process spawned per inbound CALL (see the internal
%%% macula_station_link advertise/5 — "Handlers run in a transient
%%% process spawned per CALL"). This module gives that transient
%%% process a proper shape: each inbound call starts one supervised
%%% `macula_response' child (under a `simple_one_for_one' factory
%%% this module owns), threading state through `Module:init/1' and
%%% `Module:handle_request/2', and publishing `rpc.received_v1' /
%%% `rpc.replied_v1' mesh facts around the request. This is the
%%% provider-side counterpart to `macula_request'.
%%%
%%% A crashing `Module:handle_request/2' kills the response child;
%%% that composes with the SDK's own crash mapping unchanged, since
%%% `gen_server:call/3' against a dead callee raises the same way a
%%% crashing bare handler fun already does.
%%%
%%% == Example ==
%%%
%%% ```
%%% -module(math_service).
%%% -behaviour(macula_response).
%%% -export([init/1, handle_request/2]).
%%%
%%% init(_Args) -> {ok, []}.
%%%
%%% handle_request(#{a := A, b := B}, State) ->
%%% {reply, #{result => A + B}, State}.
%%% '''
%%%
%%% ```
%%% {ok, _Sup} = macula_response:advertise(Pool, Realm,
%%% <<"math.add_v1">>, math_service, []).
%%% '''
%%%
%%% == Advertise and publish functions ==
%%%
%%% The options of `advertise/6' and `advertise_direct/7' take
%%% `advertise', the function the handler is advertised with,
%%% `macula:advertise/5' by default; and `fact_publish', the function each
%%% response announces its facts with, `macula:publish/4' by default. The
%%% other options go on to those functions without these two. A test gives its own functions this
%%% way instead of replacing a module.
%%%
%%% == Direct-dial ==
%%%
%%% `advertise/5,6' registers the handler with the pool, and since
%%% 14.1.0 that alone makes it resolvable: each link sends its station an
%%% ADVERTISE and puts the same signed `procedure_advertisement' in the
%%% DHT, signed once by the pool and naming that link's station
%%% (macula#33). A caller using `macula_request:start_link_direct/6,7' resolves it and dials in one hop,
%%% whether or not the two stations have a routing edge between them.
%%% `advertise_direct/6,7' is the same registration, kept for the 14.x API.
%%% @end
%%%-------------------------------------------------------------------
-module(macula_response).
-behaviour(gen_server).
-include_lib("kernel/include/logger.hrl").
-export([advertise/5, advertise/6, advertise_direct/6, advertise_direct/7,
unadvertise/3]).
-export([start_link/7]).
-export_type([advertise/0, advertise_opts/0]).
-export([init/1, handle_call/3, handle_cast/2, handle_info/2, terminate/2]).
-callback init(Args :: term()) ->
{ok, State :: term()} | {stop, Reason :: term()}.
-callback handle_request(Payload :: term(), State :: term()) ->
{reply, Reply :: term(), NewState :: term()} |
{error, Reason :: term(), NewState :: term()}.
-callback terminate(Reason :: term(), State :: term()) -> any().
-optional_callbacks([terminate/2]).
%% How long a response waits for its handler: `handler_timeout_ms', 30 s by
%% default, and never past the 600 s a caller may wait for any call.
-define(DEFAULT_HANDLER_TIMEOUT_MS, 30_000).
-define(MAX_HANDLER_TIMEOUT_MS, 600_000).
-define(REQUEST_RECEIVED, <<"rpc.received_v1">>).
-define(REQUEST_REPLIED, <<"rpc.replied_v1">>).
-type advertise() :: fun((macula:pool(), macula:realm(), macula:procedure(),
macula_client:handler(), map()) -> ok | {error, term()}).
-type advertise_opts() :: #{advertise => advertise(),
fact_publish => macula_lifetime_announcer:publish(),
handler_timeout_ms => 1..600_000,
atom() => term()}.
-record(rstate, {
module :: module(),
pool :: macula:pool(),
realm :: macula:realm(),
announce :: boolean(),
fact_publish :: macula_lifetime_announcer:publish(),
request_id :: binary(),
payload :: term(),
user :: term()
}).
%% @doc Advertise `Procedure' on `Pool'/`Realm'. Starts a private
%% factory supervisor for per-request response children and registers
%% a dispatch handler with `macula:advertise/5'. Returns the
%% supervisor pid so the caller can supervise it (or ignore it).
-spec advertise(macula:pool(), macula:realm(), macula:procedure(),
module(), term()) -> {ok, pid()} | {error, term()}.
advertise(Pool, Realm, Procedure, Module, Args) ->
advertise(Pool, Realm, Procedure, Module, Args, #{}).
%% @doc As `advertise/5'. `Opts' may include `announce' (default
%% `true'), `auth' (forwarded to `macula:advertise/5'),
%% `handler_timeout_ms' — how long to wait for the handler before the caller
%% is answered `temporary_relay_failure', an integer from 1 to 600000,
%% default 30000; anything else is refused as
%% `{error, {invalid_handler_timeout_ms, Value}}' — and
%% `reuse_sup' — an existing supervisor pid (as returned by a prior
%% `advertise/5,6' call) to register the handler again with, without
%% starting a new factory supervisor. Use this for a periodic
%% re-advertise (see `advertise_direct/6,7''s own doc) — calling
%% plain `advertise/5,6' on a timer would leak one orphaned
%% supervisor per tick, since each call otherwise starts a fresh one.
-spec advertise(macula:pool(), macula:realm(), macula:procedure(),
module(), term(), advertise_opts()) -> {ok, pid()} | {error, term()}.
advertise(Pool, Realm, Procedure, Module, Args, Opts) when is_map(Opts) ->
advertise_within(handler_timeout(maps:get(handler_timeout_ms, Opts, ?DEFAULT_HANDLER_TIMEOUT_MS)),
Pool, Realm, Procedure, Module, Args, Opts).
advertise_within({ok, Timeout}, Pool, Realm, Procedure, Module, Args, Opts) ->
Advertise = arity_5(maps:get(advertise, Opts, fun macula:advertise/5)),
FactPublish = arity_4(maps:get(fact_publish, Opts, fun macula:publish/4)),
Sup = existing_or_new_sup(maps:get(reuse_sup, Opts, undefined)),
Announce = maps:get(announce, Opts, true),
Handler = fun(Payload) ->
dispatch(Sup, Module, Pool, Realm, Announce, FactPublish, Args, Payload, Timeout)
end,
case Advertise(Pool, Realm, Procedure, Handler, without_functions(Opts)) of
ok -> {ok, Sup};
{error, Reason} -> {error, Reason}
end;
advertise_within({error, _} = Refused, _Pool, _Realm, _Procedure, _Module, _Args, _Opts) ->
Refused.
handler_timeout(Ms) when is_integer(Ms), Ms >= 1, Ms =< ?MAX_HANDLER_TIMEOUT_MS -> {ok, Ms};
handler_timeout(Other) -> {error, {invalid_handler_timeout_ms, Other}}.
%% A `reuse_sup' pid from a caller's prior `advertise/6' call can have
%% died since (e.g. the caller itself crashed and, being linked to the
%% factory sup it started, took it down too — see `mcl_om_capabilities'
%% for a real periodic-republish caller that does exactly this on a
%% timed-out advertise). Reusing a dead pid unconditionally used to hand
%% `dispatch/7' a `Sup' that would `noproc' on its very first
%% `supervisor:start_child' — silently breaking every inbound call for
%% that procedure until the next re-advertise happened to land. Checking
%% liveness here is pattern matching on a plain predicate, not a
%% try/catch: a dead reuse target is exactly as valid an input as an
%% absent one, and both fall through to `new_sup/0'.
existing_or_new_sup(Pid) when is_pid(Pid) ->
existing_or_new_sup(Pid, erlang:is_process_alive(Pid));
existing_or_new_sup(undefined) ->
new_sup().
existing_or_new_sup(Pid, true) -> Pid;
existing_or_new_sup(_Pid, false) -> new_sup().
new_sup() ->
{ok, Sup} = macula_response_sup:start_link(),
Sup.
%% @doc As `advertise/5'. Since 14.1.0 `advertise/5' itself makes the
%% provider resolvable: each link puts the advertisement it sends its station
%% in the DHT, signed once by the pool (macula#33), so this publishes nothing
%% of its own. `NodeIdentity' is kept for the 14.x signature and is not used.
-spec advertise_direct(macula:pool(), macula:realm(), macula:procedure(),
module(), term(), macula_node_keys:node_key()) ->
{ok, pid()} | {error, term()}.
advertise_direct(Pool, Realm, Procedure, Module, Args, NodeIdentity) ->
advertise_direct(Pool, Realm, Procedure, Module, Args, NodeIdentity, #{}).
%% @doc As `advertise_direct/6', with `Opts' forwarded to `advertise/6'.
%% `cert_chain', a 10.x option `authorization' replaces, is refused
%% with `{error, {removed_option, cert_chain}}' before the handler is
%% registered.
-spec advertise_direct(macula:pool(), macula:realm(), macula:procedure(),
module(), term(), macula_node_keys:node_key(), advertise_opts()) ->
{ok, pid()} | {error, term()}.
advertise_direct(Pool, Realm, Procedure, Module, Args, NodeIdentity, Opts) when is_map(Opts) ->
advertise_direct_unless_removed(macula_direct_dial:removed_option(advertise, Opts), Pool,
Realm, Procedure, Module, Args, NodeIdentity, Opts).
advertise_direct_unless_removed(none, Pool, Realm, Procedure, Module, Args, _NodeIdentity, Opts) ->
advertise(Pool, Realm, Procedure, Module, Args, Opts);
advertise_direct_unless_removed(Removed, _Pool, _Realm, _Procedure, _Module, _Args, _NodeIdentity,
_Opts) ->
{error, Removed}.
%% @doc Stop advertising. Does not stop the factory supervisor
%% returned by `advertise/5,6' — callers that want to tear it down
%% should `exit(Sup, shutdown)' themselves.
-spec unadvertise(macula:pool(), macula:realm(), macula:procedure()) -> ok.
unadvertise(Pool, Realm, Procedure) ->
macula:unadvertise(Pool, Realm, Procedure).
dispatch(Sup, Module, Pool, Realm, Announce, FactPublish, Args, Payload, Timeout) ->
Child = [Module, Pool, Realm, Announce, FactPublish, Args, Payload],
case supervisor:start_child(Sup, Child) of
{ok, Pid} -> gen_server:call(Pid, run, Timeout);
{error, Reason} -> {error, Reason}
end.
%% @private
-spec start_link(module(), macula:pool(), macula:realm(), boolean(),
macula_lifetime_announcer:publish(), term(), term()) ->
{ok, pid()} | {error, term()}.
start_link(Module, Pool, Realm, Announce, FactPublish, InitArgs, Payload) ->
gen_server:start_link(?MODULE,
{Module, Pool, Realm, Announce, FactPublish, InitArgs, Payload}, []).
%%%===================================================================
%%% gen_server callbacks
%%%===================================================================
%% @private
init({Module, Pool, Realm, Announce, FactPublish, InitArgs, Payload}) ->
case Module:init(InitArgs) of
{ok, UserState} ->
RequestId = crypto:strong_rand_bytes(16),
publish(FactPublish, Announce, Pool, Realm, ?REQUEST_RECEIVED,
#{request_id => RequestId}),
{ok, #rstate{module = Module, pool = Pool, realm = Realm,
announce = Announce, fact_publish = FactPublish,
request_id = RequestId,
payload = Payload, user = UserState}};
{stop, Reason} ->
{stop, Reason}
end.
%% @private
handle_call(run, _From, #rstate{module = Module, payload = Payload,
user = User} = State) ->
{Reply, NewUser} = outcome(Module:handle_request(Payload, User)),
publish_replied(State, Reply),
{stop, normal, Reply, State#rstate{user = NewUser}};
handle_call(_Request, _From, State) ->
{reply, {error, unsupported}, State}.
outcome({reply, Reply, NewUser}) -> {{ok, Reply}, NewUser};
outcome({error, Reason, NewUser}) -> {{error, Reason}, NewUser}.
%% @private
handle_cast(_Msg, State) -> {noreply, State}.
%% @private
handle_info(_Msg, State) -> {noreply, State}.
%% @private
terminate(Reason, #rstate{module = Module, user = User}) ->
maybe_terminate(Module, Reason, User).
maybe_terminate(Module, Reason, User) ->
case erlang:function_exported(Module, terminate, 2) of
true -> Module:terminate(Reason, User);
false -> ok
end.
publish_replied(#rstate{pool = Pool, realm = Realm, announce = Announce,
fact_publish = FactPublish, request_id = RequestId}, Reply) ->
publish(FactPublish, Announce, Pool, Realm, ?REQUEST_REPLIED,
outcome_fields(#{request_id => RequestId}, Reply)).
outcome_fields(Base, {ok, _}) -> Base#{outcome => replied};
outcome_fields(Base, {error, Reason}) -> Base#{outcome => failed, reason => Reason}.
publish(_FactPublish, false, _, _, _, _) -> ok;
publish(FactPublish, true, Pool, Realm, Topic, Payload) ->
_ = FactPublish(Pool, Realm, Topic, Payload), ok.
%% A function option of the wrong arity is refused with function_clause,
%% in the caller.
arity_4(Fun) when is_function(Fun, 4) -> Fun.
arity_5(Fun) when is_function(Fun, 5) -> Fun.
%% The options that are this node's own business: the functions, and the
%% handler timeout, which bounds this node's wait on its handler.
without_functions(Opts) ->
maps:without([advertise, fact_publish, handler_timeout_ms], Opts).