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
10 changes: 10 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,16 @@ All notable changes to wfpy are documented here. The format follows
describe. A value neither can write (a lock, an open file, a local class)
makes the checkpoint not resumable, naming it. The resumed run records
`resumedFrom`.
- Resuming a run at any step of its queue trace:
`wfpy run flow.py --resume-from wf-out/<run> --at-step N`. Every firing is
journalled in `run.wf-journal.jsonl` (what it took, by digest, what it put,
and the actor's state after). A resume at step N starts again from the run's
inputs and replays each actor's firings up to the step instead of running
them, then carries on live. Replaying per actor is sound with parallel
workers, where restoring a snapshot taken at a step is not. A replayed firing
whose inputs differ from the original's stops the replay (`replayStopped` in
the run record). Trace steps carry `journalSeq` and `replayed`, and the queue
trace is now written when a run fails too.

### Fixed
- A run whose tokens are arbitrary objects no longer fails writing its overlay
Expand Down
5 changes: 4 additions & 1 deletion CLAUDE.md
Original file line number Diff line number Diff line change
Expand Up @@ -63,6 +63,8 @@ src/wfpy/
├── _run_finalization.py Run finalization (overlays, artifacts, run record)
├── _checkpoint_runtime.py Run checkpoints: capture/restore a plan's state,
│ run.wf-checkpoint.json, resume_from
├── _journal_runtime.py Firing journal (run.wf-journal.jsonl) + per-actor
│ replay for resume_from + at_step
├── _workflow_instance.py Workflow instances (child(instance=...)) + run discovery
├── _context_runtime.py Shared-context (@context) runtime
├── _lsp_client.py LSP client for agent LSP integration
Expand Down Expand Up @@ -90,7 +92,8 @@ Public API is re-exported from `src/wfpy/__init__.py`.
- `cli._load_module()` prepends the workflow module dir, its package parents, and the nearest project root to `sys.path`. Prefer fixing that logic over adding per-example `sys.path` hacks.
- `RewriteEngine.export_workflow_graph()` resolves nested child workflows recursively, including `from ... import ...` imports, with recursion guards.
- `run()` defaults `queue_trace=True`; writes `run.wf-queues.json`, `run.wf-run.json`, and appends `run-log.jsonl`. `keep_intermediates=True` copies `work/` into the run output.
- A run that fails writes `run.wf-checkpoint.json` (queues, actor/control state, sub-plans, context); `run(resume_from=...)` / `wfpy run --resume-from` restores it into a freshly built plan. A firing that raises gives its tokens back (`_atomic_firing`); nested workflows and if/loop are composite, not wrapped, and mark an interrupted firing in `RuntimeActor.pending`. New per-actor runtime state must be added to `_checkpoint_runtime.capture/restore` or a resume silently loses it.
- A run that fails writes `run.wf-checkpoint.json` (queues, actor/control state, sub-plans, context); `run(resume_from=...)` / `wfpy run --resume-from` restores it into a freshly built plan. A firing that raises gives its tokens back (`_atomic_firing`); nested workflows and if/loop are composite, not wrapped, and mark an interrupted firing in `RuntimeActor.pending`. New per-actor runtime state must be added to `_checkpoint_runtime.actor_state/restore_actor_state` (shared by the checkpoint and the journal) or a resume silently loses it.
- Every wrapped firing is journalled (`run.wf-journal.jsonl`, on whenever `queue_trace` is); `run(resume_from=..., at_step=N)` replays each actor's first firings from it (`Journal.replay_firing`, called from `_step_actor`), chosen per actor by `_journal_runtime.entries_at_step` — not by one journal cut-off, which races with parallel workers. Token writes during a firing must go through `Queue.enqueue`/`dequeue` so the journal sees them.
- CLI `--agent-tools` auto-discovers `agent-tools.json` beside the workflow file. Without the flag, the registry is not auto-loaded.
- `wfpy plan --format graph --best-effort` falls back to the AST partial-graph builder in `partial.py`; plain `plan` does not.
- Do not use system `/tmp` for repo work or scratch — use the repo-local `.tmp/` directory.
Expand Down
170 changes: 91 additions & 79 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 — phases 1 (atomic firings) and 2 (checkpoint on failure, `--resume-from`) are implemented; the queue trace and the stepper are not.
**Status:** proposal — phases 1 to 3 (atomic firings, checkpoint on failure, the firing journal and `--at-step`) are implemented; the stepper's resume button 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 @@ -236,71 +236,85 @@ File tokens keep pointing into the earlier run's `work/` directory. That run
has to stay where it is, which it does: nothing deletes it. `--copy-work`
copies it into the new run when it should not have to.

