Packages

macula

10.5.2
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_publisher.erl
Raw

src/macula_publisher.erl

%%%-------------------------------------------------------------------
%%% @doc Behaviour for supervised content publishers.
%%%
%%% `macula:publish/4' is a plain blocking call — no addressable pid to
%%% cancel it from outside. This is the missing supervised counterpart
%%% to `macula_subscriber': every other primitive pair already has one
%%% on each side (`macula_request'/`macula_response' for RPC,
%%% `macula_streamer'/`macula_stream_sink' for streaming RPC,
%%% `macula_feeder'/`macula_download' for content sharing) — pubsub had
%%% only the consumer half. `start_link/5,6' returns immediately with a
%%% pid, runs `macula:publish/4' in a linked worker, delivers the
%%% outcome to `Module:handle_published/2', and publishes
%%% `pubsub.publish_started_v1' / `pubsub.publish_completed_v1' mesh
%%% facts around the publish — including `outcome => cancelled' if the
%%% publisher is stopped before the publish resolves.
%%%
%%% == Example ==
%%%
%%% ```
%%% -module(status_publisher).
%%% -behaviour(macula_publisher).
%%% -export([init/1, handle_published/2]).
%%%
%%% init(Parent) -> {ok, Parent}.
%%%
%%% handle_published(Result, Parent) ->
%%% Parent ! {published, Result},
%%% {stop, normal, Parent}.
%%% '''
%%%
%%% ```
%%% {ok, Pid} = macula_publisher:start_link(status_publisher, Pool, Realm,
%%% Topic, Payload, self()).
%%% '''
%%% @end
%%%-------------------------------------------------------------------
-module(macula_publisher).
-behaviour(gen_server).
-export([start_link/5, start_link/6]).
-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_published(Result :: ok | {error, term()}, State :: term()) ->
{noreply, NewState :: term()} | {stop, Reason :: term(), NewState :: term()}.
-define(PUBLISH_STARTED, <<"pubsub.publish_started_v1">>).
-define(PUBLISH_COMPLETED, <<"pubsub.publish_completed_v1">>).
-record(pstate, {
module :: module(),
pool :: macula:pool(),
realm :: macula:realm(),
announce :: boolean(),
publish_id :: binary(),
worker :: pid(),
completed :: boolean(),
user :: term()
}).
%% @doc Start a publisher. Publishes `Payload' on `Topic' via `Pool'.
-spec start_link(module(), macula:pool(), macula:realm(), macula:topic(),
term()) -> {ok, pid()} | {error, term()}.
start_link(Module, Pool, Realm, Topic, Payload) ->
start_link(Module, Pool, Realm, Topic, Payload, undefined).
%% @doc As `start_link/5', with `Args' passed to `Module:init/1'.
-spec start_link(module(), macula:pool(), macula:realm(), macula:topic(),
term(), term()) -> {ok, pid()} | {error, term()}.
start_link(Module, Pool, Realm, Topic, Payload, Args) ->
gen_server:start_link(?MODULE,
{Module, Pool, Realm, Topic, Payload, true, Args}, []).
%% @doc Cancel an in-flight publish. Publishes
%% `pubsub.publish_completed_v1' with `outcome => cancelled' if the
%% publish had not resolved yet.
-spec cancel(pid()) -> ok.
cancel(Pid) -> gen_server:stop(Pid).
%%%===================================================================
%%% gen_server callbacks
%%%===================================================================
%% @private
init({Module, Pool, Realm, Topic, Payload, Announce, InitArgs}) ->
process_flag(trap_exit, true),
case Module:init(InitArgs) of
{ok, UserState} ->
PublishId = crypto:strong_rand_bytes(16),
publish(Announce, Pool, Realm, ?PUBLISH_STARTED,
#{publish_id => PublishId, topic => Topic}),
Worker = spawn_worker(Pool, Realm, Topic, Payload),
{ok, #pstate{module = Module, pool = Pool, realm = Realm,
announce = Announce, publish_id = PublishId,
worker = Worker, completed = false, user = UserState}};
{stop, Reason} ->
{stop, Reason}
end.
spawn_worker(Pool, Realm, Topic, Payload) ->
Parent = self(),
spawn_link(fun() ->
Result = macula:publish(Pool, Realm, Topic, Payload),
Parent ! {publish_result, Result}
end).
%% @private
handle_call(_Request, _From, State) ->
{reply, {error, unsupported}, State}.
%% @private
handle_cast(_Msg, State) -> {noreply, State}.
%% @private
handle_info({publish_result, Result}, State) ->
NewState = announce_completed(State, Result),
#pstate{module = Module, user = User} = NewState,
deliver(Module:handle_published(Result, User), NewState);
handle_info({'EXIT', Worker, Reason}, #pstate{worker = Worker} = State)
when Reason =/= normal ->
{stop, {worker_crashed, Reason}, State};
handle_info(_Msg, State) ->
{noreply, State}.
deliver({noreply, NewUser}, State) -> {noreply, State#pstate{user = NewUser}};
deliver({stop, Reason, NewUser}, State) -> {stop, Reason, State#pstate{user = NewUser}}.
%% @private
terminate(_Reason, #pstate{worker = Worker, completed = true}) ->
unlink(Worker),
exit(Worker, kill),
ok;
terminate(_Reason, State) ->
unlink(State#pstate.worker),
exit(State#pstate.worker, kill),
_ = announce_completed(State, {error, cancelled}),
ok.
announce_completed(#pstate{completed = true} = State, _Result) ->
State;
announce_completed(#pstate{pool = Pool, realm = Realm, announce = Announce,
publish_id = PublishId} = State, Result) ->
publish(Announce, Pool, Realm, ?PUBLISH_COMPLETED,
outcome_fields(#{publish_id => PublishId}, Result)),
State#pstate{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.