Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
33 changes: 25 additions & 8 deletions docs/proposals/resume.md
Original file line number Diff line number Diff line change
@@ -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).

Expand Down Expand Up @@ -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,
Expand Down
74 changes: 72 additions & 2 deletions src/wfpy/runner.py
Original file line number Diff line number Diff line change
Expand Up @@ -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."""

Expand Down Expand Up @@ -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."""
Expand All @@ -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()})"
Expand Down Expand Up @@ -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":
Expand Down
Loading
Loading