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
Current section
Files
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.
%%%
%%% == 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]).
-export([start_link_direct/6, start_link_direct/7, start_link_direct/8]).
-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">>).
-record(qstate, {
module :: module(),
pool :: macula:pool(),
realm :: macula:realm(),
announce :: boolean(),
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) ->
gen_server:start_link(?MODULE,
{pooled, Module, Pool, Realm, Procedure, Payload, TimeoutMs, true,
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 `Opts' forwarded to
%% `macula_direct_dial:call/6' — e.g. `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(),
map()) -> {ok, pid()} | {error, term()}.
start_link_direct(Module, Pool, Realm, Procedure, Payload, TimeoutMs, Args,
Opts) ->
gen_server:start_link(?MODULE,
{direct, Module, Pool, Realm, Procedure, Payload, TimeoutMs, true,
Args, Opts}, []).
%% @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({DialMode, Module, Pool, Realm, Procedure, Payload, TimeoutMs, Announce,
InitArgs, Opts}) ->
process_flag(trap_exit, true),
case Module:init(InitArgs) of
{ok, UserState} ->
RequestId = crypto:strong_rand_bytes(16),
publish(Announce, Pool, Realm, ?REQUEST_SENT,
#{request_id => RequestId}),
Worker = spawn_worker(DialMode, Pool, Realm, Procedure, Payload,
TimeoutMs, Opts),
{ok, #qstate{module = Module, pool = Pool, realm = Realm,
announce = Announce, request_id = RequestId,
worker = Worker, completed = false, user = UserState}};
{stop, Reason} ->
{stop, Reason}
end.
spawn_worker(pooled, Pool, Realm, Procedure, Payload, TimeoutMs, _Opts) ->
Parent = self(),
spawn_link(fun() ->
Result = macula:call(Pool, Realm, Procedure, Payload, TimeoutMs),
Parent ! {request_result, Result}
end);
spawn_worker(direct, Pool, Realm, Procedure, Payload, TimeoutMs, Opts) ->
Parent = self(),
spawn_link(fun() ->
Result = macula_direct_dial:call(Pool, Realm, Procedure, Payload,
TimeoutMs, Opts),
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,
request_id = RequestId} = State, Result) ->
publish(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(false, _, _, _, _) -> ok;
publish(true, Pool, Realm, Topic, Payload) ->
_ = macula:publish(Pool, Realm, Topic, Payload), ok.