Rheo.Group (rheo v1.0.0)

Copy Markdown View Source

Local OTP coordinator for one {instance, stream, group} triple.

Owns demand, fetch scheduling, inflight leases, renewals, and drain — not durable ACK truth. Competing nodes may each run a Group for the same durable group; the backend arbitrates leases.

use Rheo.Consumer produces a child spec for this process, so a consumer's Group lives in the host application's supervision tree. Exactly one Group per {rheo, stream, group} may run on a node; a second start returns {:error, {:already_started, pid}}.

Application Supervision Tree
        |
        +-- Rheo (named instance)
        |     +-- Backend / Instance / Task.Supervisor / GroupSupervisor
        |
        +-- RiskConsumer = Rheo.Group (host-owned)
              +-- handler Tasks

Optional :partitions (:all or a list of ids) scopes fetch to a static assignment. Automatic rebalancing is not implemented (ADR 016).

Summary

Functions

Stops fetching and waits until inflight work settles or timeout elapses.

Milliseconds the supervisor should wait for terminate/2 to drain inflight work.

Functions

drain(server, timeout \\ 5000)

@spec drain(GenServer.server(), timeout()) ::
  :ok | :timeout | {:error, :already_draining}

Stops fetching and waits until inflight work settles or timeout elapses.

The group keeps serving renewals and worker results while a drain is pending. A second concurrent drain returns {:error, :already_draining}.

Arguments

  • server — Group pid or via-tuple name
  • timeout — max wait in milliseconds (default 5000)

Returns

  • :ok — nothing inflight (immediately or after workers finish)
  • :timeout — inflight work remained after timeout
  • {:error, :already_draining} — a drain is already pending

Example

# Graceful pause before a deploy or supervisor stop:
case Rheo.Group.drain(group_pid, 5_000) do
  :ok -> :ready
  :timeout -> :still_working
  {:error, :already_draining} -> :drain_in_progress
end

shutdown_ms()

@spec shutdown_ms() :: pos_integer()

Milliseconds the supervisor should wait for terminate/2 to drain inflight work.

Equals the default drain budget (5000 ms) plus a 1s cushion so terminate/2 can finish answering in-flight drain/2 calls.

Example

# Prefer this over a hard-coded shutdown when writing child specs:
%{
  id: MyApp.RiskConsumer,
  start: {MyApp.RiskConsumer, :start_link, [[]]},
  shutdown: Rheo.Group.shutdown_ms()
}