Backends implement Rheo.Backend. Application code calls Rheo; the instance
dispatches to your adapter through an opaque handle.
What you must implement
See the @callback docs on Rheo.Backend:
- Lifecycle and health:
child_spec/1,capabilities/0,ensure_indexes/1,ping/1 - Streams and groups:
create_stream/3,create_group/4 - Log:
append/4,append_batch/4,read/3,query/2 - Delivery:
fetch/4,renew/3,ack/2,retry/3,reject/3 - Group progress:
replay/4,reset_group/4,lag/4
The callbacks name semantic operations. How you store deliveries is private: a
SQL table with SKIP LOCKED, a document collection with compare-and-set, or a
native consumer group such as Redis Streams' PEL.
Capabilities
Return a validated Rheo.Backend.Capabilities struct:
@impl true
def capabilities do
Rheo.Backend.Capabilities.new(
durable: true,
distributed: true,
partitions: true,
contiguous_frontier: true,
replay: true,
atomic_compare_and_set: true,
ordered_range_scan: true
)
endGuarantees gate conformance cases; mechanisms describe how you implement
delivery. Declare yes to :distributed when several BEAM nodes may fetch
the same group through a shared store. at_least_once and lease_fencing
cannot be declared false.
Invariants
Events are immutable — ACK never deletes log rows.
Leases are fenced — every settle callback fails with
{:error, :stale_lease}unlesslease_idstill matches the current claim. If you setlease.receipt, compare it with term equality and fail with{:error, :receipt_mismatch}when it differs.At-least-once — expiry or crash may redeliver; do not claim exactly-once.
Partitions (if declared) — per-partition sequences and a contiguous frontier for
lag/4.Errors are portable — map driver failures into the
Rheo.Settlevocabulary. Common atoms:Reason Meaning :stream_not_foundStream missing (looked up before group) :group_not_foundStream exists; group does not :already_existsDuplicate stream or group :stale_leaseSettle lost the fencing token :receipt_mismatchNative receipt no longer matches :cursor_not_foundDead-letter / page :afterid absent:already_drainingSecond drain/2while one is pending:unsupportedOptional ops callback not implemented :backend_unavailableBackend unreachable / process down :confirm_requiredreset_groupwithoutconfirm: true{:ambiguous, cause}Write may have committed (e.g. timeout) {:failed, cause}Definite failure Never return driver exception structs.
Conformance
Run the shared contract against your adapter:
defmodule MyBackendContractTest do
use ExUnit.Case, async: false
use Rheo.BackendContract, backend: MyBackend, backend_opts: [name: :my_backend]
endThe suite (test/support/backend_contract.ex, ADR 012 / ADR 024 / ADR 029) is
the executable freeze for adapters. It is grouped by guarantee: lifecycle,
unavailable handles, event log, queries, consumer groups, leases and fencing,
retry and reject, replay, partitions and frontier. Correctness cases always
run; only cases for guarantees you do not declare are skipped. Shipping a
backend without the contract suite is unsupported.
Multi-node (several BEAM nodes)
Node = one BEAM VM. Rheo itself does not form a cluster: each node runs its
own Groups; a shared durable backend arbitrates leases when several nodes
fetch the same group (capabilities.guarantees.distributed).
| Backend | Several BEAM nodes, same group? | Notes |
|---|---|---|
| Redis | Yes (distributed: true) | Shared Redis; native PEL / reclaim |
| Ecto PostgreSQL | Yes (distributed: true) | Shared DB; FOR UPDATE SKIP LOCKED |
| Ecto SQLite | No (distributed: false) | Single-writer |
| Mongo | Yes (distributed: true) | Shared Mongo; compare-and-set claims |
| ETS | No | Tables die with the node |
| Mnesia | No in v0.11 (distributed: false) | Single-node disc_copies; multi-node table copies later — still not a Rheo control plane |
That is not the same as multiple partitions, multiple consumers on one node,
or several named {Rheo, name: …} instances on one node (all already supported).
Reference implementations
| Module | Role |
|---|---|
Rheo.Backend.ETS | In-memory reference (ephemeral) |
Rheo.Backend.Mnesia | Durable ETS-shaped store (OTP :mnesia, single-node) |
Rheo.Backend.Mongo | Durable document store |
Rheo.Backend.Ecto | Durable SQL via a host-owned Repo |
Rheo.Backend.Redis | Redis Streams (optional redix) |
Rheo.Backend.NativeStreamDouble (test support) | Native-stream shape: receipts, native reclaim |
Start from ETS to learn the state machine. For a native-stream store, start
from the double: it keeps the portable event.sequence and fences on both
lease_id and receipt.