fix(cron): fire one catch-up for a slot missed during a restart gap (#107485)
tick() advances a recurring job's next_run_at BEFORE dispatch so a crash mid-run cannot re-fire it on every restart (at-most-once). That leaves a window — advance persisted, fire claim not yet taken — in which the process dying (interpreter already finalizing, executor refusing new futures, SIGKILL, the Desktop idle-exit from #107485) loses the occurrence silently: the restarted scan sees only the future next_run_at, writes no execution row and no log line, and a daily job skips a day. Contract (documented in cron-internals.md): every recurring occurrence is accounted for — it runs once, or its skip is logged with a reason. - The due scan stamps `pending_slot = {scheduled_at, at, by}` in the same save that records `last_dispatch`; claim_job_for_fire and mark_job_run clear it, and an explicit schedule / next_run_at / enabled / state rewrite (edit, pause, resume, run-now) drops it. - A later scan that finds the stamp with a provably gone owner (this process and the job not in its running set, or another owner past the fire-claim lease / dead pid) restores scheduled_at as next_run_at ONCE and logs a WARNING (cron/occurrences.py::unclaimed_pending_slot). The restored instant then meets the ordinary policy — completed_occurrence() blocks a second fire of a slot that already ran, the grace window classifies it late, past-grace collapses the backlog into one run, cron.catch_up_missed: false skips it with a logged reason. Never a replay of N slots. Same store fields on both topologies: a standalone `hermes -p X gateway run` and a profile served by the default multiplexer evaluate the identical record.
This commit is contained in:
+38
-2
@@ -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
|
||||
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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"])
|
||||
Reference in New Issue
Block a user