Current section
Files
Jump to
Current section
Files
README.md
<p align="center">
<img src="guides/assets/baton-logo.svg" alt="Baton" width="320">
</p>
<p align="center">
<a href="https://hex.pm/packages/baton"><img src="https://img.shields.io/hexpm/v/baton.svg" alt="Hex.pm"></a>
<a href="https://hexdocs.pm/baton"><img src="https://img.shields.io/badge/hex-docs-lightgreen.svg" alt="Documentation"></a>
<a href="https://github.com/RudeWalrus/Baton/actions/workflows/ci.yml"><img src="https://github.com/RudeWalrus/Baton/actions/workflows/ci.yml/badge.svg" alt="CI"></a>
<a href="https://github.com/RudeWalrus/Baton/blob/main/LICENSE"><img src="https://img.shields.io/badge/license-MIT-blue.svg" alt="License"></a>
</p>
DAG-based job workflows for [Oban](https://github.com/oban-bg/oban): dependency
ordering, fan-out/fan-in (over a fixed collection or one computed mid-run),
result passing, retry idempotency, and per-step LLM cost tracking — built
entirely on Oban OSS, no Oban Pro required.
## Features
- **Directed acyclic graphs** of Oban jobs with named dependencies, validated
for cycles before insertion (Kahn's algorithm). Build one in code with
`Baton.new/1` and `Baton.add/4`, or compile one from data — see below.
- **Workflow templates** — a flow can be a portable, JSON-serializable
*definition* instead of code: nodes, dependencies, prompts, and bindings as
inert data naming no host module. Store it in a table, in Git, or in another
service; version it, diff it, or put it behind an editor, then compile it per
run with that run's input.
- **A flow compiler** ([guide](guides/flow_compiler.md)) that turns a
definition plus a run's input and context into an executable workflow —
validating node types, bindings, and cycles, expanding fan-outs, and storing
an immutable snapshot of exactly what ran. Everything live (prompts, actions,
adapters, guards) resolves through host-supplied seams, so the persisted
artifact never contains executable code. Steps can also be *seeded* with
results supplied up front, so a definition's tail runs on its own — which is
what makes a single node testable in isolation.
- **Self-gating execution** — each job checks its dependencies at runtime and
snoozes, proceeds, or cancels accordingly. No external scheduler.
- **Completion-triggered rescheduling** so downstream steps start promptly
instead of waiting out a snooze timer.
- **Result passing** between steps, stored in the engine's own table (never in
`oban_jobs.meta`).
- **Fan-out and fan-in** — a node expands into one job per item in a
collection, and a node depending on it reads the whole expansion back as an
ordered list. Two gates: `parallel` runs the expansion concurrently,
`sequential` chains it to prime a shared prompt cache or smooth a rate limit.
- **Dynamic fan-out** — when the collection is a *result* rather than an
input — a list an LLM just extracted, whose length nothing could know in
advance — the node creates its own expansion mid-run and waits for it. The
alternative is handing a whole variable-length list to one model call, which
is where degeneration and runaway token spend come from. A reader cannot
tell the two kinds apart, and a node can depend on one of each.
- **Retry idempotency** — a retried step that already produced a result returns
it without re-running side effects (important for paid LLM calls).
- **Multi-model fan-out** — run the same step across several models and
synthesize the results.
- **Structured LLM steps** (`Baton.LLMStep`) — implement `request/1` and
`handle_response/3`; the engine owns the call, the error taxonomy
(429/529 → snooze, `max_tokens` → retryable), JSON decoding, and usage
recording, so a step is a few lines instead of a hand-written transport loop.
- **Batch mode** — `use Baton.LLMStep, mode: :batch` moves a step onto the
provider's Message Batches API for roughly half the token cost, with no
change to the step's callbacks. The engine submits, polls attempt-free, and
ingests the result through the ordinary pipeline.
- **Observability** — telemetry for every transition, optional per-step token/
cost stats, optional full context-window capture, and live step events over
`Phoenix.PubSub` for building a LiveView dashboard.
## Installation
```elixir
def deps do
[
{:baton, "~> 0.1"},
{:oban, "~> 2.17"}
]
end
```
Add the schema via a migration:
```elixir
defmodule MyApp.Repo.Migrations.AddBaton do
use Ecto.Migration
# Omit :version to install the latest schema. The migration is idempotent
# (create_if_not_exists), so to upgrade an existing install you can ship a new
# migration that simply calls Baton.Migration.up/0 again.
def up, do: Baton.Migration.up()
def down, do: Baton.Migration.down()
end
```
Configure (the repo is inherited from Oban automatically):
```elixir
config :baton,
oban_name: Oban,
pubsub: MyApp.PubSub, # only for live events/dashboard
pricing: MyApp.LLMPricing # only if tracking cost
config :my_app, Oban,
plugins: [
Oban.Plugins.Pruner,
{Baton.Plugin, interval: :timer.seconds(60)}
],
queues: [default: 20]
```
### Data retention
Baton's tables (`workflow_nodes`, `workflow_step_stats`,
`workflow_debug_logs`, `workflow_completions`) have **no foreign key** to
`oban_jobs`, so Oban's `Pruner` does not clean them up — left alone they grow
without bound. Enable pruning on `Baton.Plugin` to delete Baton rows
once their backing Oban job has been pruned:
```elixir
{Baton.Plugin,
interval: :timer.seconds(60),
prune: true, # off by default
debug_log_max_age: 24 * 60 * 60} # optional: cap debug logs at 24h (seconds)
```
This piggybacks on Oban's `Pruner`, so there's a single retention policy. For
this to be safe, the `Pruner`'s `max_age` must exceed your longest workflow's
runtime (which Baton already requires for correct dependency gating) — set
it generously, e.g. `{Oban.Plugins.Pruner, max_age: 60 * 60 * 24}`.
## Upgrading
### To schema v9 — dynamic fan-out
v9 adds `workflow_nodes.fan_out_of` and `.item_index`, which identify and order
the nodes a [dynamic fan-out](guides/flow_compiler.md) creates while a workflow
is running. It also **replaces** v1's plain index on
`(workflow_id, step_name)` with a unique one — step names were unique only by
in-memory check in `Baton.add/4`, which cannot arbitrate inserts made mid-run
from a running job. Ship a migration calling `Baton.Migration.up/0` again (it is
idempotent):
```elixir
defmodule MyApp.Repo.Migrations.UpgradeBatonV9 do
use Ecto.Migration
def up, do: Baton.Migration.up()
def down, do: Baton.Migration.down(version: 8)
end
```
The columns are nullable and additive, so existing rows and every static
fan-out leave them `nil`. **The index swap is the part with an operational
cost.** Neither index is built `concurrently`, so on a large existing
`workflow_nodes` table the unique index build takes a lock that blocks writes —
run it during a quiet window, or pre-create it yourself and let v9 no-op:
```elixir
defmodule MyApp.Repo.Migrations.BatonUniqueStepNameConcurrently do
use Ecto.Migration
@disable_ddl_transaction true
@timeout :infinity
def up do
create unique_index(:workflow_nodes, [:workflow_id, :step_name],
name: :baton_workflow_nodes_workflow_id_step_name_index,
concurrently: true
)
end
def down, do: :ok
end
```
v9 uses `create_if_not_exists`, so it will skip an index you already built.
Existing data cannot violate uniqueness — a workflow insert is all-or-nothing —
but if the build does fail, this finds the offenders:
```sql
SELECT workflow_id, step_name, count(*)
FROM workflow_nodes GROUP BY 1, 2 HAVING count(*) > 1;
```
**What you do _not_ need to change:** worker code, existing definitions, or
configuration. A fan-out collection rooted at `$input.`/`$context.` compiles
exactly as before — the dynamic form (`$steps.`) was a validation error prior
to v9, so no stored definition changes meaning.
> #### One new event state {: .warning}
>
> A node that expands itself broadcasts `state: "expanded"` with
> `detail: %{count: n}` — the one event meaning the graph gained steps. If you
> pattern-match on the `state` field of `{:workflow_step_updated, _}` events,
> make sure you have a catch-all clause.
### To schema v7 — batch mode checkpoints
v7 adds `workflow_nodes.checkpoint`, where a step records progress that has to
survive between attempts — today, the provider batch id a
[batch mode](guides/batch_mode.md) step polls for. Ship a migration calling
`Baton.Migration.up/0` again (it is idempotent):
```elixir
defmodule MyApp.Repo.Migrations.UpgradeBatonV7 do
use Ecto.Migration
def up, do: Baton.Migration.up()
def down, do: Baton.Migration.down(version: 6)
end
```
**What you do _not_ need to change:** anything. The column is written only by
steps that opt into batch mode, and existing workers are untouched.
If you are also coming from before v6, that version added
`workflow_nodes.sequence_after` for the `sequential` fan-out gate's ordering
edge — the same idempotent `up/0` installs both.
> #### One new event state {: .warning}
>
> Batch steps broadcast `state: "awaiting"` while parked on a provider batch.
> If you pattern-match on the `state` field of `{:workflow_step_updated, _}`
> events, make sure you have a catch-all clause.
### To schema v4 — workflow-id index on `oban_jobs`
v4 adds a partial expression index on `oban_jobs ((meta->>'workflow_id'))` that
backs `Baton.Plugin`'s failed-workflow detection and orphan scan, keeping those
sweeps cheap as `oban_jobs` grows. **You must run a migration** — ship one that
calls `Baton.Migration.up/0` again (it is idempotent):
```elixir
defmodule MyApp.Repo.Migrations.UpgradeBatonV4 do
use Ecto.Migration
def up, do: Baton.Migration.up() # creates the workflow_id index (v4)
def down, do: Baton.Migration.down(version: 3)
end
```
Then `mix ecto.migrate`. The index is created without `concurrently`, so on a
very large existing `oban_jobs` table run it during a quiet window (or add
`@disable_ddl_transaction true` + `concurrently` in your own migration).
**What you do _not_ need to change:** worker code, configuration, or the plugin
setup — the index is used automatically. Fast dispatch and crash detection are
attached by `Baton.Application` and need no wiring.
### To schema v3 — large-result spilling
Step results larger than `inline_threshold_bytes` (default 32 KB) are now
gzipped and stored in a new `workflow_artifacts` table instead of inline on
`workflow_nodes`, keeping the hot dependency-gating table small. **You must run
a migration** to add the table — ship one that calls `Baton.Migration.up/0`
again (it is idempotent):
```elixir
defmodule MyApp.Repo.Migrations.UpgradeBaton do
use Ecto.Migration
def up, do: Baton.Migration.up() # creates workflow_artifacts (v3)
def down, do: Baton.Migration.down(version: 2)
end
```
Then `mix ecto.migrate`. Until the table exists, any result above the inline
threshold fails to store and its step retries.
**What you do _not_ need to change:**
- **Worker code is untouched.** `Baton.Results` (`store_result/2`,
`get_result/2`, `get_all_results/1`, `get_own_result/1`) keeps the same API
and semantics; large results resolve transparently. Your `perform_workflow/1`
functions don't change.
- **No data backfill.** Results already stored inline keep reading correctly.
- **No new dependencies.**
**Behavioural changes to be aware of:**
- Results whose encoded size exceeds `max_result_bytes` (default 16 MB) are now
rejected with `{:error, :result_too_large}`, failing the step. Previously they
were stored inline. If you legitimately emit larger results, raise the limit.
- A step whose result cannot be persisted now fails and retries instead of
completing silently (which previously left downstream steps waiting forever).
- `Baton.Retention.delete_orphans/2` and `delete_workflow/2` count maps gain a
`:workflow_artifacts` key — only relevant if you match the exact map shape.
- If you read `workflow_nodes.result` **directly** (bypassing `Baton.Results`),
large results now appear as a reference (`%{"__baton_artifact__" => …}`)
rather than the data — read through `Baton.Results` instead.
### New configuration (all optional, sensible defaults)
```elixir
config :baton,
# Result tiering
inline_threshold_bytes: 32_768, # spill above this (default 32 KB)
max_result_bytes: 16_777_216, # reject above this (default 16 MB)
result_store: Baton.ResultStore.Postgres, # large-result backend (default)
# Optional node-local read cache (off by default)
result_cache_enabled: false,
max_cache_bytes: 67_108_864 # 64 MB budget
```
The read cache skips the backend round-trip and gunzip/decode when a step reads
the same large upstream result more than once (fan-in / multi-model synthesis).
It is node-local and safe to drop — a miss only costs a cold read — so enabling
it never affects correctness.
## Usage
```elixir
defmodule MyApp.Steps.Fetch do
use Baton.Worker, queue: :default
@impl true
def perform_workflow(%Oban.Job{args: %{"url" => url}}) do
{:ok, %{body: fetch(url)}}
end
end
Baton.new(workflow_name: "ingest")
|> Baton.add(:fetch, MyApp.Steps.Fetch.new(%{url: "https://example.com"}))
|> Baton.add(:parse, MyApp.Steps.Parse.new(%{}), deps: [:fetch])
|> Baton.add(:store, MyApp.Steps.Store.new(%{}), deps: [:parse])
|> Baton.insert!()
```
See the [getting started guide](guides/getting_started.md), the
[building a workflow guide](guides/building_a_workflow.md) (fan-out/fan-in,
pruning, and a live LiveView), and the [multi-model guide](guides/multi_model.md).
### LLM steps
For steps that call a model, `Baton.LLMStep` owns the transport loop — the
timing, the 429/529 → snooze and `max_tokens` → retryable mapping, JSON
decoding, and usage recording — so a step only describes its request and how to
handle the decoded reply:
```elixir
defmodule MyApp.Steps.AssessQuality do
use Baton.LLMStep
@impl true
def output_schema, do: %{"type" => "object", "properties" => %{"score" => %{"type" => "integer"}}}
@impl true
def request(%Oban.Job{} = job) do
{:ok, %{"parsed" => parsed}} = Baton.Results.get_result(job, :parse)
{:ok, [%{role: "user", content: prompt(parsed)}], model: "claude-sonnet-4-20250514"}
end
@impl true
def handle_response(%{"score" => _} = quality, _ctx, _job), do: {:ok, %{"quality" => quality}}
end
```
The client is the module you set as `config :baton, llm_client: MyApp.LLM`, whose
`complete/2` returns `{:ok, %{text:, model:, stop_reason:, usage:}}` or
`{:error, reason}`. See `Baton.LLMClient` for the full contract.
### Batch mode
When nobody is waiting on a step's result — a nightly enrichment run, a
backfill, a DAG tail — it can run on the provider's Message Batches API for
roughly half the token cost, at hours-scale latency. That is one line:
```elixir
defmodule MyApp.Steps.AssessQuality do
use Baton.LLMStep, mode: :batch
# request/1, handle_response/3, output_schema/0 — all unchanged
end
```
The engine submits a batch, polls with snoozes (which cost no retry attempts,
so a step can wait a day with its budget intact), then runs the result through
the same decode → handle → attach-usage pipeline. Dependents, completion,
retries, and stats behave exactly as for a live step; the recorded usage
carries `service_tier: "batch"` so your pricing module can apply the discount.
Batch mode needs **schema v7** and three optional callbacks on your client
(`submit_batch/2`, `poll_batch/2`, `batch_results/2`). See the
[batch mode guide](guides/batch_mode.md) for a complete worked example.
## Integrating with Phoenix LiveView
Baton ships **no LiveView of its own**. Instead, every step transition is
broadcast over `Phoenix.PubSub`, so you render progress however you like. (The
same transitions are also emitted as telemetry — see `Baton.Telemetry` — if
you'd rather not use Phoenix at all.)
### 1. Point Baton at your PubSub
A Phoenix app already starts one in its supervision tree (`{Phoenix.PubSub,
name: MyApp.PubSub}`). Tell Baton to use it:
```elixir
config :baton, pubsub: MyApp.PubSub
```
If `:pubsub` is left unset, broadcasting is a no-op and the engine runs fine
without Phoenix — only telemetry is emitted.
### 2. Topics and message shape
Each transition is published on **two topics** so views can subscribe at the
granularity they need:
- `"workflow:all"` — every event from every workflow (index views)
- `"workflow:<workflow_id>"` — one workflow's events (detail views)
Don't build these strings by hand — use the helpers in `Baton.Events`. The
message is always:
```elixir
{:workflow_step_updated, %{
workflow_id: "uuid",
workflow_label: "patent:US11234567B2", # the :workflow_name you passed to new/1
step_name: "assess_quality",
worker: "MyApp.Patent.AssessQuality",
state: "completed", # see below
job_id: 123,
attempt: 1,
has_result: true,
error: nil, # an error string on failure, else nil
timestamp: ~U[2026-06-14 18:00:00Z]
}}
```
`state` is one of `"executing"`, `"snoozed"`, `"completed"`, `"retryable"`,
`"discarded"`, or `"cancelled"`. A `"snoozed"` payload's `detail` names the
wait — `%{reason: r, seconds: n}` with `r` one of `"deps"`, `"rate_budget"`,
`"provider_limit"`, `"batch_slot"`, or `"step"`, and `n` the announced wait in
seconds.
When the **last** step in a workflow settles, a single terminal event is
published on the same two topics:
```elixir
{:workflow_finished, %{
workflow_id: "uuid",
workflow_label: "patent:US11234567B2",
outcome: :completed, # or :failed
failed_steps: [], # step names that were cancelled/discarded
timestamp: ~U[2026-06-14 18:00:24Z]
}}
```
Use it to flip the page to a done state, redirect, or fire a notification
without polling. (The same signal is available as
`[:baton, :workflow, :finished]` telemetry if you're not using PubSub.)
> **Requires `Baton.Plugin` for crash-case coverage.** When a step fails by
> returning `{:error, reason}`, the finished event fires immediately. But if a
> step *hard-crashes* (raises/exits) or is killed by Oban, the worker never gets
> to announce — `Baton.Plugin`'s periodic sweep is what detects the settled
> workflow and broadcasts `{:workflow_finished, outcome: :failed}` as a backstop
> (typically within one sweep interval). Make sure the plugin is in your Oban
> `plugins:` list (see [Installation](#installation)); without it, workflows that
> die from a hard crash won't emit a terminal event.
### 3. Subscribe in a LiveView
```elixir
defmodule MyAppWeb.WorkflowLive do
use MyAppWeb, :live_view
alias Baton.Events
def mount(%{"id" => workflow_id}, _session, socket) do
if connected?(socket), do: Events.subscribe_workflow(workflow_id)
{:ok, assign(socket, workflow_id: workflow_id, steps: %{})}
end
def handle_info({:workflow_step_updated, %{step_name: name} = event}, socket) do
{:noreply, update(socket, :steps, &Map.put(&1, name, event))}
end
# ... render @steps ...
end
```
For an index of all running workflows, subscribe with `Events.subscribe_all/0`
and key your state by `event.workflow_id`. A complete, copy-paste pair of
detail and index LiveViews lives in
[`examples/my_app/live/workflow_live.ex`](examples/my_app/live/workflow_live.ex).
### Seeding initial state
PubSub only delivers events that occur *after* `mount`, so a fresh page load (or
a step that completed before the user opened the view) won't be reflected by
events alone. Seed `@steps` from the database on mount using `Baton.Query`,
then let incoming events keep it current — and handle `{:workflow_finished, _}`
to react when the whole workflow is done.
## How it compares to Oban Pro Workflow
Baton covers DAG ordering, fan-out/fan-in, dynamic workflows, result
passing, and dependency-failure cascading. It adds cycle detection, retry
idempotency, multi-model fan-out, and LLM cost tracking. The main mechanical
difference is that completion uses snooze-based gating plus an opportunistic
reschedule rather than Pro's event-driven completion; correctness does not
depend on the reschedule.
That mechanism is also what makes dynamic fan-out possible here: dependencies
live in Baton's own table and are re-resolved on every wake rather than fixed
at insert, so a running step can add nodes to its own workflow and a
dependency list that changes between wakes is handled by construction.
## License
MIT — see [LICENSE](https://github.com/RudeWalrus/Baton/blob/main/LICENSE).