# Rheo Pipelines: GenStage, Flow & Broadway

```elixir
Mix.install(
  [
    # From this repo (path). For HexDocs / standalone Livebook use {:rheo, "~> 1.0"} instead:
    {:rheo, path: Path.join(__DIR__, "..")},
    # {:rheo, "~> 1.0"},
    {:kino, "~> 0.14"},
    {:broadway, "~> 1.3"},
    {:flow, "~> 1.2"}
  ],
  config: [
    rheo: [
      start_on_application: false,
      clock: Rheo.Clock.System,
      default_lease_ms: 30_000,
      default_max_attempts: 5
    ]
  ]
)
```

## Intro

Rheo answers: *“Was this event delivered, leased, and settled for this group?”*

GenStage, Flow, and Broadway answer: *“How do we shape demand, parallelism, and
batching in the BEAM?”*

This notebook keeps that boundary sharp. You will feed the **same**
`%Rheo.Lease{}` through:

1. A plain **GenStage** consumer
2. **Flow** map / reduce / window (with settlement rules)
3. A **Broadway** pipeline

Rheo is not replaced by any of them — and none of them should invent a second
delivery identity.

Sibling notebooks: [Quickstart](quickstart.livemd) · [Concepts](concepts.livemd) ·
[Backends](backends.livemd) · [Index](rheo_demo.livemd)

<!-- livebook:{"break_markdown":true} -->

---

## Setup

We use ETS again, with **two partitions** and a handful of ticks. Each pipeline
demo uses its **own consumer group** so they do not steal work from each other.

```elixir
alias Rheo.{Lease, Producer}

case Process.whereis(Rheo) do
  nil -> :ok
  pid -> GenServer.stop(pid)
end

{:ok, _} = Rheo.start_link(backend: Rheo.Backend.ETS)
:ok = Rheo.ensure_indexes()

stream = "pipeline-events"
:ok = Rheo.create_stream(stream, partition_count: 2)
:ok = Rheo.create_group(stream, "genstage")
:ok = Rheo.create_group(stream, "broadway")

{:ok, _} =
  Rheo.append_batch(
    stream,
    for(i <- 1..12, do: %{type: "tick", n: i, key: "k-#{rem(i, 3)}"})
  )

:ok
```

> **Rule of thumb:** pick **one** consumption surface per
> `{rheo, stream, group}` — either `Rheo.Consumer` **or** `Rheo.Producer`.
> Mixing both on the same group means competing fetchers.

---

## 1. GenStage: demand becomes leases

`Rheo.Producer` is a GenStage producer. Downstream demand turns into
`Rheo.fetch/3`. Emitted events are **leases**, not bare payloads — because
settlement still belongs to Rheo.

After you finish work:

1. Settle durably (`Rheo.ack/1` or `Producer.ack/3`)
2. Free the producer's inflight slot (`confirm/2` is included in `Producer.ack/3`)

Otherwise the producer keeps renewing the lease and will not fetch more once
`:max_demand` is full.

```elixir
defmodule Pipelines.Sink do
  @moduledoc "Tiny GenStage consumer that forwards leases to the Livebook process."
  use GenStage

  def start_link(opts), do: GenStage.start_link(__MODULE__, opts)

  @impl true
  def init(opts) do
    producer = Keyword.fetch!(opts, :producer)
    owner = Keyword.fetch!(opts, :owner)
    {:consumer, owner, subscribe_to: [{producer, max_demand: 4}]}
  end

  @impl true
  def handle_events(leases, _from, owner) do
    Enum.each(leases, &send(owner, {:lease, &1}))
    {:noreply, [], owner}
  end
end

{:ok, producer} =
  Producer.start_link(
    stream: stream,
    group: "genstage",
    max_demand: 4,
    poll_ms: 50
  )

{:ok, sink} = Pipelines.Sink.start_link(producer: producer, owner: self())

genstage_leases =
  for _ <- 1..4 do
    receive do
      {:lease, %Lease{} = lease} -> lease
    after
      3_000 -> raise "timeout waiting for GenStage lease"
    end
  end

Enum.each(genstage_leases, fn lease ->
  # Prefer Producer.ack/3 in real code (settle + confirm together).
  :ok = Rheo.ack(lease)
  :ok = Producer.confirm(producer, lease.lease_id)
end)

%{
  received: length(genstage_leases),
  tick_numbers: Enum.map(genstage_leases, & &1.event.payload["n"]),
  inflight_after_confirm: Producer.inflight_count(producer)
}
```

`inflight_after_confirm` should be `0`. Stop this producer before the Flow
section so we do not leave a polling process around.

