GenStage producer that turns demand into Rheo leases.
Rheo.Producer is the demand-driven consumption surface for a durable
consumer group. It is an alternative to Rheo.Consumer / Rheo.Group, not a
layer on top of them: the producer owns fetch and lease renewal, while the
downstream pipeline owns concurrency and settles each lease. Do not point both
a Rheo.Consumer and a Rheo.Producer at the same {rheo, stream, group}
unless you intend competing consumers.
Emitted events are %Rheo.Lease{} structs. For Broadway, add
transformer: {Rheo.Broadway, :transform, []} and the lease is wrapped into a
%Broadway.Message{} whose acknowledger settles it (see Rheo.Broadway).
Broadway / GenStage demand
|
v
Rheo.Producer --fetch/renew--> Backend
|
| %Rheo.Lease{}
v
processors / batchers
|
v
Acknowledger --> Producer.ack | nack | reject
|
+--> release inflight (stops renew, frees demand)Options
:stream— required stream name:group— required consumer group:rheo— instance name (defaultRheo):max_demand— maximum unsettled leases held at once (default from:default_max_demand, typically10). This bounds leases, not pipeline concurrency, which Broadway's:concurrencyowns.:lease_ms— lease TTL; inflight leases are renewed everylease_ms / 2:poll_ms— idle poll interval when demand outruns available work (default200):consumer_id— worker identity recorded on each lease (default: generated):partitions—:all(default) or a list of partition ids to fetch from:on_failure—:nack(default) or:reject, used byRheo.Broadway.Acknowledgerwhen a message fails:name— optional process name (ignored under Broadway, which names the producer stage itself)
Broadway
producer: [
module: {Rheo.Producer, rheo: MyRheo, stream: "market-events", group: "risk"},
transformer: {Rheo.Broadway, :transform, []},
concurrency: 1
]GenStage
{:ok, producer} =
Rheo.Producer.start_link(rheo: MyRheo, stream: "market-events", group: "risk")
# In a consumer's handle_events/3:
def handle_events(leases, _from, state) do
Enum.each(leases, fn lease ->
case Risk.process(lease.event) do
:ok -> :ok = Rheo.Producer.ack(state.producer, lease, rheo: MyRheo)
{:error, reason} -> :ok = Rheo.Producer.nack(state.producer, lease, reason, rheo: MyRheo)
end
end)
{:noreply, [], state}
endSettlement boundary
A lease held by a pipeline has two pieces of state: the durable delivery in
the backend and the producer's inflight entry that drives renewal and
:max_demand. ack/3, nack/4, and reject/4 settle both in one call:
they run the durable settle and then release the inflight entry, whatever the
settle result. A fenced or unavailable settle leaves the lease to expire and
redeliver (at-least-once), so the producer must not keep renewing it.
confirm/2 is the lower-level half for callers that settle with Rheo.ack/2
directly. Flow pipelines built on Flow.from_stages/2 use the same three
helpers from on_trigger once an aggregate is durable — see the Flow
readiness spike.
Backpressure and renewal
Each handle_demand/2 accumulates demand and fetches at most
min(demand, max_demand - inflight) leases. When the backend returns fewer
leases than asked for, the producer polls every :poll_ms; when a fetch fails
it backs off exponentially (Rheo.Backoff) and emits
[:rheo, :fetch, :error]. Inflight leases are renewed on a timer and emit
[:rheo, :lease, :renew]; a lease that has gone stale is dropped so the
backend can redeliver it.
Requires {:gen_stage, "~> 1.2"} in your dependencies; Broadway's
prepare_for_draining/1 is implemented when {:broadway, "~> 1.2"} is present.
See ADR 018 and Rheo.Broadway.
Summary
Functions
Acknowledges lease durably and releases it from the producer.
Returns the configuration of the producer running in the calling process.
Tells the producer that leases were settled elsewhere.
Stops fetching and waits until inflight leases are settled or timeout elapses.
Returns the number of leases the producer holds but has not seen confirmed.
Returns lease for retry (or dead-letters it at max_attempts) and releases
it from the producer.
Dead-letters lease for its group and releases it from the producer.
Starts a producer.
Functions
@spec ack(GenServer.server(), Rheo.Lease.t(), keyword()) :: :ok | {:error, term()}
Acknowledges lease durably and releases it from the producer.
Arguments
producer— the producer that emitted the leaselease— the%Rheo.Lease{}opts—:rheoinstance name (defaultRheo)
Returns
The result of Rheo.ack/2. The inflight entry is released even on error;
see "Settlement boundary".
@spec config() :: %{ rheo: atom(), stream: Rheo.stream(), group: Rheo.group(), on_failure: :nack | :reject } | nil
Returns the configuration of the producer running in the calling process.
Returns nil outside a producer process. Broadway invokes the transformer
inside the producer, so Rheo.Broadway.transform/2 uses this to pick up
:rheo and :on_failure without repeating them in the transformer arguments.
Examples
iex> Rheo.Producer.config()
nil
@spec confirm(GenServer.server(), String.t() | [String.t()]) :: :ok
Tells the producer that leases were settled elsewhere.
Removes them from the inflight set so renewal stops and :max_demand capacity
is released. Accepts one lease_id or a list. Asynchronous. Prefer ack/3,
nack/4, and reject/4, which settle and confirm together.
@spec drain(GenServer.server(), timeout()) :: :ok | :timeout | {:error, :already_draining}
Stops fetching and waits until inflight leases are settled or timeout elapses.
Returns :ok when nothing is inflight, :timeout otherwise, or
{:error, :already_draining} if a drain is already pending. The producer
keeps serving confirmations and renewals while a drain is pending. Broadway
calls prepare_for_draining/1 on its own during shutdown; this is the
equivalent for plain GenStage pipelines.
@spec inflight_count(GenServer.server()) :: non_neg_integer()
Returns the number of leases the producer holds but has not seen confirmed.
@spec nack(GenServer.server(), Rheo.Lease.t(), term(), keyword()) :: :ok | {:error, term()}
Returns lease for retry (or dead-letters it at max_attempts) and releases
it from the producer.
Arguments
producer— the producer that emitted the leaselease— the%Rheo.Lease{}reason— stored for diagnostics (default:retry)opts—:rheoinstance name (defaultRheo)
Returns
The result of Rheo.nack/3.
@spec reject(GenServer.server(), Rheo.Lease.t(), term(), keyword()) :: :ok | {:error, term()}
Dead-letters lease for its group and releases it from the producer.
Arguments
producer— the producer that emitted the leaselease— the%Rheo.Lease{}reason— stored on the delivery (default:rejected)opts—:rheoinstance name (defaultRheo)
Returns
The result of Rheo.reject/3.
@spec start_link(keyword()) :: GenServer.on_start()
Starts a producer.
Accepts the options documented in the module doc. :name, when given, is
passed through to GenStage.start_link/3.
Returns
{:ok, pid}{:error, {:already_started, pid}}{:error, reason}