Rheo.Consumer behaviour (rheo v1.0.0)

Copy Markdown View Source

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

  • :stream — required stream name
  • :group — required consumer group
  • :rheo — Rheo instance name (default Rheo)
  • :max_demand — outstanding lease bound (default from config)
  • :concurrency — max concurrent handler tasks (default 1)
  • :lease_ms — lease TTL in milliseconds
  • :poll_ms — idle poll interval (default 200)
  • :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

handle_event(t, context)

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

setup keyword

(optional)
@callback setup(keyword()) :: {:ok, map()} | {:stop, term()}

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} — context must 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 %{}.