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 TasksOptional :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
@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 nametimeout— max wait in milliseconds (default5000)
Returns
:ok— nothing inflight (immediately or after workers finish):timeout— inflight work remained aftertimeout{: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
@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()
}