diff --git a/CHANGELOG.md b/CHANGELOG.md index 6069b5a..bc598bc 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -22,6 +22,12 @@ All notable changes to wfpy are documented here. The format follows that qualify, `createNode` takes `importFrom` for a workflow defined in another module, and the graph export marks an instance with an `instance` annotation. +- A firing that fails gives back the tokens it took: each goes back to the + head of its queue, in order, for an action, a tool, an agent and a + StreamBlocks instance. A nested workflow and an if / loop are made of other + actors' firings, each atomic on its own. A failed run's queues are + left as they were before the failing firing, the first step towards resuming + a run (`docs/proposals/resume.md`). ## [1.0.0] — 2026-07-24 diff --git a/docs/proposals/resume.md b/docs/proposals/resume.md index e488365..90f2025 100644 --- a/docs/proposals/resume.md +++ b/docs/proposals/resume.md @@ -1,6 +1,6 @@ # Proposal: resuming a run -**Status:** proposal — nothing here is implemented. +**Status:** proposal — phase 1 (atomic firings) is implemented; the rest is not. **Affects:** wfpy (the runtime and the queue trace); dialogram (the queue-trace stepper gets a resume button); wfpy-ide (one command id). @@ -85,13 +85,30 @@ Today a firing removes its tokens before it runs: dequeue on entry. When the action raises, its inputs are gone, and no checkpoint taken after the -failure can bring them back. So the first change is **atomic firings**: peek -the inputs, run, and dequeue only once the firing succeeded. - -That is safe as it stands. One actor never fires twice at once (the scheduler -keeps `running_actors`), and every queue has exactly one consumer, so nothing -else can take a peeked token in between. An internal action's guard already -works this way; the commit moves from after the guard to after the action. +failure can bring them back. So the first change is **atomic firings**: a +firing that fails gives back what it took. + +**Implemented** (phase 1). Rather than reordering each kind's step, every +`Queue.dequeue()` made during a firing is logged per thread +(`_atomic_firing` in `runner.py`, around `_step_actor`). When the firing +raises, each token goes back to the head of its queue, in reverse order, so +the queue is as it was. The same log is what phase 3 records as a step's +`consumed`. + +An `if` or a `loop` is not wrapped, any more than a nested workflow is. Its +firing runs its branch to quiescence, and each firing in the branch is atomic +on its own. Giving back the condition or the iterable after part of the +branch had run would run that part twice. Where a control node stopped is +its state, recorded by the checkpoint. + +That is safe without holding the queues for the whole firing. One actor never +fires twice at once (the scheduler keeps `running_actors`), and every queue +has exactly one consumer, so nothing else takes from the head meanwhile. A +producer appending to the tail is unaffected. + +What is not given back is the actor's own state. An action that set +`self._done = True` and then raised has still set it. Phase 2's checkpoint +records the state as the failure left it. A nested workflow is the exception: its tokens are handed to the sub-plan, which may consume them before failing. Its checkpoint is the sub-plan's own, diff --git a/src/wfpy/runner.py b/src/wfpy/runner.py index 1550b30..1bbf5fc 100644 --- a/src/wfpy/runner.py +++ b/src/wfpy/runner.py @@ -253,6 +253,48 @@ def _attach_agent_response_meta( # ═══════════════════════════════════════════════════════════════════════════ +# The tokens the firing running on this thread has taken, in the order it took +# them. Set only while `_atomic_firing` holds it; a dequeue anywhere else is not +# part of a firing and is not logged. +_firing_log = threading.local() + + +def _log_taken(queue: "Queue", values: list[Any]) -> None: + log: list[tuple[Queue, Any]] | None = getattr(_firing_log, "taken", None) + if log is not None: + log.extend((queue, value) for value in values) + + +class _atomic_firing: + """Give a firing's inputs back when it fails. + + Every actor kind takes its tokens before it runs -- an action after its + guard, a tool or an agent on entry -- so a firing that raised had already + lost them, and nothing saved after the failure could bring them back. This + puts each one back at the head of its queue, in order, before the error + propagates: the run stops where the firing had not happened. + + Safe without holding the queues for the whole firing: one actor never + fires twice at once, and every queue has one consumer, so nothing else + takes from the head meanwhile. A producer appending to the tail is + unaffected. + """ + + def __enter__(self) -> "_atomic_firing": + self._outer = getattr(_firing_log, "taken", None) + _firing_log.taken = [] + return self + + def __exit__(self, exc_type: Any, exc: Any, tb: Any) -> None: + taken: list[tuple[Queue, Any]] = _firing_log.taken + _firing_log.taken = self._outer + # A firing that succeeded keeps what it took, whatever fails after it. + if exc_type is None: + return + for queue, value in reversed(taken): + queue.give_back(value) + + class Queue: """Unbounded FIFO token queue on one connection edge. Thread-safe.""" @@ -291,7 +333,14 @@ def enqueue(self, value: Any) -> None: def dequeue(self) -> Any: with self._lock: - return self.items.popleft() + value = self.items.popleft() + _log_taken(self, [value]) + return value + + def give_back(self, value: Any) -> None: + """Return a token taken by a firing that failed, to the head.""" + with self._lock: + self.items.appendleft(value) def peek(self, n: int = 1) -> list[Any]: """Return up to *n* items without removing them.""" @@ -314,7 +363,9 @@ def try_dequeue(self, n: int = 1) -> list[Any] | None: with self._lock: if len(self.items) < n: return None - return [self.items.popleft() for _ in range(n)] + values = [self.items.popleft() for _ in range(n)] + _log_taken(self, values) + return values def __repr__(self) -> str: return f"Queue({self.id!r}, len={self.size()})" @@ -1126,6 +1177,25 @@ def _step_actor( return False if actor.kind.startswith("control-") and actor._control_running: return False + # A nested workflow and an if / loop are composite: one firing of theirs + # runs other actors' firings to completion, and each of those is atomic on + # its own. Giving back theirs as well would hand a token to a firing that + # already consumed it, and its sub-plan or branch would see it twice. + if actor.kind == "workflow": + return _step_workflow(actor, plan, out_dir, verbose) + if actor.kind in ("control-if", "control-loop"): + return _dispatch_step(actor, plan, out_dir, verbose, active_scopes) + with _atomic_firing(): + return _dispatch_step(actor, plan, out_dir, verbose, active_scopes) + + +def _dispatch_step( + actor: RuntimeActor, + plan: FifoPlan, + out_dir: Path, + verbose: bool, + active_scopes: set[str] | None, +) -> bool: if actor.kind == "internal": return _step_internal(actor, out_dir, plan, verbose) elif actor.kind == "external": diff --git a/tests/test_atomic_firings.py b/tests/test_atomic_firings.py new file mode 100644 index 0000000..cb46d09 --- /dev/null +++ b/tests/test_atomic_firings.py @@ -0,0 +1,251 @@ +"""A firing that fails keeps its inputs. + +Every actor kind takes its tokens before it runs, so a firing that raised had +lost them, and a run could not be resumed from where it stopped. Now each one +is given back to the head of its queue, in order, before the error propagates. +""" + +from __future__ import annotations + +import sys +from pathlib import Path +from typing import Any + +import pytest + +from wfpy import Port, action, agent, connect, guard, loop, task, tool, workflow +from wfpy.runner import Queue, _atomic_firing, _build_workflow_graph, build_plan, execute_plan + + +@task +class Count: + """Emits 1, 2, 3, one per firing.""" + + _n: int = 0 + + class Ports: + Out = Port[int](direction="out") + + @action(consumes={}, produces={"Out": 1}) + @guard(lambda self: self._n < 3) + def emit(self) -> int: + self._n += 1 + return self._n + + +def _plan(wf: Any) -> Any: + wf_def = wf._wfpy_workflow + return build_plan(_build_workflow_graph(wf_def), wf_def) + + +def _queued(plan: Any, actor_name: str, port: str) -> list[Any]: + actor = next(a for a in plan.actors if a.name == actor_name) + return list(actor.in_queues[port][0].items) + + +def _fail(plan: Any, tmp_path: Path, **kwargs: Any) -> None: + with pytest.raises(Exception): + execute_plan(plan, out_dir=tmp_path, max_workers=1, **kwargs) + + +def test_a_failed_action_gives_its_token_back_and_keeps_the_ones_before(tmp_path: Path) -> None: + @task + class FailsOnTwo: + class Ports: + In = Port[int](direction="in") + Out = Port[int](direction="out") + + @action(consumes={"In": 1}, produces={"Out": 1}) + def go(self, x: int) -> int: + if x == 2: + raise ValueError("two") + return x + + @workflow(outputs={"Out": int}) + def wf(): + c = Count() + f = FailsOnTwo() + connect(c.Out, f.In) + connect(f.Out, "Out") + + plan = _plan(wf) + _fail(plan, tmp_path) + + # 1 went through; 2 failed and is back at the head, ahead of 3. + assert _queued(plan, "f", "In")[:2] == [2, 3] + assert list(plan.wf_output_queues["Out"][0].items) == [1] + + +def test_a_multi_token_firing_gives_them_back_in_order(tmp_path: Path) -> None: + @task + class Pairs: + class Ports: + In = Port[int](direction="in") + Out = Port[int](direction="out") + + @action(consumes={"In": 2}, produces={"Out": 1}) + def go(self, a: int, b: int) -> int: + raise ValueError("no pairs") + + @workflow(outputs={"Out": int}) + def wf(): + c = Count() + p = Pairs() + connect(c.Out, p.In) + connect(p.Out, "Out") + + plan = _plan(wf) + _fail(plan, tmp_path) + + assert _queued(plan, "p", "In")[:2] == [1, 2] + + +def test_a_failed_tool_gives_its_token_back(tmp_path: Path) -> None: + @tool(cmd=sys.executable, args=["-c", "import sys; sys.exit(3)", "{in.In}"]) + class Fails: + class Ports: + In = Port[str](direction="in") + Out = Port[str](direction="out") + + @workflow(inputs={"In": str}, outputs={"Out": str}) + def wf(): + f = Fails() + connect("In", f.In) + connect(f.Out, "Out") + + plan = _plan(wf) + _fail(plan, tmp_path, inputs={"In": "payload"}) + + assert _queued(plan, "f", "In") == ["payload"] + + +def test_a_failed_agent_gives_its_token_back( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch +) -> None: + def boom(*_a: Any, **_k: Any) -> Any: + raise RuntimeError("the model is down") + + monkeypatch.setattr("wfpy.runner._invoke_agent", boom) + + @agent(prompt="x", useSkill=False) + class Writer: + class Ports: + In = Port[str](direction="in") + Out = Port[str](direction="out") + + @workflow(inputs={"In": str}, outputs={"Out": str}) + def wf(): + w = Writer() + connect("In", w.In) + connect(w.Out, "Out") + + plan = _plan(wf) + _fail(plan, tmp_path, inputs={"In": "brief"}) + + assert _queued(plan, "w", "In") == ["brief"] + + +def test_a_failure_inside_a_child_keeps_the_token_in_the_child_only(tmp_path: Path) -> None: + @task + class Fails: + class Ports: + In = Port[int](direction="in") + Out = Port[int](direction="out") + + @action(consumes={"In": 1}, produces={"Out": 1}) + def go(self, x: int) -> int: + raise ValueError("inside") + + @workflow(inputs={"In": int}, outputs={"Out": int}) + def child(): + f = Fails() + connect("In", f.In) + connect(f.Out, "Out") + + @workflow(inputs={"In": int}, outputs={"Out": int}) + def parent(): + c = child() + connect("In", c.In) + connect(c.Out, "Out") + + plan = _plan(parent) + _fail(plan, tmp_path, inputs={"In": 5}) + + nested = next(a for a in plan.actors if a.name == "c") + # Held once: by the child actor that failed, not also by the parent. + assert _queued(nested.sub_plan, "f", "In") == [5] + assert _queued(plan, "c", "In") == [] + + +def test_a_successful_firing_keeps_what_it_took() -> None: + q = Queue("q") + for value in (1, 2, 3): + q.enqueue(value) + + with _atomic_firing(): + assert q.dequeue() == 1 + assert q.try_dequeue(1) == [2] + + assert list(q.items) == [3] + + +def test_a_firing_that_fails_gives_back_across_queues() -> None: + a, b = Queue("a"), Queue("b") + for value in (1, 2): + a.enqueue(value) + b.enqueue("x") + + with pytest.raises(RuntimeError): + with _atomic_firing(): + a.dequeue() + b.dequeue() + a.dequeue() + a.enqueue(9) # produced meanwhile, at the tail + raise RuntimeError("fail") + + assert list(a.items) == [1, 2, 9] + assert list(b.items) == ["x"] + + +def test_inside_a_loop_only_the_failed_firing_gives_back(tmp_path: Path) -> None: + # A loop's firing runs its body: `a` succeeds on item 2, then `b` fails on + # it. Only `b`'s token goes back; `a` consumed its own and produced from + # it, so handing it back as well would run `a` on item 2 twice. + @task + class Pass: + class Ports: + In = Port[int](direction="in") + Out = Port[int](direction="out") + + @action(consumes={"In": 1}, produces={"Out": 1}) + def go(self, x: int) -> int: + return x + + @task + class FailsOnTwo: + class Ports: + In = Port[int](direction="in") + Out = Port[int](direction="out") + + @action(consumes={"In": 1}, produces={"Out": 1}) + def go(self, x: int) -> int: + if x == 2: + raise ValueError("two") + return x + + @workflow(outputs={"Out": int}) + def wf(): + lp = loop([1, 2, 3]) + with lp: + a = Pass() + b = FailsOnTwo() + connect(lp.item, a.In) + connect(a.Out, b.In) + connect(b.Out, "Out") + + plan = _plan(wf) + _fail(plan, tmp_path) + + assert _queued(plan, "b", "In") == [2] + assert _queued(plan, "a", "In") == [] + assert list(plan.wf_output_queues["Out"][0].items) == [1]