diff --git a/examples/authoring_concurrent_foreach.py b/examples/authoring_concurrent_foreach.py index cb2f701d..46959f59 100644 --- a/examples/authoring_concurrent_foreach.py +++ b/examples/authoring_concurrent_foreach.py @@ -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"]}) diff --git a/examples/demo_workflow.py b/examples/demo_workflow.py index d94a49a1..58f8343d 100644 --- a/examples/demo_workflow.py +++ b/examples/demo_workflow.py @@ -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"}, diff --git a/examples/raw_concurrent_foreach.py b/examples/raw_concurrent_foreach.py index 15cb2265..50e744a2 100644 --- a/examples/raw_concurrent_foreach.py +++ b/examples/raw_concurrent_foreach.py @@ -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}, ], diff --git a/src/wf_core/runtime/ops/flow.py b/src/wf_core/runtime/ops/flow.py index 7950aebb..2f50a70c 100644 --- a/src/wf_core/runtime/ops/flow.py +++ b/src/wf_core/runtime/ops/flow.py @@ -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" diff --git a/src/wf_core/runtime/ops/foreach.py b/src/wf_core/runtime/ops/foreach.py index b45b15ef..dec9178b 100644 --- a/src/wf_core/runtime/ops/foreach.py +++ b/src/wf_core/runtime/ops/foreach.py @@ -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 diff --git a/src/wf_core/runtime/step.py b/src/wf_core/runtime/step.py index ddea7f87..1da1491c 100644 --- a/src/wf_core/runtime/step.py +++ b/src/wf_core/runtime/step.py @@ -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 diff --git a/src/wf_core/runtime/subgraphs.py b/src/wf_core/runtime/subgraphs.py index 3292586e..013950c0 100644 --- a/src/wf_core/runtime/subgraphs.py +++ b/src/wf_core/runtime/subgraphs.py @@ -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, diff --git a/tests/authoring/test_demo_workflow.py b/tests/authoring/test_demo_workflow.py index 79c725ce..a14de31c 100644 --- a/tests/authoring/test_demo_workflow.py +++ b/tests/authoring/test_demo_workflow.py @@ -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) diff --git a/tests/core/test_concurrent_foreach.py b/tests/core/test_concurrent_foreach.py index c73afdb6..936f8ad0 100644 --- a/tests/core/test_concurrent_foreach.py +++ b/tests/core/test_concurrent_foreach.py @@ -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}), ], ) diff --git a/tests/core/test_concurrent_foreach_async.py b/tests/core/test_concurrent_foreach_async.py index 473bfdbd..f07abdb5 100644 --- a/tests/core/test_concurrent_foreach_async.py +++ b/tests/core/test_concurrent_foreach_async.py @@ -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}), ], ) diff --git a/tests/core/test_concurrent_foreach_errors.py b/tests/core/test_concurrent_foreach_errors.py index 18c1e3a4..a8d087e2 100644 --- a/tests/core/test_concurrent_foreach_errors.py +++ b/tests/core/test_concurrent_foreach_errors.py @@ -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( { diff --git a/tests/core/test_concurrent_foreach_interrupts.py b/tests/core/test_concurrent_foreach_interrupts.py index 3463e473..009f488f 100644 --- a/tests/core/test_concurrent_foreach_interrupts.py +++ b/tests/core/test_concurrent_foreach_interrupts.py @@ -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}), ], ) diff --git a/tests/core/test_foreach_back_edges.py b/tests/core/test_foreach_back_edges.py new file mode 100644 index 00000000..8a2b969e --- /dev/null +++ b/tests/core/test_foreach_back_edges.py @@ -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)