Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
81 changes: 80 additions & 1 deletion cron/jobs.py
Original file line number Diff line number Diff line change
Expand Up @@ -1187,6 +1187,13 @@ def advance_next_run(job_id: str) -> bool:

One-shot jobs are left unchanged so they can still retry on restart.

CRASH WINDOW: if the gateway crashes AFTER this call but BEFORE
mark_job_run() records completion, next_run_at points to the next period
while last_run_at remains null. On restart, the job looks "not yet due"
and is permanently skipped for that occurrence.
Mitigation: recover_crash_window_jobs() detects and resets these victims
at startup, before the first tick.

Returns True if next_run_at was advanced, False otherwise.
"""
with _jobs_lock():
Expand All @@ -1206,6 +1213,77 @@ def advance_next_run(job_id: str) -> bool:
return False


def recover_crash_window_jobs() -> int:
"""One-shot startup recovery for advance_next_run crash-window victims.

Called once by InProcessCronScheduler.start() before the tick loop begins.
Scans jobs.json for jobs where advance_next_run() wrote a future next_run_at
but mark_job_run() never followed (gateway crashed between the two), leaving
last_run_at=None with next_run_at suspiciously far in the future.

For those jobs, resets next_run_at to now so the very first tick fires them.
This isolates crash-window recovery from the normal get_due_jobs() hot-path.

Returns the number of jobs whose next_run_at was reset.
"""
with _jobs_lock():
jobs = load_jobs()
now = _hermes_now()
changed = 0
for job in jobs:
if not job.get("enabled", True):
continue
if job.get("last_run_at") is not None:
continue
next_run = job.get("next_run_at")
if not next_run:
continue
schedule = job.get("schedule", {})
kind = schedule.get("kind")
if kind not in {"cron", "interval"}:
continue
try:
next_run_dt = _ensure_aware(datetime.fromisoformat(next_run))
except (ValueError, TypeError):
continue
if next_run_dt <= now:
continue
grace = _compute_grace_seconds(schedule)
if (next_run_dt - now).total_seconds() <= grace:
continue
created_at = job.get("created_at")
if not created_at:
continue
try:
created_dt = _ensure_aware(datetime.fromisoformat(created_at))
first_run_expected = None
if kind == "interval":
first_run_expected = created_dt + timedelta(
minutes=schedule.get("minutes", 1)
)
elif kind == "cron" and HAS_CRONITER:
_fc = croniter(schedule["expr"], created_dt)
first_run_expected = _ensure_aware(_fc.get_next(datetime))
if first_run_expected is not None and now > first_run_expected:
logger.warning(
"Startup recovery: job '%s' (never ran, "
"first_run_expected=%s, next_run_at=%s, grace=%ds) "
"— advance_next_run crash window detected; "
"resetting next_run_at to fire on first tick",
job.get("name", job["id"]),
first_run_expected.isoformat(),
next_run,
grace,
)
job["next_run_at"] = now.isoformat()
changed += 1
except Exception:
pass
if changed:
save_jobs(jobs)
return changed


def _machine_id() -> str:
"""Stable-ish identifier for claim attribution/debugging (NOT correctness).

Expand Down Expand Up @@ -1384,12 +1462,13 @@ def _get_due_jobs_locked() -> List[Dict[str, Any]]:
break
continue

grace = _compute_grace_seconds(schedule)

if next_run_dt <= now:

# For recurring jobs, check if the scheduled time is stale
# (gateway was down and missed the window). Fast-forward to
# the next future occurrence instead of firing a stale run.
grace = _compute_grace_seconds(schedule)
if kind in {"cron", "interval"} and (now - next_run_dt).total_seconds() > grace:
# Job is past its catch-up grace window — skip accumulated
# missed runs but still execute once now to avoid deferring
Expand Down
59 changes: 57 additions & 2 deletions cron/scheduler_provider.py
Original file line number Diff line number Diff line change
Expand Up @@ -151,12 +151,53 @@ def resolve_cron_scheduler() -> "CronScheduler":
return InProcessCronScheduler()


def _resolve_tick_interval(default: int) -> int:
"""Resolve the cron polling interval from env/config with a safe default.

Resolution order:
1. ``HERMES_CRON_INTERVAL`` environment variable (seconds, integer >= 5)
2. ``cron.tick_interval_seconds`` in config.yaml
3. ``default`` argument (caller-supplied, normally the CLI/gateway default)

Values below 5 are clamped to 5 to prevent runaway tight loops.
Invalid values are logged and ignored.
"""
import os
import logging

