# `Rheo.Producer`
[🔗](https://github.com/thanos/rheo/blob/v1.0.0/lib/rheo/producer.ex#L3)

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

# `ack`

```elixir
@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`

```elixir
@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`

```elixir
@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`

```elixir
@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`

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

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

# `nack`

```elixir
@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`

```elixir
@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`

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

---

*Consult [api-reference.md](api-reference.md) for complete listing*
