# `Rheo`
[🔗](https://github.com/thanos/rheo/blob/v1.0.0/lib/rheo.ex#L1)

Rheo provides durable, searchable, replayable consumer-group semantics over
storage systems.

Events are immutable. Consumption records progress, leases, retries, and
dead-letters in separate consumer-group state. Delivery is at-least-once:
after a crash or lease expiry an event may be delivered again, so handlers
must be idempotent on `Rheo.Event.id`.

Streams may have several partitions. Sequences and ordering are per
partition, key routing uses `:erlang.phash2/2`, and group progress is a
contiguous ACK frontier (`Rheo.lag/3`). There is no global order across
partitions.

## Backends

A Rheo instance runs on one `Rheo.Backend`:

  * `Rheo.Backend.ETS` — in-memory, ephemeral; always available
  * `Rheo.Backend.Mnesia` — durable OTP `:mnesia` (`disc_copies`, single-node)
  * `Rheo.Backend.Mongo` — MongoDB; requires `mongodb_driver`
  * `Rheo.Backend.Ecto` — PostgreSQL or SQLite on a host-owned `Ecto.Repo`;
    requires `ecto_sql` and the repo's adapter
  * `Rheo.Backend.Redis` — Redis Streams; requires `redix`

Integration modules compile only when their dependency is present
(ADR 020).

## Supervision

    children = [
      {Rheo, name: MyRheo, backend: {Rheo.Backend.Mongo, url: "mongodb://localhost:27017/rheo"}},
      {MyApp.RiskConsumer, rheo: MyRheo, concurrency: 8, max_demand: 100}
    ]

    Supervisor.start_link(children, strategy: :one_for_one)

Ecto (the host owns the repo):

    children = [
      MyApp.Repo,
      {Rheo, name: MyRheo, backend: {Rheo.Backend.Ecto, repo: MyApp.Repo}},
      {MyApp.RiskConsumer, rheo: MyRheo, concurrency: 8, max_demand: 100}
    ]

Public APIs accept an optional `:rheo` option targeting a named instance
(default `Rheo`).

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

## Typical low-level flow

    Rheo.create_stream("market-events", partition_count: 4)
    {:ok, event} = Rheo.append("market-events", %{type: "curve_update", key: "EUR-1"})
    Rheo.create_group("market-events", "risk")
    {:ok, leases} = Rheo.fetch("market-events", "risk", limit: 10)
    :ok = Rheo.ack(hd(leases))
    {:ok, _} = Rheo.query("market-events", type: "curve_update")
    {:ok, lag} = Rheo.lag("market-events", "risk")

Search history with `query/2`, `query_page/2`, and `stream_query/2`. Replay
without copying events via `create_group/3` start cursors, `replay/3`, or
`reset_group/3` (`confirm: true`). See `Rheo.Event.Lineage` for correlation
metadata and `Rheo.Partition` for routing helpers.

See `Rheo.Consumer` for the OTP handler API, `Rheo.Producer` and
`Rheo.Broadway` for the GenStage/Broadway surface, and `Rheo.Backend` for
adapters. From 1.0 this surface follows SemVer — see the
[public API guide](public-api.html) and [ADR 030](030-semver-1-0.html).

## Ops

Inventory and health without a second settle path: `list_streams/1`,
`list_groups/2`, `dead_letters/3` (DLQ = dead-letter queue), `group_info/3`,
plus Mix inspect tasks. Optional `Rheo.LiveDashboard.Page` when
`phoenix_live_dashboard` is present (ADR 027). See the [ops guide](ops.html).

## Error atoms

Portable reasons returned by public APIs and backends:

  * `:stream_not_found`, `:group_not_found`, `:already_exists`
  * `:stale_lease`, `:receipt_mismatch`, `:cursor_not_found`
  * `:already_draining`, `:unsupported`, `:backend_unavailable`
  * `:confirm_required`
  * `{:ambiguous, cause}`, `{:failed, cause}`

`Rheo.Settle.classify/1` maps settle errors for Group / Producer policy.

# `group`

```elixir
@type group() :: String.t()
```

Name of a consumer group on a stream.

# `stream`

```elixir
@type stream() :: String.t()
```

Name of an event stream.

# `ack`

```elixir
@spec ack(Rheo.Lease.t(), keyword()) :: :ok | {:error, term()}
```

Acknowledges successful processing of a leased event.

Durable for that consumer group only. Requires the current `lease_id`.

## Arguments

  * `lease` — `%Rheo.Lease{}` returned by `fetch/3`
  * `opts` — only `:rheo` (instance name) is honoured

## Examples

    iex> stream = "doc-ack-" <> Base.encode16(:crypto.strong_rand_bytes(4), case: :lower)
    iex> :ok = Rheo.create_stream(stream)
    iex> :ok = Rheo.create_group(stream, "risk")
    iex> {:ok, _} = Rheo.append(stream, %{type: "once"})
    iex> {:ok, [lease]} = Rheo.fetch(stream, "risk", limit: 1, consumer_id: "c1")
    iex> Rheo.ack(lease)
    :ok
    iex> Rheo.ack(lease)
    {:error, :stale_lease}
    iex> Rheo.fetch(stream, "risk", limit: 1, consumer_id: "c2")
    {:ok, []}

## Returns

  * `:ok`
  * `{:error, :stale_lease}` — lease expired, already ACKed, or replaced
  * `{:error, reason}`

# `append`

```elixir
@spec append(stream(), map(), keyword()) :: {:ok, Rheo.Event.t()} | {:error, term()}
```

Appends a single event to a stream.

Allocates the next sequence number atomically and persists an immutable
`%Rheo.Event{}`.

## Arguments

  * `stream` — existing stream name
  * `payload` — map body. Common keys:
    * `:type` / `"type"` — event type (also copied to `Event.type`)
    * `:key` / `"key"` — optional key
    * `:metadata` / `"metadata"` — merged into `Event.metadata`
    * remaining keys become `Event.payload`
  * `opts` — optional keyword list:
    * `:rheo` — instance name (default `Rheo`)
    * `:id` — explicit event id (default: generated)
    * `:key` — routing / payload key; hashed into a partition when `:partition` omitted
    * `:metadata` — extra metadata map
    * `:timestamp` — `DateTime.t()` (default: clock now)
    * `:partition` — explicit partition (overrides key routing)

## Examples

    iex> stream = "doc-append-" <> Base.encode16(:crypto.strong_rand_bytes(4), case: :lower)
    iex> :ok = Rheo.create_stream(stream)
    iex> {:ok, event} = Rheo.append(stream, %{
    ...>   type: "curve_update",
    ...>   currency: "EUR",
    ...>   curve: "EUR-EURIBOR-6M",
    ...>   price: 2.913,
    ...>   metadata: %{correlation_id: "abc", producer: "pricing-v3"}
    ...> })
    iex> {event.sequence, event.type, event.payload["currency"], event.metadata["correlation_id"]}
    {1, "curve_update", "EUR", "abc"}
    iex> Rheo.append("no-such-stream", %{type: "x"})
    {:error, :stream_not_found}

## Returns

  * `{:ok, %Rheo.Event{}}`
  * `{:error, :stream_not_found}`
  * `{:error, reason}`

# `append_batch`

```elixir
@spec append_batch(stream(), [map()], keyword()) ::
  {:ok, [Rheo.Event.t()]} | {:error, term()}
```

Appends multiple events, allocating contiguous sequences.

## Arguments

  * `stream` — existing stream name
  * `payloads` — list of payload maps (same shape as `append/3`)
  * `opts` — shared options applied to each event (see `append/3`)

## Examples

    iex> stream = "doc-batch-" <> Base.encode16(:crypto.strong_rand_bytes(4), case: :lower)
    iex> :ok = Rheo.create_stream(stream)
    iex> {:ok, events} = Rheo.append_batch(stream, [
    ...>   %{type: "tick", n: 1},
    ...>   %{type: "tick", n: 2},
    ...>   %{type: "tick", n: 3}
    ...> ])
    iex> Enum.map(events, & &1.sequence)
    [1, 2, 3]
    iex> Rheo.append_batch(stream, [])
    {:ok, []}

## Returns

  * `{:ok, [%Rheo.Event{}]}` — possibly empty
  * `{:error, :stream_not_found}`
  * `{:error, reason}`

# `create_group`

```elixir
@spec create_group(stream(), group(), keyword()) :: :ok | {:error, term()}
```

Creates a consumer group on an existing stream.

Groups consume independently: an ACK in `"risk"` does not ACK `"surveillance"`.

## Arguments

  * `stream` — existing stream name
  * `group` — consumer group name
  * `opts` — optional keyword list:
    * `:rheo` — instance name (default `Rheo`)
    * `:max_attempts` — attempts before dead-letter on retry (default from
      application env, typically `5`)
    * `:start_after` — exclusive sequence (integer for all selected partitions,
      or `%{partition => sequence}` map); group begins materializing after this
    * `:start_at` — `DateTime`; begin at the first event at/after this time
    * `:partition` / `:partitions` — limit start cursors to those partitions
      (default all)
## Examples

    iex> stream = "doc-group-" <> Base.encode16(:crypto.strong_rand_bytes(4), case: :lower)
    iex> :ok = Rheo.create_stream(stream)
    iex> Rheo.create_group(stream, "risk")
    :ok
    iex> Rheo.create_group(stream, "risk")
    {:error, :already_exists}
    iex> Rheo.create_group("missing-stream", "risk")
    {:error, :stream_not_found}

## Returns

  * `:ok`
  * `{:error, :already_exists}`
  * `{:error, :stream_not_found}`
  * `{:error, reason}`

# `create_stream`

```elixir
@spec create_stream(stream(), keyword()) :: :ok | {:error, term()}
```

Creates a named stream.

## Arguments

  * `stream` — unique stream name (`t:stream/0`)
  * `opts` — optional keyword list:
    * `:rheo` — instance name (default `Rheo`)
    * `:partition_count` — number of partitions (default `1`). Sequences are
      per-partition; key routing uses `:erlang.phash2/2` when `:key` is set
      on append (ADR 016).

## Examples

    iex> stream = "doc-create-" <> Base.encode16(:crypto.strong_rand_bytes(4), case: :lower)
    iex> Rheo.create_stream(stream)
    :ok
    iex> Rheo.create_stream(stream)
    {:error, :already_exists}

## Returns

  * `:ok` when the stream is created
  * `{:error, :already_exists}` when the name is taken
  * `{:error, reason}` on backend failure

# `dead_letters`

```elixir
@spec dead_letters(stream(), group(), keyword()) ::
  {:ok, [Rheo.DeadLetter.t()]} | {:error, term()}
```

Lists **dead-lettered** (DLQ) deliveries for a group (ops inspect, v0.10+ /
ADR 027).

A dead letter is a delivery that will not be fetched again for this group
until `replay` / `reset_group` — typically after `reject/3` or max nack
attempts. See `Rheo.DeadLetter` and the [ops guide](ops.html).

## Arguments

  * `stream` — stream name
  * `group` — consumer group name
  * `opts` — optional keyword list:
    * `:rheo` — instance name (default `Rheo`)
    * `:limit` — max rows (default 100)
    * `:after` — skip until after this `event_id` (cursor)

## Examples

    iex> stream = "doc-dlq-" <> Base.encode16(:crypto.strong_rand_bytes(4), case: :lower)
    iex> :ok = Rheo.create_stream(stream)
    iex> :ok = Rheo.create_group(stream, "risk")
    iex> {:ok, _} = Rheo.append(stream, %{type: "poison"})
    iex> {:ok, [lease]} = Rheo.fetch(stream, "risk", limit: 1, consumer_id: "c1")
    iex> :ok = Rheo.reject(lease, :invalid_schema)
    iex> {:ok, [dl]} = Rheo.dead_letters(stream, "risk")
    iex> dl.event_id == lease.event_id
    true
    iex> Rheo.fetch(stream, "risk", limit: 1, consumer_id: "c2")
    {:ok, []}

## Returns

  * `{:ok, [%Rheo.DeadLetter{}]}` — empty when none
  * `{:error, :unsupported}` when the backend does not implement `dead_letters/4`
  * `{:error, reason}` on backend failure

# `ensure_indexes`

```elixir
@spec ensure_indexes(keyword()) :: :ok | {:error, term()}
```

Ensures backend indexes exist (safe to call repeatedly).

## Arguments

  * `opts` — only `:rheo` (instance name) is honoured

## Examples

    iex> Rheo.ensure_indexes()
    :ok

## Returns

  * `:ok`
  * `{:error, reason}` on index creation failure

# `fetch`

```elixir
@spec fetch(stream(), group(), keyword()) ::
  {:ok, [Rheo.Lease.t()]} | {:error, term()}
```

Fetches up to `:limit` events for a consumer group, creating leases.

Competing workers in the same group receive distinct leases. Independent
groups may lease the same events concurrently.

## Arguments

  * `stream` — stream name
  * `group` — consumer group name
  * `opts` — optional keyword list:
    * `:rheo` — instance name (default `Rheo`)
    * `:limit` — max leases / demand bound (default from config)
    * `:consumer_id` — worker identity (default: generated)
    * `:lease_ms` — lease TTL in milliseconds (default from config)

## Examples

    iex> stream = "doc-fetch-" <> Base.encode16(:crypto.strong_rand_bytes(4), case: :lower)
    iex> :ok = Rheo.create_stream(stream)
    iex> :ok = Rheo.create_group(stream, "risk")
    iex> {:ok, _} = Rheo.append(stream, %{type: "order", id: 1})
    iex> {:ok, [lease]} = Rheo.fetch(stream, "risk", limit: 10, consumer_id: "c1")
    iex> {lease.group, lease.attempt, lease.event.type}
    {"risk", 1, "order"}
    iex> Rheo.fetch(stream, "missing", limit: 1)
    {:error, :group_not_found}

## Returns

  * `{:ok, [%Rheo.Lease{}]}` — empty when no eligible work
  * `{:error, :group_not_found}`
  * `{:error, reason}`

# `group_info`

```elixir
@spec group_info(stream(), group(), keyword()) ::
  {:ok, Rheo.GroupInfo.t()} | {:error, term()}
```

Returns group health: lag plus inflight and dead-letter counts (v0.10+ / ADR 027).

## Arguments

  * `stream` — stream name
  * `group` — consumer group name
  * `opts` — optional keyword list:
    * `:rheo` — instance name (default `Rheo`)

## Examples

    iex> stream = "doc-group-info-" <> Base.encode16(:crypto.strong_rand_bytes(4), case: :lower)
    iex> :ok = Rheo.create_stream(stream)
    iex> :ok = Rheo.create_group(stream, "risk")
    iex> {:ok, _} = Rheo.append(stream, %{type: "tick"})
    iex> {:ok, info} = Rheo.group_info(stream, "risk")
    iex> {info.lag.lag, info.inflight_count, info.dead_letter_count}
    {1, 0, 0}

## Returns

  * `{:ok, %Rheo.GroupInfo{}}`
  * `{:error, :unsupported}` when the backend does not implement `group_info/4`
  * `{:error, :group_not_found}` / `{:error, :stream_not_found}`
  * `{:error, reason}` on backend failure

# `lag`

```elixir
@spec lag(stream(), group(), keyword()) :: {:ok, Rheo.Lag.t()} | {:error, term()}
```

Returns contiguous-frontier lag for a consumer group (v0.5+).

Per-partition lag is `max(high_watermark - frontier, 0)`. Aggregate `lag` is
the sum across partitions. See `Rheo.Lag` and ADR 016.

## Arguments

  * `stream` — stream name
  * `group` — consumer group name
  * `opts` — optional keyword list:
    * `:rheo` — instance name (default `Rheo`)

## Examples

    iex> stream = "doc-lag-" <> Base.encode16(:crypto.strong_rand_bytes(4), case: :lower)
    iex> :ok = Rheo.create_stream(stream)
    iex> :ok = Rheo.create_group(stream, "risk")
    iex> {:ok, _} = Rheo.append(stream, %{type: "tick"})
    iex> {:ok, lag} = Rheo.lag(stream, "risk")
    iex> {lag.lag, lag.partitions[0].high_watermark}
    {1, 1}
    iex> Rheo.lag(stream, "missing")
    {:error, :group_not_found}

## Returns

  * `{:ok, %Rheo.Lag{}}`
  * `{:error, :group_not_found}` / `{:error, :stream_not_found}`
  * `{:error, reason}` on backend failure

# `list_groups`

```elixir
@spec list_groups(stream(), keyword()) :: {:ok, [group()]} | {:error, term()}
```

Lists consumer group names for a stream (ops inspect, v0.10+ / ADR 027).

## Arguments

  * `stream` — stream name
  * `opts` — optional keyword list:
    * `:rheo` — instance name (default `Rheo`)

## Examples

    iex> stream = "doc-list-groups-" <> Base.encode16(:crypto.strong_rand_bytes(4), case: :lower)
    iex> :ok = Rheo.create_stream(stream)
    iex> :ok = Rheo.create_group(stream, "risk")
    iex> Rheo.list_groups(stream)
    {:ok, ["risk"]}

## Returns

  * `{:ok, [group_name]}` — empty list when the stream has no groups
  * `{:error, :unsupported}` when the backend does not implement `list_groups/3`
  * `{:error, reason}` on backend failure

# `list_streams`

```elixir
@spec list_streams(keyword()) :: {:ok, [stream()]} | {:error, term()}
```

Lists registered stream names (ops inspect, v0.10+ / ADR 027).

## Arguments

  * `opts` — optional keyword list:
    * `:rheo` — instance name (default `Rheo`)

## Examples

    iex> stream = "doc-list-streams-" <> Base.encode16(:crypto.strong_rand_bytes(4), case: :lower)
    iex> :ok = Rheo.create_stream(stream)
    iex> {:ok, names} = Rheo.list_streams()
    iex> stream in names
    true

## Returns

  * `{:ok, [stream_name]}`
  * `{:error, :unsupported}` when the backend does not implement `list_streams/2`
  * `{:error, reason}` on backend failure

# `nack`

```elixir
@spec nack(Rheo.Lease.t(), term(), keyword()) :: :ok | {:error, term()}
```

Returns a leased event for retry / redelivery according to group policy.

Sets delivery status back to available (or dead-letters when
`attempt >= max_attempts`).

## Arguments

  * `lease` — `%Rheo.Lease{}` from `fetch/3`
  * `reason` — any term stored for diagnostics (default `:retry`)
  * `opts` — only `:rheo` (instance name) is honoured

## Examples

    iex> stream = "doc-nack-" <> Base.encode16(:crypto.strong_rand_bytes(4), case: :lower)
    iex> :ok = Rheo.create_stream(stream)
    iex> :ok = Rheo.create_group(stream, "risk", max_attempts: 5)
    iex> {:ok, _} = Rheo.append(stream, %{type: "tmp"})
    iex> {:ok, [lease]} = Rheo.fetch(stream, "risk", limit: 1, consumer_id: "c1")
    iex> Rheo.nack(lease, :temporary_error)
    :ok
    iex> {:ok, [again]} = Rheo.fetch(stream, "risk", limit: 1, consumer_id: "c2")
    iex> again.attempt
    2

## Returns

  * `:ok`
  * `{:error, :stale_lease}`
  * `{:error, reason}`

# `ping`

```elixir
@spec ping(keyword()) :: :ok | {:error, term()}
```

Verifies connectivity to the configured backend.

## Arguments

  * `opts` — only `:rheo` (instance name) is honoured

## Examples

    iex> Rheo.ping()
    :ok

## Returns

  * `:ok`
  * `{:error, reason}` when the backend is unreachable

# `query`

```elixir
@spec query(Rheo.Query.t() | stream(), keyword()) ::
  {:ok, [Rheo.Event.t()]} | {:error, term()}
```

Queries historical events. Consumption never removes events from the log.

Accepts a `%Rheo.Query{}` or a stream name plus keyword filters (converted via
`Rheo.Query.new/2`).

## Arguments

  * `query_or_stream` — `%Rheo.Query{}` or stream name
  * `opts` — when the first argument is a stream, filter options (see
    `Rheo.Query`); always accepts `:rheo`

## Examples

    iex> stream = "doc-query-" <> Base.encode16(:crypto.strong_rand_bytes(4), case: :lower)
    iex> :ok = Rheo.create_stream(stream)
    iex> {:ok, _} = Rheo.append(stream, %{
    ...>   type: "curve_update",
    ...>   currency: "EUR",
    ...>   curve: "EUR-EURIBOR-6M",
    ...>   price: 2.913
    ...> })
    iex> {:ok, _} = Rheo.append(stream, %{type: "curve_update", currency: "USD", curve: "USD-SOFR"})
    iex> {:ok, [event]} = Rheo.query(stream, type: "curve_update", currency: "EUR")
    iex> event.payload["curve"]
    "EUR-EURIBOR-6M"
    iex> q = Rheo.Query.new(stream, where: [type: "curve_update", currency: "USD"])
    iex> {:ok, [usd]} = Rheo.query(q)
    iex> usd.payload["currency"]
    "USD"

## Returns

  * `{:ok, [%Rheo.Event{}]}`
  * `{:error, reason}` on backend failure

# `query_page`

```elixir
@spec query_page(Rheo.Query.t() | stream(), keyword()) ::
  {:ok, Rheo.Page.t()} | {:error, term()}
```

Queries one page of historical events.

Returns `{:ok, %Rheo.Page{}}`. When more results may exist, `page.next_cursor`
is a `%{partition => after_sequence}` map; pass it as `cursor:` on the next
call (or embed on `%Rheo.Query{}`). Cursor pagination is defined for
ascending sequence order. The legacy `%{after_sequence: n}` cursor is still
accepted as a global lower bound.

## Arguments

  * `query_or_stream` — `%Rheo.Query{}` or stream name
  * `opts` — when the first argument is a stream, filter options (see
    `Rheo.Query`); always accepts `:rheo` (instance name, default `Rheo`)

## Examples

    iex> stream = "doc-query-page-" <> Base.encode16(:crypto.strong_rand_bytes(4), case: :lower)
    iex> :ok = Rheo.create_stream(stream)
    iex> {:ok, _} = Rheo.append_batch(stream, [
    ...>   %{type: "a"}, %{type: "b"}, %{type: "c"}
    ...> ])
    iex> {:ok, page} = Rheo.query_page(stream, limit: 2)
    iex> Enum.map(page.events, & &1.sequence)
    [1, 2]
    iex> page.next_cursor
    %{0 => 2}
    iex> {:ok, page2} = Rheo.query_page(stream, limit: 2, cursor: page.next_cursor)
    iex> {Enum.map(page2.events, & &1.sequence), page2.next_cursor}
    {[3], nil}

## Returns

  * `{:ok, %Rheo.Page{}}` — `events` may be empty; `next_cursor` is `nil` when done
  * `{:error, reason}` on backend failure

# `read`

```elixir
@spec read(stream(), keyword()) :: {:ok, [Rheo.Event.t()]} | {:error, term()}
```

Reads events by sequence without affecting consumer-group state.

## Arguments

  * `stream` — stream name
  * `opts` — optional keyword list:
    * `:rheo` — instance name (default `Rheo`)
    * `:after` — return sequences strictly greater than this (default `0`)
    * `:limit` — max events (default `100`)
    * `:partition` — partition (default `0`)

## Examples

    iex> stream = "doc-read-" <> Base.encode16(:crypto.strong_rand_bytes(4), case: :lower)
    iex> :ok = Rheo.create_stream(stream)
    iex> {:ok, _} = Rheo.append_batch(stream, [%{type: "a"}, %{type: "b"}, %{type: "c"}])
    iex> {:ok, [first, second]} = Rheo.read(stream, after: 0, limit: 2)
    iex> {first.sequence, second.sequence}
    {1, 2}
    iex> {:ok, [third]} = Rheo.read(stream, after: 2, limit: 10)
    iex> third.type
    "c"

## Returns

  * `{:ok, [%Rheo.Event{}]}` — empty list when nothing matches

# `reject`

```elixir
@spec reject(Rheo.Lease.t(), term(), keyword()) :: :ok | {:error, term()}
```

Permanently rejects a leased event for this consumer group (dead-letter).

Other groups are unaffected. The immutable event remains queryable.

## Arguments

  * `lease` — `%Rheo.Lease{}` from `fetch/3`
  * `reason` — any term stored on the delivery (default `:rejected`)
  * `opts` — only `:rheo` (instance name) is honoured

## Examples

    iex> stream = "doc-reject-" <> Base.encode16(:crypto.strong_rand_bytes(4), case: :lower)
    iex> :ok = Rheo.create_stream(stream)
    iex> :ok = Rheo.create_group(stream, "risk")
    iex> {:ok, _} = Rheo.append(stream, %{type: "poison"})
    iex> {:ok, [lease]} = Rheo.fetch(stream, "risk", limit: 1, consumer_id: "c1")
    iex> Rheo.reject(lease, :invalid_schema)
    :ok
    iex> Rheo.fetch(stream, "risk", limit: 1, consumer_id: "c2")
    {:ok, []}

## Returns

  * `:ok`
  * `{:error, :stale_lease}`
  * `{:error, reason}`

# `renew`

```elixir
@spec renew(Rheo.Lease.t(), keyword()) :: {:ok, Rheo.Lease.t()} | {:error, term()}
```

Extends an active lease when `lease_id` still matches.

## Arguments

  * `lease` — `%Rheo.Lease{}` from `fetch/3`
  * `opts` — optional keyword list:
    * `:rheo` — instance name (default `Rheo`)
    * `:lease_ms` — new TTL from now (default from config)

## Examples

    iex> stream = "doc-renew-" <> Base.encode16(:crypto.strong_rand_bytes(4), case: :lower)
    iex> :ok = Rheo.create_stream(stream)
    iex> :ok = Rheo.create_group(stream, "risk")
    iex> {:ok, _} = Rheo.append(stream, %{type: "hold"})
    iex> {:ok, [lease]} = Rheo.fetch(stream, "risk", limit: 1, consumer_id: "c1")
    iex> {:ok, renewed} = Rheo.renew(lease, lease_ms: 60_000)
    iex> renewed.lease_id == lease.lease_id
    true
    iex> :ok = Rheo.ack(lease)
    iex> Rheo.renew(lease)
    {:error, :stale_lease}

## Returns

  * `{:ok, %Rheo.Lease{}}` with updated `expires_at`
  * `{:error, :stale_lease}`
  * `{:error, :backend_unavailable}`
  * `{:error, reason}`

# `replay`

```elixir
@spec replay(stream(), group(), keyword()) :: :ok | {:error, term()}
```

Replays history for a consumer group without copying events.

Prefer a new group with `:start_after` / `:start_at` when isolating
replay from production consumers. Design details:
[ADR 015](https://github.com/thanos/rheo/blob/main/docs/adr/015-replay-semantics.md).

## Arguments

  * `stream` — stream name
  * `group` — consumer group name
  * `opts` — optional keyword list (exactly one of the replay selectors below,
    plus optional scope / instance keys):
    * `:rheo` — instance name (default `Rheo`)
    * `:from_sequence` — exclusive lower bound; deliveries from the next
      sequence become available again
    * `:from` — `DateTime`; resolved to a sequence then same as `:from_sequence`
    * `:query` — `%Rheo.Query{}` or keyword filters on the stream; matching
      events are re-opened for this group
    * `:partition` / `:partitions` — limit replay scope

## Examples

    iex> stream = "doc-replay-" <> Base.encode16(:crypto.strong_rand_bytes(4), case: :lower)
    iex> :ok = Rheo.create_stream(stream)
    iex> :ok = Rheo.create_group(stream, "risk")
    iex> {:ok, event} = Rheo.append(stream, %{type: "once"})
    iex> {:ok, [lease]} = Rheo.fetch(stream, "risk", limit: 1, consumer_id: "c1")
    iex> :ok = Rheo.ack(lease)
    iex> Rheo.fetch(stream, "risk", limit: 1, consumer_id: "c2")
    {:ok, []}
    iex> Rheo.replay(stream, "risk", from_sequence: 0)
    :ok
    iex> {:ok, [again]} = Rheo.fetch(stream, "risk", limit: 1, consumer_id: "c3")
    iex> again.event_id == event.id
    true
    iex> Rheo.replay(stream, "risk")
    {:error, :invalid_replay_opts}

## Returns

  * `:ok`
  * `{:error, :invalid_replay_opts}` when no replay selector is given
  * `{:error, :group_not_found}` / `{:error, :stream_not_found}`
  * `{:error, :no_events_in_range}` when `:from` finds no matching event
  * `{:error, reason}` on backend failure

# `reset_group`

```elixir
@spec reset_group(stream(), group(), keyword()) :: :ok | {:error, term()}
```

Destructively clears deliveries for one group and resets its cursor.

Requires `confirm: true`. Never deletes immutable events. Other groups are
unaffected. Optional `:start_after` sets the post-reset materialization cursor
(exclusive), default `0` (replay from the beginning).

## Arguments

  * `stream` — stream name
  * `group` — consumer group name
  * `opts` — optional keyword list:
    * `:rheo` — instance name (default `Rheo`)
    * `:confirm` — must be `true` or the call returns `{:error, :confirm_required}`
    * `:start_after` — exclusive materialization cursor after reset (default `0`)
    * `:partition` / `:partitions` — limit reset scope when supported

## Examples

    iex> stream = "doc-reset-" <> Base.encode16(:crypto.strong_rand_bytes(4), case: :lower)
    iex> :ok = Rheo.create_stream(stream)
    iex> :ok = Rheo.create_group(stream, "risk")
    iex> {:ok, event} = Rheo.append(stream, %{type: "keep"})
    iex> {:ok, [lease]} = Rheo.fetch(stream, "risk", limit: 1, consumer_id: "c1")
    iex> :ok = Rheo.ack(lease)
    iex> Rheo.reset_group(stream, "risk")
    {:error, :confirm_required}
    iex> Rheo.reset_group(stream, "risk", confirm: true)
    :ok
    iex> {:ok, [still]} = Rheo.read(stream, after: 0, limit: 1)
    iex> still.id == event.id
    true
    iex> {:ok, [again]} = Rheo.fetch(stream, "risk", limit: 1, consumer_id: "c2")
    iex> again.event_id == event.id
    true

## Returns

  * `:ok`
  * `{:error, :confirm_required}` when `confirm: true` is missing
  * `{:error, :group_not_found}` / `{:error, :stream_not_found}`
  * `{:error, reason}` on backend failure

# `start_link`

```elixir
@spec start_link(keyword()) :: Supervisor.on_start()
```

Starts a Rheo instance supervisor.

## Options

  * `:name` — instance name (default `Rheo`). Also used as the supervisor name.
  * `:backend` — `{module, opts}` or a module. Required unless `:url` is
    given, in which case `Rheo.Backend.Mongo` is used with the remaining
    options (`:url`, `:pool_size`) when that module is available.

## Examples

    iex> name = String.to_atom("rheo_doc_189444")
    iex> url = System.get_env("RHEO_MONGO_URL", "mongodb://localhost:27017/rheo_test")
    iex> {:ok, pid} = Rheo.start_link(name: name, backend: {Rheo.Backend.Mongo, url: url})
    iex> is_pid(pid)
    true

## Returns

  * `{:ok, pid}` on success
  * `{:error, {:already_started, pid}}` if the instance name is taken
  * `{:error, reason}` on backend start failure

## Errors

Raises `ArgumentError` when no `:backend` is given and the Mongo shorthand
cannot be used.

# `stream_query`

```elixir
@spec stream_query(Rheo.Query.t() | stream(), keyword()) :: Enumerable.t()
```

Lazily streams query results page by page.

Each element is a `%Rheo.Event{}`. Uses `query_page/2` internally. On a
`query_page/2` error this **raises**
`RuntimeError` with message `"Rheo.stream_query failed: …"` rather than
returning `{:error, reason}`.

## Arguments

  * `query_or_stream` — `%Rheo.Query{}` or stream name
  * `opts` — when the first argument is a stream, filter options (see
    `Rheo.Query`); always accepts `:rheo` (instance name, default `Rheo`)

## Examples

    iex> stream = "doc-stream-query-" <> Base.encode16(:crypto.strong_rand_bytes(4), case: :lower)
    iex> :ok = Rheo.create_stream(stream)
    iex> {:ok, _} = Rheo.append_batch(stream, [
    ...>   %{type: "a"}, %{type: "b"}, %{type: "c"}, %{type: "d"}
    ...> ])
    iex> stream |> Rheo.stream_query(limit: 2) |> Enum.map(& &1.sequence)
    [1, 2, 3, 4]

## Returns

  * an `Enumerable` of `%Rheo.Event{}` (possibly empty)

## Errors

Raises `RuntimeError` when an underlying `query_page/2` returns
`{:error, reason}` (`raise "Rheo.stream_query failed: #{inspect(reason)}"`).

---

*Consult [api-reference.md](api-reference.md) for complete listing*