_log = logging.getLogger("cron.scheduler_provider")

env_val = os.getenv("HERMES_CRON_INTERVAL", "").strip()
if env_val:
try:
return max(5, int(env_val))
except (ValueError, TypeError):
_log.warning("Invalid HERMES_CRON_INTERVAL=%r; using config/default", env_val)

try:
from hermes_cli.config import load_config
cfg = load_config() or {}
cfg_val = (cfg.get("cron") or {}).get("tick_interval_seconds")
if cfg_val is not None:
return max(5, int(cfg_val))
except Exception:
pass

return default


class InProcessCronScheduler(CronScheduler):
"""Default provider: the historical in-process 60s ticker.
"""Default provider: the historical in-process ticker.

``start()`` blocks in the tick loop until ``stop_event`` is set, identical
to the pre-refactor ``_start_cron_ticker`` core loop. The caller runs it in
a daemon thread.

Default polling interval is 60 seconds. Override via the
``HERMES_CRON_INTERVAL`` env var or ``cron.tick_interval_seconds`` in
config.yaml (minimum 5 s). On startup, ``recover_crash_window_jobs()``
runs once to reset any jobs that were lost in an advance_next_run crash
window so they fire on the very first tick.
"""

@property
Expand All @@ -166,10 +207,24 @@ def name(self) -> str:
def start(self, stop_event, *, adapters=None, loop=None, interval=60):
import logging
from cron.scheduler import tick as cron_tick
from cron.jobs import record_ticker_heartbeat
from cron.jobs import record_ticker_heartbeat, recover_crash_window_jobs

interval = _resolve_tick_interval(interval)

logger = logging.getLogger("cron.scheduler_provider")
logger.info("In-process cron scheduler started (interval=%ds)", interval)

# One-shot startup recovery: resets any jobs whose next_run_at was
# pre-written by advance_next_run but whose mark_job_run never ran
# (gateway crashed in the window between the two). After this call,
# those jobs have next_run_at=now and will fire on the first tick.
_recovered = recover_crash_window_jobs()
if _recovered:
logger.info(
"Startup: reset %d crash-window job(s) to fire on first tick",
_recovered,
)

# Heartbeat once before the first sleep so `hermes cron status` sees a
# live ticker immediately after startup, not only after the first tick.
record_ticker_heartbeat()
Expand Down
146 changes: 146 additions & 0 deletions tests/cron/test_jobs.py
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@
mark_job_run,
advance_next_run,
get_due_jobs,
recover_crash_window_jobs,
save_job_output,
)

Expand Down Expand Up @@ -1049,6 +1050,151 @@ def test_offset_migration_at_wall_clock_equal_now_falls_through(self, tmp_cron_d
assert nr is None or datetime.fromisoformat(nr).utcoffset() == now.utcoffset() or "+10:00" in nr


# ── crash-window catch-up (advance_next_run survived but mark_job_run didn't) ──

def _crash_window_job(self, job_id, created_iso, next_run_iso, schedule_dict):
"""Build a minimal job dict simulating advance_next_run crash window."""
return {
"id": job_id,
"name": job_id,
"prompt": "...",
"schedule": schedule_dict,
"schedule_display": schedule_dict.get("display", ""),
"repeat": {"times": None, "completed": 0},
"enabled": True,
"state": "scheduled",
"paused_at": None,
"paused_reason": None,
"created_at": created_iso,
"next_run_at": next_run_iso,
"last_run_at": None,
"last_status": None,
"last_error": None,
"deliver": "local",
"origin": None,
}

def test_crash_window_interval_job_resets_next_run_at_on_startup(
self, tmp_cron_dir, monkeypatch
):
"""recover_crash_window_jobs() resets next_run_at to now for an interval
job whose advance_next_run crashed before mark_job_run ran.

