From cd537187611769ebb6a1aa9460265e9ef5694606 Mon Sep 17 00:00:00 2001 From: morluto <76467478+morluto@users.noreply.github.com> Date: Sun, 12 Jul 2026 20:23:04 +0000 Subject: [PATCH] fix(cron): prevent long-running scheduled scripts from running twice --- cron/scheduler.py | 73 +++++++++- tests/cron/test_script_claim_heartbeat.py | 165 ++++++++++++++++++++++ 2 files changed, 234 insertions(+), 4 deletions(-) create mode 100644 tests/cron/test_script_claim_heartbeat.py diff --git a/cron/scheduler.py b/cron/scheduler.py index 9d8dc7507757..a50704ea10be 100644 --- a/cron/scheduler.py +++ b/cron/scheduler.py @@ -1976,6 +1976,7 @@ def _deliver_result(job: dict, content: str, adapters=None, loop=None) -> Option _DEFAULT_SCRIPT_TIMEOUT = 3600 # seconds (1 hour) # Backward-compatible module override used by tests and emergency monkeypatches. _SCRIPT_TIMEOUT = _DEFAULT_SCRIPT_TIMEOUT +_RUN_CLAIM_HEARTBEAT_SECONDS = 60.0 def _get_script_timeout() -> int: @@ -2135,6 +2136,71 @@ def _run_job_script(script_path: str) -> tuple[bool, str]: return False, f"Script execution failed: {exc}" +def _run_job_script_with_claim_heartbeat( + job: dict, script_path: str +) -> tuple[bool, str]: + """Run a cron script while keeping its owned one-shot claim fresh. + + Script execution is synchronous and may legitimately outlive the stale + claim TTL. Without a concurrent heartbeat, another scheduler process can + mistake the live run for a dead owner and dispatch the same one-shot again. + Recurring jobs and unclaimed/manual runs have no durable one-shot claim and + therefore use the ordinary script path without starting a thread. + + The claim owner is captured from the dispatched job and never re-read from + storage. ``heartbeat_run_claim`` compares that stable owner before every + refresh, so a stale runner cannot extend a replacement owner's claim. + """ + schedule = job.get("schedule") + claim = job.get("run_claim") + owner = str(claim.get("by") or "") if isinstance(claim, dict) else "" + if not ( + isinstance(schedule, dict) + and schedule.get("kind") == "once" + and owner + ): + return _run_job_script(script_path) + + job_id = str(job.get("id") or "") + stop = threading.Event() + heartbeat_context = contextvars.copy_context() + + def _heartbeat_loop() -> None: + while not stop.wait(_RUN_CLAIM_HEARTBEAT_SECONDS): + try: + heartbeat_run_claim(job_id, expected_owner=owner) + except Exception: + logger.debug( + "Job '%s': script run_claim heartbeat failed", + job_id, + exc_info=True, + ) + + heartbeat_thread = threading.Thread( + target=heartbeat_context.run, + args=(_heartbeat_loop,), + name="cron-script-claim-heartbeat", + daemon=True, + ) + try: + heartbeat_thread.start() + except Exception: + logger.debug( + "Job '%s': could not start script run_claim heartbeat", + job_id, + exc_info=True, + ) + return _run_job_script(script_path) + + try: + return _run_job_script(script_path) + finally: + stop.set() + # Event.wait() wakes immediately. Keep completion bounded if the + # heartbeat is already waiting on another process's jobs-file lock. + heartbeat_thread.join(timeout=1.0) + + def _parse_wake_gate(script_output: str) -> bool: """Parse the last non-empty stdout line of a cron job's pre-check script as a wake gate. @@ -2542,7 +2608,7 @@ def run_job( _prior_cwd = None try: - ok, output = _run_job_script(script_path) + ok, output = _run_job_script_with_claim_heartbeat(job, script_path) finally: if _prior_cwd is not None: try: @@ -2689,7 +2755,7 @@ def run_job( prerun_script = None script_path = job.get("script") if script_path: - prerun_script = _run_job_script(script_path) + prerun_script = _run_job_script_with_claim_heartbeat(job, script_path) _ran_ok, _script_output = prerun_script if _ran_ok and not _parse_wake_gate(_script_output): logger.info( @@ -3171,7 +3237,6 @@ def run_job( _run_claim_owner = ( str(_run_claim.get("by") or "") if isinstance(_run_claim, dict) else "" ) - _CLAIM_HEARTBEAT_SECONDS = 60.0 _last_claim_heartbeat = time.monotonic() def _heartbeat_run_claim_if_due(): @@ -3179,7 +3244,7 @@ def run_job( if not _is_oneshot or not _run_claim_owner: return _mono = time.monotonic() - if _mono - _last_claim_heartbeat < _CLAIM_HEARTBEAT_SECONDS: + if _mono - _last_claim_heartbeat < _RUN_CLAIM_HEARTBEAT_SECONDS: return _last_claim_heartbeat = _mono try: diff --git a/tests/cron/test_script_claim_heartbeat.py b/tests/cron/test_script_claim_heartbeat.py new file mode 100644 index 000000000000..0f4c3fd863c9 --- /dev/null +++ b/tests/cron/test_script_claim_heartbeat.py @@ -0,0 +1,165 @@ +"""Regression coverage for one-shot claims during blocking cron scripts.""" + +from datetime import datetime, timedelta, timezone +import threading +from unittest.mock import MagicMock, patch + +import pytest + + +@pytest.mark.parametrize( + ("no_agent", "script_output"), + [ + (True, "watchdog complete"), + (False, '{"wakeAgent": false}'), + ], + ids=("script-only-job", "pre-agent-script"), +) +def test_long_running_script_refreshes_owned_claim_in_profile_store( + tmp_path, monkeypatch, no_agent, script_output +): + """Both blocking script paths keep their one-shot claim alive. + + The real store update runs on the heartbeat thread. A second store holds + the same job ID, proving the thread inherited the active profile's + ContextVar instead of falling back to another profile's default paths. + """ + import cron.jobs as jobs + import cron.scheduler as scheduler + + profile_home = tmp_path / "profile" + default_cron = tmp_path / "default" / "cron" + default_cron.mkdir(parents=True) + profile_home.mkdir() + + monkeypatch.setattr(jobs, "CRON_DIR", default_cron) + monkeypatch.setattr(jobs, "JOBS_FILE", default_cron / "jobs.json") + monkeypatch.setattr(jobs, "OUTPUT_DIR", default_cron / "output") + monkeypatch.setattr(scheduler, "_RUN_CLAIM_HEARTBEAT_SECONDS", 0.01) + + original_timestamp = "2026-07-12T12:00:00+00:00" + original_time = datetime.fromisoformat(original_timestamp) + claim_ttl = jobs._oneshot_run_claim_ttl_seconds() + current_time = [original_time + timedelta(seconds=claim_ttl - 60)] + monkeypatch.setattr(jobs, "_hermes_now", lambda: current_time[0]) + + def _job() -> dict: + return { + "id": "long-script", + "name": "long script", + "prompt": "inspect the script output", + "script": "watchdog.py", + "no_agent": no_agent, + "schedule": { + "kind": "once", + "run_at": original_timestamp, + }, + "next_run_at": original_timestamp, + "enabled": True, + "run_claim": { + "at": original_timestamp, + "by": "dispatch-owner", + }, + } + + # Safe fallback store: if ContextVars are not propagated to the heartbeat + # thread, this record would be modified instead of the profile record. + jobs.save_jobs([_job()]) + with jobs.use_cron_store(profile_home): + jobs.save_jobs([_job()]) + claimed_job = jobs.get_job("long-script") + + heartbeat_seen = threading.Event() + real_heartbeat = jobs.heartbeat_run_claim + second_scheduler_scan = {} + + def _observed_heartbeat(job_id: str, *, expected_owner: str) -> bool: + updated = real_heartbeat(job_id, expected_owner=expected_owner) + # A different scheduler scans after the ORIGINAL claim's TTL while the + # script is still blocked. The refreshed claim must keep the job out of + # the due set and preserve its durable record. + current_time[0] = original_time + timedelta(seconds=claim_ttl + 10) + second_scheduler_scan["due"] = jobs.get_due_jobs() + second_scheduler_scan["record_present"] = jobs.get_job(job_id) is not None + heartbeat_seen.set() + return updated + + def _blocking_script(_script_path: str) -> tuple[bool, str]: + assert heartbeat_seen.wait(timeout=2), ( + "claim was not refreshed while script blocked" + ) + return True, script_output + + monkeypatch.setattr(scheduler, "heartbeat_run_claim", _observed_heartbeat) + monkeypatch.setattr(scheduler, "_run_job_script", _blocking_script) + + with ( + jobs.use_cron_store(profile_home), + patch("hermes_state.SessionDB", return_value=MagicMock()), + ): + success, _doc, _response, error = scheduler.run_job(claimed_job) + profile_claim = jobs.get_job("long-script")["run_claim"] + + assert success is True + assert error is None + assert profile_claim["at"] != original_timestamp + assert profile_claim["by"] == "dispatch-owner" + assert second_scheduler_scan == {"due": [], "record_present": True} + assert jobs.get_job("long-script")["run_claim"] == { + "at": original_timestamp, + "by": "dispatch-owner", + } + + +def test_script_heartbeat_uses_captured_claim_owner(tmp_path, monkeypatch): + """A stale script runner cannot refresh a replacement owner's claim.""" + import cron.jobs as jobs + import cron.scheduler as scheduler + + profile_home = tmp_path / "profile" + profile_home.mkdir() + original_timestamp = "2026-07-12T12:00:00+00:00" + replacement_timestamp = "2026-07-12T12:00:30+00:00" + job = { + "id": "reclaimed-script", + "script": "watchdog.py", + "schedule": {"kind": "once", "run_at": original_timestamp}, + "run_claim": {"at": original_timestamp, "by": "original-owner"}, + } + + with jobs.use_cron_store(profile_home): + jobs.save_jobs([ + { + **job, + "run_claim": { + "at": replacement_timestamp, + "by": "replacement-owner", + }, + } + ]) + + heartbeat_seen = threading.Event() + real_heartbeat = jobs.heartbeat_run_claim + + def _observed_heartbeat(job_id: str, *, expected_owner: str) -> bool: + updated = real_heartbeat(job_id, expected_owner=expected_owner) + heartbeat_seen.set() + return updated + + def _blocking_script(_script_path: str) -> tuple[bool, str]: + assert heartbeat_seen.wait(timeout=2) + return True, "done" + + monkeypatch.setattr(scheduler, "_RUN_CLAIM_HEARTBEAT_SECONDS", 0.01) + monkeypatch.setattr(scheduler, "heartbeat_run_claim", _observed_heartbeat) + monkeypatch.setattr(scheduler, "_run_job_script", _blocking_script) + + with jobs.use_cron_store(profile_home): + assert scheduler._run_job_script_with_claim_heartbeat(job, "watchdog.py") == ( + True, + "done", + ) + assert jobs.get_job("reclaimed-script")["run_claim"] == { + "at": replacement_timestamp, + "by": "replacement-owner", + }