## Resuming from a completed step, through the queue trace
## Resuming from a completed step: replaying a firing journal

Restoring the state **after** firing N needs that state, and the failure
checkpoint has only the last one. Writing a full checkpoint after every firing
is too expensive for a run with large tokens. The queue trace is already
written once per firing; it lacks only what changed **in** the firing.
**Implemented** (phase 3), in `_journal_runtime.py`, differently from the
first version of this section. That version extended the queue trace with
each step's consumed and produced tokens, and rebuilt the state after step N
by applying steps 1 to N. That is unsound with parallel workers. A step is
recorded when a firing completes, while other firings are in flight, so the
state "after step N" can hold half of a firing that completes later. A
consumer can even complete, and be recorded, before the firing that produced
its input. Applying steps 1 to N then takes a token that was never put.

### The trace becomes a journal (version 2)
### Replay per actor, not per state

Each step keeps what it has now and gains what the firing did:
Every firing that completes appends one line to `run.wf-journal.jsonl`:

```json
{"step": 12, "actorInstanceName": "review", "actorKind": "agent", "actorFireCount": 3,
"queueSizes": [ "…as today…" ],
"consumed": {"q:draft.Out->review.In": 1},
"produced": {"q:review.Out->wf:output:Report": [{"$wfType": "file", "path": "work/review__Out__3.md"}]},
"state": {"_round": 4},
"runtime": {"scheduleState": "done", "chatHistory": "…appended turns only…"}}
{"seq": 12, "actor": "child/review", "fireCount": 3,
"consumed": [["draft.Out-->review.In", 1, "sha256 of the tokens"]],
"produced": [["review.Out-->WF.Report", [{"$wfType": "file", "path": "…"}]]],
"state": {"_round": 4}, "chatHistory": ["…"]}
```

- `consumed` is a count per queue. The values are already in an earlier
step's `produced`, or in the workflow inputs, so they are not written twice.
- `produced` is the tokens the firing enqueued, encoded as checkpoint tokens
are.
- `state` and `runtime` are the actor's state after the firing, written only
when it changed.

A step is recorded under the scheduler's lock, in completion order. The trace
is therefore a serialization of the run, even with parallel workers.

The state after step N is computed without running anything. Start from the
inputs, then for each step up to N: dequeue `consumed`, enqueue `produced`,
and replace the actor's state. `queueSizes` is kept as it is. The stepper
still reads it, and on resume it is a check: replayed sizes that differ from
the recorded ones mean a corrupt or edited trace, and the resume is refused.

Two changes to when the trace is written:

- **On failure too.** Today it is built only in `_finalize_run`. A failed run
is the run most worth stepping through. Its last step is the last firing
that completed, and the error overlay names the one that did not.
- **Appended as the run goes**, to `run.wf-queues.jsonl`, one step per line,
then assembled into `run.wf-queues.json` at the end as now. A run killed
outright still leaves every completed step on disk, and the IDE can step
through a run that is still going.

Values make the trace larger. Version 2 is on by default only when it stays
under a size (see Open questions). `--queue-trace=sizes` keeps version 1.
Version 1 is still read: a run with a version 1 trace can be stepped through,
but resumed only from its failure checkpoint.

### Resume

```bash
wfpy run flows/pipeline.py --resume-from wf-out/<run> --at-step 12
```

`--at-step` uses the stepper's numbering, so the number a person reads in the
IDE is the number they pass.

Step N is a firing of the top-level plan. A nested workflow's firing is one
step, as the stepper shows it today: resuming lands before or after a child,
not inside it. Resuming inside a child needs the child's trace too. That is
left for later, since the failure checkpoint already covers a failure inside
a child.
A header line holds the run's inputs and graph fingerprint. Tokens and state
are encoded as checkpoints encode them.

`wfpy run flow.py --resume-from wf-out/<run> --at-step N` starts the run
again from that run's inputs. Each actor's first firings, those the run had
completed by step N, are **replayed** instead of run: the recorded tokens are
taken and put, and the recorded state is set. Then the run carries on live.

