Building your own backend

Copy Markdown View Source

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
  )
end

Guarantees 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

  1. Events are immutable — ACK never deletes log rows.

  2. Leases are fenced — every settle callback fails with {:error, :stale_lease} unless lease_id still matches the current claim. If you set lease.receipt, compare it with term equality and fail with {:error, :receipt_mismatch} when it differs.

  3. At-least-once — expiry or crash may redeliver; do not claim exactly-once.

  4. Partitions (if declared) — per-partition sequences and a contiguous frontier for lag/4.

  5. Errors are portable — map driver failures into the Rheo.Settle vocabulary. Common atoms:

    ReasonMeaning
    :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 :after id absent
    :already_drainingSecond drain/2 while one is pending
    :unsupportedOptional ops callback not implemented
    :backend_unavailableBackend unreachable / process down
    :confirm_requiredreset_group without confirm: 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]
end

The 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).

BackendSeveral BEAM nodes, same group?Notes
RedisYes (distributed: true)Shared Redis; native PEL / reclaim
Ecto PostgreSQLYes (distributed: true)Shared DB; FOR UPDATE SKIP LOCKED
Ecto SQLiteNo (distributed: false)Single-writer
MongoYes (distributed: true)Shared Mongo; compare-and-set claims
ETSNoTables die with the node
MnesiaNo 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

ModuleRole
Rheo.Backend.ETSIn-memory reference (ephemeral)
Rheo.Backend.MnesiaDurable ETS-shaped store (OTP :mnesia, single-node)
Rheo.Backend.MongoDurable document store
Rheo.Backend.EctoDurable SQL via a host-owned Repo
Rheo.Backend.RedisRedis 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.

ADRs: 005, 021, 023, 024, 028, 029.