Rheo.Backend behaviour (rheo v1.0.0)

Copy Markdown View Source

Behaviour for Rheo storage backends.

A backend owns storage for streams, events, groups, and deliveries behind an opaque handle (a process name, pid, table reference, repo, …). Application code calls Rheo, which resolves the handle for the named instance and dispatches here.

Part of the SemVer-frozen surface (ADR 029 / ADR 030). Required callbacks and the fencing / settle vocabulary will not change meaning without a major version. Ops inspect callbacks remain optional.

Architecture

+------------------+
|  Application     |
|  (Rheo facade)   |
+--------+---------+
         |
         | resolve name -> {backend, handle}
         v
+--------+---------+       +---------------------------+
| Rheo.Instance    |------>| Rheo.Backend callbacks    |
| (supervisor)     |       | child_spec / append / …   |
+------------------+       +-------------+-------------+
                                         |
               +-------------------------+-------------------------+
               |                         |                         |
               v                         v                         v
        +------------+            +------------+            +------------+
        | ETS/Mnesia |            | Mongo/Ecto |            | Redis /    |
        | (tables)   |            | (docs/SQL) |            | custom     |
        +------------+            +------------+            +------------+

Delivery rows, SQL locking, and native consumer-group commands stay inside the adapter. A backend built on Redis Streams may implement fetch with XREADGROUP, ack with a fenced XACK, and lag with XINFO GROUPS without pretending to keep a deliveries table.

Semantic contract

Callbacks describe what Rheo needs, not how storage does it:

  • 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
  • ops inspect (optional) — list_streams/2, list_groups/3, dead_letters/4, group_info/4 (ADR 027)

Fencing and receipts

fetch/4 returns leases carrying a unique lease_id. Every settle callback (renew/3, ack/2, retry/3, reject/3) must fail with {:error, :stale_lease} when the delivery is no longer held under that lease_id. A backend may also set lease.receipt to an opaque native claim identity (ADR 021); if it does, settle callbacks must compare it with term equality and fail with {:error, :receipt_mismatch} when it differs. Database backends usually leave receipt equal to lease_id.

Error vocabulary

Map driver failures into portable atoms at the adapter boundary. Do not return driver exception structs directly.

  • :stream_not_found — stream missing (checked before group)
  • :group_not_found — stream exists; group does not
  • :already_exists — duplicate stream or group
  • :stale_lease — settle lost the fencing token
  • :receipt_mismatch — native receipt no longer matches
  • :unsupported — optional ops callback not implemented
  • :backend_unavailable — backend unreachable / process down
  • {:failed, cause} — definite failure
  • {:ambiguous, cause} — outcome unknown (timeout, disconnect mid-write)

See Rheo.Settle for how Group / Producer classify settle errors.

Capabilities

capabilities/0 returns a Rheo.Backend.Capabilities struct. The conformance suite (test/support/backend_contract.ex) gates optional cases on declared guarantees; it never skips fencing or at-least-once checks.

Implementing a backend

Prefer the Rheo facade in application code. A minimal adapter is a GenServer (or borrowed connection) that implements the required callbacks and declares honest capabilities. Abbreviated ETS-style sketch:

defmodule MyApp.MemoryBackend do
  @behaviour Rheo.Backend
  use GenServer

  @impl true
  def capabilities do
    Rheo.Backend.Capabilities.new(%{
      durable: false,
      distributed: false,
      batch_writes: true,
      ordered_range_scan: true,
      replay: true,
      partitions: true,
      contiguous_frontier: true
    })
  end

  @impl true
  def child_spec(opts) do
    name = Keyword.get(opts, :name, __MODULE__)

    %{
      id: {__MODULE__, name},
      start: {__MODULE__, :start_link, [opts]},
      type: :worker,
      restart: :permanent
    }
  end

  @impl true
  def ping(handle), do: GenServer.call(handle, :ping)

  @impl true
  def ensure_indexes(_handle), do: :ok

  @impl true
  def create_stream(handle, stream, opts),
    do: GenServer.call(handle, {:create_stream, stream, opts})

  # create_group, append, append_batch, read, query,
  # fetch, renew, ack, retry, reject, replay, reset_group, lag …
  #
  # Settle must return {:error, :stale_lease} when lease_id no longer
  # matches. Map driver failures to :backend_unavailable,
  # {:failed, cause}, or {:ambiguous, cause}.
