Workflow: downstream dependency on a graft is dropped when the grafting workflow is itself grafted (nested graft)

Summary

A downstream job that deps on a graft correctly waits for the grafted jobs at one graft level, but stops waiting when the workflow containing that graft is itself grafted into a parent workflow (a nested / graft-within-a-graft). The downstream job cascades as soon as the grafter node finishes, before the grafted jobs run.

This looks adjacent to the 1.7.7 fix “Wait for sub-workflows grafted via add_workflow/4 — that fix holds for a single graft level but does not extend to a second level of nesting.

Environment

  • oban_pro 1.7.7
  • oban 2.22.1
  • ecto_sql 3.13.5, postgrex 0.22.2
  • Elixir 1.20.0 / OTP 27, PostgreSQL 16
  • engine: Oban.Pro.Engines.Smart

Expected vs actual

A workflow with an inner graft node :inner (which grafts a :leaf job) and a :consumer with deps: [:inner]:

  • Expected: :consumer waits for the grafted :leaf in all cases.
  • Actual: it waits when the workflow is inserted directly (1 graft level), but not when the same workflow is grafted in by an outer grafter (2 graft levels) — :consumer runs before :leaf.

Minimal reproduction

Single-file script (own repo + Oban instance, no app code). Needs a Postgres via DATABASE_URL (defaults to a local one). Deps: {:oban, "~> 2.22"}, {:oban_pro, "~> 1.7", repo: "oban"}, {:ecto_sql, "~> 3.12"}, {:postgrex, ">= 0.0.0"}. Run with mix run --no-start repro.exs.

Application.put_env(:repro, Repro.Repo,
  url: System.get_env("DATABASE_URL") || "postgres://postgres:postgres@localhost/repro_dev",
  pool_size: 10
)

defmodule Repro.Repo do
  use Ecto.Repo, otp_app: :repro, adapter: Ecto.Adapters.Postgres
end

defmodule Repro.Migration do
  use Ecto.Migration

  def up do
    Oban.Migration.up()
    Oban.Pro.Migration.up()
  end

  def down do
    Oban.Pro.Migration.down()
    Oban.Migration.down()
  end
end

# Shared recorder to observe execution order across worker processes.
defmodule Repro.Recorder do
  def start, do: Agent.start_link(fn -> %{} end, name: __MODULE__)
  def put(key, val), do: Agent.update(__MODULE__, &Map.put(&1, key, val))
  def get(key), do: Agent.get(__MODULE__, &Map.get(&1, key))
end

# The grafted job. Sleeps so a consumer with no real dependency runs before it finishes.
defmodule Repro.Leaf do
  use Oban.Pro.Worker, queue: :default, max_attempts: 1

  @impl Oban.Pro.Worker
  def process(%Oban.Job{args: %{"run" => run}}) do
    Process.sleep(400)
    Repro.Recorder.put({run, :leaf_done}, true)
    :ok
  end
end

# Downstream consumer. Records whether :leaf had finished at the moment it ran.
defmodule Repro.Consumer do
  use Oban.Pro.Worker, queue: :default, max_attempts: 1

  @impl Oban.Pro.Worker
  def process(%Oban.Job{args: %{"run" => run}}) do
    Repro.Recorder.put({run, :consumer_saw_leaf}, Repro.Recorder.get({run, :leaf_done}) == true)
    :ok
  end
end

# Inner grafter: grafts :leaf into the workflow it is running inside.
defmodule Repro.InnerGrafter do
  use Oban.Pro.Worker, queue: :default, max_attempts: 1

  @impl Oban.Pro.Worker
  def process(%Oban.Job{args: %{"run" => run}}) do
    Oban.Pro.Workflow.new()
    |> Oban.Pro.Workflow.add(:leaf, Repro.Leaf.new(%{run: run}))
    |> Oban.Pro.Workflow.apply_graft()
    |> Oban.insert_all()

    :ok
  end
end

# "gen" workflow, built identically in both scenarios: an inner graft that spawns :leaf,
# and a :consumer that depends on that graft.
defmodule Repro.Gen do
  def workflow(run) do
    Oban.Pro.Workflow.new()
    |> Oban.Pro.Workflow.add_graft(:inner, Repro.InnerGrafter.new(%{run: run}))
    |> Oban.Pro.Workflow.add(:consumer, Repro.Consumer.new(%{run: run}), deps: [:inner])
  end
end

# Outer grafter: grafts the WHOLE gen workflow in — this is the second graft level.
defmodule Repro.OuterGrafter do
  use Oban.Pro.Worker, queue: :default, max_attempts: 1

  @impl Oban.Pro.Worker
  def process(%Oban.Job{args: %{"run" => run}}) do
    Repro.Gen.workflow(run)
    |> Oban.Pro.Workflow.apply_graft()
    |> Oban.insert_all()

    :ok
  end
end

# ---- boot ----
{:ok, _} = Application.ensure_all_started(:postgrex)
{:ok, _} = Application.ensure_all_started(:ecto_sql)
{:ok, _} = Application.ensure_all_started(:oban)
{:ok, _} = Application.ensure_all_started(:oban_pro)

{:ok, _} = Repro.Repo.start_link()
Ecto.Migrator.run(Repro.Repo, [{0, Repro.Migration}], :up, all: true)

{:ok, _} =
  Oban.start_link(
    repo: Repro.Repo,
    engine: Oban.Pro.Engines.Smart,
    queues: [default: 10],
    plugins: [Oban.Pro.Plugins.DynamicLifeline]
  )

{:ok, _} = Repro.Recorder.start()

# CONTROL — one graft level: the gen workflow is inserted directly.
Repro.Gen.workflow("control") |> Oban.insert_all()

# NESTED — two graft levels: the gen workflow is itself grafted in by an outer grafter.
Oban.Pro.Workflow.new()
|> Oban.Pro.Workflow.add_graft(:outer, Repro.OuterGrafter.new(%{run: "nested"}))
|> Oban.insert_all()

Enum.reduce_while(1..100, nil, fn _, _ ->
  if Repro.Recorder.get({"control", :consumer_saw_leaf}) != nil and
       Repro.Recorder.get({"nested", :consumer_saw_leaf}) != nil do
    {:halt, :ok}
  else
    Process.sleep(100)
    {:cont, nil}
  end
end)

control = Repro.Recorder.get({"control", :consumer_saw_leaf})
nested = Repro.Recorder.get({"nested", :consumer_saw_leaf})

IO.puts("control (1 graft level): consumer saw :leaf finished? #{inspect(control)}   (expected true)")
IO.puts("nested  (2 graft levels): consumer saw :leaf finished? #{inspect(nested)}   (expected true)")

Output:

control (1 graft level): consumer saw :leaf finished? true   (expected true)
nested  (2 graft levels): consumer saw :leaf finished? false   (expected true)

Question

Is nested grafting (a graft inside a workflow that is itself grafted) intended to be supported? If so, the root_id computation in apply_graft/2 seems to need to preserve the current graft’s graft_id for the grafted jobs (or reparent_grafted/3 needs to tie them back to it) rather than collapsing to the outermost grafted_* id, so that a downstream dependency on the inner graft still resolves.