Packages
macula
13.0.0
13.0.1
13.0.0
12.12.0
12.11.1
12.11.0
12.10.0
12.9.1
12.9.0
12.8.0
12.7.0
12.6.0
12.5.1
12.5.0
12.4.0
12.3.0
12.2.1
12.2.0
12.1.0
12.0.0
11.5.0
11.4.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
Current section
Files
src/macula_feeder.erl
%%%-------------------------------------------------------------------
%%% @doc Behaviour for supervised content feeders (the share side).
%%%
%%% `start_link/4,5,6' returns immediately with a pid, shares the bytes from
%%% this node (`macula:share_content/4', D27: the node keeps and serves what
%%% it shares), delivers the root content id to `Module:handle_fed/2', and
%%% publishes `sharing.put_started_v1' / `sharing.put_completed_v1' mesh facts
%%% around it, including `outcome => cancelled' when `cancel/1' lands before
%%% the share resolves. A cancelled share is withdrawn: the root content id
%%% is known from the bytes before anything is sent, and it is unshared, so a
%%% cancel never leaves content shared behind the caller's back.
%%%
%%% This is content sharing, not general-purpose RPC streaming; see
%%% `macula_streamer' / `macula_stream_sink' for that.
%%%
%%% == Start options ==
%%%
%%% `share' and `unshare', as `macula:share_content/4' and
%%% `macula:unshare_content/3' (the defaults); `share_opts', passed to
%%% `share' (`name', `org'); `fact_publish', `macula:publish/4' by default. A
%%% function of another shape is refused with `function_clause', in the caller.
%%%
%%% == Example ==
%%%
%%% ```
%%% -module(doc_feeder).
%%% -behaviour(macula_feeder).
%%% -export([init/1, handle_fed/2]).
%%%
%%% init(Parent) -> {ok, Parent}.
%%%
%%% handle_fed(Result, Parent) ->
%%% Parent ! {fed, Result},
%%% {stop, normal, Parent}.
%%% '''
%%%
%%% ```
%%% {ok, Pid} = macula_feeder:start_link(doc_feeder, Pool, Realm,
%%% Bytes, self()).
%%% '''
%%% @end
%%%-------------------------------------------------------------------
-module(macula_feeder).
-behaviour(gen_server).
-export([start_link/4, 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_fed(Result :: {ok, macula:mcid()} | {error, term()}, State :: term()) ->
{noreply, NewState :: term()} | {stop, Reason :: term(), NewState :: term()}.
-define(PUT_STARTED, <<"sharing.put_started_v1">>).
-define(PUT_COMPLETED, <<"sharing.put_completed_v1">>).
-export_type([start_opts/0]).
-type start_opts() :: #{share => fun((macula:pool(), macula:realm(), binary(), map()) ->
{ok, macula:mcid()} | {error, term()}),
unshare => fun((macula:pool(), macula:realm(), macula:mcid()) -> ok),
share_opts => map(),
fact_publish => macula_lifetime_announcer:publish()}.
-record(fstate, {
module :: module(),
pool :: macula:pool(),
realm :: macula:realm(),
mcid :: macula:mcid(),
unshare :: fun((macula:pool(), macula:realm(), macula:mcid()) -> ok),
fact_publish :: macula_lifetime_announcer:publish(),
share_id :: binary(),
worker :: pid(),
completed :: boolean(),
user :: term()
}).
%% @doc Start a feeder sharing `Bytes' in `Realm' through `Pool'.
-spec start_link(module(), macula:pool(), macula:realm(), binary()) -> {ok, pid()} | {error, term()}.
start_link(Module, Pool, Realm, Bytes) ->
start_link(Module, Pool, Realm, Bytes, undefined).
%% @doc As `start_link/4', with `Args' passed to `Module:init/1'.
-spec start_link(module(), macula:pool(), macula:realm(), binary(), term()) -> {ok, pid()} | {error, term()}.
start_link(Module, Pool, Realm, Bytes, Args) ->
start_link(Module, Pool, Realm, Bytes, Args, #{}).
%% @doc As `start_link/5', with start options (see "Start options" above).
-spec start_link(module(), macula:pool(), macula:realm(), binary(), term(), start_opts()) ->
{ok, pid()} | {error, term()}.
start_link(Module, Pool, Realm, Bytes, Args, Opts) when is_binary(Bytes), is_map(Opts) ->
Functions = #{share => arity_4(maps:get(share, Opts, fun macula:share_content/4)),
unshare => arity_3(maps:get(unshare, Opts, fun macula:unshare_content/3)),
share_opts => maps:get(share_opts, Opts, #{}),
fact_publish => arity_4(maps:get(fact_publish, Opts, fun macula:publish/4))},
gen_server:start_link(?MODULE, {Module, Pool, Realm, Bytes, Args, Functions}, []).
%% @doc Cancel an in-flight feed. Publishes `sharing.put_completed_v1' with
%% `outcome => cancelled' and withdraws the share if it had not resolved yet.
-spec cancel(pid()) -> ok.
cancel(Pid) -> gen_server:stop(Pid).
arity_3(Fun) when is_function(Fun, 3) -> Fun.
arity_4(Fun) when is_function(Fun, 4) -> Fun.
%%%===================================================================
%%% gen_server callbacks
%%%===================================================================
%% @private
init({Module, Pool, Realm, Bytes, InitArgs, Functions}) ->
process_flag(trap_exit, true),
started(Module:init(InitArgs), Module, Pool, Realm, Bytes, Functions).
started({ok, User}, Module, Pool, Realm, Bytes,
#{share := Share, unshare := Unshare, share_opts := ShareOpts, fact_publish := Publish}) ->
ShareId = crypto:strong_rand_bytes(16),
_ = Publish(Pool, Realm, ?PUT_STARTED, #{share_id => ShareId, size => byte_size(Bytes)}),
%% The root the share will name, known before anything is sent, so a cancel can withdraw it.
{MCID, _} = macula_content_store:added(Bytes, ShareOpts, macula_content_store:new()),
Parent = self(),
Worker = spawn_link(fun() -> Parent ! {feed_result, Share(Pool, Realm, Bytes, ShareOpts)} end),
{ok, #fstate{module = Module, pool = Pool, realm = Realm, mcid = MCID, unshare = Unshare,
fact_publish = Publish, share_id = ShareId, worker = Worker, completed = false, user = User}};
started({stop, Reason}, _Module, _Pool, _Realm, _Bytes, _Functions) ->
{stop, Reason}.
%% @private
handle_call(_Request, _From, State) ->
{reply, {error, unsupported}, State}.
%% @private
handle_cast(_Msg, State) ->
{noreply, State}.
%% @private
handle_info({feed_result, Result}, #fstate{module = Module, user = User} = State) ->
deliver(Module:handle_fed(Result, User), completed(Result, State));
handle_info({'EXIT', Worker, Reason}, #fstate{worker = Worker} = State) when Reason =/= normal ->
{stop, {worker_crashed, Reason}, State};
handle_info(_Msg, State) ->
{noreply, State}.
deliver({noreply, User}, State) -> {noreply, State#fstate{user = User}};
deliver({stop, Reason, User}, State) -> {stop, Reason, State#fstate{user = User}}.
%% @private
terminate(_Reason, #fstate{completed = true, worker = Worker}) ->
unlink(Worker),
exit(Worker, kill),
ok;
terminate(_Reason, #fstate{worker = Worker, pool = Pool, realm = Realm, mcid = MCID, unshare = Unshare} = State) ->
unlink(Worker),
exit(Worker, kill),
_ = Unshare(Pool, Realm, MCID),
_ = completed({error, cancelled}, State),
ok.
completed(_Result, #fstate{completed = true} = State) ->
State;
completed(Result, #fstate{pool = Pool, realm = Realm, fact_publish = Publish, share_id = ShareId} = State) ->
_ = Publish(Pool, Realm, ?PUT_COMPLETED, outcome(#{share_id => ShareId}, Result)),
State#fstate{completed = true}.
outcome(Base, {ok, MCID}) -> Base#{outcome => completed, mcid => MCID, chunked => is_chunked(MCID)};
outcome(Base, {error, cancelled}) -> Base#{outcome => cancelled};
outcome(Base, {error, Reason}) -> Base#{outcome => failed, reason => Reason}.
is_chunked(<<2, 16#56, _/binary>>) -> true;
is_chunked(_) -> false.