Rheo is an Elixir/OTP library: supervise a Rheo instance, append events to a stream, and consume them with a durable consumer group. Delivery is at-least-once — use event IDs for idempotency.

Install

def deps do
  [{:rheo, "~> 1.0"}]
end

Pick a backend when you start Rheo:

BackendWhen
Rheo.Backend.ETSTests, Livebook, ephemeral apps
{Rheo.Backend.Mongo, url: ...}Durable MongoDB
{Rheo.Backend.Ecto, repo: MyApp.Repo}Durable PostgreSQL or SQLite

Minimal ETS example

children = [
  {Rheo, name: MyRheo, backend: Rheo.Backend.ETS}
]

Supervisor.start_link(children, strategy: :one_for_one)

:ok = Rheo.create_stream("orders", rheo: MyRheo)
:ok = Rheo.create_group("orders", "fulfillment", rheo: MyRheo)

{:ok, _event} =
  Rheo.append("orders", %{type: "order_created", key: "cust-1"}, rheo: MyRheo)

{:ok, [lease]} = Rheo.fetch("orders", "fulfillment", limit: 1, rheo: MyRheo)
:ok = Rheo.ack(lease, rheo: MyRheo)

Idiomatic consumer

defmodule MyApp.FulfillmentConsumer do
  use Rheo.Consumer, stream: "orders", group: "fulfillment"

  @impl true
  def handle_event(event, _context) do
    :ok = MyApp.Fulfillment.process(event)
    :ack
  end
end

children = [
  {Rheo, name: MyRheo, backend: Rheo.Backend.ETS},
  {MyApp.FulfillmentConsumer, rheo: MyRheo, concurrency: 4, max_demand: 50}
]

Next steps