feat: return foreach items through owner back-edges

This commit is contained in:
lda
2026-09-04 07:40:49 +07:00 Verified
parent f69c4502cf
commit f23e760b5f
13 changed files with 902 additions and 21 deletions
+2 -2
View File
@@ -114,7 +114,7 @@ def build_concurrent_foreach_workflow(
)
builder.set_entry_point(each)
builder.connect(each, "loop", record)
builder.connect(record, "ok", END)
builder.connect(record, "ok", each)
builder.connect(each, "done", END)
if _item_error_action(item_error) in {"collect", "skip"}:
builder.connect(each, "completed_with_errors", END)
@@ -165,7 +165,7 @@ def run_replace_conflict_example() -> None:
)
builder.set_entry_point(each)
builder.connect(each, "loop", record)
builder.connect(record, "ok", END)
builder.connect(record, "ok", each)
builder.connect(each, "done", END)
try:
builder.execute({"items": ["a", "b"]})
+1 -1
View File
@@ -255,7 +255,7 @@ def build_demo_workflow() -> Workflow:
"outcome": "done",
"to": "combine_summaries",
},
{"from": "summarize_one", "outcome": "ok", "to": END},
{"from": "summarize_one", "outcome": "ok", "to": "summarize_each"},
{"from": "combine_summaries", "outcome": "ok", "to": "should_email"},
{"from": "should_email", "outcome": "true", "to": "approve_email"},
{"from": "should_email", "outcome": "false", "to": "skip_email"},
+1 -1
View File
@@ -97,7 +97,7 @@ def build_raw_concurrent_foreach_workflow() -> Workflow:
],
"edges": [
{"from": "each", "outcome": "loop", "to": "record"},
{"from": "record", "outcome": "ok", "to": END},
{"from": "record", "outcome": "ok", "to": "each"},
{"from": "each", "outcome": "done", "to": END},
{"from": "each", "outcome": "completed_with_errors", "to": END},
],
+58
View File
@@ -2,6 +2,7 @@ from __future__ import annotations
from typing import Any
from wf_core.errors import WorkflowExecutionError
from wf_core.models.workflow import Workflow
from wf_core.run_state import (
ExecutionFrame,
@@ -76,6 +77,39 @@ def advance_frame(
next_node_id: str,
front: bool = False,
) -> None:
# Foreach back-edge return is an ownership check, not generic cycle
# detection. Only the frame's immediate recorded owner completes the item;
# a root frame targeting the same foreach enters it normally.
from wf_core.runtime.foreach_state import item_frame_owner
owner = item_frame_owner(frame)
if owner is not None:
if next_node_id == END:
raise WorkflowExecutionError(
f"foreach item frame {frame.id!r} cannot target workflow END; "
f"return to owning foreach {owner.foreach_node_id!r}"
)
if next_node_id == owner.foreach_node_id:
source_node_id = frame.node_id
frame.prior_outcome = outcome
frame.activated_incoming_edge = source_node_id
frame.node_id = owner.foreach_node_id
frame.status = FrameStatus.COMPLETED
frame.finished_at_node_id = owner.foreach_node_id
# The child does not execute the controller again; the blocked
# parent activation consumes the result and admits the next item
# or emits done. The owner location stays inspectable in trace
# and checkpoint state.
wake_parent_for_child_progress(run, frame.id)
run.sync_from_current_frame()
return
ancestors = _foreach_ancestor_ids(run, frame)
if next_node_id in ancestors[1:]:
raise WorkflowExecutionError(
f"foreach item frame {frame.id!r} targets non-immediate "
f"ancestor {next_node_id!r}; only {owner.foreach_node_id!r} "
"can complete this item"
)
frame.prior_outcome = outcome
frame.activated_incoming_edge = frame.node_id
frame.node_id = next_node_id
@@ -93,6 +127,30 @@ def advance_frame(
run.sync_from_current_frame()
def _foreach_ancestor_ids(run: RunState, frame: ExecutionFrame) -> list[str]:
"""Derive active foreach owners from frame ancestry for fail-closed checks.
The first entry is the frame's immediate owner; later entries are older
ancestors. A target naming an older ancestor is a non-local return, while
a target naming an inactive foreach is an ordinary nested entry.
"""
from wf_core.runtime.foreach_state import item_frame_owner
ancestors: list[str] = []
cursor: ExecutionFrame | None = frame
seen: set[str] = set()
while cursor is not None:
owner = item_frame_owner(cursor)
if owner is not None:
if owner.foreach_node_id in seen:
break
seen.add(owner.foreach_node_id)
ancestors.append(owner.foreach_node_id)
parent_id = cursor.parent_frame_id
cursor = run.frames.get(parent_id) if parent_id is not None else None
return ancestors
def finalize_run(workflow: Workflow, run: RunState) -> RunState:
if run.outcome is None:
run.outcome = "ok"
+17
View File
@@ -12,6 +12,7 @@ from wf_core.runtime.foreach_state import (
ForeachBarrierState,
ItemErrorRecord,
PendingItemResult,
close_foreach_activation,
load_or_begin_foreach_activation,
save_foreach_activation,
)
@@ -87,6 +88,9 @@ def _step_foreach_serial(
state_changes={},
),
)
# Close the visit before following `done` so a self-looping completion
# edge or a later revisit starts a fresh activation.
close_foreach_activation(frame, activation)
advance_frame(run, frame, outcome=outcome, next_node_id=next_node_id)
return run
@@ -96,6 +100,15 @@ def _step_foreach_serial(
save_foreach_activation(frame, activation)
child_id = _child_frame_id(activation, loop_index)
child_lineage_id = _child_lineage_id(activation, loop_index)
# Serial items still own a lineage so nested subgraph/boundary commits have
# a parent lineage to buffer into; top-level serial writes commit through
# the parent scope root.
add_lineage(
run,
scope_id=frame.scope_id,
lineage_id=child_lineage_id,
parent_id=frame.lineage_id,
)
add_frame(
run,
ExecutionFrame(
@@ -168,6 +181,7 @@ def _step_foreach_concurrent(
frame=frame,
step=step,
index=index,
activation=activation,
barrier=barrier,
reducers=reducers,
)
@@ -318,6 +332,7 @@ def _finish_concurrent_foreach(
frame: ExecutionFrame,
step: ForeachNode,
index: WorkflowIndex,
activation: ForeachActivationState,
barrier: ForeachBarrierState,
reducers: Mapping[str, ReducerDefinition] | None = None,
) -> RunState:
@@ -373,6 +388,8 @@ def _finish_concurrent_foreach(
state_changes=state_changes,
),
)
# Close the visit before following completion so later revisits start fresh.
close_foreach_activation(frame, activation)
advance_frame(run, frame, outcome=outcome, next_node_id=next_node_id)
return run
+5
View File
@@ -87,6 +87,11 @@ def complete_end_step(
"""Record an explicit workflow terminal and complete the active frame."""
result = StepExecutionResult(outcome=outcome)
frame = run.frames[frame_id]
if item_frame_owner(frame) is not None:
raise WorkflowExecutionError(
f"foreach item frame {frame.id!r} cannot target explicit end node "
f"{node_id!r}; return to its owning foreach"
)
frame.metadata["workflow_outcome"] = outcome
if frame.parent_frame_id is None:
run.outcome = outcome
+25 -1
View File
@@ -236,7 +236,31 @@ def _finish_subgraph(
reducers=reducers,
missing_field_message="subgraph output did not include required field {field}",
)
state_changes = commit_patch_for_frame(run, frame, patch)
# Match node execution: serial item writes commit through the parent
# scope so top-level serial subgraphs land in root state; concurrent
# item writes stay buffered in the item lineage for barrier merge.
from wf_core.runtime.foreach_state import (
item_frame_owner,
load_foreach_activation,
)
commit_frame = frame
try:
owner = item_frame_owner(frame)
except Exception:
owner = None
if owner is not None:
parent_frame = run.frames.get(owner.parent_frame_id)
if parent_frame is not None:
foreach_activation = load_foreach_activation(
parent_frame, owner.foreach_node_id, owner.activation_id
)
if (
foreach_activation is not None
and foreach_activation.barrier.mode == "serial"
):
commit_frame = parent_frame
state_changes = commit_patch_for_frame(run, commit_frame, patch)
return StepExecutionResult(
outcome=child_outcome,
resolved_input=activation.child_input,
+1 -1
View File
@@ -210,7 +210,7 @@ def build_authoring_demo_workflow():
builder.connect(list_files, "ok", summarize_each)
builder.connect(summarize_each, "loop", summarize_one)
builder.connect(summarize_each, "done", combine_summaries)
builder.connect(summarize_one, "ok", END)
builder.connect(summarize_one, "ok", summarize_each)
builder.connect(combine_summaries, "ok", should_email)
builder.connect(should_email, "true", approve_email)
builder.connect(should_email, "false", skip_email)
+31 -11
View File
@@ -286,12 +286,14 @@ def test_sync_concurrent_foreach_barrier_replays_add_reducer_inputs() -> None:
assert run.output["number"] == 6
assert run.lineages["root:each#0[0]"].writes[0].incoming_value == 3
assert run.lineages["root:each#0[1]"].writes[0].incoming_value == 1
active = load_or_begin_foreach_activation(
# The visit closed before `done`; lineage history remains while a new
# load starts fresh barrier state with a new activation id.
fresh = load_or_begin_foreach_activation(
run.frames["root"], "each", mode="concurrent"
)
assert active.id == "root:each#0"
assert active.barrier.pending_results[0].lineage_id == "root:each#0[0]"
assert active.barrier.pending_results[0].patch.writes == []
assert fresh.id == "root:each#1"
assert fresh.barrier.next_index == 0
assert fresh.barrier.pending_results == {}
foreach_entries = [entry for entry in run.trace if entry.step_type == "foreach"]
assert foreach_entries[-1].state_changes["state.number"] == 6
@@ -407,7 +409,7 @@ def _sum_items_workflow() -> Workflow:
],
edges=[
Edge.model_validate({"from": "each", "outcome": "loop", "to": "add_item"}),
Edge.model_validate({"from": "add_item", "outcome": "ok", "to": END}),
Edge.model_validate({"from": "add_item", "outcome": "ok", "to": "each"}),
Edge.model_validate({"from": "each", "outcome": "done", "to": END}),
],
)
@@ -517,7 +519,7 @@ def _same_item_reducer_visibility_workflow() -> Workflow:
"to": "read_number",
}
),
Edge.model_validate({"from": "read_number", "outcome": "ok", "to": END}),
Edge.model_validate({"from": "read_number", "outcome": "ok", "to": "each"}),
Edge.model_validate({"from": "each", "outcome": "done", "to": END}),
],
)
@@ -593,6 +595,15 @@ def _nested_foreach_lineage_workflow() -> Workflow:
"output": [{"source": "seen", "target": "state.seen"}],
}
),
NodeUse.model_validate(
{
"id": "tail",
"type": "node",
"node": "record",
"input": [{"target": "seen", "path": "context.outer"}],
"output": [{"source": "seen", "target": "state.seen"}],
}
),
],
edges=[
Edge.model_validate(
@@ -609,8 +620,13 @@ def _nested_foreach_lineage_workflow() -> Workflow:
"to": "record",
}
),
Edge.model_validate({"from": "record", "outcome": "ok", "to": END}),
Edge.model_validate({"from": "inner_each", "outcome": "done", "to": END}),
Edge.model_validate(
{"from": "record", "outcome": "ok", "to": "inner_each"}
),
Edge.model_validate(
{"from": "inner_each", "outcome": "done", "to": "tail"}
),
Edge.model_validate({"from": "tail", "outcome": "ok", "to": "outer_each"}),
Edge.model_validate({"from": "outer_each", "outcome": "done", "to": END}),
],
)
@@ -624,7 +640,7 @@ def _workflow(
) -> Workflow:
edges = [
Edge.model_validate({"from": "each", "outcome": "loop", "to": "record"}),
Edge.model_validate({"from": "record", "outcome": "ok", "to": END}),
Edge.model_validate({"from": "record", "outcome": "ok", "to": "each"}),
Edge.model_validate({"from": "each", "outcome": "done", "to": END}),
]
if include_completed_with_errors:
@@ -794,7 +810,9 @@ def _multi_step_overlay_workflow() -> Workflow:
"to": "read_scratch",
}
),
Edge.model_validate({"from": "read_scratch", "outcome": "ok", "to": END}),
Edge.model_validate(
{"from": "read_scratch", "outcome": "ok", "to": "each"}
),
Edge.model_validate({"from": "each", "outcome": "done", "to": END}),
],
)
@@ -861,7 +879,9 @@ def _same_path_replace_workflow() -> Workflow:
"to": "write_winner",
}
),
Edge.model_validate({"from": "write_winner", "outcome": "ok", "to": END}),
Edge.model_validate(
{"from": "write_winner", "outcome": "ok", "to": "each"}
),
Edge.model_validate({"from": "each", "outcome": "done", "to": END}),
],
)
+1 -1
View File
@@ -123,7 +123,7 @@ def _workflow(*, max_active: int) -> Workflow:
],
edges=[
Edge.model_validate({"from": "each", "outcome": "loop", "to": "record"}),
Edge.model_validate({"from": "record", "outcome": "ok", "to": END}),
Edge.model_validate({"from": "record", "outcome": "ok", "to": "each"}),
Edge.model_validate({"from": "each", "outcome": "done", "to": END}),
],
)
+1 -1
View File
@@ -144,7 +144,7 @@ def _workflow(*, item_error: dict[str, object]) -> Workflow:
],
edges=[
Edge.model_validate({"from": "each", "outcome": "loop", "to": "record"}),
Edge.model_validate({"from": "record", "outcome": "ok", "to": END}),
Edge.model_validate({"from": "record", "outcome": "ok", "to": "each"}),
Edge.model_validate({"from": "each", "outcome": "done", "to": END}),
Edge.model_validate(
{
@@ -52,8 +52,15 @@ async def test_resume_prioritizes_interrupted_item_before_siblings() -> None:
resume_payload={},
)
from wf_core.runtime.foreach_state import item_frame_owner
interrupted_owner = item_frame_owner(run.frames["root:each#0:1"])
assert interrupted_owner is not None
assert resumed.status is RunStatus.COMPLETED
assert resumed.state["seen"] == ["a", "b", "c"]
resumed_owner = item_frame_owner(resumed.frames["root:each#0:1"])
assert resumed_owner is not None
assert resumed_owner.activation_id == interrupted_owner.activation_id
assert resumed.trace[interrupted_trace_len].frame_id == "root:each#0:1"
assert resumed.trace[interrupted_trace_len].step_type == "interrupt"
assert resumed.trace[interrupted_trace_len].outcome == "submitted"
@@ -132,7 +139,7 @@ def _workflow() -> Workflow:
],
edges=[
Edge.model_validate({"from": "each", "outcome": "loop", "to": "route"}),
Edge.model_validate({"from": "route", "outcome": "ok", "to": END}),
Edge.model_validate({"from": "route", "outcome": "ok", "to": "each"}),
Edge.model_validate(
{
"from": "route",
@@ -140,7 +147,7 @@ def _workflow() -> Workflow:
"to": "ask",
}
),
Edge.model_validate({"from": "ask", "outcome": "submitted", "to": END}),
Edge.model_validate({"from": "ask", "outcome": "submitted", "to": "each"}),
Edge.model_validate({"from": "each", "outcome": "done", "to": END}),
],
)
+750
View File
@@ -0,0 +1,750 @@
from __future__ import annotations
from typing import Any
import pytest
from wf_core import (
END,
ConditionNode,
Edge,
ForeachNode,
NodeDef,
NodeUse,
ReducerRef,
SchemaRef,
StateField,
StateSchema,
SubgraphNode,
Workflow,
WorkflowExecutionError,
execute_workflow,
)
from wf_core.run_state import ExecutionFrame, FrameStatus, RunState, RunStatus
from wf_core.runtime.foreach_state import item_frame_owner
from wf_core.runtime.scheduler import add_frame
def _node_use(node_id: str, *, node: str = "record") -> NodeUse:
return NodeUse.model_validate(
{
"id": node_id,
"type": "node",
"node": node,
"input": [{"target": "value", "path": "context.item"}],
"output": [{"source": "seen", "target": "state.seen"}],
}
)
def _serial_workflow() -> Workflow:
foreach = ForeachNode.model_validate(
{
"id": "each",
"type": "foreach",
"over": "state.items",
"as": "item",
"mode": "serial",
}
)
return Workflow(
name="foreach_back_edge",
input_schema=SchemaRef(type="object", properties={"items": {"type": "array"}}),
state_schema=StateSchema.from_field_map(
{
"items": StateField(type="array"),
"seen": StateField(
type="array", reducer=ReducerRef(name="wf.std.append")
),
}
),
output_schema=SchemaRef(type="object", properties={"seen": {"type": "array"}}),
node_defs=[
NodeDef(
name="record",
input_schema=SchemaRef(
type="object", properties={"value": {}}, required=["value"]
),
output_schema=SchemaRef(
type="object", properties={"seen": {}}, required=["seen"]
),
outcomes=["ok"],
)
],
start="each",
nodes=[
foreach,
NodeUse.model_validate(
{
"id": "work",
"type": "node",
"node": "record",
"input": [{"target": "value", "path": "context.item"}],
"output": [{"source": "seen", "target": "state.seen"}],
}
),
],
edges=[
Edge.model_validate({"from": "each", "outcome": "loop", "to": "work"}),
Edge.model_validate({"from": "work", "outcome": "ok", "to": "each"}),
Edge.model_validate({"from": "each", "outcome": "done", "to": END}),
],
)
def test_serial_item_return_wakes_parent_and_admits_next_item() -> None:
workflow = _serial_workflow()
run = execute_workflow(
workflow,
{"items": ["a", "b"]},
{
"record": lambda payload, _ctx: {
"outcome": "ok",
"output": {"seen": payload["value"]},
}
},
)
assert run.status == RunStatus.COMPLETED
assert run.state["seen"] == ["a", "b"]
item_frames = [
frame for frame in run.frames.values() if frame.kind == "foreach_iteration"
]
assert len(item_frames) == 2
assert all(frame.status == FrameStatus.COMPLETED for frame in item_frames)
assert all(frame.finished_at_node_id == "each" for frame in item_frames)
assert run.output["seen"] == ["a", "b"]
def test_foreach_body_cycle_can_repeat_then_return() -> None:
foreach = ForeachNode.model_validate(
{
"id": "each",
"type": "foreach",
"over": "state.items",
"as": "item",
"mode": "serial",
}
)
workflow = Workflow(
name="foreach_body_cycle",
input_schema=SchemaRef(type="object", properties={"items": {"type": "array"}}),
state_schema=StateSchema.from_field_map(
{
"items": StateField(type="array"),
"seen": StateField(
type="array", reducer=ReducerRef(name="wf.std.append")
),
}
),
output_schema=SchemaRef(type="object", properties={"seen": {"type": "array"}}),
node_defs=[
NodeDef(
name="step_a",
input_schema=SchemaRef(
type="object", properties={"value": {}}, required=["value"]
),
output_schema=SchemaRef(
type="object", properties={"seen": {}}, required=["seen"]
),
outcomes=["again", "done"],
),
NodeDef(
name="step_b",
input_schema=SchemaRef(
type="object", properties={"seen": {}}, required=["seen"]
),
output_schema=SchemaRef(
type="object", properties={"seen": {}}, required=["seen"]
),
outcomes=["ok"],
),
],
start="each",
nodes=[
foreach,
NodeUse.model_validate(
{
"id": "a",
"type": "node",
"node": "step_a",
"input": [{"target": "value", "path": "context.item"}],
"output": [{"source": "seen", "target": "state.seen"}],
}
),
NodeUse.model_validate(
{
"id": "b",
"type": "node",
"node": "step_b",
"input": [{"target": "seen", "path": "context.item"}],
"output": [],
}
),
],
edges=[
Edge.model_validate({"from": "each", "outcome": "loop", "to": "a"}),
Edge.model_validate({"from": "a", "outcome": "again", "to": "b"}),
Edge.model_validate({"from": "b", "outcome": "ok", "to": "a"}),
Edge.model_validate({"from": "a", "outcome": "done", "to": "each"}),
Edge.model_validate({"from": "each", "outcome": "done", "to": END}),
],
)
calls = {"count": 0}
def step_a(payload: dict[str, Any], _ctx: object) -> dict[str, Any]:
calls["count"] += 1
outcome = "again" if calls["count"] == 1 else "done"
return {"outcome": outcome, "output": {"seen": payload["value"]}}
run = execute_workflow(
workflow,
{"items": ["x"]},
{
"step_a": step_a,
"step_b": lambda payload, _ctx: {
"outcome": "ok",
"output": {"seen": payload["seen"]},
},
},
)
assert run.status == RunStatus.COMPLETED
# `b` writes nothing; `a` writes once per visit (again + done).
assert run.state["seen"] == ["x", "x"]
def test_conditional_body_can_return_on_either_outcome() -> None:
foreach = ForeachNode.model_validate(
{
"id": "each",
"type": "foreach",
"over": "state.items",
"as": "item",
"mode": "serial",
}
)
workflow = Workflow(
name="foreach_conditional_return",
input_schema=SchemaRef(type="object", properties={"items": {"type": "array"}}),
state_schema=StateSchema.from_field_map(
{
"items": StateField(type="array"),
"seen": StateField(
type="array", reducer=ReducerRef(name="wf.std.append")
),
}
),
output_schema=SchemaRef(type="object", properties={"seen": {"type": "array"}}),
node_defs=[
NodeDef(
name="decide",
input_schema=SchemaRef(
type="object", properties={"value": {}}, required=["value"]
),
output_schema=SchemaRef(
type="object", properties={"seen": {}}, required=["seen"]
),
outcomes=["true", "false"],
),
NodeDef(
name="work",
input_schema=SchemaRef(
type="object", properties={"value": {}}, required=["value"]
),
output_schema=SchemaRef(
type="object", properties={"seen": {}}, required=["seen"]
),
outcomes=["ok"],
),
],
start="each",
nodes=[
foreach,
NodeUse.model_validate(
{
"id": "condition",
"type": "node",
"node": "decide",
"input": [{"target": "value", "path": "context.item"}],
"output": [],
}
),
NodeUse.model_validate(
{
"id": "work",
"type": "node",
"node": "work",
"input": [{"target": "value", "path": "context.item"}],
"output": [{"source": "seen", "target": "state.seen"}],
}
),
NodeUse.model_validate(
{
"id": "work_false",
"type": "node",
"node": "work",
"input": [{"target": "value", "path": "context.item"}],
"output": [{"source": "seen", "target": "state.seen"}],
}
),
],
edges=[
Edge.model_validate({"from": "each", "outcome": "loop", "to": "condition"}),
Edge.model_validate({"from": "condition", "outcome": "true", "to": "work"}),
Edge.model_validate(
{"from": "condition", "outcome": "false", "to": "work_false"}
),
Edge.model_validate({"from": "work", "outcome": "ok", "to": "each"}),
Edge.model_validate({"from": "work_false", "outcome": "ok", "to": "each"}),
Edge.model_validate({"from": "each", "outcome": "done", "to": END}),
],
)
def decide(payload: dict[str, Any], _ctx: object) -> dict[str, Any]:
outcome = "true" if payload["value"] == "a" else "false"
return {"outcome": outcome, "output": {"seen": payload["value"]}}
def _work(payload: dict[str, Any], _ctx: object) -> dict[str, Any]:
return {"outcome": "ok", "output": {"seen": payload["value"]}}
run = execute_workflow(
workflow,
{"items": ["a", "b"]},
{"decide": decide, "work": _work},
)
assert run.status == RunStatus.COMPLETED
assert sorted(run.state["seen"]) == ["a", "b"]
def test_nested_foreach_returns_inner_then_outer() -> None:
outer = ForeachNode.model_validate(
{
"id": "outer",
"type": "foreach",
"over": "state.items",
"as": "outer_item",
"mode": "concurrent",
"concurrent": {"max_active": 2, "max_outstanding": 2},
}
)
inner = ForeachNode.model_validate(
{
"id": "inner",
"type": "foreach",
"over": "state.inner_items",
"as": "inner_item",
"mode": "concurrent",
"concurrent": {"max_active": 2, "max_outstanding": 2},
}
)
workflow = Workflow(
name="nested_foreach_return",
input_schema=SchemaRef(type="object", properties={}),
state_schema=StateSchema.from_field_map(
{
"items": StateField(type="array"),
"inner_items": StateField(type="array"),
"seen": StateField(
type="array", reducer=ReducerRef(name="wf.std.append")
),
}
),
output_schema=SchemaRef(type="object", properties={"seen": {"type": "array"}}),
node_defs=[
NodeDef(
name="work",
input_schema=SchemaRef(
type="object", properties={"value": {}}, required=["value"]
),
output_schema=SchemaRef(
type="object", properties={"seen": {}}, required=["seen"]
),
outcomes=["ok"],
)
],
start="outer",
nodes=[
outer,
inner,
NodeUse.model_validate(
{
"id": "work",
"type": "node",
"node": "work",
"input": [{"target": "value", "path": "context.inner_item"}],
"output": [{"source": "seen", "target": "state.seen"}],
}
),
NodeUse.model_validate(
{
"id": "tail",
"type": "node",
"node": "work",
"input": [{"target": "value", "path": "context.outer_item"}],
"output": [{"source": "seen", "target": "state.seen"}],
}
),
NodeUse.model_validate(
{
"id": "after",
"type": "node",
"node": "work",
"input": [{"target": "value", "path": "state.seen"}],
"output": [],
}
),
],
edges=[
Edge.model_validate({"from": "outer", "outcome": "loop", "to": "inner"}),
Edge.model_validate({"from": "inner", "outcome": "loop", "to": "work"}),
Edge.model_validate({"from": "work", "outcome": "ok", "to": "inner"}),
Edge.model_validate({"from": "inner", "outcome": "done", "to": "tail"}),
Edge.model_validate({"from": "tail", "outcome": "ok", "to": "outer"}),
Edge.model_validate({"from": "outer", "outcome": "done", "to": "after"}),
Edge.model_validate({"from": "after", "outcome": "ok", "to": END}),
],
)
run = execute_workflow(
workflow,
{"items": ["a"], "inner_items": [1, 2]},
{
"work": lambda payload, _ctx: {
"outcome": "ok",
"output": {"seen": payload["value"]},
}
},
)
assert run.status == RunStatus.COMPLETED
# Inner items 1, 2 plus outer tail "a" plus after echo.
assert run.state["seen"][:3] == [1, 2, "a"]
def test_reentering_foreach_uses_fresh_activation_and_item_frames() -> None:
workflow = Workflow(
name="foreach_reentry",
input_schema=SchemaRef(type="object", properties={}),
state_schema=StateSchema.from_field_map(
{
"items": StateField(type="array"),
"count": StateField(type="integer", default=0),
"seen": StateField(
type="array", reducer=ReducerRef(name="wf.std.append")
),
}
),
output_schema=SchemaRef(type="object", properties={"seen": {"type": "array"}}),
node_defs=[
NodeDef(
name="record",
input_schema=SchemaRef(
type="object", properties={"value": {}}, required=["value"]
),
output_schema=SchemaRef(
type="object", properties={"seen": {}}, required=["seen"]
),
outcomes=["ok"],
),
NodeDef(
name="bump",
input_schema=SchemaRef(
type="object", properties={"count": {"type": "integer"}}
),
output_schema=SchemaRef(
type="object", properties={"count": {"type": "integer"}}
),
outcomes=["ok"],
),
],
start="again",
nodes=[
ConditionNode.model_validate(
{
"id": "again",
"type": "condition",
"check": {
"op": "lt",
"left": {"path": "state.count"},
"right": {"value": 2},
},
}
),
ForeachNode.model_validate(
{
"id": "each",
"type": "foreach",
"over": "state.items",
"as": "item",
"mode": "serial",
}
),
NodeUse.model_validate(
{
"id": "work",
"type": "node",
"node": "record",
"input": [{"target": "value", "path": "context.item"}],
"output": [{"source": "seen", "target": "state.seen"}],
}
),
NodeUse.model_validate(
{
"id": "bump",
"type": "node",
"node": "bump",
"input": [{"target": "count", "path": "state.count"}],
"output": [{"source": "count", "target": "state.count"}],
}
),
],
edges=[
Edge.model_validate({"from": "again", "outcome": "true", "to": "each"}),
Edge.model_validate({"from": "each", "outcome": "loop", "to": "work"}),
Edge.model_validate({"from": "work", "outcome": "ok", "to": "each"}),
Edge.model_validate({"from": "each", "outcome": "done", "to": "bump"}),
Edge.model_validate({"from": "bump", "outcome": "ok", "to": "again"}),
Edge.model_validate({"from": "again", "outcome": "false", "to": END}),
],
)
def bump(payload: dict[str, Any], _ctx: object) -> dict[str, Any]:
count = payload.get("count", 0)
assert isinstance(count, int)
return {"outcome": "ok", "output": {"count": count + 1}}
run = execute_workflow(
workflow,
{"items": ["a"]},
{
"record": lambda payload, _ctx: {
"outcome": "ok",
"output": {"seen": payload["value"]},
},
"bump": bump,
},
)
assert run.status == RunStatus.COMPLETED
assert run.state["seen"] == ["a", "a"]
item_frames = [
frame for frame in run.frames.values() if frame.kind == "foreach_iteration"
]
assert len(item_frames) == 2
owners = [item_frame_owner(frame) for frame in item_frames]
assert all(owner is not None for owner in owners)
assert owners[0] is not None and owners[1] is not None
assert owners[0].activation_id != owners[1].activation_id
assert item_frames[0].id != item_frames[1].id
# Both visits run item zero, but activation-qualified frame ids differ.
assert all(frame.id.endswith(":0") for frame in item_frames)
def test_subgraph_end_returns_to_subgraph_node_then_foreach_owner() -> None:
from wf_core import PreparedSubgraph
child = Workflow(
name="child",
input_schema=SchemaRef(type="object", properties={"value": {}}),
state_schema=StateSchema.from_field_map({"seen": StateField(type="string")}),
output_schema=SchemaRef(type="object", properties={"seen": {}}),
node_defs=[
NodeDef(
name="inner_record",
input_schema=SchemaRef(
type="object", properties={"value": {}}, required=["value"]
),
output_schema=SchemaRef(
type="object", properties={"seen": {}}, required=["seen"]
),
outcomes=["ok"],
)
],
start="inner_record",
nodes=[
NodeUse.model_validate(
{
"id": "inner_record",
"type": "node",
"node": "inner_record",
"input": [{"target": "value", "path": "input.value"}],
"output": [{"source": "seen", "target": "state.seen"}],
}
)
],
edges=[
Edge.model_validate({"from": "inner_record", "outcome": "ok", "to": END})
],
)
foreach = ForeachNode.model_validate(
{
"id": "each",
"type": "foreach",
"over": "state.items",
"as": "item",
"mode": "serial",
}
)
workflow = Workflow(
name="foreach_subgraph",
input_schema=SchemaRef(type="object", properties={"items": {"type": "array"}}),
state_schema=StateSchema.from_field_map(
{
"items": StateField(type="array"),
"seen": StateField(
type="array", reducer=ReducerRef(name="wf.std.append")
),
}
),
output_schema=SchemaRef(type="object", properties={"seen": {"type": "array"}}),
node_defs=[],
start="each",
nodes=[
foreach,
SubgraphNode.model_validate(
{
"id": "child",
"type": "subgraph",
"workflow": "child.workflow",
"input_schema": {"type": "object", "properties": {"value": {}}},
"output_schema": {"type": "object", "properties": {"seen": {}}},
"input": [{"target": "value", "path": "context.item"}],
"output": [{"source": "seen", "target": "state.seen"}],
"outcomes": ["ok"],
}
),
],
edges=[
Edge.model_validate({"from": "each", "outcome": "loop", "to": "child"}),
Edge.model_validate({"from": "child", "outcome": "ok", "to": "each"}),
Edge.model_validate({"from": "each", "outcome": "done", "to": END}),
],
)
run = execute_workflow(
workflow,
{"items": ["a", "b"]},
{},
subgraphs={
"child.workflow": PreparedSubgraph(
workflow=child,
registry={
"inner_record": lambda payload, _ctx: {"seen": payload["value"]}
},
)
},
)
assert run.status == RunStatus.COMPLETED
assert run.state["seen"] == ["a", "b"]
def test_nonlocal_runtime_return_fails_closed_when_validation_is_bypassed() -> None:
run = RunState(
workflow_name="nonlocal",
status=RunStatus.RUNNING,
workflow_input={},
state={},
frames={},
)
add_frame(
run,
ExecutionFrame(id="root", kind="workflow", node_id="outer"),
)
add_frame(
run,
ExecutionFrame(
id="outer-item",
kind="foreach_iteration",
node_id="inner",
parent_frame_id="root",
metadata={
"foreach_node_id": "outer",
"activation_id": "root:outer#0",
"loop_index": 0,
"loop_item": "a",
"loop_alias": "outer_item",
},
),
)
add_frame(
run,
ExecutionFrame(
id="inner-item",
kind="foreach_iteration",
node_id="work",
parent_frame_id="outer-item",
metadata={
"foreach_node_id": "inner",
"activation_id": "outer-item:inner#0",
"loop_index": 0,
"loop_item": 1,
"loop_alias": "inner_item",
},
),
)
run.current_frame_id = "inner-item"
run.sync_from_current_frame()
from wf_core.runtime.ops.flow import advance_frame
with pytest.raises(WorkflowExecutionError, match="non-local|ancestor|immediate"):
advance_frame(run, run.frames["inner-item"], outcome="ok", next_node_id="outer")
def test_completed_activation_cannot_consume_later_activation_result_or_wake() -> None:
"""A closed visit rejects buffered results and wake-ups from other visits."""
from wf_core.runtime.foreach_state import (
close_foreach_activation,
load_or_begin_foreach_activation,
require_foreach_activation,
)
from wf_core.runtime.scheduler import (
block_frame_on_children,
wake_parent_for_child_progress,
)
parent = ExecutionFrame(id="root", kind="workflow", node_id="each")
first = load_or_begin_foreach_activation(parent, "each", mode="concurrent")
first.barrier.next_index = 1
first.barrier.start_child(f"{first.id}:0")
from wf_core.runtime.foreach_state import save_foreach_activation
save_foreach_activation(parent, first)
close_foreach_activation(parent, first)
second = load_or_begin_foreach_activation(parent, "each", mode="concurrent")
save_foreach_activation(parent, second)
assert second.id != first.id
with pytest.raises(WorkflowExecutionError, match="closed|superseded"):
require_foreach_activation(parent, "each", first.id)
run = RunState(
workflow_name="activation_isolation",
status=RunStatus.RUNNING,
workflow_input={},
state={},
frames={parent.id: parent},
)
stale_child = ExecutionFrame(
id=f"{first.id}:0",
kind="foreach_iteration",
node_id="work",
parent_frame_id="root",
metadata={
"foreach_node_id": "each",
"activation_id": first.id,
"loop_index": 0,
"loop_item": "a",
"loop_alias": "item",
},
)
run.frames[stale_child.id] = stale_child
block_frame_on_children(run, "root", (stale_child.id,))
stale_child.status = FrameStatus.COMPLETED
with pytest.raises(WorkflowExecutionError, match="closed activation"):
wake_parent_for_child_progress(run, stale_child.id)