diff --git a/backend/apps/workflows/executor.py b/backend/apps/workflows/executor.py index 2be409c3..0989da90 100644 --- a/backend/apps/workflows/executor.py +++ b/backend/apps/workflows/executor.py @@ -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: diff --git a/backend/apps/workflows/scheduler.py b/backend/apps/workflows/scheduler.py index 7aafe8ed..bec374a6 100644 --- a/backend/apps/workflows/scheduler.py +++ b/backend/apps/workflows/scheduler.py @@ -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 diff --git a/backend/apps/workflows/storage.py b/backend/apps/workflows/storage.py index 6aab9dbd..49facb3c 100644 --- a/backend/apps/workflows/storage.py +++ b/backend/apps/workflows/storage.py @@ -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 diff --git a/backend/tests/test_workflows_semantics.py b/backend/tests/test_workflows_semantics.py index 9bae4e13..62d21fb9 100644 --- a/backend/tests/test_workflows_semantics.py +++ b/backend/tests/test_workflows_semantics.py @@ -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)