In Livebook, stop the **consumer first**, then the producer. If you only stop the
producer, GenStage logs a cancel notice on the sink *after* the cell returns —
Livebook's stdout is already closed, so Logger's Writer dies with `:epipe`. The
stop itself already succeeded (`:normal`); the epipe is I/O noise, not a failed
shutdown.

```elixir
:ok = GenStage.stop(sink)
:ok = GenStage.stop(producer)
```

---

## 2. Flow: when is it safe to settle?

Flow is **not** a Rheo Hex package. You compose it yourself with
`Flow.from_stages/1`. The lease type and settlement API stay the same as
Broadway (and work against Redis the same way — see [Backends](backends.livemd)).

| Stage       | When to ACK                                       |
| ----------- | ------------------------------------------------- |
| `map` (1:1) | After the map side-effect succeeds                |
| `partition` | Still holding the lease; settle downstream        |
| `reduce`    | **After** the aggregate is durable (`on_trigger`) |
| `window`    | **After** the window closes                       |

Crashing before settle ⇒ lease expiry ⇒ at-least-once redelivery. Never ACK
inside `reduce` “because Flow accepted the event”.

### Livebook pitfall (read this)

`Enum.take(n)` only **keeps** `n` results. Upstream stages may still demand and
run side effects for more events. If you `ack` inside `Flow.map` on a shared
group, you can drain the whole group before the next cell runs — and the next
cell appears to “hang forever”.

Each demo below therefore uses:

* a **fresh group**
* a **finite** number of appended events
* a helper that **times out** and stops the producer

```elixir
run_flow! = fn producer, flow, n, timeout ->
  task = Task.async(fn -> flow |> Enum.take(n) end)

  case Task.yield(task, timeout) || Task.shutdown(task, :brutal_kill) do
    {:ok, result} ->
      # Wait for Flow stages to unsubscribe before stopping the producer so
      # GenStage's cancel notice is not logged after Livebook closes the cell I/O
      # (that race shows up as Logger Writer `:epipe`, not a failed stop).
      Process.sleep(50)
      _ = GenStage.stop(producer)
      result

    nil ->
      Process.sleep(50)
      _ = GenStage.stop(producer)
      raise "Flow did not finish in #{timeout}ms"
  end
end

:ok
```

### Map 1:1 — settle per lease

```elixir
:ok = Rheo.create_group(stream, "flow-map")

{:ok, _} =
  Rheo.append_batch(
    stream,
    for(i <- 1..4, do: %{type: "flow-map", n: i, key: "m-#{i}"})
  )

{:ok, flow_map_producer} =
  Producer.start_link(
    stream: stream,
    group: "flow-map",
    max_demand: 4,
    poll_ms: 50
  )

mapped =
  run_flow!.(
    flow_map_producer,
    Flow.from_stages([flow_map_producer])
    |> Flow.map(fn %Lease{} = lease ->
      :ok = Producer.ack(flow_map_producer, lease)
      {lease.event.payload["n"], lease.receipt}
    end),
    4,
    5_000
  )

%{mapped: mapped}
```

### Reduce — settle only on trigger

We keep leases in the reducer accumulator and ACK them in `on_trigger` after a
batch of three. That mimics “materialize the aggregate, then settle inputs”.

```elixir
:ok = Rheo.create_group(stream, "flow-reduce")

{:ok, _} =
  Rheo.append_batch(
    stream,
    for(i <- 1..6, do: %{type: "flow-reduce", n: i, key: "r-#{rem(i, 2)}"})
  )

{:ok, flow_reduce_producer} =
  Producer.start_link(
    stream: stream,
    group: "flow-reduce",
    max_demand: 6,
    poll_ms: 50
  )

window = Flow.Window.global() |> Flow.Window.trigger_every(3)

[batch_size] =
  run_flow!.(
    flow_reduce_producer,
    Flow.from_stages([flow_reduce_producer])
    |> Flow.partition(key: fn _ -> :one end, window: window, stages: 1)
    |> Flow.reduce(fn -> [] end, fn %Lease{} = lease, acc -> [lease | acc] end)
    |> Flow.on_trigger(fn leases ->
      Enum.each(leases, fn lease ->
        :ok = Producer.ack(flow_reduce_producer, lease)
      end)

      {[length(leases)], []}
    end),
    1,
    5_000
  )

%{reduce_batch_size: batch_size}
```

`reduce_batch_size` should be `3` — the trigger size, not the full group.

### Window — two flushes of two

