Powered by AppSignal & Oban Pro

Rheo Livebook Demo (v0.3.0)

notebooks/rheo_demo.livemd

Rheo Livebook Demo (v0.3.0)

Mix.install(
  [
    {:rheo, path: Path.join(__DIR__, "..")},
    {:kino, "~> 0.14"}
  ],
  config: [
    rheo: [
      start_on_application: false,
      clock: Rheo.Clock.Frozen,
      default_lease_ms: 5_000,
      default_max_attempts: 3
    ]
  ]
)

Prerequisites

  1. Open this notebook from the repo (notebooks/rheo_demo.livemd) so Mix.install can resolve {:rheo, path: ...}
  2. Evaluate cells top to bottom
  3. No Docker required — this demo uses Rheo.Backend.ETS by default

Rheo v0.3.0 is an embedded Elixir/OTP library: consumer groups over a searchable event log (Mongo durable, ETS ephemeral). Delivery is at-least-once — use event IDs for idempotency.


1. Start Rheo (ETS)

case Rheo.start_link(backend: Rheo.Backend.ETS) do
  {:ok, _} -> :ok
  {:error, {:already_started, _}} -> :ok
end

:ok = Rheo.ensure_indexes()
:ok = Rheo.ping()

Optional Mongo path: start Docker (docker compose up -d), then Rheo.start_link(url: System.get_env("RHEO_MONGO_URL", "mongodb://localhost:27017/rheo_livebook")) and clear collections before continuing.


2. Create a stream and publish market events

stream = "market-events"
:ok = Rheo.create_stream(stream)

events =
  for i <- 1..30 do
    %{
      type: "curve_update",
      currency: Enum.at(["EUR", "USD", "GBP"], rem(i, 3)),
      curve: Enum.at(["EUR-EURIBOR-6M", "USD-SOFR", "GBP-SONIA"], rem(i, 3)),
      price: Float.round(1.0 + :rand.uniform() * 3.0, 4),
      metadata: %{
        correlation_id: "corr-#{rem(i, 7)}",
        producer: "pricing-service-v3"
      }
    }
  end

{:ok, written} = Rheo.append_batch(stream, events)

Kino.DataTable.new(
  Enum.map(written, fn e ->
    %{
      sequence: e.sequence,
      id: e.id,
      type: e.type,
      currency: e.payload["currency"],
      curve: e.payload["curve"],
      price: e.payload["price"]
    }
  end)
)

3. Independent consumer groups

risk and surveillance each consume the same immutable stream.

:ok = Rheo.create_group(stream, "risk")
:ok = Rheo.create_group(stream, "surveillance")

{:ok, risk_leases} = Rheo.fetch(stream, "risk", limit: 5, consumer_id: "risk-1")
{:ok, surv_leases} = Rheo.fetch(stream, "surveillance", limit: 5, consumer_id: "surv-1")

%{
  risk_count: length(risk_leases),
  surveillance_count: length(surv_leases),
  same_event_ids?: Enum.map(risk_leases, & &1.event_id) == Enum.map(surv_leases, & &1.event_id)
}

Acknowledge surveillance; leave risk leases hanging (simulate a crash before ACK).

Enum.each(surv_leases, &Rheo.ack/1)

stale_risk = risk_leases
:ok

4. Lease expiry and redelivery

Advance the frozen clock past the lease TTL. A different risk worker can claim the same events; the stale lease can no longer ACK.

Rheo.Clock.Frozen.set(DateTime.utc_now())
Rheo.Clock.Frozen.advance(6_000)

{:ok, redelivered} = Rheo.fetch(stream, "risk", limit: 5, consumer_id: "risk-2")

stale_ack = Rheo.ack(hd(stale_risk))
fresh_ack = Rheo.ack(hd(redelivered))

Enum.each(tl(redelivered), &Rheo.ack/1)

%{
  redelivered: length(redelivered),
  stale_ack: stale_ack,
  fresh_ack: fresh_ack,
  attempts: Enum.map(redelivered, & &1.attempt)
}

5. Search history after consumption

Events stay in Mongo after ACK — that is Rheo's differentiator.

{:ok, eur} =
  Rheo.query(stream,
    type: "curve_update",
    currency: "EUR",
    curve: "EUR-EURIBOR-6M",
    limit: 10
  )

{:ok, by_corr} = Rheo.query(stream, correlation_id: "corr-1", limit: 10)

Kino.Layout.grid(
  [
    Kino.DataTable.new(
      Enum.map(eur, fn e ->
        %{sequence: e.sequence, price: e.payload["price"], id: e.id}
      end)
    ),
    Kino.DataTable.new(
      Enum.map(by_corr, fn e ->
        %{
          sequence: e.sequence,
          currency: e.payload["currency"],
          correlation_id: e.metadata["correlation_id"]
        }
      end)
    )
  ],
  columns: 2
)

6. Idiomatic Rheo.Consumer

defmodule Livebook.RiskConsumer do
  use Rheo.Consumer, stream: "market-events", group: "risk"

  @impl true
  def setup(opts) do
    {:ok, %{seen: [], agent: Keyword.fetch!(opts, :agent)}}
  end

  @impl true
  def handle_event(event, %{agent: agent} = state) do
    Agent.update(agent, fn ids -> [event.id | ids] end)
    {:ack, state}
  end
end

{:ok, agent} = Agent.start_link(fn -> [] end)

{:ok, _pid} =
  Livebook.RiskConsumer.start_link(
    stream: stream,
    group: "risk",
    agent: agent,
    name: Livebook.RiskConsumer,
    max_demand: 10,
    poll_ms: 100
  )

Process.sleep(1_500)
seen = Agent.get(agent, &Enum.reverse/1)

%{processed: length(seen), sample_ids: Enum.take(seen, 5)}

Takeaways

  • Immutable log — consumption never deletes events
  • Independent groups — risk and surveillance progress separately
  • Leases + fencing — stale ACKs fail; expired work is redelivered (at-least-once)
  • Search — investigate history with Rheo.query/2 after processing
  • OTP — supervise {Rheo, ...} and Rheo.Consumer in your app

CLI alternative from the repo root: mix rheo.demo