Packages

macula

9.8.1
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_stream_sink.erl
Raw

src/macula_stream_sink.erl

%%%-------------------------------------------------------------------
%%% @doc Behaviour for supervised streaming RPC consumers.
%%%
%%% `call_stream/5' hands back a raw stream pid; a real consumer has to
%%% hand-write a `recv/2' loop around it — the provider side already
%%% gets this for free via `advertise_stream/5''s callback handler, this
%%% is the missing consumer-side half. `macula_stream_sink' opens the
%%% stream for you, drives the `recv/2' loop in a linked reader process
%%% (so a slow or stuck `recv' never blocks your gen_server's own
%%% mailbox), and calls `Module:handle_chunk/2' once per item against
%%% state your module owns, `Module:handle_close/2' when the stream ends
%%% or errors.
%%%
%%% This is the general-purpose RPC streaming feature (`call_stream/5',
%%% e.g. a `logs.tail_v1'-style procedure) — unrelated to content
%%% sharing's own chunked-transfer protocol; see `macula_feeder' /
%%% `macula_download' for that.
%%%
%%% Publishes `streaming.started_v1' / `streaming.completed_v1' mesh
%%% facts around the stream's lifetime, from the consumer's own
%%% perspective — the provider side (`macula_streamer') publishes its
%%% own copy from its perspective; the two are not deduplicated,
%%% mirroring how `macula_feeder' / `macula_download' each announce
%%% their own side of a content transfer.
%%%
%%% == Direct-dial ==
%%%
%%% `start_link/5,6' opens through the pool's existing links — the same
%%% gossip-propagated routing `call_stream/5' always used.
%%% `start_link_direct/5,6' is the direct-dial counterpart: it resolves
%%% the procedure's `procedure_advertisement' from the DHT (published by
%%% `macula_streamer:advertise_direct/6,7' on the provider side) and
%%% opens the stream there directly, in one hop, instead of depending on
%%% advertise-gossip having propagated a route between arbitrary
%%% stations. Requires the provider to have advertised via
%%% `advertise_direct/6,7', not plain `advertise/5,6'. See
%%% `macula_direct_dial''s module doc, "Trust model".
%%%
%%% == Example ==
%%%
%%% ```
%%% -module(log_tailer).
%%% -behaviour(macula_stream_sink).
%%% -export([init/1, handle_chunk/2, handle_close/2]).
%%%
%%% init(_Args) -> {ok, []}.
%%%
%%% handle_chunk(Line, Lines) ->
%%% io:format("~s", [Line]),
%%% {noreply, [Line | Lines]}.
%%%
%%% handle_close(_Reason, _Lines) -> ok.
%%% '''
%%%
%%% ```
%%% {ok, Pid} = macula_stream_sink:start_link(log_tailer, Pool, Realm,
%%% <<"logs.tail_v1">>, []).
%%% '''
%%% @end
%%%-------------------------------------------------------------------
-module(macula_stream_sink).
-behaviour(gen_server).
-export([start_link/5, start_link/6]).
-export([start_link_direct/5, start_link_direct/6]).
-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_chunk(Chunk :: term(), State :: term()) ->
{noreply, NewState :: term()} | {stop, Reason :: term(), NewState :: term()}.
-callback handle_close(Reason :: normal | term(), State :: term()) -> any().
-optional_callbacks([handle_close/2]).
-define(RECV_TIMEOUT, 30_000).
-define(STREAMING_STARTED, <<"streaming.started_v1">>).
-define(STREAMING_COMPLETED, <<"streaming.completed_v1">>).
-record(kstate, {
module :: module(),
pool :: macula:pool(),
realm :: macula:realm(),
announce :: boolean(),
stream_id :: binary(),
stream :: pid(),
reader :: pid(),
user :: term()
}).
%% @doc Start a sink. Opens a stream to `Procedure' on `(Realm)' via
%% `Pool' and passes `Args' to `Module:init/1'.
-spec start_link(module(), macula:pool(), macula:realm(), macula:procedure(),
term()) -> {ok, pid()} | {error, term()}.
start_link(Module, Pool, Realm, Procedure, Args) ->
start_link(Module, Pool, Realm, Procedure, Args, #{}).
%% @doc As `start_link/5', with `CallArgs' passed to `call_stream/5' as
%% the RPC argument payload.
-spec start_link(module(), macula:pool(), macula:realm(), macula:procedure(),
term(), term()) -> {ok, pid()} | {error, term()}.
start_link(Module, Pool, Realm, Procedure, Args, CallArgs) ->
gen_server:start_link(?MODULE,
{pooled, Module, Pool, Realm, Procedure, Args, CallArgs}, []).
%% @doc As `start_link/5', but resolves and dials the procedure's
%% provider 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()) ->
{ok, pid()} | {error, term()}.
start_link_direct(Module, Pool, Realm, Procedure, Args) ->
start_link_direct(Module, Pool, Realm, Procedure, Args, undefined).
%% @doc As `start_link_direct/5', with `CallArgs' passed to
%% `macula_direct_dial:call_stream/5' as the RPC argument payload.
-spec start_link_direct(module(), macula:pool(), macula:realm(),
macula:procedure(), term(), term()) ->
{ok, pid()} | {error, term()}.
start_link_direct(Module, Pool, Realm, Procedure, Args, CallArgs) ->
gen_server:start_link(?MODULE,
{direct, Module, Pool, Realm, Procedure, Args, CallArgs}, []).
%%%===================================================================
%%% gen_server callbacks
%%%===================================================================
%% @private
init({DialMode, Module, Pool, Realm, Procedure, InitArgs, CallArgs}) ->
process_flag(trap_exit, true),
case Module:init(InitArgs) of
{ok, UserState} ->
open_stream(DialMode, Module, Pool, Realm, Procedure, CallArgs,
UserState);
{stop, Reason} ->
{stop, Reason}
end.
open_stream(DialMode, Module, Pool, Realm, Procedure, CallArgs, UserState) ->
case dial_stream(DialMode, Pool, Realm, Procedure, CallArgs) of
{ok, Stream} ->
Reader = spawn_reader(Stream),
StreamId = crypto:strong_rand_bytes(16),
publish(true, Pool, Realm, ?STREAMING_STARTED,
#{stream_id => StreamId}),
{ok, #kstate{module = Module, pool = Pool, realm = Realm,
announce = true, stream_id = StreamId,
stream = Stream, reader = Reader, user = UserState}};
{error, Reason} ->
{stop, Reason}
end.
dial_stream(pooled, Pool, Realm, Procedure, CallArgs) ->
macula:call_stream(Pool, Realm, Procedure, CallArgs, #{});
dial_stream(direct, Pool, Realm, Procedure, CallArgs) ->
macula_direct_dial:call_stream(Pool, Realm, Procedure, CallArgs, #{}).
spawn_reader(Stream) ->
Parent = self(),
spawn_link(fun() -> reader_loop(Parent, Stream) end).
reader_loop(Parent, Stream) ->
dispatch_recv(macula:recv(Stream, ?RECV_TIMEOUT), Parent, Stream).
dispatch_recv({chunk, Data}, Parent, Stream) ->
Parent ! {stream_item, Data}, reader_loop(Parent, Stream);
dispatch_recv({data, Data}, Parent, Stream) ->
Parent ! {stream_item, Data}, reader_loop(Parent, Stream);
dispatch_recv(eof, Parent, _Stream) ->
Parent ! stream_eof;
dispatch_recv({error, Reason}, Parent, _Stream) ->
Parent ! {stream_error, Reason}.
%% @private
handle_call(_Request, _From, State) ->
{reply, {error, unsupported}, State}.
%% @private
handle_cast(_Msg, State) ->
{noreply, State}.
%% @private
handle_info({stream_item, Data}, #kstate{module = Module, user = User} = State) ->
deliver(Module:handle_chunk(Data, User), State);
handle_info(stream_eof, State) ->
{stop, normal, State};
handle_info({stream_error, Reason}, State) ->
{stop, Reason, State};
handle_info({'EXIT', Reader, Reason}, #kstate{reader = Reader} = State)
when Reason =/= normal ->
{stop, {reader_crashed, Reason}, State};
handle_info(_Msg, State) ->
{noreply, State}.
deliver({noreply, NewUser}, State) ->
{noreply, State#kstate{user = NewUser}};
deliver({stop, Reason, NewUser}, State) ->
{stop, Reason, State#kstate{user = NewUser}}.
%% @private
terminate(Reason, #kstate{module = Module, pool = Pool, realm = Realm,
announce = Announce, stream_id = StreamId,
stream = Stream, reader = Reader, user = User}) ->
%% A `normal'-reason exit does not propagate across a link to a
%% non-trapping process, so a clean stop (eof, or the callback
%% returning {stop, normal, _}) would otherwise leave the reader
%% looping on `recv/2' forever against a stream nobody is reading
%% for anymore. Stop it unconditionally.
unlink(Reader),
exit(Reader, kill),
try macula:close_stream(Stream) catch _:_ -> ok end,
publish(Announce, Pool, Realm, ?STREAMING_COMPLETED,
outcome_fields(#{stream_id => StreamId}, Reason)),
maybe_close(Module, Reason, User).
outcome_fields(Base, normal) -> Base#{outcome => completed};
outcome_fields(Base, Reason) -> Base#{outcome => failed, reason => Reason}.
maybe_close(Module, Reason, User) ->
case erlang:function_exported(Module, handle_close, 2) of
true -> Module:handle_close(Reason, User);
false -> ok
end.
publish(false, _, _, _, _) -> ok;
publish(true, Pool, Realm, Topic, Payload) ->
_ = macula:publish(Pool, Realm, Topic, Payload), ok.