Skip to content

Commit e5a4223

Browse files
author
Markus
committed
fix(workflows): let interrupts during checkpoints pause the run
RunState.save() caught BaseException around the checkpoint write, so a KeyboardInterrupt during any checkpoint became CheckpointError, marked the state as failed to checkpoint, and bypassed the KeyboardInterrupt handler in execute()/resume(). The run stayed "running" and could not be resumed. On main the interrupt reaches that handler and the run pauses. Introduced in 7aa424d. Convert only Exception to CheckpointError. Write failures still stop further writes; interrupts propagate unchanged, and the pause re-saves the tree from memory. Tests interrupt a checkpoint at the root, inside a workflow call, and during resume; all failed before this change. Assisted-by: OpenCode (model: claude-opus-5.5, autonomous)
1 parent 3e5e66e commit e5a4223

2 files changed

Lines changed: 84 additions & 1 deletion

File tree

‎src/specify_cli/workflows/engine.py‎

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -774,7 +774,10 @@ def save(self) -> None:
774774
try:
775775
self._atomic_write_json(runs_dir / "state.json", state_data)
776776
self._atomic_write_json(runs_dir / "inputs.json", {"inputs": self.inputs})
777-
except BaseException as exc:
777+
except Exception as exc:
778+
# Only write failures poison the instance. A graceful interrupt
779+
# (KeyboardInterrupt) propagates unchanged so execute()/resume()
780+
# pause the run, as on main; the pause re-saves from memory.
778781
from ._execution import CheckpointError
779782

780783
self._checkpoint_failed = True

‎tests/workflows/test_composition_execution.py‎

Lines changed: 80 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -681,6 +681,86 @@ def write(path, data):
681681
assert probe["work"] == 1
682682

683683

684+
def _interrupt_state_write(monkeypatch, when):
685+
"""Raise ``KeyboardInterrupt`` from the first ``state.json`` write matching *when*."""
686+
original = RunState._atomic_write_json
687+
fired = []
688+
689+
def write(path, data):
690+
if path.name == "state.json" and not fired and when(data):
691+
fired.append(path)
692+
raise KeyboardInterrupt
693+
original(path, data)
694+
695+
monkeypatch.setattr(RunState, "_atomic_write_json", staticmethod(write))
696+
return fired
697+
698+
699+
def _logged_events(state):
700+
log = (state.runs_dir / "log.jsonl").read_text(encoding="utf-8").splitlines()
701+
return [json.loads(line)["event"] for line in log]
702+
703+
704+
@pytest.mark.parametrize("scope", ["root", "workflow-call"])
705+
def test_interrupt_during_checkpoint_pauses_run(tmp_path, monkeypatch, probe, scope):
706+
steps = [{"id": "first", "type": "probe"}, {"id": "second", "type": "probe"}]
707+
if scope == "workflow-call":
708+
install(tmp_path, definition("child", steps))
709+
steps = [call()]
710+
# Interrupt the checkpoint that marks ``second`` active, after ``first`` ran.
711+
fired = _interrupt_state_write(
712+
monkeypatch, lambda data: data.get("current_step_id") == "second"
713+
)
714+
715+
state = WorkflowEngine(tmp_path).execute(definition("parent", steps))
716+
717+
# A graceful interrupt is not a checkpoint failure: it reaches the engine's
718+
# pause path, as on main, even inside a called workflow.
719+
assert fired
720+
assert state.status == RunStatus.PAUSED
721+
assert RunState.load(state.run_id, tmp_path).status == RunStatus.PAUSED
722+
assert _logged_events(state)[-1] == "workflow_interrupted"
723+
assert probe == {"first": 1}
724+
725+
state = WorkflowEngine(tmp_path).resume(state.run_id)
726+
727+
assert state.status == RunStatus.COMPLETED
728+
assert probe == {"first": 1, "second": 1}
729+
730+
731+
def test_interrupt_during_resume_checkpoint_pauses_run(tmp_path, monkeypatch, probe):
732+
root = definition(
733+
"parent",
734+
[
735+
{"id": "wait", "type": "probe", "await": True},
736+
{"id": "next", "type": "probe"},
737+
],
738+
inputs={"approve": {"type": "boolean", "default": False}},
739+
)
740+
state = WorkflowEngine(tmp_path).execute(root)
741+
assert state.status == RunStatus.PAUSED
742+
# Interrupt the checkpoint that records ``wait`` as completed on resume.
743+
fired = _interrupt_state_write(
744+
monkeypatch,
745+
lambda data: data.get("step_results", {}).get("wait", {}).get("status")
746+
== "completed",
747+
)
748+
749+
state = WorkflowEngine(tmp_path).resume(state.run_id, {"approve": True})
750+
751+
assert fired
752+
assert state.status == RunStatus.PAUSED
753+
assert RunState.load(state.run_id, tmp_path).status == RunStatus.PAUSED
754+
assert _logged_events(state)[-1] == "workflow_interrupted"
755+
756+
state = WorkflowEngine(tmp_path).resume(state.run_id)
757+
758+
# The interrupted checkpoint's transition is saved by the pause, so the
759+
# completed ``wait`` is not run a third time.
760+
assert state.status == RunStatus.COMPLETED
761+
assert probe == {"wait": 2, "next": 1}
762+
763+
684764
def test_checkpoint_failure_stops_concurrent_fan_out_logging(tmp_path, monkeypatch):
685765
import specify_cli.workflows._execution as execution
686766
from specify_cli.workflows._execution import CheckpointError

0 commit comments

Comments
 (0)