Packages
macula
11.3.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 put/share side).
%%%
%%% `start_link/4,5' returns immediately with a pid, delivers the
%%% outcome to `Module:handle_fed/2', and publishes
%%% `sharing.put_started_v1' / `sharing.put_completed_v1' mesh facts
%%% around the transfer — including `outcome => cancelled' if the
%%% feeder is stopped before the put resolves.
%%%
%%% This is content sharing, not general-purpose RPC streaming — see
%%% `macula_streamer' / `macula_stream_sink' for that (`streaming.*'
%%% facts belong to that pair).
%%%
%%% == Real cancel, real underneath ==
%%%
%%% Internally this drives `macula_content_transfer' (PLAN_PUSH_UPLOAD.md
%%% Phase 4) rather than a blocking `macula:put_content/2' call run in
%%% a linked worker — that blocking shape had no addressable handle to
%%% the actual transfer, so `cancel/1' (`gen_server:stop/1') could only
%%% ever kill the local worker process waiting on it, never touch the
%%% underlying stream. A `macula_content_transfer' cancelled that way
%%% doesn't even notice: nothing links a `gen_server:call' caller's
%%% death to the callee, so it would run to completion, or sit
%%% resolved-but-never-reaped, forever — orphaned, leaking its
%%% `content_stream_bufs' entry on the link and its
%%% `macula_content_transfer_registry' entry, for no purpose. This
%%% module now holds the `macula_content_transfer' pid directly (a
%%% `content_transfer' state field, alongside the lightweight resolve
%%% + await proxy `worker', which asks this process to start the
%%% transfer) so `cancel/1' can call
%%% `macula_content_transfer:cancel/1' on it for real — the same
%%% peer-visible QUIC RESET_STREAM abort described there, not a local
%%% kill with nothing downstream the wiser. The share_id this module
%%% already minted for its own `sharing.*' mesh facts is threaded
%%% through as `macula_content_transfer''s own `share_id' too, so both
%%% layers resolve to the same id.
%%%
%%% == Direct-dial ==
%%%
%%% `start_link/4,5' puts through the pool's own connected link
%%% (whichever `pick_connected_link/1' picks). `start_link_direct/4,5'
%%% is the direct-dial counterpart: unlike `macula_download''s (which
%%% resolves an MCID to find out WHO has it), a PUT already knows its
%%% own target — the caller names `Station' directly, and it is
%%% resolved to a dialable endpoint via that station's own signed
%%% `station_endpoint' record (`macula_direct_dial:resolve_station_endpoint/2',
%%% a fast, non-addressable DHT lookup that stays a plain blocking call
%%% inside the resolve+await proxy — nothing has ever needed to cancel
%%% mid-resolve) and dialed in one hop via `macula_content_transfer:
%%% start_put_station/5', deliberately seeding that specific station
%%% instead of whichever the pool picks. See `macula_direct_dial''s
%%% module doc, "Content" section, for the trust model.
%%%
%%% == Transfer I/O ==
%%%
%%% A feeder starts, awaits and cancels its transfer with `start_put/3',
%%% `start_put_station/5', `await/1' and `cancel/1', the
%%% `macula_content_transfer' ones by default; resolves its station for
%%% `start_link_direct' with `resolve_station_endpoint/2',
%%% `macula_direct_dial''s by default; and announces its facts with
%%% `fact_publish', `macula:publish/4' by default. `start_link/6' and
%%% `start_link_direct/7' take them in their start options, the four
%%% transfer functions as `transfer_io', checked by
%%% `macula_content_transfer:transfer_io/2', and pass a `link_io' option on
%%% to the transfer they start. 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([start_link_direct/5, start_link_direct/6, start_link_direct/7]).
-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">>).
%% Bounds only the QUIC handshake wait when `start_link_direct/5,6'
%% must dial a fresh link — matches `macula_client:connect/2''s own
%% `connect_timeout_ms' default. The block/manifest transfer that
%% follows has its own separate, internal timeouts regardless.
-define(DIRECT_DIAL_CONNECT_TIMEOUT_MS, 30_000).
-export_type([start_opts/0]).
-type start_opts() :: #{transfer_io => macula_content_transfer:transfer_io(),
resolve_station_endpoint => fun((macula:pool(), macula_identity:pubkey()) ->
{ok, binary()} | {error, term()}),
fact_publish => macula_lifetime_announcer:publish(),
link_io => macula_content_transfer:link_io()}.
-record(fstate, {
module :: module(),
pool :: macula:pool(),
realm :: macula:realm(),
announce :: boolean(),
transfer_io :: macula_content_transfer:transfer_io(),
fact_publish :: macula_lifetime_announcer:publish(),
share_id :: binary(),
worker :: pid(),
content_transfer :: pid() | undefined,
completed :: boolean(),
user :: term()
}).
%% @doc Start a feeder. Puts `Bytes' into content storage via `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 "Transfer I/O" 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_map(Opts) ->
gen_server:start_link(?MODULE,
{pooled, Module, Pool, Realm, Bytes, true, Args, functions(Opts)}, []).
%% @doc As `start_link/4', but resolves `Station''s own
%% `station_endpoint' and dials it directly instead of putting through
%% the pool's existing links. See the "Direct-dial" section above.
-spec start_link_direct(module(), macula:pool(), <<_:256>>,
macula:realm(), binary()) ->
{ok, pid()} | {error, term()}.
start_link_direct(Module, Pool, Station, Realm, Bytes) ->
start_link_direct(Module, Pool, Station, Realm, Bytes, undefined).
%% @doc As `start_link_direct/5', with `Args' passed to `Module:init/1'.
-spec start_link_direct(module(), macula:pool(), <<_:256>>,
macula:realm(), binary(), term()) ->
{ok, pid()} | {error, term()}.
start_link_direct(Module, Pool, Station, Realm, Bytes, Args) ->
start_link_direct(Module, Pool, Station, Realm, Bytes, Args, #{}).
%% @doc As `start_link_direct/6', with start options (see "Transfer I/O"
%% above).
-spec start_link_direct(module(), macula:pool(), macula_identity:pubkey(),
macula:realm(), binary(), term(), start_opts()) ->
{ok, pid()} | {error, term()}.
start_link_direct(Module, Pool, Station, Realm, Bytes, Args, Opts) when is_map(Opts) ->
gen_server:start_link(?MODULE,
{direct, Module, Pool, Station, Realm, Bytes, true, Args, functions(Opts)}, []).
%% @doc Cancel an in-flight feed. Publishes `sharing.put_completed_v1'
%% with `outcome => cancelled' if the put had not resolved yet.
-spec cancel(pid()) -> ok.
cancel(Pid) -> gen_server:stop(Pid).
%% The functions a feeder runs on, from its start options or else the
%% defaults; one of another shape is refused with function_clause, in the
%% caller. `link_io' goes on to the transfer the feeder starts.
functions(Opts) ->
TransferIo = macula_content_transfer:transfer_io(default_transfer_io(),
maps:get(transfer_io, Opts, undefined)),
#{transfer_io => TransferIo,
resolve_station_endpoint =>
arity_2(maps:get(resolve_station_endpoint, Opts,
fun macula_direct_dial:resolve_station_endpoint/2)),
fact_publish => arity_4(maps:get(fact_publish, Opts, fun macula:publish/4)),
transfer_opts => maps:with([link_io], Opts)}.
default_transfer_io() ->
#{start_put => fun macula_content_transfer:start_put/3,
start_put_station => fun macula_content_transfer:start_put_station/5,
await => fun macula_content_transfer:await/1,
cancel => fun macula_content_transfer:cancel/1}.
arity_2(Fun) when is_function(Fun, 2) -> Fun.
arity_4(Fun) when is_function(Fun, 4) -> Fun.
%%%===================================================================
%%% gen_server callbacks
%%%===================================================================
%% @private
init({pooled, Module, Pool, Realm, Bytes, Announce, InitArgs, Functions}) ->
start_feeder(Module, InitArgs, Pool, Realm, Bytes, Announce, Functions,
fun(ShareId) -> spawn_worker(pooled, Functions, Pool, Bytes, ShareId) end);
init({direct, Module, Pool, Station, Realm, Bytes, Announce, InitArgs, Functions}) ->
start_feeder(Module, InitArgs, Pool, Realm, Bytes, Announce, Functions,
fun(ShareId) -> spawn_worker(direct, Functions, Pool, Station, Bytes, ShareId) end).
start_feeder(Module, InitArgs, Pool, Realm, Bytes, Announce,
#{transfer_io := TransferIo, fact_publish := FactPublish}, SpawnFun) ->
process_flag(trap_exit, true),
case Module:init(InitArgs) of
{ok, UserState} ->
ShareId = crypto:strong_rand_bytes(16),
publish(Announce, FactPublish, Pool, Realm, ?PUT_STARTED,
#{share_id => ShareId, size => byte_size(Bytes)}),
Worker = SpawnFun(ShareId),
{ok, #fstate{module = Module, pool = Pool, realm = Realm,
announce = Announce, transfer_io = TransferIo,
fact_publish = FactPublish, share_id = ShareId,
worker = Worker, content_transfer = undefined,
completed = false, user = UserState}};
{stop, Reason} ->
{stop, Reason}
end.
%% The lightweight proxy: have this feeder start the addressable
%% transfer (see `run_transfer/2'), block for the outcome, reap the
%% transfer (a no-op if it's already being cancelled from outside, see
%% `reap_content_transfer/1'), report the outcome.
spawn_worker(pooled, #{transfer_io := TransferIo, transfer_opts := TransferOpts}, Pool, Bytes,
ShareId) ->
Parent = self(),
Start = pooled_put(TransferIo, Pool, Bytes, TransferOpts#{share_id => ShareId}),
spawn_link(fun() -> run_transfer(Parent, TransferIo, Start) end).
%% Resolving `Station''s endpoint stays a plain blocking DHT lookup
%% here (matches what `macula_direct_dial:put_content/4' already did);
%% only the transfer itself becomes addressable.
spawn_worker(direct, Functions, Pool, Station, Bytes, ShareId) ->
Parent = self(),
spawn_link(fun() -> direct_worker_run(Functions, Pool, Station, Bytes, ShareId, Parent) end).
direct_worker_run(#{transfer_io := TransferIo, transfer_opts := TransferOpts,
resolve_station_endpoint := ResolveStationEndpoint},
Pool, Station, Bytes, ShareId, Parent) ->
case ResolveStationEndpoint(Pool, Station) of
{ok, DialUrl} ->
Opts = TransferOpts#{share_id => ShareId, expected_node_id => Station,
pin_tls_cert => false, verify => none},
run_transfer(Parent, TransferIo, station_put(TransferIo, Pool, DialUrl, Bytes, Opts));
{error, Reason} ->
Parent ! {feed_result, {error, {unresolved, Reason}}}
end.
pooled_put(#{start_put := StartPut}, Pool, Bytes, Opts) ->
fun() -> StartPut(Pool, Bytes, Opts) end.
station_put(#{start_put_station := StartPutStation}, Pool, DialUrl, Bytes, Opts) ->
fun() -> StartPutStation(Pool, DialUrl, Bytes, ?DIRECT_DIAL_CONNECT_TIMEOUT_MS, Opts) end.
%% This feeder starts the transfer itself, in a call from the worker, so
%% the transfer's pid is in the state before any `cancel/1' is handled:
%% a cancel handled first finds nothing started, and a cancel handled
%% after finds the pid.
run_transfer(Parent, #{await := Await} = TransferIo, Start) ->
{ok, CTPid} = gen_server:call(Parent, {start_transfer, Start}, infinity),
Result = Await(CTPid),
reap_content_transfer(TransferIo, CTPid),
Parent ! {feed_result, Result}.
%% @private
handle_call({start_transfer, Start}, _From, State) ->
Started = Start(),
{reply, Started, record_transfer(Started, State)};
handle_call(_Request, _From, State) ->
{reply, {error, unsupported}, State}.
record_transfer({ok, CTPid}, State) -> State#fstate{content_transfer = CTPid};
record_transfer(_NotStarted, State) -> State.
%% @private
handle_cast(_Msg, State) -> {noreply, State}.
%% @private
handle_info({feed_result, Result}, State) ->
NewState = announce_completed(State, Result),
#fstate{module = Module, user = User} = NewState,
deliver(Module:handle_fed(Result, User), NewState#fstate{content_transfer = undefined});
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, NewUser}, State) -> {noreply, State#fstate{user = NewUser}};
deliver({stop, Reason, NewUser}, State) -> {stop, Reason, State#fstate{user = NewUser}}.
%% @private
terminate(_Reason, #fstate{worker = Worker, completed = true}) ->
unlink(Worker),
exit(Worker, kill),
ok;
terminate(_Reason, #fstate{transfer_io = TransferIo, content_transfer = CTPid} = State) ->
unlink(State#fstate.worker),
exit(State#fstate.worker, kill),
reap_content_transfer(TransferIo, CTPid),
_ = announce_completed(State, {error, cancelled}),
ok.
%% Killing the proxy `worker' does NOT cascade into stopping the
%% `macula_content_transfer' it waits on — that gen_server traps exits
%% and doesn't recognize this feeder or the proxy as one of ITS OWN
%% tracked workers, so an incoming exit signal just falls through its
%% catch-all `handle_info' clause, unnoticed. This is the actual fix:
%% reach in and cancel it explicitly. `undefined' covers a cancel
%% handled before the worker asked for the transfer to start (still
%% resolving, for direct-dial): nothing was started. `catch' covers the
%% benign race where the proxy's own natural reap (in `run_transfer/2')
%% and an external `cancel/1' land at the same time — the second
%% `cancel/1' call reaches an already-dead pid.
reap_content_transfer(_TransferIo, undefined) -> ok;
reap_content_transfer(#{cancel := Cancel}, CTPid) ->
try Cancel(CTPid) catch _:_ -> ok end,
ok.
announce_completed(#fstate{completed = true} = State, _Result) ->
State;
announce_completed(#fstate{pool = Pool, realm = Realm, announce = Announce,
fact_publish = FactPublish, share_id = ShareId} = State,
Result) ->
publish(Announce, FactPublish, Pool, Realm, ?PUT_COMPLETED,
outcome_fields(#{share_id => ShareId}, Result)),
State#fstate{completed = true}.
outcome_fields(Base, {ok, Mcid}) ->
Base#{outcome => completed, mcid => Mcid, chunked => is_chunked_mcid(Mcid)};
outcome_fields(Base, {error, cancelled}) ->
Base#{outcome => cancelled};
outcome_fields(Base, {error, Reason}) ->
Base#{outcome => failed, reason => Reason}.
is_chunked_mcid(<<2, 16#56, _/binary>>) -> true;
is_chunked_mcid(_) -> false.
publish(false, _FactPublish, _, _, _, _) -> ok;
publish(true, FactPublish, Pool, Realm, Topic, Payload) ->
_ = FactPublish(Pool, Realm, Topic, Payload), ok.