This is sound under any scheduling, because in a dataflow network each actor
sees the same sequence of tokens however its firings interleave with others'.
An actor's k-th firing can therefore be replayed whenever its inputs arrive.
It does not matter which other firings happen to be in flight, or which
completed first.

### What "by step N" means

Which firings a step covers is chosen per actor, not by one journal cut-off
(that is the same race again: another actor's firing can be journalled
between a firing's entry and its step):

- an actor the trace shows replays its entries up to its own last step at or
before N;
- an actor inside a nested workflow (path `child/leaf`) follows that
workflow's steps;
- an actor the trace never shows fires only inside an if or a loop, which run
with nothing else in flight, so the cut-off of step N is exact for it.

Each trace step carries `journalSeq`, the journal's length when it was
recorded, and `replayed: true` when the firing was replayed.

### When the run goes another way

A replayed firing checks that the tokens it takes are the ones the original
took (by digest). The firings feeding a replayed one are replayed too, so in
a straight resume they always match. They do not when something upstream runs
live and behaves differently: a firing of the race above, or a journal that
was cut short. Then replay stops, the reason goes into the run record
(`replayStopped`), and the run carries on live from there. Replay also stops if
the run goes quiet with firings left to replay, and at a firing whose values
could not be saved (`notReplayable`).

What a replayed firing returns is what it returned then. Resuming at step N
keeps everything before N as it happened, including what an agent answered.
A fix to an actor's code applies to its firings after N.

### Also

- The queue trace is now written when a run fails too: the last step is the
last firing that completed.
- A resumed run keeps its own journal, replayed firings included, so it can be
resumed at a step in turn. A run resumed from a checkpoint (phase 2) cannot,
because it did not start from inputs. Resume at a step of the run it came
from instead.
- `queue_trace=False` keeps no journal either.
- Not done: the trace is not appended as the run goes (`run.wf-queues.jsonl`).
The journal is, so a run killed outright can still be resumed at a step if
its trace survives, which today it does not.

## The stepper resumes

Expand All @@ -312,9 +326,11 @@ cluster:
```

**⟲ Resume from here** runs `--resume-from <that run> --at-step <shown
step>`, through the profile's run driver as ▶ Run does. It is enabled when the
run's trace is version 2. On a failed run's last step it resumes from the
failure checkpoint, so "fix, then ⟲" is the whole loop.
step>`, through the profile's run driver as ▶ Run does. The number shown is
the trace's `step`, which is the number `--at-step` takes. It is enabled when
the run's trace steps carry `journalSeq`. On a failed run's last step it
replays every firing that completed and runs the failed one live, so "fix,
then ⟲" is the whole loop.

On a node, **Rerun from here** finds the step of that node's last firing and
resumes from the step before it, so the node fires again.
Expand All @@ -338,20 +354,16 @@ flags.

## Phasing

1. **Atomic firings.** Peek, run, then commit, for every actor kind. This is
useful alone: a failure no longer loses a token. Tests: each kind fails
once, and its inputs are still queued afterwards.
2. **Checkpoint on failure, and `--resume-from`.** Serialization, graph
fingerprint, restore. Tests: a workflow that fails at a known firing is
fixed and resumed. Firings before the failure do not run again (counted),
and the outputs equal those of a clean run. Also a failure inside a nested
workflow, inside a loop, and in an agent with chat history.
3. **Queue trace version 2.** Written on failure, appended as it goes,
`consumed`, `produced` and state per step, and `--at-step`. Tests: for
every step of a run, the replayed state equals the state captured live at
that step. The stepper keeps working on version 1 and version 2 traces.
1. **Atomic firings.** *Done.*
2. **Checkpoint on failure, and `--resume-from`.** *Done.*
3. **The firing journal, and `--at-step`.** *Done.* Tests resume at every
step of a run, with one worker and with parallel ones. They check that the
firings before the step are replayed and not run, that those after it run,
and that the outputs equal a clean run's. Also covered: a failed run at its
last step, a loop over a child workflow, a resumed run resumed again,
replay stopping when inputs differ, and the CLI.
4. **The stepper's ⟲ button**, and "Rerun from here" on a node (dialogram +
wfpy-ide).
wfpy-ide). The button passes the step it shows to `--at-step`.

