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 availableRheo.Backend.Mnesia— durable OTP:mnesia(disc_copies, single-node)Rheo.Backend.Mongo— MongoDB; requiresmongodb_driverRheo.Backend.Ecto— PostgreSQL or SQLite on a host-ownedEcto.Repo; requiresecto_sqland the repo's adapterRheo.Backend.Redis— Redis Streams; requiresredix
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
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
Functions
@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 byfetch/3opts— 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}
@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 namepayload— map body. Common keys::type/"type"— event type (also copied toEvent.type):key/"key"— optional key:metadata/"metadata"— merged intoEvent.metadata- remaining keys become
Event.payload
opts— optional keyword list::rheo— instance name (defaultRheo):id— explicit event id (default: generated):key— routing / payload key; hashed into a partition when:partitionomitted: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}
@spec append_batch(stream(), [map()], keyword()) :: {:ok, [Rheo.Event.t()]} | {:error, term()}
Appends multiple events, allocating contiguous sequences.
Arguments
stream— existing stream namepayloads— list of payload maps (same shape asappend/3)opts— shared options applied to each event (seeappend/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}
Creates a consumer group on an existing stream.
Groups consume independently: an ACK in "risk" does not ACK "surveillance".
Arguments
stream— existing stream namegroup— consumer group nameopts— optional keyword list::rheo— instance name (defaultRheo):max_attempts— attempts before dead-letter on retry (default from application env, typically5):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}
Creates a named stream.
Arguments
stream— unique stream name (stream/0)opts— optional keyword list::rheo— instance name (defaultRheo):partition_count— number of partitions (default1). Sequences are per-partition; key routing uses:erlang.phash2/2when:keyis 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
:okwhen the stream is created{:error, :already_exists}when the name is taken{:error, reason}on backend failure
@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 namegroup— consumer group nameopts— optional keyword list::rheo— instance name (defaultRheo):limit— max rows (default 100):after— skip until after thisevent_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 implementdead_letters/4{:error, reason}on backend failure
Ensures backend indexes exist (safe to call repeatedly).
Arguments
opts— only:rheo(instance name) is honoured
Examples
iex> Rheo.ensure_indexes()
:okReturns
:ok{:error, reason}on index creation failure
@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 namegroup— consumer group nameopts— optional keyword list::rheo— instance name (defaultRheo):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}
@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 namegroup— consumer group nameopts— optional keyword list::rheo— instance name (defaultRheo)
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 implementgroup_info/4{:error, :group_not_found}/{:error, :stream_not_found}{:error, reason}on backend failure
@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 namegroup— consumer group nameopts— optional keyword list::rheo— instance name (defaultRheo)
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
Lists consumer group names for a stream (ops inspect, v0.10+ / ADR 027).
Arguments
stream— stream nameopts— optional keyword list::rheo— instance name (defaultRheo)
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 implementlist_groups/3{:error, reason}on backend failure
Lists registered stream names (ops inspect, v0.10+ / ADR 027).
Arguments
opts— optional keyword list::rheo— instance name (defaultRheo)
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
trueReturns
{:ok, [stream_name]}{:error, :unsupported}when the backend does not implementlist_streams/2{:error, reason}on backend failure
@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{}fromfetch/3reason— 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
2Returns
:ok{:error, :stale_lease}{:error, reason}
Verifies connectivity to the configured backend.
Arguments
opts— only:rheo(instance name) is honoured
Examples
iex> Rheo.ping()
:okReturns
:ok{:error, reason}when the backend is unreachable
@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 nameopts— when the first argument is a stream, filter options (seeRheo.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
@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 nameopts— when the first argument is a stream, filter options (seeRheo.Query); always accepts:rheo(instance name, defaultRheo)
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{}}—eventsmay be empty;next_cursorisnilwhen done{:error, reason}on backend failure
@spec read(stream(), keyword()) :: {:ok, [Rheo.Event.t()]} | {:error, term()}
Reads events by sequence without affecting consumer-group state.
Arguments
stream— stream nameopts— optional keyword list::rheo— instance name (defaultRheo):after— return sequences strictly greater than this (default0):limit— max events (default100):partition— partition (default0)
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
@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{}fromfetch/3reason— 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}
@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{}fromfetch/3opts— optional keyword list::rheo— instance name (defaultRheo):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 updatedexpires_at{:error, :stale_lease}{:error, :backend_unavailable}{:error, reason}
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 namegroup— consumer group nameopts— optional keyword list (exactly one of the replay selectors below, plus optional scope / instance keys)::rheo— instance name (defaultRheo):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:fromfinds no matching event{:error, reason}on backend failure
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 namegroup— consumer group nameopts— optional keyword list::rheo— instance name (defaultRheo):confirm— must betrueor the call returns{:error, :confirm_required}:start_after— exclusive materialization cursor after reset (default0):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
trueReturns
:ok{:error, :confirm_required}whenconfirm: trueis missing{:error, :group_not_found}/{:error, :stream_not_found}{:error, reason}on backend failure
@spec start_link(keyword()) :: Supervisor.on_start()
Starts a Rheo instance supervisor.
Options
:name— instance name (defaultRheo). Also used as the supervisor name.:backend—{module, opts}or a module. Required unless:urlis given, in which caseRheo.Backend.Mongois 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)
trueReturns
{: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.
@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 nameopts— when the first argument is a stream, filter options (seeRheo.Query); always accepts:rheo(instance name, defaultRheo)
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
Enumerableof%Rheo.Event{}(possibly empty)
Errors
Raises RuntimeError when an underlying query_page/2 returns
{:error, reason} (raise "Rheo.stream_query failed: #{inspect(reason)}").