Packages

macula

9.5.0
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, []).
%%% '''
%%%
%%% == Direct-dial ==
%%%
%%% `advertise/5,6' registers the handler with the pool's advertise-
%%% gossip mechanism only — nothing published lets a caller on another
%%% station find this procedure without a route having propagated
%%% between the two stations first. `advertise_direct/6' does that
%%% AND publishes a signed `procedure_advertisement' DHT record naming
%%% this pool's currently-connected station as the server, so a caller
%%% using `macula_request:start_link_direct/6,7' can resolve and dial
%%% here directly, in one hop, regardless of whether the two stations
%%% have a routing edge between them.
%%% @end
%%%-------------------------------------------------------------------
-module(macula_response).
-behaviour(gen_server).
-export([advertise/5, advertise/6, advertise_direct/6, advertise_direct/7,
unadvertise/3]).
-export([start_link/6]).
-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]).
-define(CALL_TIMEOUT, 30_000).
-define(REQUEST_RECEIVED, <<"rpc.received_v1">>).
-define(REQUEST_REPLIED, <<"rpc.replied_v1">>).
-record(rstate, {
module :: module(),
pool :: macula:pool(),
realm :: macula:realm(),
announce :: boolean(),
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') and `auth' (forwarded to `macula:advertise/5').
-spec advertise(macula:pool(), macula:realm(), macula:procedure(),
module(), term(), map()) -> {ok, pid()} | {error, term()}.
advertise(Pool, Realm, Procedure, Module, Args, Opts) ->
{ok, Sup} = macula_response_sup:start_link(),
Announce = maps:get(announce, Opts, true),
Handler = fun(Payload) ->
dispatch(Sup, Module, Pool, Realm, Announce, Args, Payload)
end,
case macula:advertise(Pool, Realm, Procedure, Handler, Opts) of
ok -> {ok, Sup};
{error, Reason} -> {error, Reason}
end.
%% @doc As `advertise/5', and additionally publishes a signed
%% `procedure_advertisement' DHT record naming this pool's connected
%% station as the server, so `macula_request:start_link_direct/6,7'
%% can resolve and dial here directly. `Identity' signs the
%% advertisement — reuse the same one across re-advertises so each one
%% doesn't mint a fresh advertiser identity.
%%
%% The DHT publish is best-effort: if it fails (e.g. no healthy link
%% at that instant), the handler is still advertised and reachable via
%% the ordinary pooled path — direct-dial callers just won't be able
%% to resolve it until a later publish succeeds.
-spec advertise_direct(macula:pool(), macula:realm(), macula:procedure(),
module(), term(), macula_identity:key_pair()) ->
{ok, pid()} | {error, term()}.
advertise_direct(Pool, Realm, Procedure, Module, Args, Identity) ->
advertise_direct(Pool, Realm, Procedure, Module, Args, Identity, #{}).
%% @doc As `advertise_direct/6', with `Opts' forwarded to
%% `macula_direct_dial:publish_advertisement/5' — e.g. `cert_chain =>
%% ChainPem' (leaf ++ org CA, PEM), so a verifying consumer's
%% `verify_cert_chain' opt can check this advertiser's org/realm
%% authorization (Slice 7c Direction B, managed realms only. See
%% `macula_direct_dial''s module doc, "Trust model").
-spec advertise_direct(macula:pool(), macula:realm(), macula:procedure(),
module(), term(), macula_identity:key_pair(), map()) ->
{ok, pid()} | {error, term()}.
advertise_direct(Pool, Realm, Procedure, Module, Args, Identity, Opts) ->
case advertise(Pool, Realm, Procedure, Module, Args) of
{ok, Sup} ->
_ = macula_direct_dial:publish_advertisement(Pool, Realm,
Procedure, Identity,
Opts),
{ok, Sup};
{error, _} = Error ->
Error
end.
%% @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, Args, Payload) ->
case supervisor:start_child(Sup, [Module, Pool, Realm, Announce, Args, Payload]) of
{ok, Pid} -> gen_server:call(Pid, run, ?CALL_TIMEOUT);
{error, Reason} -> {error, Reason}
end.
%% @private
-spec start_link(module(), macula:pool(), macula:realm(), boolean(),
term(), term()) -> {ok, pid()} | {error, term()}.
start_link(Module, Pool, Realm, Announce, InitArgs, Payload) ->
gen_server:start_link(?MODULE,
{Module, Pool, Realm, Announce, InitArgs, Payload}, []).
%%%===================================================================
%%% gen_server callbacks
%%%===================================================================
%% @private
init({Module, Pool, Realm, Announce, InitArgs, Payload}) ->
case Module:init(InitArgs) of
{ok, UserState} ->
RequestId = crypto:strong_rand_bytes(16),
publish(Announce, Pool, Realm, ?REQUEST_RECEIVED,
#{request_id => RequestId}),
{ok, #rstate{module = Module, pool = Pool, realm = Realm,
announce = Announce, 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,
request_id = RequestId}, Reply) ->
publish(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(false, _, _, _, _) -> ok;
publish(true, Pool, Realm, Topic, Payload) ->
_ = macula:publish(Pool, Realm, Topic, Payload), ok.