diff --git a/cron/jobs.py b/cron/jobs.py index 386599c310..f0e790f353 100644 --- a/cron/jobs.py +++ b/cron/jobs.py @@ -1989,6 +1989,10 @@ def update_job(job_id: str, updates: Dict[str, Any]) -> Optional[Dict[str, Any]] if "schedule" in updates: _apply_schedule_update(updated, updates, job_id) + if {"schedule", "next_run_at", "enabled", "state"}.intersection(updates): + # An explicit schedule/lifecycle rewrite supersedes any occurrence the dispatcher + # left unclaimed — pause/resume/edit must not resurrect a slot from before the edit. + updated.pop("pending_slot", None) if inference_fields_changed: snapshots = _compute_provider_model_snapshots( provider=updated.get("provider"), @@ -2234,6 +2238,7 @@ def _record_run_outcome( job["last_delivery_error"] = delivery_error # Clear both claims: the run is over, so the job is claimable again. job["fire_claim"] = None + job.pop("pending_slot", None) if job.get("run_claim") is not None: # keep key absence for legacy records job["run_claim"] = None @@ -2581,6 +2586,8 @@ def claim_job_for_fire( # Per-acquisition token: a process may legitimately reclaim its own stale lease, and the # previous runner must not heartbeat the new claim merely because hostname + PID match. job["fire_claim"] = {"at": now.isoformat(), "by": f"{_machine_id()}:{uuid.uuid4().hex}"} + # Claimed: the occurrence is now owned by a run (its ledger row + fire claim carry it). + job.pop("pending_slot", None) if job.get("schedule", {}).get("kind") in {"cron", "interval"}: nxt = compute_next_run(job["schedule"], now.isoformat()) if nxt: @@ -2969,6 +2976,29 @@ def _oneshot_dispatch_limit_reached(job: Dict[str, Any], scan: _DueScan) -> bool return True +def _restore_unclaimed_slot(job: Dict[str, Any], scan: _DueScan) -> Optional[str]: + """Put an occurrence the dispatcher advanced past but never claimed back on the schedule + (#107485); returns the restored ``next_run_at`` or None. Restored ONCE: the stamp is dropped + here, so the slot then meets the ordinary late / fast-forward / ``cron.catch_up_missed`` + policy like any other overdue instant — never a replay of every missed slot.""" + from cron.occurrences import unclaimed_pending_slot + + slot = unclaimed_pending_slot(job, scan.now) + if slot is None: + return None + logger.warning( + "Job '%s' (%s): occurrence %s was taken off the schedule but never claimed " + "(scheduler stopped before dispatch); restoring it as the due instant (was %s).", + job.get("name", job.get("id")), job.get("id"), slot, job.get("next_run_at")) + job["next_run_at"] = slot + job.pop("pending_slot", None) + rj = scan.find(job["id"]) + if rj is not None: + rj.pop("pending_slot", None) + scan.persist(job["id"], next_run_at=slot) + return slot + + def _evaluate_due_job(job: Dict[str, Any], scan: _DueScan, run_claim_ttl: float) -> bool: """Decide whether one enabled, non-terminal job fires this tick, persisting any repairs. Ordering matters: recover missing next_run_at, repair timezone shifts, re-arm stale-error @@ -2984,7 +3014,7 @@ def _evaluate_due_job(job: Dict[str, Any], scan: _DueScan, run_claim_ttl: float) ): return False - next_run = job.get("next_run_at") or _recover_missing_next_run(job, scan) + next_run = _restore_unclaimed_slot(job, scan) or job.get("next_run_at") or _recover_missing_next_run(job, scan) if not next_run: return False raw_next_run_dt = datetime.fromisoformat(next_run) @@ -3040,7 +3070,13 @@ def _evaluate_due_job(job: Dict[str, Any], scan: _DueScan, run_claim_ttl: float) "kind": _classify_dispatch_lateness(lateness, grace), } job["last_dispatch"] = dispatch_stamp - scan.persist(job["id"], last_dispatch=dispatch_stamp) + # The tick advances next_run_at past this occurrence before any fire claim exists; the + # stamp survives a process death in that window so the slot is restored, not lost. + from cron.occurrences import pending_slot_stamp + + scan.persist( + job["id"], last_dispatch=dispatch_stamp, + pending_slot=pending_slot_stamp(next_run, now)) return True diff --git a/cron/occurrences.py b/cron/occurrences.py index 8b86a706ae..9b598c5f93 100644 --- a/cron/occurrences.py +++ b/cron/occurrences.py @@ -34,3 +34,52 @@ def completed_occurrence(job, instant): except Exception: logger.warning("Cannot check completed occurrence for job %s", job['id'], exc_info=True) return False + + +# --- Pending slot: the occurrence a tick took off the schedule but has not yet claimed --- +# +# The tick advances a recurring job's ``next_run_at`` BEFORE dispatch (at-most-once: a crash +# mid-run must not re-fire on every restart). That leaves a window — advance done, fire claim +# not yet taken (interpreter finalizing, executor refusing work, process killed) — in which the +# process exiting loses the occurrence silently: the restarted scan sees a future ``next_run_at`` +# and nothing ever ran (#107485). ``pending_slot`` is the durable record of that window: the due +# scan stamps it with the exact stored instant plus the stamping owner, ``claim_job_for_fire`` +# (the point after which side effects may exist) clears it, and any explicit rewrite of the +# schedule drops it. A slot still pending once its owner is provably gone (or its lease has +# expired) was never claimed, so it is restored ONCE as the due instant and then flows through +# the ordinary late / fast-forward / ``cron.catch_up_missed`` policy — never N replays. + +def pending_slot_stamp(next_run, now): + """Store value for a recurring occurrence about to be handed to the dispatcher.""" + from cron.jobs import _machine_id + + return {"scheduled_at": next_run, "at": now.isoformat(), "by": _machine_id()} + + +def unclaimed_pending_slot(job, now): + """Stored instant of a slot the dispatcher never claimed, or None. + + None for non-recurring jobs and malformed stamps (never a fire) and for a job still in this + process's running set (its queued worker will claim and clear the slot itself). A stamp by + THIS process on a job not running here is orphaned (dispatch refused). A stamp by another + process is honoured while that owner may still be alive within the fire-claim lease — a + second live gateway on the same store is mid-dispatch, not dead.""" + from cron.jobs import ( + FIRE_CLAIM_TTL_SECONDS, _claim_is_live, _job_running_in_this_process, _machine_id, + ) + + pending = job.get("pending_slot") + if not isinstance(pending, dict): + return None + slot = pending.get("scheduled_at") + if job.get("schedule", {}).get("kind") not in {"cron", "interval"} or not isinstance(slot, str): + return None + try: + datetime.fromisoformat(slot) + except ValueError: + return None + if _job_running_in_this_process(str(job.get("id", ""))): + return None + if pending.get("by") != _machine_id() and _claim_is_live(pending, now, FIRE_CLAIM_TTL_SECONDS): + return None + return slot diff --git a/tests/cron/test_missed_window_catchup.py b/tests/cron/test_missed_window_catchup.py new file mode 100644 index 0000000000..dfc063d71b --- /dev/null +++ b/tests/cron/test_missed_window_catchup.py @@ -0,0 +1,115 @@ +"""A recurring occurrence that a tick took off the schedule but never dispatched must fire once +after the scheduler restarts — never be silently lost, never fire twice (#107485). + +The tick advances ``next_run_at`` BEFORE dispatch (at-most-once across a crash mid-run). When the +process dies in the window between that advance and the fire claim — the interpreter was already +finalizing, the executor refused work, SIGKILL — the restarted scan used to see only the future +``next_run_at`` and the occurrence vanished: no execution row, no log line. The store now carries a +``pending_slot`` stamp across that window and a later scan restores it as the due instant. + +Drives the REAL ``tick()`` against a throwaway HERMES_HOME with a ``no_agent`` script job that +appends one line per fire. Process 1 is a real subprocess that dies inside the window, so the +restarted scan sees a provably dead owner exactly as a gateway restart does. +""" +from __future__ import annotations + +import os +import subprocess +import sys +from datetime import timedelta +from pathlib import Path + +import pytest + +REPO = Path(__file__).resolve().parents[2] + +_CRASH_BEFORE_DISPATCH = """ +import os, sys +sys.path.insert(0, sys.argv[1]) +import cron.scheduler as S +S.create_execution = lambda *a, **k: os._exit(137) # SIGKILL in the advance→claim window +S.tick(verbose=False, sync=True) +""" + + +@pytest.fixture +def slot_env(tmp_path, monkeypatch): + home = tmp_path / ".hermes" + (home / "cron" / "output").mkdir(parents=True) + (home / "scripts").mkdir() + monkeypatch.setenv("HERMES_HOME", str(home)) + monkeypatch.delenv("HERMES_MACHINE_ID", raising=False) + + import cron.executions as E + import cron.jobs as J + import cron.scheduler as S + + monkeypatch.setattr(J, "HERMES_DIR", home) + monkeypatch.setattr(J, "CRON_DIR", home / "cron") + monkeypatch.setattr(J, "JOBS_FILE", home / "cron" / "jobs.json") + monkeypatch.setattr(J, "OUTPUT_DIR", home / "cron" / "output") + monkeypatch.setattr(E, "EXECUTIONS_FILE", home / "cron" / "executions.db") + monkeypatch.setattr(S, "_hermes_home", home) + S._running_job_ids.clear() + S._running_since.clear() + S._running_futures.clear() + + counter = home / "fires.txt" + (home / "scripts" / "fire.sh").write_text( + f"#!/bin/sh\necho fired >> {counter}\necho fired\n", encoding="utf-8") + (home / "scripts" / "fire.sh").chmod(0o755) + job = J.create_job(prompt=None, schedule="every 1h", name="slot", script="fire.sh", + no_agent=True, deliver="local") + slot = (J._hermes_now() - timedelta(minutes=1)).replace(microsecond=0).isoformat() + stored = J.load_jobs() + next(r for r in stored if r["id"] == job["id"])["next_run_at"] = slot + J.save_jobs(stored) + + def fires() -> int: + return counter.read_text(encoding="utf-8").count("\n") if counter.exists() else 0 + + def crash_before_dispatch() -> None: + """Process 1: its tick advances the schedule, then it dies before any fire claim.""" + env = dict(os.environ, HERMES_HOME=str(home)) + proc = subprocess.run([sys.executable, "-c", _CRASH_BEFORE_DISPATCH, str(REPO)], + env=env, cwd=str(REPO), capture_output=True, text=True, timeout=120) + assert proc.returncode == 137, proc.stderr[-2000:] + + yield {"job_id": job["id"], "slot": slot, "fires": fires, "crash": crash_before_dispatch, + "S": S, "J": J, "E": E} + S._shutdown_parallel_pool() + + +class TestMissedWindowCatchUp: + def test_slot_lost_before_dispatch_fires_once_after_restart(self, slot_env): + S, J, E = slot_env["S"], slot_env["J"], slot_env["E"] + job_id, slot = slot_env["job_id"], slot_env["slot"] + + slot_env["crash"]() + after_crash = J.get_job(job_id) + assert J._ensure_aware(J.datetime.fromisoformat(after_crash["next_run_at"])) > J._hermes_now() + assert slot_env["fires"]() == 0 + + # Restarted scheduler: the occurrence must come back and run exactly once. + S.tick(verbose=False, sync=True) + assert slot_env["fires"]() == 1, "occurrence lost in the restart gap must fire once" + row = E.latest_execution(job_id) + assert row["status"] == "completed" + from cron.occurrences import scheduled_instant + assert row["scheduled_instant"] == scheduled_instant(slot), "restored slot keeps its identity" + + S.tick(verbose=False, sync=True) + assert slot_env["fires"]() == 1 + rec = J.get_job(job_id) + assert "pending_slot" not in rec + assert J._ensure_aware(J.datetime.fromisoformat(rec["next_run_at"])) > J._hermes_now() + + def test_fired_slot_is_not_replayed_after_restart(self, slot_env): + """Contract half two: a slot that DID run before the restart stays run.""" + S, J = slot_env["S"], slot_env["J"] + S.tick(verbose=False, sync=True) + assert slot_env["fires"]() == 1 + S._running_job_ids.clear() # what a restart forgets + S.tick(verbose=False, sync=True) + assert slot_env["fires"]() == 1 + assert "pending_slot" not in J.get_job(slot_env["job_id"])