Skip to content

Commit 109efc3

Browse files
perf(workflows): copy a fixed pre-fan-out snapshot per concurrent item
Concurrent fan-out items each copied the live shared steps dict, which completed items keep growing with their namespaced results. With a small worker count that made every later item's copy larger, O(K * N^2) work overall, and what a late-starting item saw depended on scheduling. The concurrent path now takes one snapshot before any worker starts and each item copies that. Also corrects the run_item comment that still described applying only the last item's aliases (the post-join code folds every collected item's aliases in item order), and adds an end-to-end resume() test for the reserved-id collision guard, since resume() builds its own run context separately from execute(). Assisted-by: Claude Code (model: Claude Opus 5.5, autonomous) Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
1 parent a8b1030 commit 109efc3

2 files changed

Lines changed: 137 additions & 5 deletions

File tree

‎src/specify_cli/workflows/engine.py‎

Lines changed: 31 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -1791,7 +1791,11 @@ def item_id(idx: int) -> str:
17911791
return f"{step_id}:{base_id}:{idx}"
17921792

17931793
def run_item(
1794-
idx: int, item_ctx: StepContext, *, local_only: bool
1794+
idx: int,
1795+
item_ctx: StepContext,
1796+
*,
1797+
local_only: bool,
1798+
base_steps: dict[str, Any] | None = None,
17951799
) -> tuple[Any, dict[str, dict[str, Any]]]:
17961800
# Namespace every id in the template's subtree (not just the
17971801
# template's own top-level id) so a step nested inside e.g. an
@@ -1820,9 +1824,21 @@ def run_item(
18201824
# shared steps dict that other concurrently-running items also
18211825
# read from — the actual race Copilot flagged: every worker
18221826
# writing the same bare-id key could otherwise expose another
1823-
# item's value to a sibling read. The caller applies exactly
1824-
# one item's aliases to shared state — deterministically the
1825-
# last item in item order — once every item has finished.
1827+
# item's value to a sibling read. Each item returns its aliases
1828+
# instead; the concurrent caller folds every collected item's
1829+
# aliases in item order once the pool has joined (so an id
1830+
# resolves to the latest item that actually wrote it), and the
1831+
# sequential caller publishes them after each item.
1832+
#
1833+
# ``base_steps``, when given, is the dict to snapshot instead of
1834+
# the live ``original_steps``. The concurrent path passes one
1835+
# snapshot taken before any item starts: completed items keep
1836+
# publishing their namespaced keys into ``original_steps``
1837+
# (below), so copying it per item would cost O(items x keys)
1838+
# per copy and O(K * N^2) overall, and what a late-starting item
1839+
# saw would depend on scheduling. Copying the fixed pre-fan-out
1840+
# snapshot keeps each copy the same size and makes every
1841+
# concurrent item see the same namespace.
18261842
#
18271843
# A plain dict copy, not a ``ChainMap`` overlay: expression
18281844
# interpolation (``{{ steps.x.output... }}``, via
@@ -1834,7 +1850,12 @@ def run_item(
18341850
# the shared dict made after this item started, but fan-out
18351851
# items were never entitled to see those anyway.
18361852
original_steps = item_ctx.steps
1837-
item_steps = dict(original_steps) if local_only else original_steps
1853+
if local_only:
1854+
item_steps = dict(
1855+
original_steps if base_steps is None else base_steps
1856+
)
1857+
else:
1858+
item_steps = original_steps
18381859
item_ctx.steps = item_steps
18391860

18401861
# Sequential items (not local_only) execute directly against the
@@ -1954,6 +1975,10 @@ def run_item(
19541975
n = len(items)
19551976
slots: list[Any] = [None] * n
19561977
alias_slots: list[dict[str, dict[str, Any]]] = [{}] * n
1978+
# Taken once, before any worker starts, so every item copies the
1979+
# same fixed-size namespace rather than the live dict that completed
1980+
# items keep growing (see ``base_steps`` in run_item).
1981+
base_steps = dict(context.steps)
19571982

19581983
def run_isolated(idx: int) -> tuple[Any, dict[str, dict[str, Any]]]:
19591984
# Each item runs against its own context copy so context.item is not
@@ -1968,6 +1993,7 @@ def run_isolated(idx: int) -> tuple[Any, dict[str, dict[str, Any]]]:
19681993
inside_fan_out=True,
19691994
),
19701995
local_only=True,
1996+
base_steps=base_steps,
19711997
)
19721998

19731999
def item_halt_status(idx: int) -> RunStatus | None:

‎tests/test_workflows.py‎

Lines changed: 106 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -7385,6 +7385,112 @@ def test_fan_out_reserved_id_collision_end_to_end(
73857385
)
73867386
assert state.step_results["after"]["output"]["stdout"].strip() == "outside"
73877387

7388+
@pytest.mark.parametrize("max_concurrency", [1, 2])
7389+
def test_fan_out_reserved_id_collision_end_to_end_on_resume(
7390+
self, project_dir, max_concurrency
7391+
):
7392+
"""`resume()` builds its own run context, separately from
7393+
`execute()`, so it needs its own end-to-end coverage of the
7394+
`reserved_step_ids` wiring.
7395+
7396+
The outside `leaf` step completes, a gate pauses the run, and the
7397+
colliding fan-out only runs after `engine.resume()`. On resume
7398+
`context.steps` IS `state.step_results`, so a missing reserved set
7399+
there would let the fan-out's bare `leaf` alias overwrite the
7400+
persisted outside result directly.
7401+
"""
7402+
from unittest.mock import patch
7403+
7404+
from specify_cli.workflows.base import RunStatus, StepResult
7405+
from specify_cli.workflows.engine import WorkflowDefinition, WorkflowEngine
7406+
7407+
yaml_str = f"""
7408+
schema_version: "1.0"
7409+
workflow:
7410+
id: "fan-out-reserved-collision-resume"
7411+
name: "Fan Out Reserved Collision Resume"
7412+
version: "1.0.0"
7413+
steps:
7414+
- id: leaf
7415+
type: shell
7416+
run: "echo outside"
7417+
- id: approve
7418+
type: gate
7419+
message: "Approve?"
7420+
- id: fan
7421+
type: fan-out
7422+
items: "{{{{ ['a', 'b', 'c'] }}}}"
7423+
max_concurrency: {max_concurrency}
7424+
step:
7425+
id: leaf
7426+
type: shell
7427+
run: "echo {{{{ item }}}}"
7428+
- id: after
7429+
type: shell
7430+
run: "echo {{{{ steps.leaf.output.stdout }}}}"
7431+
"""
7432+
definition = WorkflowDefinition.from_string(yaml_str)
7433+
engine = WorkflowEngine(project_dir)
7434+
state = engine.execute(definition)
7435+
assert state.status == RunStatus.PAUSED
7436+
7437+
with patch(
7438+
"specify_cli.workflows.step.gate.GateStep.execute",
7439+
return_value=StepResult(output={"approved": True}),
7440+
):
7441+
state = engine.resume(state.run_id)
7442+
7443+
assert state.status == RunStatus.COMPLETED
7444+
assert state.step_results["leaf"]["output"]["stdout"] == "outside\n"
7445+
for idx, item in enumerate(["a", "b", "c"]):
7446+
assert (
7447+
state.step_results[f"fan:leaf:{idx}"]["output"]["stdout"]
7448+
== f"{item}\n"
7449+
)
7450+
assert state.step_results["after"]["output"]["stdout"].strip() == "outside"
7451+
7452+
def test_concurrent_fan_out_items_copy_fixed_pre_fan_out_snapshot(
7453+
self, tmp_path
7454+
):
7455+
"""Every concurrent item must copy the namespace as it was before
7456+
the fan-out started, not the live steps dict.
7457+
7458+
Completed items publish their namespaced keys back into the shared
7459+
steps dict. With 2 workers, item 2 is only submitted once item 0
7460+
has been collected, so if items copied the live dict, item 2 would
7461+
see `fan:probe:0`, item 4 would see three earlier keys, and so on.
7462+
That makes the copies grow with the item count (quadratic overall)
7463+
and makes what an item sees depend on scheduling.
7464+
"""
7465+
from specify_cli.workflows.base import RunStatus, StepBase, StepContext, StepResult, StepStatus
7466+
from specify_cli.workflows.engine import RunState, WorkflowEngine
7467+
7468+
seen: dict[int, list[str]] = {}
7469+
7470+
class _ProbeStep(StepBase):
7471+
type_key = "probe"
7472+
7473+
def execute(self, config, context):
7474+
seen[context.item] = sorted(
7475+
k for k in context.steps if k.startswith("fan:")
7476+
)
7477+
return StepResult(status=StepStatus.COMPLETED, output={})
7478+
7479+
engine = WorkflowEngine(project_root=tmp_path)
7480+
context = StepContext()
7481+
context.steps["before"] = {"output": {}}
7482+
state = RunState(run_id="r", workflow_id="w", project_root=tmp_path)
7483+
state.status = RunStatus.RUNNING
7484+
registry = {"probe": _ProbeStep()}
7485+
template = {"id": "probe", "type": "probe"}
7486+
items = list(range(6))
7487+
engine._run_fan_out(items, template, "fan", context, state, registry, 2)
7488+
7489+
assert seen == {i: [] for i in items}
7490+
# Every item's namespaced result is still published afterwards.
7491+
for i in items:
7492+
assert f"fan:probe:{i}" in context.steps
7493+
73887494
@pytest.mark.parametrize("inner_concurrency", [1, 2])
73897495
def test_reserved_alias_visible_after_nested_fan_out_in_sequential_item(
73907496
self, project_dir, inner_concurrency

0 commit comments

Comments
 (0)