end

See the building your own backend guide and run the shared contract suite against the adapter before shipping.

Summary

Types

Consumer group name.

Opaque backend connection handle (process name, pid, table ref, …).

Backend-specific options (keyword list).

Stream name.

Callbacks

Acknowledges a lease when lease_id (and receipt) still match.

Appends one event to a stream, allocating the next per-partition sequence.

Appends many events with contiguous per-partition sequences.

Returns the backend's static capability declaration.

Returns a child spec that starts backend resources under a Rheo instance.

Registers a consumer group on an existing stream.

Creates a stream registry entry and its per-partition sequence allocators.

Lists dead-lettered deliveries for a group (ops inspect, ADR 027).

Creates indexes / schema needed for correct operation.

Claims up to opts[:limit] events for a group, returning fenced leases.

Returns lag plus inflight and dead-letter counts (ops inspect, ADR 027).

Returns committed frontier vs high-watermark lag for a group.

Lists consumer group names for a stream (ops inspect, ADR 027).

Lists registered stream names (ops inspect, ADR 027).

Health-checks the backend connection.

Queries historical events using a portable Rheo.Query.

Reads events by sequence without mutating consumer-group state.

Dead-letters a lease for this group only.

Extends an active lease's expiry when lease_id (and receipt) still match.

Re-opens deliveries for a group for replay without copying events.

Clears deliveries for a group and resets its materialization cursor.

Marks a lease for redelivery, or dead-letters it at the group's max attempts.

Types

group()

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

Consumer group name.

handle()

@type handle() :: term()

Opaque backend connection handle (process name, pid, table ref, …).

opts()

@type opts() :: keyword()

Backend-specific options (keyword list).

stream()

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

Stream name.

Callbacks

ack(handle, t)

@callback ack(handle(), Rheo.Lease.t()) :: :ok | {:error, term()}

Acknowledges a lease when lease_id (and receipt) still match.

Durable for that consumer group only. Never deletes immutable log events.

Returns

  • :ok
  • {:error, :stale_lease}
  • {:error, :receipt_mismatch}
  • {:error, :backend_unavailable}
  • {:error, {:failed, cause}} / {:error, {:ambiguous, cause}}
  • {:error, term()}

append(handle, stream, map, opts)

@callback append(handle(), stream(), map(), opts()) ::
  {:ok, Rheo.Event.t()} | {:error, term()}

Appends one event to a stream, allocating the next per-partition sequence.

Options

  • :id — explicit event id (default generated)
  • :key — routing / payload key
  • :metadata — extra metadata map
  • :timestamp — DateTime.t()
  • :partition — explicit partition (overrides key routing)

Returns

  • {:ok, Event.t()}
  • {:error, :stream_not_found}
  • {:error, :backend_unavailable}
  • {:error, {:failed, cause}} / {:error, {:ambiguous, cause}}
  • {:error, term()}

append_batch(handle, stream, list, opts)

@callback append_batch(handle(), stream(), [map()], opts()) ::
  {:ok, [Rheo.Event.t()]} | {:error, term()}

Appends many events with contiguous per-partition sequences.

Shared options apply to each payload (same keys as append/4).

Returns

  • {:ok, [Event.t()]} — possibly empty
  • {:error, :stream_not_found}
  • {:error, :backend_unavailable}
  • {:error, term()}

capabilities()

@callback capabilities() :: Rheo.Backend.Capabilities.t()

Returns the backend's static capability declaration.

Returns

A Rheo.Backend.Capabilities.t(). Guarantees gate conformance cases; mechanisms select optional optimizations. at_least_once and lease_fencing cannot be declared false.

child_spec(opts)

@callback child_spec(opts()) :: Supervisor.child_spec()

Returns a child spec that starts backend resources under a Rheo instance.

Backends that borrow host-owned resources (an Ecto.Repo) still return a small configuration process so the instance has a handle.

Returns

A Supervisor.child_spec/0 map. Started by Rheo under the instance supervisor; the registered name (or equivalent) becomes the opaque handle.

create_group(handle, stream, group, opts)

@callback create_group(handle(), stream(), group(), opts()) :: :ok | {:error, term()}

Registers a consumer group on an existing stream.

