Packages
macula
11.3.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_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.
%%%
%%% The sink's `macula_lifetime_announcer' publishes these facts, in
%%% order, from a process of its own, so a pool that is gone or slow
%%% never fails or holds up the sink or its stream; a publish that fails
%%% is logged. A sink killed before it could hand over its stream's end
%%% has `streaming.completed_v1' published for it, with outcome `failed'
%%% and the reason it went down for. An end fact and an abort message
%%% name their reason, `killed' or `timeout' for example, and carry none
%%% of the reason's terms, which go to the local log only.
%%%
%%% == Cancel ==
%%%
%%% Stopping this gen_server for any non-`normal' reason (a `recv'
%%% error, the reader crashing, `Module:handle_chunk/2' returning a
%%% non-normal stop) sends the provider an explicit `macula:abort/3'
%%% STREAM_ERROR instead of an ordinary close — the provider learns
%%% the consumer cancelled or failed rather than mistaking it for a
%%% clean end-of-stream. A `normal' stop closes both sides cleanly
%%% instead, same as before.
%%%
%%% == 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".
%%%
%%% == Stream I/O ==
%%%
%%% A sink opens, reads and ends its stream through four functions,
%%% `call_stream/5', `recv/2', `close_stream/1' and `abort/3', and
%%% announces its facts through a fifth, `fact_publish/4'. They are the
%%% `macula' facade's by default, and a direct-dial sink dials with
%%% `macula_direct_dial:call_stream/5'. `start_link/7' and
%%% `start_link_direct/7' take a `stream_io' start option, checked by
%%% `macula_stream:stream_io/2', with the four stream functions and any
%%% other `macula_stream:stream_io()' ones, and a `fact_publish' start
%%% option, to run a sink on something else, such as a test's scripted
%%% stream.
%%%
%%% == 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, start_link/7]).
-export([start_link_direct/5, start_link_direct/6, start_link_direct/7]).
-export([init/1, handle_call/3, handle_cast/2, handle_info/2, terminate/2]).
-export_type([start_opts/0]).
-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">>).
-define(CANCEL_CODE, <<"cancelled">>).
-type start_opts() :: #{stream_io => macula_stream:stream_io(),
fact_publish => macula_lifetime_announcer:publish()}.
-record(kstate, {
io :: macula_stream:stream_io(),
module :: module(),
pool :: macula:pool(),
realm :: macula:realm(),
announcer :: pid(),
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) ->
start_link(Module, Pool, Realm, Procedure, Args, CallArgs, #{}).
%% @doc As `start_link/6', with start options: `stream_io' gives the
%% functions the sink runs its stream on, and `fact_publish' the one it
%% announces its facts with (see "Stream I/O" above).
-spec start_link(module(), macula:pool(), macula:realm(), macula:procedure(),
term(), term(), start_opts()) -> {ok, pid()} | {error, term()}.
start_link(Module, Pool, Realm, Procedure, Args, CallArgs, Opts) when is_map(Opts) ->
start(pooled, Opts, {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) ->
start_link_direct(Module, Pool, Realm, Procedure, Args, CallArgs, #{}).
%% @doc As `start_link_direct/6', with start options: `stream_io' gives
%% the functions the sink runs its stream on, and `fact_publish' the one
%% it announces its facts with (see "Stream I/O" above).
-spec start_link_direct(module(), macula:pool(), macula:realm(),
macula:procedure(), term(), term(), start_opts()) ->
{ok, pid()} | {error, term()}.
start_link_direct(Module, Pool, Realm, Procedure, Args, CallArgs, Opts) when is_map(Opts) ->
start(direct, Opts, {Module, Pool, Realm, Procedure, Args, CallArgs}).
%% A sink starts on stream functions macula_stream:stream_io/2 accepts
%% and a fact_publish of arity 4; any other is refused with
%% function_clause, in the caller.
start(DialMode, Opts, Start) ->
StreamIo = macula_stream:stream_io(default_stream_io(DialMode),
maps:get(stream_io, Opts, undefined)),
FactPublish = arity_4(maps:get(fact_publish, Opts, fun macula:publish/4)),
gen_server:start_link(?MODULE, {StreamIo, FactPublish, Start}, []).
default_stream_io(DialMode) ->
#{call_stream => dial(DialMode),
recv => fun macula:recv/2,
close_stream => fun macula:close_stream/1,
abort => fun macula:abort/3}.
arity_4(Fun) when is_function(Fun, 4) -> Fun.
dial(pooled) -> fun macula:call_stream/5;
dial(direct) -> fun macula_direct_dial:call_stream/5.
%%%===================================================================
%%% gen_server callbacks
%%%===================================================================
%% @private
init({StreamIo, FactPublish, {Module, Pool, Realm, Procedure, InitArgs, CallArgs}}) ->
process_flag(trap_exit, true),
case Module:init(InitArgs) of
{ok, UserState} ->
open_stream(StreamIo, FactPublish, Module, Pool, Realm, Procedure, CallArgs,
UserState);
{stop, Reason} ->
{stop, Reason}
end.
open_stream(#{call_stream := CallStream, recv := Recv} = StreamIo, FactPublish, Module, Pool,
Realm, Procedure, CallArgs, UserState) ->
case CallStream(Pool, Realm, Procedure, CallArgs, #{}) of
{ok, Stream} ->
StreamId = crypto:strong_rand_bytes(16),
Announcer = macula_lifetime_announcer:start(
true, facts(FactPublish, Pool, Realm, StreamId)),
Reader = spawn_reader(Recv, Stream),
{ok, #kstate{io = StreamIo, module = Module, pool = Pool, realm = Realm,
announcer = Announcer, stream_id = StreamId,
stream = Stream, reader = Reader, user = UserState}};
{error, Reason} ->
{stop, Reason}
end.
spawn_reader(Recv, Stream) ->
Parent = self(),
spawn_link(fun() -> reader_loop(Parent, Recv, Stream) end).
reader_loop(Parent, Recv, Stream) ->
dispatch_recv(Recv(Stream, ?RECV_TIMEOUT), Parent, Recv, Stream).
dispatch_recv({chunk, Data}, Parent, Recv, Stream) ->
Parent ! {stream_item, Data}, reader_loop(Parent, Recv, Stream);
dispatch_recv({data, Data}, Parent, Recv, Stream) ->
Parent ! {stream_item, Data}, reader_loop(Parent, Recv, Stream);
dispatch_recv(eof, Parent, _Recv, _Stream) ->
Parent ! stream_eof;
dispatch_recv({error, Reason}, Parent, _Recv, _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{io = StreamIo, module = Module, announcer = Announcer,
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, and wait until it has
%% exited: a kill arrives asynchronously, so without the wait the
%% sink could be gone while its reader still calls `recv/2'.
stop_reader(Reader),
finish_stream(StreamIo, Reason, Stream),
ok = macula_lifetime_announcer:announce_end(
Announcer, outcome_fields(#{stream_id => StreamId}, Reason)),
maybe_close(Module, Reason, User).
stop_reader(Reader) ->
Ref = monitor(process, Reader),
unlink(Reader),
exit(Reader, kill),
receive
{'DOWN', Ref, process, Reader, _} -> ok
end.
%% @private A `normal' reason (eof, or the callback choosing to stop
%% cleanly) closes both sides. Anything else sends the provider an
%% explicit abort instead of an ordinary close, so it learns this was
%% a cancellation/failure rather than a clean end-of-stream. `Stream'
%% may already be dead by the time this runs (e.g. `{stream_error,_}'
%% means the provider already tore it down) — harmless, caught below.
finish_stream(#{close_stream := CloseStream}, normal, Stream) ->
try CloseStream(Stream) catch _:_ -> ok end;
finish_stream(#{abort := Abort}, Reason, Stream) ->
try Abort(Stream, ?CANCEL_CODE, macula_reason_name:text(Reason))
catch _:_ -> ok end.
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.
%% The facts the sink's announcer publishes: the stream's start, and its
%% end, both with the stream's id.
facts(FactPublish, Pool, Realm, StreamId) ->
#{publish => FactPublish, pool => Pool, realm => Realm,
started => {?STREAMING_STARTED, #{stream_id => StreamId}},
ended => {?STREAMING_COMPLETED, #{stream_id => StreamId}}}.