Packages

macula

11.1.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
macula src macula_streamer.erl
Raw

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/2', 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/2' 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.
%%%
%%% For `client_stream' mode — a consumer pushing chunks INTO the
%%% provider, e.g. a batch upload — export the optional
%%% `handle_chunk/2' callback (mirroring `macula_stream_sink''s
%%% consumer-side callback exactly) and this module drives the same
%%% linked-reader `recv/2' loop for you, on the provider side. A
%%% `server_stream'-mode module that doesn't export it is unaffected.
%%%
%%% A `client_stream' provider that also needs to hand the consumer a
%%% terminal result (not just accept chunks) exports the optional
%%% `handle_eof/1' callback: called once, when the consumer's own
%%% `close_send/1' surfaces here as end-of-stream, in place of the
%%% default unconditional `{stop, normal, State}'. Returning
%%% `{reply, Result, NewState}' sets the stream's terminal reply
%%% (`macula_stream:set_reply/2' for `{ok, Value}', `set_error/2' for
%%% `{error, Reason}') so the consumer's own `macula:await_reply/1,2'
%%% unblocks with it, before stopping. A module that doesn't export
%%% `handle_eof/1' keeps the exact prior behavior — no reply is ever
%%% set, eof just stops the stream.
%%%
%%% 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.
%%%
%%% == Cancel ==
%%%
%%% Stopping this gen_server for any non-`normal' reason (a crash, the
%%% underlying stream dying, `Module:handle_open/2'/`handle_chunk/2'
%%% returning a non-normal stop) sends the peer an explicit
%%% `macula_stream:abort/3' STREAM_ERROR, not just a graceful close —
%%% the peer learns the transfer was cancelled/failed instead of
%%% mistaking it for an ordinary end-of-stream. A `normal' stop closes
%%% both sides cleanly instead.
%%%
%%% == Direct-dial ==
%%%
%%% `advertise/5,6' registers the handler with the pool's advertise-
%%% gossip mechanism only — nothing published lets a caller on another
%%% station find this procedure without a route having propagated
%%% between the two stations first. `advertise_direct/6,7' does that
%%% AND publishes a signed `procedure_advertisement' DHT record naming
%%% this pool's currently-connected station as the server — the exact
%%% same record type and publish function `macula_response:advertise_direct/6,7'
%%% uses for plain RPC (a `procedure_advertisement' does not distinguish
%%% RPC from streaming), so a caller using
%%% `macula_stream_sink:start_link_direct/5,6' can resolve and dial
%%% here directly, in one hop, regardless of whether the two stations
%%% have a routing edge between them.
%%%
%%% == Stream I/O ==
%%%
%%% A streamer advertises its procedure with `advertise_stream', a
%%% function of arity 6, `macula:advertise_stream/6' by default, which
%%% gets the procedure's `auth' policy and no other option;
%%% `advertise_direct/6,7' publishes its DHT record with
%%% `publish_advertisement', `macula_direct_dial:publish_advertisement/5'
%%% by default; and each streamer announces its facts with
%%% `fact_publish', `macula:publish/4' by default. Each streamer runs its
%%% stream on seven `macula_stream:stream_io()' functions, `recv/2',
%%% `send/3', `close_send/1', `close/1', `abort/3', `set_reply/2' and
%%% `set_error/2', which are `macula:recv/2' and the `macula_stream' ones
%%% by default. `advertise/6' and `advertise_direct/7' take all of these
%%% in their options, the seven as `stream_io', checked by
%%% `macula_stream:stream_io/2', and refuse a function of another arity
%%% with `function_clause' before anything is advertised.
%%%
%%% == Example ==
%%%
%%% ```
%%% -module(log_tailer_provider).
%%% -behaviour(macula_streamer).
%%% -export([init/1, handle_open/2]).
%%%
%%% init(Registry) -> {ok, Registry}.
%%%
%%% handle_open(#{topic := Topic}, Registry) ->
%%% Registry ! {tailer_ready, Topic, self()},
%%% {ok, Registry}.
%%% '''
%%%
%%% ```
%%% {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">>).
%%% '''
%%%
%%% A `client_stream'-mode provider exports `handle_chunk/2' instead,
%%% and never calls `send/2,3' itself — the consumer is the one
%%% pushing:
%%%
%%% ```
%%% -module(batch_upload_provider).
%%% -behaviour(macula_streamer).
%%% -export([init/1, handle_open/2, handle_chunk/2]).
%%%
%%% init(Parent) -> {ok, {Parent, []}}.
%%%
%%% handle_open(_StreamArgs, State) -> {ok, State}.
%%%
%%% handle_chunk(Chunk, {Parent, Acc}) ->
%%% {noreply, {Parent, [Chunk | Acc]}}.
%%% '''
%%%
%%% ```
%%% {ok, _Sup} = macula_streamer:advertise(Pool, Realm,
%%% <<"bulk.ingest">>, batch_upload_provider, self(),
%%% #{mode => client_stream}).
%%% '''
%%% @end
%%%-------------------------------------------------------------------
-module(macula_streamer).
-behaviour(gen_server).
-include_lib("kernel/include/logger.hrl").
-export([advertise/5, advertise/6, advertise_direct/6, advertise_direct/7,
unadvertise/3]).
-export([send/2, send/3, close/1]).
-export([start_link/8]).
-export([init/1, handle_call/3, handle_cast/2, handle_info/2, terminate/2]).
-export_type([advertise_stream/0, publish_advertisement/0, advertise_opts/0]).
-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 handle_chunk(Chunk :: term(), State :: term()) ->
{noreply, NewState :: term()} | {stop, Reason :: term(), NewState :: term()}.
-callback handle_eof(State :: term()) ->
{noreply, NewState :: term()}
| {reply, {ok, term()} | {error, term()}, NewState :: term()}
| {stop, Reason :: term(), NewState :: term()}.
-callback terminate(Reason :: term(), State :: term()) -> any().
-optional_callbacks([terminate/2, handle_chunk/2, handle_eof/1]).
-define(STREAMING_STARTED, <<"streaming.started_v1">>).
-define(STREAMING_COMPLETED, <<"streaming.completed_v1">>).
-define(RECV_TIMEOUT, 30_000).
-define(CANCEL_CODE, <<"cancelled">>).
-type advertise_stream() :: fun((macula:pool(), macula:realm(), macula:procedure(),
macula_stream:mode(), fun((pid(), term()) -> ok), map()) ->
ok | {error, term()}).
-type publish_advertisement() :: fun((macula:pool(), macula:realm(), macula:procedure(),
macula_node_keys:node_key(), map()) ->
ok | {error, term()}).
-type advertise_opts() :: #{advertise_stream => advertise_stream(),
publish_advertisement => publish_advertisement(),
fact_publish => macula_lifetime_announcer:publish(),
stream_io => macula_stream:stream_io(),
atom() => term()}.
%% What each streamer runs its stream on and announces its facts with.
-type functions() :: #{stream_io := macula_stream:stream_io(),
fact_publish := macula_lifetime_announcer:publish()}.
-record(tstate, {
module :: module(),
pool :: macula:pool(),
realm :: macula:realm(),
announce :: boolean(),
io :: macula_stream:stream_io(),
fact_publish :: macula_lifetime_announcer:publish(),
stream_id :: binary(),
stream :: pid(),
reader :: pid() | undefined,
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,6'. 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'), `mode' (default `server_stream'), `auth' (the procedure's
%% auth policy, default `open', see `macula:advertise_stream/6'), and
%% `reuse_sup' — an
%% existing supervisor pid (as returned by a prior `advertise/5,6'
%% call) to register the handler again with, without starting a new
%% factory supervisor. Use this for a periodic re-advertise (see
%% `advertise_direct/6,7''s own doc) — calling plain
%% `advertise/5,6' on a timer would leak one orphaned supervisor per
%% tick, since each call otherwise starts a fresh one. The functions a
%% streamer runs on come from `Opts' too; see "Stream I/O" above.
-spec advertise(macula:pool(), macula:realm(), macula:procedure(),
module(), term(), advertise_opts()) -> {ok, pid()} | {error, term()}.
advertise(Pool, Realm, Procedure, Module, Args, Opts) ->
AdvertiseStream = arity_6(maps:get(advertise_stream, Opts, fun macula:advertise_stream/6)),
Functions = functions(Opts),
Sup = existing_or_new_sup(maps:get(reuse_sup, Opts, undefined)),
Announce = maps:get(announce, Opts, true),
Mode = maps:get(mode, Opts, server_stream),
Handler = fun(StreamPid, StreamArgs) ->
dispatch(Sup, Module, Pool, Realm, Announce, Args, Functions, StreamPid, StreamArgs)
end,
%% The advertise function gets the procedure's auth policy, when the
%% options give one, and no other option.
case AdvertiseStream(Pool, Realm, Procedure, Mode, Handler, maps:with([auth], Opts)) of
ok -> {ok, Sup};
{error, Reason} -> {error, Reason}
end.
%% The functions each streamer runs on, from the options or else the
%% defaults. Stream functions macula_stream:stream_io/2 does not accept,
%% or a fact_publish of another arity, are refused with function_clause,
%% in the caller, before anything is advertised.
functions(Opts) ->
StreamIo = macula_stream:stream_io(default_stream_io(), maps:get(stream_io, Opts, undefined)),
#{stream_io => StreamIo,
fact_publish => arity_4(maps:get(fact_publish, Opts, fun macula:publish/4))}.
default_stream_io() ->
#{recv => fun macula:recv/2,
controlling_process => fun macula_stream:controlling_process/2,
send => fun macula_stream:send/3,
close_send => fun macula_stream:close_send/1,
close => fun macula_stream:close/1,
abort => fun macula_stream:abort/3,
set_reply => fun macula_stream:set_reply/2,
set_error => fun macula_stream:set_error/2}.
%% The options the advertisement publish gets: all but the functions.
without_functions(Opts) ->
maps:without([advertise_stream, publish_advertisement, fact_publish, stream_io], Opts).
arity_4(Fun) when is_function(Fun, 4) -> Fun.
arity_5(Fun) when is_function(Fun, 5) -> Fun.
arity_6(Fun) when is_function(Fun, 6) -> Fun.
%% See `macula_response:existing_or_new_sup/1' for why a dead `reuse_sup'
%% pid must fall through to a fresh one rather than being handed to
%% `dispatch/9' as-is.
existing_or_new_sup(Pid) when is_pid(Pid) ->
existing_or_new_sup(Pid, erlang:is_process_alive(Pid));
existing_or_new_sup(undefined) ->
new_sup().
existing_or_new_sup(Pid, true) -> Pid;
existing_or_new_sup(_Pid, false) -> new_sup().
new_sup() ->
{ok, Sup} = macula_streamer_sup:start_link(),
Sup.
%% @doc As `advertise/5', and additionally publishes a signed
%% `procedure_advertisement' DHT record naming this pool's connected
%% station as the server, so `macula_stream_sink:start_link_direct/5,6'
%% can resolve and dial here directly. `NodeIdentity' signs it and must be
%% the node identity key `Pool' was started with: a caller targets that
%% node_id, and the station knows the pool's connection by it.
%%
%% The DHT publish is best-effort: if it fails, the handler is still
%% advertised and reachable via the ordinary pooled path — direct-dial
%% callers just won't be able to resolve it until a later publish
%% succeeds. "Best-effort" still means the failure is logged, not
%% silently discarded — a caller that only ever calls this once (never
%% retries) has no other way to learn its handler is pooled-only, and
%% "a later publish succeeds" cannot happen if nothing ever tries again.
-spec advertise_direct(macula:pool(), macula:realm(), macula:procedure(),
module(), term(), macula_node_keys:node_key()) ->
{ok, pid()} | {error, term()}.
advertise_direct(Pool, Realm, Procedure, Module, Args, NodeIdentity) ->
advertise_direct(Pool, Realm, Procedure, Module, Args, NodeIdentity, #{}).
%% @doc As `advertise_direct/6', with `Opts' forwarded BOTH to
%% `advertise/6' (so `mode'/`announce'/`reuse_sup' and the functions
%% apply here too, e.g. `mode => client_stream') and, without the
%% functions, to the advertisement publish, `publish_advertisement' in
%% `Opts' or `macula_direct_dial:publish_advertisement/5' (e.g.
%% `authorization', the provider authorization an org namespaced
%% procedure needs): each side reads only the keys it recognizes, so one
%% `Opts' map serves both.
%% `reuse_sup' matters here specifically: the procedure's DHT record
%% expires with its TTL, and callers reach the provider only through
%% that record, so the provider republishes it — a periodic
%% re-advertise with `reuse_sup => Sup' (the pid this function
%% returned the first time) registers the handler again and
%% republishes the DHT record without leaking a new supervisor per
%% tick. `cert_chain', a 10.x option `authorization' replaces, is refused
%% with `{error, {removed_option, cert_chain}}' before the handler is
%% registered.
-spec advertise_direct(macula:pool(), macula:realm(), macula:procedure(),
module(), term(), macula_node_keys:node_key(), advertise_opts()) ->
{ok, pid()} | {error, term()}.
advertise_direct(Pool, Realm, Procedure, Module, Args, NodeIdentity, Opts) ->
advertise_direct_unless_removed(macula_direct_dial:removed_option(advertise, Opts), Pool,
Realm, Procedure, Module, Args, NodeIdentity, Opts).
advertise_direct_unless_removed(none, Pool, Realm, Procedure, Module, Args, NodeIdentity, Opts) ->
PublishAdvertisement = arity_5(maps:get(publish_advertisement, Opts,
fun macula_direct_dial:publish_advertisement/5)),
case advertise(Pool, Realm, Procedure, Module, Args, Opts) of
{ok, Sup} ->
log_publish_result(
PublishAdvertisement(Pool, Realm, Procedure, NodeIdentity, without_functions(Opts)),
Procedure),
{ok, Sup};
{error, _} = Error ->
Error
end;
advertise_direct_unless_removed(Removed, _Pool, _Realm, _Procedure, _Module, _Args, _NodeIdentity,
_Opts) ->
{error, Removed}.
log_publish_result(ok, _Procedure) ->
ok;
log_publish_result({error, Reason}, Procedure) ->
?LOG_WARNING("[macula_streamer] direct-dial advertisement publish "
"failed for ~s: ~p -- handler stays reachable via the "
"pooled path only until a later publish succeeds",
[Procedure, Reason]).
%% @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, Functions, StreamPid, StreamArgs) ->
case supervisor:start_child(Sup, [Module, Pool, Realm, Announce, Args,
StreamPid, StreamArgs, Functions]) of
{ok, Pid} -> hand_stream_to(Functions, StreamPid, Pid);
{error, _Reason} -> ok
end.
%% @private The process running this dispatch owns the stream, and the
%% streamer it started takes the stream over, so the stream ends when the
%% streamer ends rather than when this dispatch returns.
hand_stream_to(#{stream_io := #{controlling_process := HandOver}}, StreamPid, StreamerPid) ->
_ = HandOver(StreamPid, StreamerPid),
ok.
%% @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(), functions()) -> {ok, pid()} | {error, term()}.
start_link(Module, Pool, Realm, Announce, InitArgs, StreamPid, StreamArgs, Functions) ->
gen_server:start_link(?MODULE,
{Module, Pool, Realm, Announce, InitArgs, StreamPid, StreamArgs, Functions}, []).
%%%===================================================================
%%% gen_server callbacks
%%%===================================================================
%% @private
init({Module, Pool, Realm, Announce, InitArgs, StreamPid, StreamArgs, Functions}) ->
process_flag(trap_exit, true),
case Module:init(InitArgs) of
{ok, UserState} ->
open(Module, Pool, Realm, Announce, Functions, StreamPid, StreamArgs, UserState);
{stop, Reason} ->
{stop, Reason}
end.
open(Module, Pool, Realm, Announce, #{stream_io := StreamIo, fact_publish := FactPublish},
StreamPid, StreamArgs, UserState) ->
case Module:handle_open(StreamArgs, UserState) of
{ok, NewUserState} ->
link(StreamPid),
Reader = maybe_spawn_reader(Module, StreamIo, StreamPid),
StreamId = crypto:strong_rand_bytes(16),
publish(Announce, FactPublish, Pool, Realm, ?STREAMING_STARTED,
#{stream_id => StreamId}),
{ok, #tstate{module = Module, pool = Pool, realm = Realm,
announce = Announce, io = StreamIo, fact_publish = FactPublish,
stream_id = StreamId, stream = StreamPid, reader = Reader,
user = NewUserState}};
{stop, Reason, _NewUserState} ->
abort_rejected_stream(StreamIo, Reason, StreamPid),
{stop, Reason}
end.
%% @private A rejected open (`handle_open/2' returning `{stop, Reason, _}')
%% never links `StreamPid', so `terminate/2' never runs on it and the peer
%% that opened the stream would otherwise be stranded until its own `recv'
%% timeout. Abort it explicitly so the peer gets an immediate signal
%% instead of silence, naming the reason and carrying none of its terms.
abort_rejected_stream(#{abort := Abort}, Reason, StreamPid) ->
Message = macula_reason_name:text(Reason),
try Abort(StreamPid, ?CANCEL_CODE, Message) catch _:_ -> ok end.
%% @private For `client_stream'-mode providers that export
%% `handle_chunk/2': spawn the same linked-reader `recv/2' loop
%% `macula_stream_sink' drives on the consumer side, applied here to
%% the provider's own stream. A `server_stream'-mode module has no
%% reason to export `handle_chunk/2', so this is a no-op for it.
maybe_spawn_reader(Module, #{recv := Recv}, Stream) ->
case erlang:function_exported(Module, handle_chunk, 2) of
true -> spawn_reader(Recv, Stream);
false -> undefined
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({send, Chunk}, _From, #tstate{io = #{send := Send}, stream = Stream} = State) ->
{reply, Send(Stream, Chunk, raw), State};
handle_call({send, Chunk, Encoding}, _From,
#tstate{io = #{send := Send}, stream = Stream} = State) ->
{reply, Send(Stream, Chunk, Encoding), State};
handle_call(close, _From, #tstate{io = #{close_send := CloseSend}, stream = Stream} = State) ->
{reply, CloseSend(Stream), State};
handle_call(_Request, _From, State) ->
{reply, {error, unsupported}, State}.
%% @private
handle_cast(_Msg, State) -> {noreply, State}.
%% @private
handle_info({stream_item, Data}, #tstate{module = Module, user = User} = State) ->
deliver(Module:handle_chunk(Data, User), State);
handle_info(stream_eof, State) ->
handle_eof(State);
handle_info({stream_error, Reason}, State) ->
{stop, Reason, State};
handle_info({'EXIT', Reader, Reason}, #tstate{reader = Reader} = State)
when Reason =/= normal ->
{stop, {reader_crashed, Reason}, State};
handle_info({'EXIT', Stream, Reason}, #tstate{stream = Stream} = State) ->
{stop, Reason, State};
%% The stream's session ended (`macula_stream:controlling_process/2'): a
%% streamer has nothing left to serve, whether or not its module would stop
%% by itself.
handle_info({macula_stream, ended, Stream, closed}, #tstate{stream = Stream} = State) ->
{stop, normal, State};
handle_info({macula_stream, ended, Stream, _How}, #tstate{stream = Stream} = State) ->
{stop, {shutdown, session_ended}, State};
handle_info(_Msg, State) ->
{noreply, State}.
deliver({noreply, NewUser}, State) ->
{noreply, State#tstate{user = NewUser}};
deliver({stop, Reason, NewUser}, State) ->
{stop, Reason, State#tstate{user = NewUser}}.
%% @private Default (no `handle_eof/1' exported): unchanged prior
%% behavior, eof just stops the stream. Otherwise gives the callback
%% one last chance to set a terminal reply before stopping.
handle_eof(#tstate{module = Module, user = User} = State) ->
case erlang:function_exported(Module, handle_eof, 1) of
true -> deliver_eof(Module:handle_eof(User), State);
false -> {stop, normal, State}
end.
deliver_eof({noreply, NewUser}, State) ->
{stop, normal, State#tstate{user = NewUser}};
deliver_eof({reply, {ok, Value}, NewUser},
#tstate{io = #{set_reply := SetReply}, stream = Stream} = State) ->
_ = SetReply(Stream, Value),
{stop, normal, State#tstate{user = NewUser}};
deliver_eof({reply, {error, Reason}, NewUser},
#tstate{io = #{set_error := SetError}, stream = Stream} = State) ->
_ = SetError(Stream, Reason),
{stop, normal, State#tstate{user = NewUser}};
deliver_eof({stop, Reason, NewUser}, State) ->
{stop, Reason, State#tstate{user = NewUser}}.
%% @private
terminate(Reason, #tstate{module = Module, pool = Pool, realm = Realm,
announce = Announce, io = StreamIo, fact_publish = FactPublish,
stream_id = StreamId, stream = Stream, reader = Reader,
user = User}) ->
stop_reader(Reader),
finish_stream(StreamIo, Reason, Stream),
publish(Announce, FactPublish, Pool, Realm, ?STREAMING_COMPLETED,
outcome_fields(#{stream_id => StreamId}, Reason)),
maybe_terminate(Module, Reason, User).
stop_reader(undefined) -> ok;
stop_reader(Reader) ->
unlink(Reader),
exit(Reader, kill).
%% @private A `normal' reason closes both sides cleanly. Anything else
%% (a crash, the underlying stream dying, a non-normal stop from
%% `handle_open/2'/`handle_chunk/2') sends the peer an explicit
%% `STREAM_ERROR' abort, whose message is the reason's name, instead of
%% leaving it to infer cancellation from the connection simply going
%% away. `Stream' may already be dead by the time this runs (e.g. its
%% own exit is what triggered this termination), which is harmless and
%% caught below.
finish_stream(#{close := Close}, normal, Stream) ->
try Close(Stream) catch _:_ -> ok end;
finish_stream(#{abort := Abort}, Reason, Stream) ->
Message = macula_reason_name:text(Reason),
try Abort(Stream, ?CANCEL_CODE, Message) catch _:_ -> ok end.
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, _FactPublish, _, _, _, _) -> ok;
publish(true, FactPublish, Pool, Realm, Topic, Payload) ->
_ = FactPublish(Pool, Realm, Topic, Payload), ok.