OTP consumer behaviour for Rheo streams.
use Rheo.Consumer turns a module into a child spec for one local
Rheo.Group. The Group owns demand, concurrency, lease renewal, and
settlement; the module implements handle_event/2.
Part of the SemVer-frozen surface (ADR 029 / ADR 030) — handler outcomes and read-only context will not change meaning without a major version.
Application Supervision Tree
|
+-- Rheo (named instance)
| +-- Backend / Instance / Task.Supervisor / …
|
+-- MyApp.RiskConsumer = Rheo.Group (host-owned)
+-- handler Tasks (up to :concurrency)Handler contract
handle_event/2 receives the event and a read-only context map built once by
setup/1 (or %{}). It returns one outcome:
:ack
{:retry, reason}
{:reject, reason}Handlers run as tasks, up to :concurrency at a time. There is no callback
state: shared mutable state belongs in an Agent, ETS table, or GenServer owned
by the application and referenced from the context (ADR 022).
Options (use and start_link/1)
:stream— required stream name:group— required consumer group:rheo— Rheo instance name (defaultRheo):max_demand— outstanding lease bound (default from config):concurrency— max concurrent handler tasks (default1):lease_ms— lease TTL in milliseconds:poll_ms— idle poll interval (default200):consumer_id— worker identity (default: generated):partitions—:all(default) or a list of partition ids:id— supervisor child id (default{module, rheo, stream, group})
Every option is also passed to setup/1, so application-specific keys
(an Agent pid, a repo, …) can travel through the child spec.
Example
defmodule MyApp.RiskConsumer do
use Rheo.Consumer,
stream: "market-events",
group: "risk",
concurrency: 8,
max_demand: 100
@impl true
def setup(opts) do
{:ok, %{counter: Keyword.fetch!(opts, :counter)}}
end
@impl true
def handle_event(event, %{counter: counter}) do
case Risk.process(event) do
:ok ->
:ok = MyApp.Counter.increment(counter)
:ack
{:temporary_error, reason} ->
{:retry, reason}
{:permanent_error, reason} ->
{:reject, reason}
end
end
end
children = [
{Rheo, name: MyRheo, backend: Rheo.Backend.ETS},
{MyApp.Counter, name: MyApp.Counter},
{MyApp.RiskConsumer, rheo: MyRheo, counter: MyApp.Counter}
]Local ownership
One Rheo.Group per {rheo, stream, group} may run on a BEAM node, and the
supervisor that starts the consumer owns it. Scale local work with
:concurrency; scale across nodes by running one consumer per node. Starting
a second consumer for the same identity returns
{:error, {:already_started, pid}}.
Summary
Callbacks
Handles a single leased event.
Builds the read-only handler context when the Group starts.
Callbacks
@callback handle_event(Rheo.Event.t(), context :: map()) :: :ack | {:retry, term()} | {:reject, term()}
Handles a single leased event.
context is the read-only map returned by setup/1 (or %{}). The Group
settles the lease according to the returned outcome.
Returns
:ack— acknowledge successful processing{:retry, reason}— nack for redelivery (may dead-letter at max attempts){:reject, reason}— permanently dead-letter for this group
Errors
Raising, throwing, exiting, or returning any other value is treated as a
handler failure: the Group logs, emits telemetry, and nacks the lease for
redelivery ({:handler_error, …} or {:invalid_outcome, …}). The callback
itself should not raise for expected business failures — return
{:retry, reason} or {:reject, reason} instead.
Builds the read-only handler context when the Group starts.
Receives every option given to the child spec (:stream, :group, :rheo,
plus any application-specific keys). The returned map is passed unchanged to
every handle_event/2 invocation.
Returns
{:ok, context}—contextmust be a map{:stop, reason}— refuse to start; the Group stops with that reason
Returning {:ok, other} where other is not a map stops the Group with
{:invalid_context, other}. When setup/1 is not implemented, the Group
uses %{}.