Rheo.Producer is a plain GenStage producer. Emitted events are %Rheo.Lease{} structs. Broadway is optional sugar on top.

Start a producer

{:ok, producer} =
  Rheo.Producer.start_link(
    rheo: MyRheo,
    stream: "market-events",
    group: "risk",
    max_demand: 20,
    poll_ms: 200
  )

Ensure the group exists (Rheo.create_group/3).

Consumer that settles leases

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

Rheo.Producer.ack/3, nack/4, and reject/4 settle the lease in the backend and release the producer's inflight entry in one call, so renewal stops and a :max_demand slot is freed. Settling with Rheo.ack/2 directly requires a separate Rheo.Producer.confirm/2. Rheo.Broadway.Acknowledger uses the same helpers.

Options

OptionRole
:max_demandCap on unsettled leases
:lease_msTTL; renewals every lease_ms / 2
:poll_msIdle poll when demand exceeds available work
:partitions:all or partition id list
:on_failureDefault Broadway failure settle (:nack / :reject)

When to use GenStage vs Consumer

SurfaceUse when
Rheo.ConsumerSimple OTP handlers, Rheo owns concurrency
Rheo.Producer + GenStageCustom demand topology
Rheo.Producer + BroadwayProcessors, batchers, rate limits

See Broadway and ADR 018.