Packages
macula
9.4.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_streamer.erl
%%%-------------------------------------------------------------------
%%% @doc Behaviour for supervised streaming RPC providers.
%%%
%%% `advertise_stream/5' on the raw SDK takes a bare handler fun
%%% invoked as `Handler(StreamPid, Args)' in a transient process
%%% spawned per inbound STREAM_OPEN (see the internal
%%% macula_station_link advertise_stream/5 — "this link spawns a
%%% server-side macula_stream and dispatches Handler(StreamPid, Args)
%%% in a transient process"). This is the provider-side counterpart to
%%% `macula_stream_sink': each inbound stream starts one supervised
%%% `macula_streamer' child (under a `simple_one_for_one' factory this
%%% module owns), threading state through `Module:init/1' and
%%% `Module:handle_open/3', and publishing `streaming.started_v1' /
%%% `streaming.completed_v1' mesh facts around the stream's lifetime.
%%%
%%% Sending is push-based and driven from outside the callback: once
%%% `Module:handle_open/3' has done whatever registration it needs
%%% (e.g. stashing `self()' in a registry keyed by some connection id),
%%% any process holding this streamer's pid can call `send/2,3' /
%%% `close/1' on it. This module does not prescribe the discovery
%%% mechanism.
%%%
%%% This is the general-purpose RPC streaming feature (`call_stream/5',
%%% `advertise_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.
%%%
%%% == Example ==
%%%
%%% ```
%%% -module(log_tailer_provider).
%%% -behaviour(macula_streamer).
%%% -export([init/1, handle_open/3]).
%%%
%%% init(Registry) -> {ok, Registry}.
%%%
%%% handle_open(#{topic := Topic}, Registry, State) ->
%%% Registry ! {tailer_ready, Topic, self()},
%%% {ok, State}.
%%% '''
%%%
%%% ```
%%% {ok, _Sup} = macula_streamer:advertise(Pool, Realm,
%%% <<"logs.tail_v1">>, log_tailer_provider, self()).
%%%
%%% %% elsewhere, once the provider has announced its pid:
%%% ok = macula_streamer:send(TailerPid, <<"a log line\n">>).
%%% '''
%%% @end
%%%-------------------------------------------------------------------
-module(macula_streamer).
-behaviour(gen_server).
-export([advertise/5, advertise/6, unadvertise/3]).
-export([send/2, send/3, close/1]).
-export([start_link/7]).
-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_open(StreamArgs :: term(), State :: term()) ->
{ok, NewState :: term()} | {stop, Reason :: term(), NewState :: term()}.
-callback terminate(Reason :: term(), State :: term()) -> any().
-optional_callbacks([terminate/2]).
-define(STREAMING_STARTED, <<"streaming.started_v1">>).
-define(STREAMING_COMPLETED, <<"streaming.completed_v1">>).
-record(tstate, {
module :: module(),
pool :: macula:pool(),
realm :: macula:realm(),
announce :: boolean(),
stream_id :: binary(),
stream :: pid(),
user :: term()
}).
%% @doc Advertise `Procedure' on `Pool'/`Realm'. Starts a private
%% factory supervisor for per-stream provider children and registers
%% a dispatch handler with `macula:advertise_stream/5'. Returns the
%% supervisor pid so the caller can supervise it (or ignore it).
-spec advertise(macula:pool(), macula:realm(), macula:procedure(),
module(), term()) -> {ok, pid()} | {error, term()}.
advertise(Pool, Realm, Procedure, Module, Args) ->
advertise(Pool, Realm, Procedure, Module, Args, #{}).
%% @doc As `advertise/5'. `Opts' may include `announce' (default
%% `true') and `mode' (default `server_stream').
-spec advertise(macula:pool(), macula:realm(), macula:procedure(),
module(), term(), map()) -> {ok, pid()} | {error, term()}.
advertise(Pool, Realm, Procedure, Module, Args, Opts) ->
{ok, Sup} = macula_streamer_sup:start_link(),
Announce = maps:get(announce, Opts, true),
Mode = maps:get(mode, Opts, server_stream),
Handler = fun(StreamPid, StreamArgs) ->
dispatch(Sup, Module, Pool, Realm, Announce, Args, StreamPid, StreamArgs)
end,
case macula:advertise_stream(Pool, Realm, Procedure, Mode, Handler) of
ok -> {ok, Sup};
{error, Reason} -> {error, Reason}
end.
%% @doc Stop advertising. Does not stop the factory supervisor
%% returned by `advertise/5,6' — callers that want to tear it down
%% should `exit(Sup, shutdown)' themselves.
-spec unadvertise(macula:pool(), macula:realm(), macula:procedure()) -> ok.
unadvertise(Pool, Realm, Procedure) ->
macula:unadvertise_stream(Pool, Realm, Procedure).
dispatch(Sup, Module, Pool, Realm, Announce, Args, StreamPid, StreamArgs) ->
case supervisor:start_child(Sup, [Module, Pool, Realm, Announce, Args,
StreamPid, StreamArgs]) of
{ok, _Pid} -> ok;
{error, _Reason} -> ok
end.
%% @doc Send a chunk out on the stream this streamer owns.
-spec send(pid(), binary()) -> ok | {error, term()}.
send(Pid, Chunk) -> gen_server:call(Pid, {send, Chunk}).
%% @doc As `send/2', with an explicit encoding.
-spec send(pid(), binary() | term(), macula_stream:encoding()) ->
ok | {error, term()}.
send(Pid, Chunk, Encoding) -> gen_server:call(Pid, {send, Chunk, Encoding}).
%% @doc Close the send side of the stream.
-spec close(pid()) -> ok.
close(Pid) -> gen_server:call(Pid, close).
%% @private
-spec start_link(module(), macula:pool(), macula:realm(), boolean(),
term(), pid(), term()) -> {ok, pid()} | {error, term()}.
start_link(Module, Pool, Realm, Announce, InitArgs, StreamPid, StreamArgs) ->
gen_server:start_link(?MODULE,
{Module, Pool, Realm, Announce, InitArgs, StreamPid, StreamArgs}, []).
%%%===================================================================
%%% gen_server callbacks
%%%===================================================================
%% @private
init({Module, Pool, Realm, Announce, InitArgs, StreamPid, StreamArgs}) ->
process_flag(trap_exit, true),
case Module:init(InitArgs) of
{ok, UserState} ->
open(Module, Pool, Realm, Announce, StreamPid, StreamArgs, UserState);
{stop, Reason} ->
{stop, Reason}
end.
open(Module, Pool, Realm, Announce, StreamPid, StreamArgs, UserState) ->
case Module:handle_open(StreamArgs, UserState) of
{ok, NewUserState} ->
link(StreamPid),
StreamId = crypto:strong_rand_bytes(16),
publish(Announce, Pool, Realm, ?STREAMING_STARTED,
#{stream_id => StreamId}),
{ok, #tstate{module = Module, pool = Pool, realm = Realm,
announce = Announce, stream_id = StreamId,
stream = StreamPid, user = NewUserState}};
{stop, Reason, _NewUserState} ->
{stop, Reason}
end.
%% @private
handle_call({send, Chunk}, _From, #tstate{stream = Stream} = State) ->
{reply, macula_stream:send(Stream, Chunk), State};
handle_call({send, Chunk, Encoding}, _From, #tstate{stream = Stream} = State) ->
{reply, macula_stream:send(Stream, Chunk, Encoding), State};
handle_call(close, _From, #tstate{stream = Stream} = State) ->
{reply, macula_stream:close_send(Stream), State};
handle_call(_Request, _From, State) ->
{reply, {error, unsupported}, State}.
%% @private
handle_cast(_Msg, State) -> {noreply, State}.
%% @private
handle_info({'EXIT', Stream, Reason}, #tstate{stream = Stream} = State) ->
{stop, Reason, State};
handle_info(_Msg, State) ->
{noreply, State}.
%% @private
terminate(Reason, #tstate{module = Module, pool = Pool, realm = Realm,
announce = Announce, stream_id = StreamId,
user = User}) ->
publish(Announce, Pool, Realm, ?STREAMING_COMPLETED,
outcome_fields(#{stream_id => StreamId}, Reason)),
maybe_terminate(Module, Reason, User).
outcome_fields(Base, normal) -> Base#{outcome => completed};
outcome_fields(Base, Reason) -> Base#{outcome => failed, reason => Reason}.
maybe_terminate(Module, Reason, User) ->
case erlang:function_exported(Module, terminate, 2) of
true -> Module:terminate(Reason, User);
false -> ok
end.
publish(false, _, _, _, _) -> ok;
publish(true, Pool, Realm, Topic, Payload) ->
_ = macula:publish(Pool, Realm, Topic, Payload), ok.