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:

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 and ADR 030.

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.

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.

Summary

Types

Name of a consumer group on a stream.

Name of an event stream.

Functions

Acknowledges successful processing of a leased event.

Appends a single event to a stream.

Appends multiple events, allocating contiguous sequences.

Creates a consumer group on an existing stream.

Creates a named stream.

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

Ensures backend indexes exist (safe to call repeatedly).

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

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

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

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

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

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

Verifies connectivity to the configured backend.

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

Queries one page of historical events.

Reads events by sequence without affecting consumer-group state.

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

Extends an active lease when lease_id still matches.

Replays history for a consumer group without copying events.

Destructively clears deliveries for one group and resets its cursor.

Starts a Rheo instance supervisor.

Lazily streams query results page by page.

Types

group()

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

Name of a consumer group on a stream.

stream()

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

Name of an event stream.

Functions

ack(lease, opts \\ [])

@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(stream, payload, opts \\ [])

@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(stream, payloads, opts \\ [])

@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(stream, group, opts \\ [])

@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(stream, opts \\ [])

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

Creates a named stream.

Arguments

  • stream — unique stream name (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(stream, group, opts \\ [])

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

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(opts \\ [])

@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(stream, group, opts \\ [])

@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(stream, group, opts \\ [])

@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(stream, group, opts \\ [])

@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(stream, opts \\ [])

@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(opts \\ [])

@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(lease, reason \\ :retry, opts \\ [])

@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(opts \\ [])

@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(query_or_stream, opts \\ [])

@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(query_or_stream, opts \\ [])

@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(stream, opts \\ [])

@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(lease, reason \\ :rejected, opts \\ [])

@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(lease, opts \\ [])

@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(stream, group, opts \\ [])

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

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(stream, group, opts \\ [])

@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(opts \\ [])

@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(query_or_stream, opts \\ [])

@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

Errors

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