Packages

macula

11.0.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_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 ==
%%%
%%% Both starts reach a provider the same way. `macula:call/5' 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 the advertised provider
%%% there in one hop via `macula:call_station/8'. `start_link_direct'
%%% also hands its options to that resolution. 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(), 1..600_000) -> {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(), 1..600_000, 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(), 1..600_000, term(), start_opts()) ->
{ok, pid()} | {error, term()}.
start_link(Module, Pool, Realm, Procedure, Payload, TimeoutMs, Args, Opts)
when is_integer(TimeoutMs), TimeoutMs > 0, TimeoutMs =< 600_000, 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(), 1..600_000) ->
{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(), 1..600_000, 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. An org namespaced
%% procedure's authorization is checked against the realm key the pool
%% pinned (see `macula_direct_dial''s module doc, "Trust model").
-spec start_link_direct(module(), macula:pool(), macula:realm(),
macula:procedure(), term(), 1..600_000, term(),
direct_opts()) -> {ok, pid()} | {error, term()}.
start_link_direct(Module, Pool, Realm, Procedure, Payload, TimeoutMs, Args, Opts)
when is_integer(TimeoutMs), TimeoutMs > 0, TimeoutMs =< 600_000, 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.