Rheo.Producer (rheo v1.0.0)

Copy Markdown View Source

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 (default Rheo)
  • :max_demand — maximum unsettled leases held at once (default from :default_max_demand, typically 10). This bounds leases, not pipeline concurrency, which Broadway's :concurrency owns.
  • :lease_ms — lease TTL; inflight leases are renewed every lease_ms / 2
  • :poll_ms — idle poll interval when demand outruns available work (default 200)
  • :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 by Rheo.Broadway.Acknowledger when 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}
end

Settlement 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

ack(producer, lease, opts \\ [])

@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 lease
  • lease — the %Rheo.Lease{}
  • opts — :rheo instance name (default Rheo)

Returns

The result of Rheo.ack/2. The inflight entry is released even on error; see "Settlement boundary".

config()

@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

confirm(producer, lease_id)

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

drain(producer, timeout \\ 5000)

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

inflight_count(producer)

@spec inflight_count(GenServer.server()) :: non_neg_integer()

Returns the number of leases the producer holds but has not seen confirmed.

nack(producer, lease, reason \\ :retry, opts \\ [])

@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 lease
  • lease — the %Rheo.Lease{}
  • reason — stored for diagnostics (default :retry)
  • opts — :rheo instance name (default Rheo)

Returns

The result of Rheo.nack/3.

reject(producer, lease, reason \\ :rejected, opts \\ [])

@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 lease
  • lease — the %Rheo.Lease{}
  • reason — stored on the delivery (default :rejected)
  • opts — :rheo instance name (default Rheo)

Returns

The result of Rheo.reject/3.

start_link(opts)

@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}