From cff9728587da4f3c0beed0786f9bea528e489f13 Mon Sep 17 00:00:00 2001 From: kshitijk4poor <82637225+kshitijk4poor@users.noreply.github.com> Date: Wed, 29 Jul 2026 19:32:09 +0500 Subject: [PATCH] fix(cron): record failure for BaseException escapes and leave a diagnostic when removing wedged one-shots Fixes #73973. A finite one-shot whose dispatch was claimed (claim_dispatch increments repeat.completed BEFORE execution) but whose run died before mark_job_run was left permanently wedged: completed==times, last_run_at null, state 'scheduled'. The run-claim TTL blocked re-dispatch, and once it expired the dispatch-limit guard silently removed the job with no output and no error. Two complementary fixes: - cron/scheduler.py run_one_job: the outer handler now catches BaseException, not just Exception. The inner run_job handler re-raises CancelledError/KeyboardInterrupt/SystemExit after agent teardown, and none of those are Exception subclasses, so the outer 'except Exception' missed them and mark_job_run(False) was never called. Failures are now recorded first (mark_job_run + finish_execution, each independently guarded), then non-Exception BaseExceptions are re-raised to preserve teardown semantics. Plain Exceptions keep the existing behavior (recorded, return False, no re-raise). Empty str(e) (bare CancelledError) falls back to the exception class name. - cron/jobs.py: when either removal site (claim_dispatch or the get_due_jobs dispatch-limit guard) drops a one-shot whose claimed run never completed (last_run_at null), _write_wedged_oneshot_diagnostic now writes an operator-visible .md into cron/output// instead of vanishing silently. Best-effort: diagnostics can never break the removal. No diagnostic when last_run_at is set (normal completion race). --- cron/jobs.py | 52 ++++++++++++++++++++++++++++++++++ cron/scheduler.py | 35 +++++++++++++++++++---- tests/cron/test_jobs.py | 51 +++++++++++++++++++++++++++++++++ tests/cron/test_run_one_job.py | 51 +++++++++++++++++++++++++++++++++ 4 files changed, 184 insertions(+), 5 deletions(-) diff --git a/cron/jobs.py b/cron/jobs.py index 7d3860997ff..57c77dc16a0 100644 --- a/cron/jobs.py +++ b/cron/jobs.py @@ -1740,6 +1740,50 @@ def mark_job_run(job_id: str, success: bool, error: Optional[str] = None, logger.warning("mark_job_run: job_id %s not found, skipping save", job_id) +def _write_wedged_oneshot_diagnostic(job: Dict[str, Any]) -> None: + """Leave an operator-visible trace when a wedged one-shot is removed. + + A finite one-shot whose dispatch was claimed (``repeat.completed`` >= + ``repeat.times``) but which never reached ``mark_job_run`` (``last_run_at`` + is null) was interrupted mid-run — scheduler restart, gateway kill, or a + non-Exception escape (#73973). The recovery guards remove such jobs so + they stop appearing due, but a silent removal leaves the user with no + output, no error, and no job record. Write a small diagnostic file into + the job's output directory so the removal is observable and debuggable. + + Best-effort: diagnostics must never break the removal itself. + """ + if job.get("last_run_at") is not None: + return # a prior run was recorded — normal completion race, not a wedge + try: + repeat = job.get("repeat") or {} + claim = job.get("run_claim") or {} + text = ( + "# Cron job removed without producing output\n\n" + f"- job id: {job.get('id')}\n" + f"- name: {job.get('name')}\n" + f"- dispatch claimed: {repeat.get('completed', '?')}/{repeat.get('times', '?')}\n" + f"- run claimed at: {claim.get('at', 'unknown')} by {claim.get('by', 'unknown')}\n" + f"- removed at: {_hermes_now().isoformat()}\n\n" + "This one-shot job's dispatch was claimed, but the run never " + "completed (`last_run_at` was never written) — the scheduler " + "process was most likely killed or restarted mid-execution. The " + "job has been removed to stop it re-firing; recreate it to run " + "again.\n" + ) + save_job_output(job.get("id", ""), text) + logger.warning( + "Job '%s': removed without a completed run — diagnostic written to " + "its output directory", + job.get("name", job.get("id", "?")), + ) + except Exception as e: + logger.debug( + "Failed to write wedged-oneshot diagnostic for job %r: %s", + job.get("id"), e, + ) + + def claim_dispatch(job_id: str) -> bool: """Atomically claim a finite one-shot job dispatch BEFORE execution. @@ -1777,6 +1821,9 @@ def claim_dispatch(job_id: str) -> bool: # Clean up so it stops appearing as due on every tick. jobs.pop(i) save_jobs(jobs) + # If the claimed run never completed (#73973), leave an + # operator-visible diagnostic instead of vanishing silently. + _write_wedged_oneshot_diagnostic(job) logger.info( "Job '%s': dispatch limit reached (%d/%d) — removing", job.get("name", job.get("id", "?")), @@ -2249,6 +2296,11 @@ def _get_due_jobs_locked() -> List[Dict[str, Any]]: raw_jobs.remove(rj) needs_save = True break + # The claimed run never completed here by + # definition (last_run_at unwritten is what made + # the entry look due) — leave an operator-visible + # diagnostic instead of vanishing silently (#73973). + _write_wedged_oneshot_diagnostic(job) continue # Durably claim a one-shot for the DURATION of its run before diff --git a/cron/scheduler.py b/cron/scheduler.py index 2428d500a73..dfcef0ec724 100644 --- a/cron/scheduler.py +++ b/cron/scheduler.py @@ -4028,11 +4028,36 @@ def run_one_job(job: dict, *, adapters=None, loop=None, verbose: bool = False) - finish_execution(execution_id, success=success, error=error) return True - except Exception as e: - logger.error("Error processing job %s: %s", job['id'], e) - if not _consume_interrupted_flag(job["id"]): - mark_job_run(job["id"], False, str(e)) - finish_execution(execution_id, success=False, error=str(e)) + except BaseException as e: # noqa: BLE001 — deliberate: see below + # BaseException, not Exception (#73973): the inner run_job handler + # re-raises CancelledError / KeyboardInterrupt / SystemExit after agent + # teardown, and none of those are Exception subclasses. If they escape + # without mark_job_run(False), a finite one-shot is left wedged — + # claim_dispatch() already consumed repeat.completed, but last_run_at + # is never written, so the job sits in state "scheduled" until the + # run-claim TTL expires and the dispatch-limit guard removes it with + # no output and no error. Record the failure first, then re-raise + # anything that isn't a plain Exception. + _err_text = str(e) or type(e).__name__ + logger.error("Error processing job %s: %s", job['id'], _err_text) + try: + if not _consume_interrupted_flag(job["id"]): + mark_job_run(job["id"], False, _err_text) + except Exception as record_err: + # Never let bookkeeping mask the original interruption. + logger.error( + "Failed to record interrupted run for job %s: %s", + job["id"], record_err, + ) + try: + finish_execution(execution_id, success=False, error=_err_text) + except Exception as record_err: + logger.error( + "Failed to finish execution record for job %s: %s", + job["id"], record_err, + ) + if not isinstance(e, Exception): + raise return False diff --git a/tests/cron/test_jobs.py b/tests/cron/test_jobs.py index 09fc7d2faab..c263bffa22d 100644 --- a/tests/cron/test_jobs.py +++ b/tests/cron/test_jobs.py @@ -1909,6 +1909,57 @@ class TestClaimDispatch: assert due == [] assert load_jobs() == [] # cleaned up + def test_get_due_jobs_stale_removal_writes_diagnostic(self, tmp_cron_dir, monkeypatch): + """#73973: when the due-scan removes a wedged one-shot (dispatch claimed, + run never completed), it must leave an operator-visible diagnostic file + in the job's output dir instead of vanishing silently.""" + import cron.jobs as jobs_mod + from cron.jobs import _hermes_now, _oneshot_run_claim_ttl_seconds + monkeypatch.delenv("HERMES_CRON_TIMEOUT", raising=False) + ttl = _oneshot_run_claim_ttl_seconds() + stale = (_hermes_now() - timedelta(seconds=ttl + 300)).isoformat() + save_jobs([{ + "id": "wedged1", + "name": "wedged one-shot", + "enabled": True, + "schedule": {"kind": "once", "run_at": stale}, + "repeat": {"times": 1, "completed": 1}, + "run_claim": {"at": stale, "by": "dead-process:123"}, + "next_run_at": stale, + }]) + assert get_due_jobs() == [] + assert load_jobs() == [] + out_dir = jobs_mod._job_output_dir("wedged1") + files = list(out_dir.glob("*.md")) + assert files, "expected a diagnostic file in the output dir" + text = files[0].read_text(encoding="utf-8") + assert "removed without producing output" in text + assert "wedged1" in text + + def test_claim_dispatch_stale_removal_writes_diagnostic(self, tmp_cron_dir): + """#73973: the claim_dispatch removal path for an already-maxed one-shot + with no completed run also leaves a diagnostic.""" + import cron.jobs as jobs_mod + save_jobs([self._oneshot(times=1, completed=1)]) + assert claim_dispatch("os1") is False + assert load_jobs() == [] + out_dir = jobs_mod._job_output_dir("os1") + files = list(out_dir.glob("*.md")) + assert files, "expected a diagnostic file in the output dir" + assert "removed without producing output" in files[0].read_text(encoding="utf-8") + + def test_no_diagnostic_when_run_completed(self, tmp_cron_dir): + """A one-shot that DID complete a run (last_run_at set) is a normal + completion race, not a wedge — no diagnostic file is written.""" + import cron.jobs as jobs_mod + job = self._oneshot(times=1, completed=1) + job["last_run_at"] = "2026-01-01T00:05:00+00:00" + save_jobs([job]) + assert claim_dispatch("os1") is False + assert load_jobs() == [] + out_dir = jobs_mod._job_output_dir("os1") + assert not out_dir.exists() or not list(out_dir.glob("*.md")) + def test_bad_schedule_does_not_crash_or_block_sibling_jobs(self, tmp_cron_dir): """Regression for a job with non-dict 'schedule' (null / string / etc. diff --git a/tests/cron/test_run_one_job.py b/tests/cron/test_run_one_job.py index 4fdf46ff710..71506758616 100644 --- a/tests/cron/test_run_one_job.py +++ b/tests/cron/test_run_one_job.py @@ -119,6 +119,57 @@ def test_run_one_job_exception_marks_failure(monkeypatch): assert marks == [("j6", False)] +def test_run_one_job_base_exception_records_failure_then_reraises(monkeypatch): + """#73973: a BaseException escaping run_job (CancelledError re-raised by the + inner teardown handler, KeyboardInterrupt, SystemExit) must still record the + failure via mark_job_run — otherwise a claim_dispatch()-consumed one-shot is + left wedged with completed==times but last_run_at never written. The + BaseException itself is re-raised after recording so shutdown semantics are + preserved.""" + import asyncio + + import pytest + + for exc in (asyncio.CancelledError(), KeyboardInterrupt(), SystemExit(1)): + def boom(job, *, defer_agent_teardown=None, _exc=exc): + raise _exc + + monkeypatch.setattr(s, "run_job", boom) + marks = [] + monkeypatch.setattr( + s, "mark_job_run", + lambda jid, ok, err=None, delivery_error=None: marks.append((jid, ok, err)), + ) + + with pytest.raises(type(exc)): + s.run_one_job({"id": "jbase", "name": "t"}) + + assert marks and marks[0][0] == "jbase" and marks[0][1] is False, ( + f"{type(exc).__name__}: failure was not recorded" + ) + # Empty str(exc) (e.g. bare CancelledError) falls back to the class name. + assert marks[0][2], f"{type(exc).__name__}: error text must be non-empty" + + +def test_run_one_job_plain_exception_still_swallowed(monkeypatch): + """The BaseException widening must not change plain-Exception behavior: + recorded, returns False, NOT re-raised.""" + def boom(job, *, defer_agent_teardown=None): + raise ValueError("plain failure") + + monkeypatch.setattr(s, "run_job", boom) + marks = [] + monkeypatch.setattr( + s, "mark_job_run", + lambda jid, ok, err=None, delivery_error=None: marks.append((jid, ok)), + ) + + ok = s.run_one_job({"id": "jplain", "name": "t"}) + + assert ok is False + assert marks == [("jplain", False)] + + def test_run_one_job_installs_secret_scope_under_multiplex(monkeypatch, tmp_path): """Regression: under profile isolation (multiplex active), run_one_job must execute run_job inside a profile secret scope so credential reads