[aidan] fix/schedule: harden run recovery and storage writes against crashes

This commit is contained in:
abccodes
2026-06-23 02:56:15 -07:00
parent 4c62b8a336
commit 3cf3a820ad
4 changed files with 68 additions and 44 deletions
+26 -24
View File
@@ -25,7 +25,7 @@ _running: dict[str, str] = {}
_running_lock = asyncio.Lock()
def _ran_late(started_at: datetime, scheduled_for: datetime) -> bool:
def p_ran_late(started_at: datetime, scheduled_for: datetime) -> bool:
"""Late means the run STARTED well after its slot (app was closed, event
loop backed up), not that it ran long. Measured from started_at so a
punctual run that simply takes a while isn't mislabeled. Both sides
@@ -215,30 +215,32 @@ async def execute(
return run
_running[wf.id] = run.id
wf.last_run_at = run.started_at
wf.last_run_status = "running"
wf.last_run_id = run.id
_persist_run_fields(wf, {
"last_run_at": run.started_at,
"last_run_status": "running",
"last_run_id": run.id,
})
# Announce the run as running the instant it claims execution, not at the
# first step. Without this a run that fails fast (e.g. no runnable steps) or
# hasn't streamed yet never hits the Home "Ongoing runs" list, so a batch of
# missed re-runs only ever surfaced whichever one reached its first step.
try:
from backend.apps.agents.core.ws_manager import ws_manager as _wsm_start
await _wsm_start.broadcast_global("workflow:run", {
"workflow_id": wf.id,
"run": run.model_dump(mode="json"),
})
except Exception:
pass
session = None
try:
wf.last_run_at = run.started_at
wf.last_run_status = "running"
wf.last_run_id = run.id
_persist_run_fields(wf, {
"last_run_at": run.started_at,
"last_run_status": "running",
"last_run_id": run.id,
})
# Announce the run as running the instant it claims execution, not at
# the first step. Without this a run that fails fast (e.g. no runnable
# steps) or hasn't streamed yet never hits the Home "Ongoing runs" list.
# Both this and the persist above sit inside the try whose finally frees
# _running, so a persist/broadcast failure can't strand the workflow as
# permanently "running" (which would block every future fire).
try:
from backend.apps.agents.core.ws_manager import ws_manager as _wsm_start
await _wsm_start.broadcast_global("workflow:run", {
"workflow_id": wf.id,
"run": run.model_dump(mode="json"),
})
except Exception:
pass
steps = [s for s in wf.steps if s.enabled and s.text and s.text.strip()]
if not steps:
raise ValueError("Workflow has no steps")
@@ -381,7 +383,7 @@ async def execute(
run.status = "failure"
run.error = step_error
wf.last_run_status = "failure"
elif scheduled_for is not None and _ran_late(run.started_at, scheduled_for):
elif scheduled_for is not None and p_ran_late(run.started_at, scheduled_for):
run.status = "ran_late"
wf.last_run_status = "ran_late"
else:
+9 -4
View File
@@ -113,7 +113,7 @@ def is_schedule_configured(sched: ScheduleConfig) -> bool:
return True
def _first_after(anchor: datetime, ref: datetime, step: timedelta) -> datetime:
def p_first_after(anchor: datetime, ref: datetime, step: timedelta) -> datetime:
"""First instant on the grid {anchor + k*step} strictly after ref."""
if anchor > ref:
return anchor
@@ -139,12 +139,12 @@ def _next_fire_after(
if sched.repeat_unit == "minute":
step = max(15, sched.repeat_every)
grid = anchor_local.replace(second=0, microsecond=0)
return _first_after(grid, ref_local, timedelta(minutes=step)).astimezone(timezone.utc)
return p_first_after(grid, ref_local, timedelta(minutes=step)).astimezone(timezone.utc)
if sched.repeat_unit == "hour":
step = max(1, sched.repeat_every)
grid = anchor_local.replace(minute=sched.minute, second=0, microsecond=0)
return _first_after(grid, ref_local, timedelta(hours=step)).astimezone(timezone.utc)
return p_first_after(grid, ref_local, timedelta(hours=step)).astimezone(timezone.utc)
candidate = base.replace(hour=sched.hour, minute=sched.minute)
@@ -328,6 +328,11 @@ async def _fire(wf: Workflow, scheduled_for: Optional[datetime]) -> None:
def _seconds_until_next() -> float:
# While globally paused, _tick no-ops and never rolls next_run_at forward,
# so an overdue slot would otherwise spin this loop at the 1s floor. Resume
# calls kick(), so idling the full interval here costs nothing.
if storage.get_paused():
return 60.0
now_utc = datetime.now(timezone.utc)
soonest: Optional[datetime] = None
for wf in storage.list_workflows():
@@ -372,7 +377,7 @@ def _mark_stuck_runs_failed() -> None:
storage.update_run(
r.id,
status="failure",
error="OpenSwarm closed before this run finished.",
error="Interrupted: OpenSwarm or your computer shut down before this run finished.",
finished_at=now,
)
# The run row is fixed, but the workflow still summarizes this
+27 -10
View File
@@ -11,6 +11,7 @@ last_run_* / next_run_at summary fields; full history lives in the runs file.
import json
import os
import tempfile
from threading import Lock
from typing import Optional
@@ -65,6 +66,27 @@ def _runs_path(wid: str) -> str:
return os.path.join(RUNS_DIR, f"{wid}.json")
def p_atomic_write_json(path: str, data: object) -> None:
"""Write JSON crash-safely: a power-off mid-write must not leave a truncated
file, because the loader drops a workflow whose JSON fails to parse (the
record would silently vanish). Write a unique sibling temp file, fsync it,
then os.replace (atomic on POSIX and Windows) so readers only ever see the
old complete file or the new complete one."""
fd, tmp = tempfile.mkstemp(dir=os.path.dirname(path), suffix=".tmp")
try:
with os.fdopen(fd, "w") as f:
json.dump(data, f, indent=2)
f.flush()
os.fsync(f.fileno())
os.replace(tmp, path)
except Exception:
try:
os.unlink(tmp)
except OSError:
pass
raise
def _load_all_from_disk() -> None:
global _cache_loaded, _paused
_ensure_dirs()
@@ -144,8 +166,7 @@ def save_workflow(wf: Workflow) -> Workflow:
with _io_lock:
_ensure_dirs()
_workflow_cache[wf.id] = wf
with open(_wf_path(wf.id), "w") as f:
json.dump(wf.model_dump(mode="json"), f, indent=2)
p_atomic_write_json(_wf_path(wf.id), wf.model_dump(mode="json"))
return wf
@@ -197,8 +218,7 @@ def record_run(run: WorkflowRun) -> WorkflowRun:
# Bound the per-workflow history to keep disk + memory cheap.
if len(arr) > RUNS_PER_WORKFLOW:
del arr[: len(arr) - RUNS_PER_WORKFLOW]
with open(_runs_path(run.workflow_id), "w") as f:
json.dump([r.model_dump(mode="json") for r in arr], f, indent=2)
p_atomic_write_json(_runs_path(run.workflow_id), [r.model_dump(mode="json") for r in arr])
return run
@@ -213,14 +233,12 @@ def set_paused(value: bool) -> bool:
with _io_lock:
_ensure_dirs()
_paused = bool(value)
with open(PAUSED_FILE, "w") as f:
json.dump({"paused": _paused}, f)
p_atomic_write_json(PAUSED_FILE, {"paused": _paused})
return _paused
def _write_missed() -> None:
with open(MISSED_FILE, "w") as f:
json.dump([m.model_dump(mode="json") for m in _missed_cache], f, indent=2)
p_atomic_write_json(MISSED_FILE, [m.model_dump(mode="json") for m in _missed_cache])
def list_missed() -> list[MissedRun]:
@@ -272,7 +290,6 @@ def update_run(run_id: str, **fields) -> Optional[WorkflowRun]:
updated = r.model_copy(update=fields)
arr[i] = updated
with _io_lock:
with open(_runs_path(updated.workflow_id), "w") as f:
json.dump([x.model_dump(mode="json") for x in arr], f, indent=2)
p_atomic_write_json(_runs_path(updated.workflow_id), [x.model_dump(mode="json") for x in arr])
return updated
return None
+6 -6
View File
@@ -317,17 +317,17 @@ def test_weekly_every_n_weeks_phase_is_stable_across_recompute():
def test_ran_late_is_measured_from_start_not_finish():
"""A run that STARTS on time is 'success' no matter how long it runs; a run
that starts >5min after its slot is 'ran_late'."""
from backend.apps.workflows.executor import _ran_late
from backend.apps.workflows.executor import p_ran_late
slot = datetime(2026, 6, 22, 9, 0, tzinfo=timezone.utc)
# Started on time -> not late (even though such a run might finish much later).
assert _ran_late(slot, slot) is False
assert _ran_late(slot + timedelta(minutes=4), slot) is False
assert p_ran_late(slot, slot) is False
assert p_ran_late(slot + timedelta(minutes=4), slot) is False
# Started well after the slot -> late.
assert _ran_late(slot + timedelta(minutes=6), slot) is True
assert p_ran_late(slot + timedelta(minutes=6), slot) is True
# Naive started_at (host-local, as datetime.now() produces) is normalized
# to UTC rather than subtracted across the offset.
naive_on_time = slot.astimezone().replace(tzinfo=None)
assert _ran_late(naive_on_time, slot) is False
assert p_ran_late(naive_on_time, slot) is False
def test_frozen_empty_tool_set_does_not_fall_back_to_defaults():
@@ -864,7 +864,7 @@ def test_killed_by_restart_message_is_friendly():
storage.record_run(WorkflowRun(workflow_id=wf.id, status="running"))
scheduler._mark_stuck_runs_failed()
runs = storage.list_runs(wf.id, limit=10)
assert any(r.status == "failure" and "OpenSwarm closed" in (r.error or "") for r in runs)
assert any(r.status == "failure" and "Interrupted" in (r.error or "") and "shut down" in (r.error or "") for r in runs)
assert not any("Killed by restart" in (r.error or "") for r in runs)