GenStage.stream/1.
Subscribes the current process to one or more gen_stage producers and
delivers their events to this process' mailbox, the Erlang
equivalent of Elixir's GenStage.stream/1.
Like the Elixir stream, this module "hijacks" the inbox of the calling
process while the stream is open: it guarantees it will not leave
unwanted messages in the mailbox after close/1 (unless a producer
comes from a remote node).
## Example
{ok, Stream} = gen_stage_stream:subscribe([{Producer, max_demand: 100}]), loop(Stream).
loop({MonitorPid, MonitorRef, Subscriptions}) -> receive {'$gen_consumer', {Pid, {MonitorRef, InnerRef}} = From, Events} when is_list(Events) -> %% ... process Events ... gen_stage:ask(From, 50), %% ask for more loop({MonitorPid, MonitorRef, Subscriptions}); {MonitorRef, {DOWN, _InnerRef, _Reason}} -> ok %% a producer went down end.
When done:
gen_stage_stream:close(Stream)
## Messages received by the caller
*{$gen_consumer', {Producer, {MonitorRef, InnerRef}}, Events}-
a batch of events from `Producer. The tuple
{Producer, {MonitorRef, InnerRef}}` is the `From to be used with
gen_stage:ask/3 (to ask for more events) and gen_stage:cancel/3.
* {$gen_consumer', {Producer, {MonitorRef, InnerRef}},
{cancel, Reason}}- the producer cancelled the subscription.
* `{MonitorRef, {DOWN, InnerRef, Reason}}` - a producer went down.
If a producer process exits, the stream reacts according to the
`cancel subscription option (permanent, the default, makes the
calling process exit with the same reason; transient only for
abnormal exits; temporary never).
cancel_mode() = permanent | transient | temporary
stream() = {pid(), reference(), #{reference() => subscription()}}
The type returned by subscribe/2.
subscription() = {subscribed, pid(), cancel_mode(), non_neg_integer(), non_neg_integer(), non_neg_integer()} | {cancel, pid()}
| ask/2 | Asks the producer of a subscription for more events. |
| ask/3 | |
| close/1 | Cancels all subscriptions of the stream and kills the helper process. |
| close/2 | |
| subscribe/1 | Subscribes the current process to the given producers with default options. |
| subscribe/2 | Subscribes the current process to the given producers. |
ask(From::gen_stage:from(), Demand::non_neg_integer()) -> ok | noconnect | nosuspend
Asks the producer of a subscription for more events. From is the
{Producer, {MonitorRef, InnerRef}} tuple received with the events.
ask(From::gen_stage:from(), Demand::non_neg_integer(), Opts::[noconnect | nosuspend]) -> ok | noconnect | nosuspend
close(Stream::stream()) -> ok
Cancels all subscriptions of the stream and kills the helper process.
In-flight events may still arrive after this call; the call drains the caller's mailbox of all messages belonging to the stream before returning. Pass aTimeout to stop draining early.
close(X1::stream(), Timeout::timeout()) -> ok
subscribe(Subscriptions::[gen_stage:stage() | {gen_stage:stage(), gen_stage:subscription_options()}]) -> {ok, stream()}
Subscribes the current process to the given producers with default
options. See subscribe/2.
subscribe(Subscriptions::[gen_stage:stage() | {gen_stage:stage(), gen_stage:subscription_options()}], Opts::[{atom(), any()}]) -> {ok, stream()}
Subscribes the current process to the given producers.
Subscriptions is a list of producers or {Producer, Opts} tuples,
where Opts are the subscription options of gen_stage:sync_subscribe/3
(max_demand, min_demand, cancel, ...).
Options:
* demand` - sets the demand mode of the producers to `forward or
accumulate after subscription. Defaults to forward so the
stream can receive items.
* producers` - the processes to set the demand mode on
initialization. Defaults to the processes being subscribed to.
Sometimes the stream subscribes to a `producer_consumer instead of
a producer; in such cases set this option to an empty list or to
the list of actual producers so their demand is properly set.
close/1, the helper
process exits with it and the producers clean up their subscriptions
via monitoring.
Generated by EDoc