prop/vendor/oban_pro/usage-rules/workflows.md
Graham McIntire e99bf06eb4
deps: re-vendor oban_pro 1.7.0 (revert hex-repo dep)
Previous commit (3c988a5f) switched oban_pro to the licensed hex
repo at deploy time. Reverting that — keep all three Pro packages
vendored so prod image builds don't depend on oban.pro reachability
or auth on every CI run.

oban_pro 1.7.0 dropped into vendor/oban_pro via
`mix hex.package fetch oban_pro 1.7.0 --repo=oban --unpack`. mix.exs
goes back to `path: "vendor/oban_pro"` (matching oban_met / oban_web,
which were already vendored — and stay vendored since the licensed
repo only has older versions of those: 0.1.11 / 2.10.6 vs the
1.1.0 / 2.12.1 we vendor).

Schema migration from 3c988a5f stays — 1.7.0 tables/indexes are
already applied. No code changes.
2026-04-30 09:29:55 -05:00

1.9 KiB

Workflows

Workflows compose jobs with dependency relationships: sequential, fan-out, and fan-in patterns.

alias Oban.Pro.Workflow

Workflow.new()
|> Workflow.add(:extract, ExtractWorker.new(%{source: "db"}))
|> Workflow.add(:transform, TransformWorker.new(%{}), deps: :extract)
|> Workflow.add(:notify, NotifyWorker.new(%{}), deps: :extract)  # fan-out
|> Workflow.add(:load, LoadWorker.new(%{}), deps: [:transform, :notify])  # fan-in
|> Oban.insert_all()

Patterns

  • Sequential: deps: :previous_job
  • Fan-out: Multiple jobs depend on the same parent
  • Fan-in: One job depends on multiple parents with deps: [:a, :b, :c]

Accessing upstream results

Use recorded: true workers:

def process(job) do
  {:ok, data} = Workflow.fetch_recorded(job, :extract)
end

Cascading Functions

Function captures that automatically receive context and upstream results:

Workflow.new()
|> Workflow.put_context(%{source: source})
|> Workflow.add_cascade(:extract, &extract/1)
|> Workflow.add_cascade(:transform, &transform/1, deps: :extract)
|> Workflow.add_cascade(:load, &load/1, deps: :transform)
|> Oban.insert_all()

def extract(%{source: source}), do: %{records: fetch_records(source)}
def transform(%{extract: %{records: records}}), do: %{data: process(records)}
def load(%{transform: %{data: data}}), do: insert_all(data)

Fan-out with {enumerable, function/2}:

Workflow.add_cascade(:accounts, {account_ids, &sync_account/2})

Sub-workflows

Compose from reusable pieces:

extract_flow =
  Workflow.new()
  |> Workflow.add(:fetch, FetchWorker.new(%{}))
  |> Workflow.add(:parse, ParseWorker.new(%{}), deps: :fetch)

Workflow.new()
|> Workflow.add(:setup, SetupWorker.new(%{}))
|> Workflow.add_workflow(:extract, extract_flow, deps: :setup)
|> Workflow.add(:cleanup, CleanupWorker.new(%{}), deps: :extract)
|> Oban.insert_all()