Scenario: job created at 09:00, interval 60 min → first_run_expected 10:00.
advance_next_run wrote next_run_at = 23:00 before the crash.
At startup (10:30), recover_crash_window_jobs() resets next_run_at to
now so the very first tick fires the job.
"""
from datetime import timezone

now = datetime(2026, 6, 22, 10, 30, 0, tzinfo=timezone.utc)
monkeypatch.setattr("cron.jobs._hermes_now", lambda: now)

save_jobs([self._crash_window_job(
job_id="crash-interval",
created_iso="2026-06-22T09:00:00+00:00",
next_run_iso="2026-06-22T23:00:00+00:00", # far future, > grace
schedule_dict={"kind": "interval", "minutes": 60, "display": "every 60m"},
)])

recovered = recover_crash_window_jobs()
assert recovered == 1, "expected exactly 1 job to be reset"
# next_run_at must now be <= now so the first tick fires it
nr = get_job("crash-interval")["next_run_at"]
from cron.jobs import _ensure_aware
assert _ensure_aware(datetime.fromisoformat(nr)) <= now, (
"recover_crash_window_jobs did not reset next_run_at to now"
)

def test_crash_window_does_not_reset_when_first_run_still_future(
self, tmp_cron_dir, monkeypatch
):
"""recover_crash_window_jobs() leaves job untouched when first_run_expected
hasn't arrived yet — this is a legitimately new job, not a crash victim.

Scenario: job created at 09:00, interval 60 min → first_run_expected 10:00.
advance_next_run wrote next_run_at = 23:00 but it's only 09:45 now.
"""
from datetime import timezone

now = datetime(2026, 6, 22, 9, 45, 0, tzinfo=timezone.utc)
monkeypatch.setattr("cron.jobs._hermes_now", lambda: now)

save_jobs([self._crash_window_job(
job_id="crash-not-yet",
created_iso="2026-06-22T09:00:00+00:00",
next_run_iso="2026-06-22T23:00:00+00:00",
schedule_dict={"kind": "interval", "minutes": 60, "display": "every 60m"},
)])

recovered = recover_crash_window_jobs()
assert recovered == 0, (
"recover_crash_window_jobs reset a job whose first_run_expected is still in the future"
)

def test_crash_window_does_not_reset_when_job_already_ran(
self, tmp_cron_dir, monkeypatch
):
"""recover_crash_window_jobs() ignores jobs where last_run_at is set —
those ran successfully before; next_run_at is legitimately in the future.
"""
from datetime import timezone

now = datetime(2026, 6, 22, 10, 30, 0, tzinfo=timezone.utc)
monkeypatch.setattr("cron.jobs._hermes_now", lambda: now)

job = self._crash_window_job(
job_id="crash-already-ran",
created_iso="2026-06-22T09:00:00+00:00",
next_run_iso="2026-06-22T23:00:00+00:00",
schedule_dict={"kind": "interval", "minutes": 60, "display": "every 60m"},
)
job["last_run_at"] = "2026-06-22T10:00:00+00:00"
job["last_status"] = "ok"
save_jobs([job])

recovered = recover_crash_window_jobs()
assert recovered == 0, (
"recover_crash_window_jobs reset a job that already ran (last_run_at is set)"
)

def test_crash_window_recovery_resets_multiple_victims(
self, tmp_cron_dir, monkeypatch
):
"""recover_crash_window_jobs() resets all qualifying jobs in a single pass,
matching the 11-job scenario from issue #51038.
"""
from datetime import timezone

now = datetime(2026, 6, 22, 15, 45, 0, tzinfo=timezone.utc)
monkeypatch.setattr("cron.jobs._hermes_now", lambda: now)

victims = [
self._crash_window_job(
job_id=f"crash-{i}",
created_iso="2026-06-22T09:00:00+00:00",
next_run_iso="2026-06-23T09:00:00+00:00", # tomorrow
schedule_dict={"kind": "interval", "minutes": 60, "display": "every 60m"},
)
for i in range(3)
]
# One legitimate job that ran before — must NOT be reset
ran_before = self._crash_window_job(
job_id="ran-before",
created_iso="2026-06-22T09:00:00+00:00",
next_run_iso="2026-06-23T09:00:00+00:00",
schedule_dict={"kind": "interval", "minutes": 60, "display": "every 60m"},
)
ran_before["last_run_at"] = "2026-06-22T10:00:00+00:00"
save_jobs(victims + [ran_before])

recovered = recover_crash_window_jobs()
assert recovered == 3, f"expected 3 resets, got {recovered}"
assert get_job("ran-before")["next_run_at"] == "2026-06-23T09:00:00+00:00", (
"job that already ran was incorrectly reset"
)


class TestEnabledToolsets:
def test_enabled_toolsets_stored(self, tmp_cron_dir):
job = create_job(prompt="monitor", schedule="every 1h", enabled_toolsets=["web", "terminal"])
Expand Down
Loading
Loading