## Open questions

Expand All @@ -362,7 +374,7 @@ flags.
stateless agent asked again after a resume may answer differently than it
would have the first time. Is that acceptable, or should a resume replay the
recorded answer for a firing that is in the journal?
- **Trace size.** Version 2 writes every produced token, so a large non-file
token (a big string, a list) appears inline in the step that produced it.
What threshold sends it to a side file instead, as File outputs already
are? And which default: version 2 always, or only when asked for?
- **Journal size.** The journal writes every produced token, so a large
non-file token (a big string, a list) appears inline in the firing that
produced it. What threshold sends it to a side file instead, as File outputs
already are? It is on whenever the queue trace is (the default).
52 changes: 35 additions & 17 deletions src/wfpy/_checkpoint_runtime.py
Original file line number Diff line number Diff line change
Expand Up @@ -299,6 +299,39 @@ def _instance_state(actor: Any) -> dict[str, Any]:
return {name: getattr(actor.instance, name) for name in names if hasattr(actor.instance, name)}


def actor_state(actor: Any, enc: Any, at: str) -> dict[str, Any]:
"""An actor's own state -- not a control node's -- written with *enc*.

Its task fields and the runtime attributes wfpy keeps on the instance, and
an agent's budget, conversation and CLI sessions. Shared by the checkpoint
and the firing journal, so both carry the same state.
"""
entry: dict[str, Any] = {
"state": {
name: enc(value, f"{at}.{name}") for name, value in _instance_state(actor).items()
}
}
if actor.agent_fire_budget is not None:
entry["agentFireBudget"] = actor.agent_fire_budget
if actor.chat_history:
entry["chatHistory"] = enc(actor.chat_history, f"{at}.chatHistory")
if actor.agent_cli_session_ids:
entry["agentCliSessionIds"] = dict(actor.agent_cli_session_ids)
return entry


def restore_actor_state(actor: Any, entry: dict[str, Any]) -> None:
"""Put back what :func:`actor_state` wrote."""
for field, value in (entry.get("state") or {}).items():
setattr(actor.instance, field, decode(value))
if "agentFireBudget" in entry:
actor.agent_fire_budget = entry["agentFireBudget"]
if "chatHistory" in entry:
actor.chat_history = decode(entry["chatHistory"])
if "agentCliSessionIds" in entry:
actor.agent_cli_session_ids = dict(entry["agentCliSessionIds"])


def capture(plan: Any, problems: list[str], where: str = "") -> dict[str, Any]:
"""The state of *plan* and its sub-plans; what cannot be saved goes to *problems*."""

Expand Down Expand Up @@ -333,15 +366,7 @@ def enc(value: Any, at: str) -> Any:
control["loopIndex"] = actor._loop_index
entry["control"] = control
else:
entry["state"] = {
name: enc(value, f"{at}.{name}") for name, value in _instance_state(actor).items()
}
if actor.agent_fire_budget is not None:
entry["agentFireBudget"] = actor.agent_fire_budget
if actor.chat_history:
entry["chatHistory"] = enc(actor.chat_history, f"{at}.chatHistory")
if actor.agent_cli_session_ids:
entry["agentCliSessionIds"] = dict(actor.agent_cli_session_ids)
entry.update(actor_state(actor, enc, at))
if actor.sub_plan is not None:
entry["subPlan"] = capture(actor.sub_plan, problems, f"{at}/")
actors[actor.name] = entry
Expand Down Expand Up @@ -379,14 +404,7 @@ def restore(plan: Any, saved: dict[str, Any]) -> None:
actor._loop_source = source
actor._loop_index = control["loopIndex"]
actor._loop_iter = iterator
for field, value in (entry.get("state") or {}).items():
setattr(actor.instance, field, decode(value))
if "agentFireBudget" in entry:
actor.agent_fire_budget = entry["agentFireBudget"]
if "chatHistory" in entry:
actor.chat_history = decode(entry["chatHistory"])
if "agentCliSessionIds" in entry:
actor.agent_cli_session_ids = dict(entry["agentCliSessionIds"])
restore_actor_state(actor, entry)
if "subPlan" in entry and actor.sub_plan is not None:
restore(actor.sub_plan, entry["subPlan"])

Expand Down
Loading
Loading