Packages
macula
9.13.8
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' that reports it back) 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.
%%%
%%% == 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]).
-export([start_link_direct/5, start_link_direct/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">>).
%% 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).
-record(fstate, {
module :: module(),
pool :: macula:pool(),
realm :: macula:realm(),
announce :: boolean(),
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) ->
gen_server:start_link(?MODULE,
{pooled, Module, Pool, Realm, Bytes, true, Args}, []).
%% @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(), macula_identity:pubkey(),
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(), macula_identity:pubkey(),
macula:realm(), binary(), term()) ->
{ok, pid()} | {error, term()}.
start_link_direct(Module, Pool, Station, Realm, Bytes, Args) ->
gen_server:start_link(?MODULE,
{direct, Module, Pool, Station, Realm, Bytes, true, Args}, []).
%% @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).
%%%===================================================================
%%% gen_server callbacks
%%%===================================================================
%% @private
init({pooled, Module, Pool, Realm, Bytes, Announce, InitArgs}) ->
start_feeder(Module, InitArgs, Pool, Realm, Bytes, Announce,
fun(ShareId) -> spawn_worker(pooled, Pool, Bytes, ShareId) end);
init({direct, Module, Pool, Station, Realm, Bytes, Announce, InitArgs}) ->
start_feeder(Module, InitArgs, Pool, Realm, Bytes, Announce,
fun(ShareId) -> spawn_worker(direct, Pool, Station, Bytes, ShareId) end).
start_feeder(Module, InitArgs, Pool, Realm, Bytes, Announce, SpawnFun) ->
process_flag(trap_exit, true),
case Module:init(InitArgs) of
{ok, UserState} ->
ShareId = crypto:strong_rand_bytes(16),
publish(Announce, Pool, Realm, ?PUT_STARTED,
#{share_id => ShareId, size => byte_size(Bytes)}),
Worker = SpawnFun(ShareId),
{ok, #fstate{module = Module, pool = Pool, realm = Realm,
announce = Announce, share_id = ShareId,
worker = Worker, content_transfer = undefined,
completed = false, user = UserState}};
{stop, Reason} ->
{stop, Reason}
end.
%% The lightweight proxy: start the addressable transfer, report its
%% pid back immediately (so `terminate/2' can reach it even if this
%% proxy itself gets killed mid-flight), 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, Pool, Bytes, ShareId) ->
Parent = self(),
spawn_link(fun() ->
{ok, CTPid} = macula_content_transfer:start_put(Pool, Bytes, #{share_id => ShareId}),
Parent ! {content_transfer, CTPid},
Result = macula_content_transfer:await(CTPid),
catch macula_content_transfer:cancel(CTPid),
Parent ! {feed_result, Result}
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, Pool, Station, Bytes, ShareId) ->
Parent = self(),
spawn_link(fun() ->
case macula_direct_dial:resolve_station_endpoint(Pool, Station) of
{ok, DialUrl} ->
Opts = #{share_id => ShareId, expected_node_id => Station,
pin_tls_cert => false, verify => none},
{ok, CTPid} = macula_content_transfer:start_put_station(
Pool, DialUrl, Bytes, ?DIRECT_DIAL_CONNECT_TIMEOUT_MS, Opts),
Parent ! {content_transfer, CTPid},
Result = macula_content_transfer:await(CTPid),
catch macula_content_transfer:cancel(CTPid),
Parent ! {feed_result, Result};
{error, Reason} ->
Parent ! {feed_result, {error, {unresolved, Reason}}}
end
end).
%% @private
handle_call(_Request, _From, State) ->
{reply, {error, unsupported}, State}.
%% @private
handle_cast(_Msg, State) -> {noreply, State}.
%% @private
handle_info({content_transfer, CTPid}, State) ->
{noreply, State#fstate{content_transfer = CTPid}};
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{content_transfer = CTPid} = State) ->
unlink(State#fstate.worker),
exit(State#fstate.worker, kill),
reap_content_transfer(CTPid),
_ = announce_completed(State, {error, cancelled}),
ok.
%% Killing the proxy `worker' does NOT cascade into stopping the
%% `macula_content_transfer' it started — that gen_server traps exits
%% and doesn't recognize the proxy as one of ITS OWN tracked workers,
%% so an incoming `{'EXIT', Proxy, killed}' just falls through its
%% catch-all `handle_info' clause, unnoticed. This is the actual fix:
%% reach in and cancel it explicitly. `undefined' covers the window
%% before `{content_transfer, CTPid}' has arrived yet (still resolving,
%% for direct-dial) — nothing addressable exists to cancel there,
%% same as before this phase. `catch' covers the benign race where the
%% proxy's own natural reap (in `spawn_worker/4') and an external
%% `cancel/1' land at the same time — the second `cancel/1' call
%% reaches an already-dead pid.
reap_content_transfer(undefined) -> ok;
reap_content_transfer(CTPid) -> catch macula_content_transfer:cancel(CTPid), ok.
announce_completed(#fstate{completed = true} = State, _Result) ->
State;
announce_completed(#fstate{pool = Pool, realm = Realm, announce = Announce,
share_id = ShareId} = State, Result) ->
publish(Announce, 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(<<1, 16#56, _/binary>>) -> true;
is_chunked_mcid(_) -> false.
publish(false, _, _, _, _) -> ok;
publish(true, Pool, Realm, Topic, Payload) ->
_ = macula:publish(Pool, Realm, Topic, Payload), ok.