Options

  • :max_attempts — attempts before dead-letter on retry (default from env)
  • :start_after — exclusive sequence (integer or %{partition => sequence})
  • :start_at — DateTime; begin at the first event at/after this time
  • :partition / :partitions — limit start cursors (default all)

Returns

  • :ok
  • {:error, :stream_not_found}
  • {:error, :already_exists}
  • {:error, :backend_unavailable}
  • {:error, term()}

create_stream(handle, stream, opts)

@callback create_stream(handle(), stream(), opts()) :: :ok | {:error, term()}

Creates a stream registry entry and its per-partition sequence allocators.

Options

  • :partition_count — number of partitions (default 1, must be >= 1)

Returns

  • :ok
  • {:error, :already_exists}
  • {:error, :invalid_partition_count}
  • {:error, :backend_unavailable}
  • {:error, term()}

dead_letters(handle, stream, group, opts)

(optional)
@callback dead_letters(handle(), stream(), group(), opts()) ::
  {:ok, [Rheo.DeadLetter.t()]} | {:error, term()}

Lists dead-lettered deliveries for a group (ops inspect, ADR 027).

Options

  • :limit — max rows (default 100)
  • :after — event id cursor for pagination

Optional — when unimplemented, Rheo.dead_letters/3 returns {:error, :unsupported}.

Returns

  • {:ok, [DeadLetter.t()]}
  • {:error, :group_not_found}
  • {:error, :cursor_not_found}
  • {:error, :unsupported}
  • {:error, :backend_unavailable}
  • {:error, term()}

ensure_indexes(handle)

@callback ensure_indexes(handle()) :: :ok | {:error, term()}

Creates indexes / schema needed for correct operation.

Must be idempotent. Called from Rheo.ensure_indexes/1.

Returns

  • :ok
  • {:error, :backend_unavailable}
  • {:error, {:failed, cause}} / {:error, {:ambiguous, cause}}
  • {:error, term()}

fetch(handle, stream, group, opts)

@callback fetch(handle(), stream(), group(), opts()) ::
  {:ok, [Rheo.Lease.t()]} | {:error, term()}

Claims up to opts[:limit] events for a group, returning fenced leases.

Options

  • :limit — max leases (default from env / default_max_demand)
  • :consumer_id — claimant identity (default generated)
  • :lease_ms — lease TTL in milliseconds
  • :partition / :partitions — claim only those partitions (default all)

Returns

  • {:ok, [Lease.t()]} — empty when nothing is claimable
  • {:error, :stream_not_found}
  • {:error, :group_not_found}
  • {:error, :backend_unavailable}
  • {:error, term()}

group_info(handle, stream, group, opts)

(optional)
@callback group_info(handle(), stream(), group(), opts()) ::
  {:ok, Rheo.GroupInfo.t()} | {:error, term()}

Returns lag plus inflight and dead-letter counts (ops inspect, ADR 027).

Options

  • :partition / :partitions — scope lag (default all)

Optional — when unimplemented, Rheo.group_info/3 returns {:error, :unsupported}.

Returns

  • {:ok, GroupInfo.t()}
  • {:error, :stream_not_found}
  • {:error, :group_not_found}
  • {:error, :unsupported}
  • {:error, :backend_unavailable}
  • {:error, term()}

lag(handle, stream, group, opts)

@callback lag(handle(), stream(), group(), opts()) ::
  {:ok, Rheo.Lag.t()} | {:error, term()}

Returns committed frontier vs high-watermark lag for a group.

Options

  • :partition / :partitions — report lag for those partitions (default all)

Returns

  • {:ok, Rheo.Lag.t()}
  • {:error, :stream_not_found}
  • {:error, :group_not_found}
  • {:error, :backend_unavailable}
  • {:error, term()}

list_groups(handle, stream, opts)

(optional)
@callback list_groups(handle(), stream(), opts()) :: {:ok, [group()]} | {:error, term()}

Lists consumer group names for a stream (ops inspect, ADR 027).

Optional — when unimplemented, Rheo.list_groups/2 returns {:error, :unsupported}.

Returns

  • {:ok, [group()]}
  • {:error, :stream_not_found}
  • {:error, :unsupported}
  • {:error, :backend_unavailable}
  • {:error, term()}

list_streams(handle, opts)

