Module gen_stage_stream

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.

Description

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).

Data Types

cancel_mode()

cancel_mode() = permanent | transient | temporary

stream()

stream() = {pid(), reference(), #{reference() => subscription()}}

The type returned by subscribe/2.

subscription()

subscription() = {subscribed, pid(), cancel_mode(), non_neg_integer(), non_neg_integer(), non_neg_integer()} | {cancel, pid()}

Function Index

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.

Function Details

ask/2

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/3

ask(From::gen_stage:from(), Demand::non_neg_integer(), Opts::[noconnect | nosuspend]) -> ok | noconnect | nosuspend

close/1

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 a Timeout to stop draining early.

close/2

close(X1::stream(), Timeout::timeout()) -> ok

subscribe/1

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/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.

If the caller process exits before calling close/1, the helper process exits with it and the producers clean up their subscriptions via monitoring.


Generated by EDoc