Same idea with `trigger_every(2)` and `Enum.take(2)` so we observe both
window emissions.

```elixir
:ok = Rheo.create_group(stream, "flow-window")

{:ok, _} =
  Rheo.append_batch(
    stream,
    for(i <- 1..4, do: %{type: "flow-window", n: i, key: "same"})
  )

{:ok, flow_window_producer} =
  Producer.start_link(
    stream: stream,
    group: "flow-window",
    max_demand: 4,
    poll_ms: 50
  )

window = Flow.Window.global() |> Flow.Window.trigger_every(2)

batches =
  run_flow!.(
    flow_window_producer,
    Flow.from_stages([flow_window_producer])
    |> Flow.partition(key: fn _ -> :one end, window: window, stages: 1)
    |> Flow.reduce(fn -> [] end, fn lease, acc -> [lease | acc] end)
    |> Flow.on_trigger(fn leases ->
      Enum.each(leases, fn lease ->
        :ok = Producer.ack(flow_window_producer, lease)
      end)

      {[length(leases)], []}
    end),
    2,
    5_000
  )

%{window_batches: batches}
```

Expect `window_batches == [2, 2]`.

---

## 3. Broadway: processors without owning delivery

Broadway is excellent at concurrency, batchers, and rate limits. Rheo remains
the source of truth for leases.

`Rheo.Broadway.transform/2` wraps each lease into a `%Broadway.Message{}`.
Successful messages ACK; failed messages NACK (or reject if configured).

This uses the `"broadway"` group created in setup — it does not compete with
GenStage/Flow groups above.

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

defmodule Pipelines.RiskBroadway do
  use Broadway

  def start_link(opts) do
    Broadway.start_link(__MODULE__,
      name: __MODULE__,
      context: %{collector: Keyword.fetch!(opts, :collector)},
      producer: [
        module:
          {Rheo.Producer,
           stream: Keyword.fetch!(opts, :stream),
           group: "broadway",
           max_demand: 20,
           poll_ms: 50},
        transformer: {Rheo.Broadway, :transform, []},
        concurrency: 1
      ],
      processors: [default: [concurrency: 4]]
    )
  end

  @impl true
  def handle_message(_processor, message, %{collector: collector}) do
    Agent.update(collector, fn seen ->
      [{message.data.partition, message.data.sequence, message.metadata.attempt} | seen]
    end)

    message
  end
end

{:ok, _pipeline} =
  Pipelines.RiskBroadway.start_link(stream: stream, collector: collector)

Process.sleep(1_500)

{:ok, lag} = Rheo.lag(stream, "broadway")

%{
  processed: collector |> Agent.get(&Enum.reverse/1) |> length(),
  sample: collector |> Agent.get(&Enum.reverse/1) |> Enum.take(5),
  frontiers: Map.new(lag.partitions, fn {p, info} -> {p, info.frontier} end),
  lag: lag.lag
}
```

### Failure modes (for your handlers)

```elixir
# Broadway.Message.failed(message, :reason)
#   → Rheo.nack/3  (retry until the group's max_attempts, then dead-letter)
#
# Broadway.Message.configure_ack(message, on_failure: :reject)
#   → Rheo.reject/3 (dead-letter immediately)
:ok = Broadway.stop(Pipelines.RiskBroadway)
```

Two knobs people confuse:

* **`:max_demand`** on the producer — caps unsettled **leases**
* **Broadway `:concurrency`** — caps parallel **handler** work

---

## Choosing a surface

| Surface                    | Use when                                   |
| -------------------------- | ------------------------------------------ |
| `Rheo.Consumer`            | Simple OTP handlers, most apps             |
| `Rheo.Producer` + GenStage | Custom demand topology                     |
| `Rheo.Producer` + Flow     | Parallel map / partition / reduce / window |
| `Rheo.Producer` + Broadway | Processors, batchers, rate limits          |

Rheo keeps leases, fencing, frontier, and receipts. Pipelines must not invent a
second delivery identity.

Next: [Backends](backends.livemd)

HexDocs: [GenStage](https://hexdocs.pm/rheo/genstage.html) ·
[Broadway](https://hexdocs.pm/rheo/broadway.html) ·
[Flow spike](https://github.com/thanos/rheo/blob/main/docs/design/flow-readiness-spike.md) ·
[Article 14](https://github.com/thanos/rheo/blob/main/docs/tutorials/14-rheo-is-not-broadway-it-feeds-broadway.md) ·
[Article 16](https://github.com/thanos/rheo/blob/main/docs/tutorials/16-rheo-on-redis-streams-portable-sequence-native-pel.md)