(optional)
@callback list_streams(handle(), opts()) :: {:ok, [stream()]} | {:error, term()}

Lists registered stream names (ops inspect, ADR 027).

Optional — when unimplemented, Rheo.list_streams/1 returns {:error, :unsupported}.

Returns

  • {:ok, [stream()]}
  • {:error, :unsupported}
  • {:error, :backend_unavailable}
  • {:error, term()}

ping(handle)

@callback ping(handle()) :: :ok | {:error, term()}

Health-checks the backend connection.

Returns

  • :ok
  • {:error, :backend_unavailable}
  • {:error, term()}

query(handle, t)

@callback query(handle(), Rheo.Query.t()) :: {:ok, [Rheo.Event.t()]} | {:error, term()}

Queries historical events using a portable Rheo.Query.

Filtering may use secondary indexes when declared, or walk and filter in the adapter when secondary_indexes: false.

Returns

  • {:ok, [Event.t()]}
  • {:error, :backend_unavailable}
  • {:error, term()}

read(handle, stream, opts)

@callback read(handle(), stream(), opts()) :: {:ok, [Rheo.Event.t()]} | {:error, term()}

Reads events by sequence without mutating consumer-group state.

Options

  • :after — return sequences strictly greater than this (default 0)
  • :limit — max events (default 100)
  • :partition — partition (default 0)

Returns

  • {:ok, [Event.t()]} — empty list when nothing matches
  • {:error, :stream_not_found} when the backend distinguishes missing streams
  • {:error, :backend_unavailable}
  • {:error, term()}

reject(handle, t, term)

@callback reject(handle(), Rheo.Lease.t(), term()) :: :ok | {:error, term()}

Dead-letters a lease for this group only.

Must fence on lease_id (and receipt). Does not delete the event from the log.

Returns

  • :ok
  • {:error, :stale_lease}
  • {:error, :receipt_mismatch}
  • {:error, :backend_unavailable}
  • {:error, term()}

renew(handle, t, opts)

@callback renew(handle(), Rheo.Lease.t(), opts()) ::
  {:ok, Rheo.Lease.t()} | {:error, term()}

Extends an active lease's expiry when lease_id (and receipt) still match.

Options

  • :lease_ms — new TTL from now (default from config)

Returns

  • {:ok, Lease.t()} with updated expires_at
  • {:error, :stale_lease}
  • {:error, :receipt_mismatch}
  • {:error, :backend_unavailable}
  • {:error, term()}

replay(handle, stream, group, opts)

@callback replay(handle(), stream(), group(), opts()) :: :ok | {:error, term()}

Re-opens deliveries for a group for replay without copying events.

Options

Backends typically accept one of:

  • :from_sequence — reopen from this exclusive sequence onward
  • :event_ids — reopen specific event ids
  • :events — reopen from concrete %Rheo.Event{} structs
  • :partition / :partitions — scope the rewind

The Rheo.replay/3 facade may translate :from (DateTime) or :query before calling this callback.

Returns

  • :ok
  • {:error, :stream_not_found}
  • {:error, :group_not_found}
  • {:error, :invalid_replay_opts}
  • {:error, :backend_unavailable}
  • {:error, term()}

reset_group(handle, stream, group, opts)

@callback reset_group(handle(), stream(), group(), opts()) :: :ok | {:error, term()}

Clears deliveries for a group and resets its materialization cursor.

Never deletes immutable events. Rheo.reset_group/3 requires confirm: true before calling this.

Options

  • :start_after / :start_at — new start cursors (same as create_group/4)
  • :partition / :partitions — scope the reset (default all)

Returns

  • :ok
  • {:error, :stream_not_found}
  • {:error, :group_not_found}
  • {:error, :backend_unavailable}
  • {:error, term()}

retry(handle, t, term)

@callback retry(handle(), Rheo.Lease.t(), term()) :: :ok | {:error, term()}

Marks a lease for redelivery, or dead-letters it at the group's max attempts.

Must fence on lease_id (and receipt). When lease.attempt >= max_attempts, dead-letter instead of reopening.

Returns

  • :ok
  • {:error, :stale_lease}
  • {:error, :receipt_mismatch}
  • {:error, :group_not_found}
  • {:error, :backend_unavailable}
  • {:error, term()}