Packages

macula

10.25.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_request.erl
Raw

src/macula_request.erl

%%%-------------------------------------------------------------------
%%% @doc Behaviour for supervised RPC requests.
%%%
%%% `call/5' is a plain blocking call in the caller's own process —
%%% there is no addressable pid to cancel it from outside. This is the
%%% consumer-side counterpart to `macula_response': `start_link/6,7'
%%% returns immediately with a pid, runs `macula:call/5' in a linked
%%% worker, delivers the outcome to `Module:handle_reply/2', and
%%% publishes `rpc.sent_v1' / `rpc.completed_v1' mesh facts around the
%%% request — including `outcome => cancelled' if the request is
%%% cancelled before a reply arrives.
%%%
%%% == Call and publish functions ==
%%%
%%% `start_link/8' takes `call', the function the request calls with,
%%% `macula:call/5' by default. `start_link_direct/8' takes `direct_call',
%%% `macula_direct_dial:call/6' by default, which gets the other options.
%%% Both take `fact_publish', the function the request announces its facts
%%% with, `macula:publish/4' by default. A test gives its own functions
%%% this way instead of replacing a module.
%%%
%%% == Example ==
%%%
%%% ```
%%% -module(add_caller).
%%% -behaviour(macula_request).
%%% -export([init/1, handle_reply/2]).
%%%
%%% init(Parent) -> {ok, Parent}.
%%%
%%% handle_reply(Result, Parent) ->
%%% Parent ! {add_result, Result},
%%% {stop, normal, Parent}.
%%% '''
%%%
%%% ```
%%% {ok, Pid} = macula_request:start_link(add_caller, Pool, Realm,
%%% <<"math.add_v1">>, #{a => 2, b => 3}, 30_000, self()).
%%% '''
%%%
%%% == Direct-dial ==
%%%
%%% `start_link/6,7' routes through the pool's existing links — first
%%% success across whichever are healthy, the same gossip-propagated
%%% routing `call/5' always used. `start_link_direct/6,7' is the
%%% direct-dial counterpart: it resolves the procedure's
%%% `procedure_advertisement' from the DHT (published by
%%% `macula_response:advertise_direct/6' on the provider side),
%%% resolves that record's `serving_station' to a dialable endpoint via
%%% the station's own `station_endpoint' record (every macula-station
%%% publishes its own automatically), and calls there in one hop via
%%% `macula:call_station/6' — instead of depending on advertise-gossip
%%% having propagated a route between arbitrary stations. Requires the
%%% provider to have advertised via `advertise_direct/6', not plain
%%% `advertise/5' — a plain advertise publishes no discoverable record.
%%% @end
%%%-------------------------------------------------------------------
-module(macula_request).
-behaviour(gen_server).
-export([start_link/6, start_link/7, start_link/8]).
-export([start_link_direct/6, start_link_direct/7, start_link_direct/8]).
-export_type([call/0, direct_call/0, start_opts/0, direct_opts/0]).
-export([cancel/1]).
-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_reply(Result :: {ok, term()} | {error, term()}, State :: term()) ->
{noreply, NewState :: term()} | {stop, Reason :: term(), NewState :: term()}.
-optional_callbacks([]).
-define(REQUEST_SENT, <<"rpc.sent_v1">>).
-define(REQUEST_COMPLETED, <<"rpc.completed_v1">>).
-type call() :: fun((macula:pool(), macula:realm(), macula:procedure(), term(),
pos_integer()) -> {ok, term()} | {error, term()}).
-type direct_call() :: fun((macula:pool(), macula:realm(), macula:procedure(), term(),
pos_integer(), map()) -> {ok, term()} | {error, term()}).
-type start_opts() :: #{call => call(), fact_publish => macula_lifetime_announcer:publish()}.
-type direct_opts() :: #{direct_call => direct_call(),
fact_publish => macula_lifetime_announcer:publish(),
atom() => term()}.
-record(qstate, {
module :: module(),
pool :: macula:pool(),
realm :: macula:realm(),
announce :: boolean(),
fact_publish :: macula_lifetime_announcer:publish(),
request_id :: binary(),
worker :: pid(),
completed :: boolean(),
user :: term()
}).
%% @doc Start a request. Calls `Procedure' on `(Pool, Realm)' with
%% `Payload', timing out after `TimeoutMs'; `Args' is passed to
%% `Module:init/1'.
-spec start_link(module(), macula:pool(), macula:realm(), macula:procedure(),
term(), pos_integer()) -> {ok, pid()} | {error, term()}.
start_link(Module, Pool, Realm, Procedure, Payload, TimeoutMs) ->
start_link(Module, Pool, Realm, Procedure, Payload, TimeoutMs, undefined).
%% @doc As `start_link/6', with `Args' passed to `Module:init/1'.
-spec start_link(module(), macula:pool(), macula:realm(), macula:procedure(),
term(), pos_integer(), term()) -> {ok, pid()} | {error, term()}.
start_link(Module, Pool, Realm, Procedure, Payload, TimeoutMs, Args) ->
start_link(Module, Pool, Realm, Procedure, Payload, TimeoutMs, Args, #{}).
%% @doc As `start_link/7', with options: `call' and `fact_publish' give
%% the functions the request calls and announces with (see "Call and
%% publish functions" above).
-spec start_link(module(), macula:pool(), macula:realm(), macula:procedure(),
term(), pos_integer(), term(), start_opts()) ->
{ok, pid()} | {error, term()}.
start_link(Module, Pool, Realm, Procedure, Payload, TimeoutMs, Args, Opts) when is_map(Opts) ->
start(arity_5(maps:get(call, Opts, fun macula:call/5)), Opts,
{Module, Pool, Realm, Procedure, Payload, TimeoutMs, Args}).
%% @doc As `start_link/6', but resolves and dials the serving station
%% directly instead of routing through the pool's existing links. See
%% the "Direct-dial" section above.
-spec start_link_direct(module(), macula:pool(), macula:realm(),
macula:procedure(), term(), pos_integer()) ->
{ok, pid()} | {error, term()}.
start_link_direct(Module, Pool, Realm, Procedure, Payload, TimeoutMs) ->
start_link_direct(Module, Pool, Realm, Procedure, Payload, TimeoutMs, undefined).
%% @doc As `start_link_direct/6', with `Args' passed to `Module:init/1'.
-spec start_link_direct(module(), macula:pool(), macula:realm(),
macula:procedure(), term(), pos_integer(), term()) ->
{ok, pid()} | {error, term()}.
start_link_direct(Module, Pool, Realm, Procedure, Payload, TimeoutMs, Args) ->
start_link_direct(Module, Pool, Realm, Procedure, Payload, TimeoutMs,
Args, #{}).
%% @doc As `start_link_direct/7', with options: `direct_call' and
%% `fact_publish' give the functions the request calls and announces with
%% (see "Call and publish functions" above), and the other options go to
%% the call as `macula_direct_dial:call/6' takes them, for example
%% `verify_cert_chain => {RealmCaPem, Org}' (Slice 7c Direction B, managed
%% realms only; see `macula_direct_dial''s module doc, "Trust model").
-spec start_link_direct(module(), macula:pool(), macula:realm(),
macula:procedure(), term(), pos_integer(), term(),
direct_opts()) -> {ok, pid()} | {error, term()}.
start_link_direct(Module, Pool, Realm, Procedure, Payload, TimeoutMs, Args,
Opts) when is_map(Opts) ->
DirectCall = arity_6(maps:get(direct_call, Opts, fun macula_direct_dial:call/6)),
DialOpts = maps:without([direct_call, fact_publish], Opts),
Call = fun(CallPool, CallRealm, CallProcedure, CallPayload, CallTimeoutMs) ->
DirectCall(CallPool, CallRealm, CallProcedure, CallPayload, CallTimeoutMs,
DialOpts)
end,
start(Call, Opts, {Module, Pool, Realm, Procedure, Payload, TimeoutMs, Args}).
%% A request starts with its call function and the options' fact publish
%% function, or macula:publish/4 without one. A function option of the
%% wrong arity is refused with function_clause, in the caller.
start(Call, Opts, {Module, Pool, Realm, Procedure, Payload, TimeoutMs, Args}) ->
FactPublish = arity_4(maps:get(fact_publish, Opts, fun macula:publish/4)),
gen_server:start_link(?MODULE,
{Call, FactPublish, Module, Pool, Realm, Procedure, Payload,
TimeoutMs, true, Args}, []).
arity_4(Fun) when is_function(Fun, 4) -> Fun.
arity_5(Fun) when is_function(Fun, 5) -> Fun.
arity_6(Fun) when is_function(Fun, 6) -> Fun.
%% @doc Cancel an in-flight request. Publishes `rpc.completed_v1' with
%% `outcome => cancelled' if no reply had arrived yet.
-spec cancel(pid()) -> ok.
cancel(Pid) -> gen_server:stop(Pid).
%%%===================================================================
%%% gen_server callbacks
%%%===================================================================
%% @private
init({Call, FactPublish, Module, Pool, Realm, Procedure, Payload, TimeoutMs, Announce,
InitArgs}) ->
process_flag(trap_exit, true),
case Module:init(InitArgs) of
{ok, UserState} ->
RequestId = crypto:strong_rand_bytes(16),
publish(FactPublish, Announce, Pool, Realm, ?REQUEST_SENT,
#{request_id => RequestId}),
Worker = spawn_worker(Call, Pool, Realm, Procedure, Payload, TimeoutMs),
{ok, #qstate{module = Module, pool = Pool, realm = Realm,
announce = Announce, fact_publish = FactPublish,
request_id = RequestId,
worker = Worker, completed = false, user = UserState}};
{stop, Reason} ->
{stop, Reason}
end.
spawn_worker(Call, Pool, Realm, Procedure, Payload, TimeoutMs) ->
Parent = self(),
spawn_link(fun() ->
Result = Call(Pool, Realm, Procedure, Payload, TimeoutMs),
Parent ! {request_result, Result}
end).
%% @private
handle_call(_Request, _From, State) ->
{reply, {error, unsupported}, State}.
%% @private
handle_cast(_Msg, State) -> {noreply, State}.
%% @private
handle_info({request_result, Result}, State) ->
NewState = announce_completed(State, Result),
#qstate{module = Module, user = User} = NewState,
deliver(Module:handle_reply(Result, User), NewState);
handle_info({'EXIT', Worker, Reason}, #qstate{worker = Worker} = State)
when Reason =/= normal ->
{stop, {worker_crashed, Reason}, State};
handle_info(_Msg, State) ->
{noreply, State}.
deliver({noreply, NewUser}, State) -> {noreply, State#qstate{user = NewUser}};
deliver({stop, Reason, NewUser}, State) -> {stop, Reason, State#qstate{user = NewUser}}.
%% @private
terminate(_Reason, #qstate{worker = Worker, completed = true}) ->
unlink(Worker),
exit(Worker, kill),
ok;
terminate(_Reason, State) ->
unlink(State#qstate.worker),
exit(State#qstate.worker, kill),
_ = announce_completed(State, {error, cancelled}),
ok.
announce_completed(#qstate{completed = true} = State, _Result) ->
State;
announce_completed(#qstate{pool = Pool, realm = Realm, announce = Announce,
fact_publish = FactPublish,
request_id = RequestId} = State, Result) ->
publish(FactPublish, Announce, Pool, Realm, ?REQUEST_COMPLETED,
outcome_fields(#{request_id => RequestId}, Result)),
State#qstate{completed = true}.
outcome_fields(Base, {ok, _}) -> Base#{outcome => completed};
outcome_fields(Base, {error, cancelled}) -> Base#{outcome => cancelled};
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.