From 45b42202fac97d47c3aba3d5d2549232d766870c Mon Sep 17 00:00:00 2001 From: teknium1 <127238744+teknium1@users.noreply.github.com> Date: Mon, 14 Sep 2026 10:57:42 -0700 Subject: [PATCH] fix(cron): a raising profile gate ticks nothing; housekeeping respawns a dead ticker (#111010) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Follow-up on the salvaged #111034: - `_start_multiplex` published the enumerated home list before the gate had filtered it, so a raising `profile_gate` (the Desktop stand-down probe from #100489) kept the thread alive but ticked every profile UNGATED — racing the gateway that owns them for the same cron store. The list is now assigned only after gating; a gate failure yields zero ticks for that cycle. - `cron/scheduler_thread.py::SupervisedTickerThread` wraps the gateway ticker thread; `_start_gateway_housekeeping` gets a per-tick "Cron ticker supervisor" chore that respawns a ticker that ended without a stop request and logs the outage at ERROR. Every guard inside `start()` keeps the loop alive, but nothing outside it could notice a thread that had already ended. - Tests trimmed to the invariants, proven red on origin/main: a REAL corrupt `executions.db` (no patched recover) no longer kills the ticker; a raising gate keeps the thread alive with zero ticks; housekeeping restarts a dead ticker and leaves a stopped one alone. Root cause: the unguarded pre-loop `recover_interrupted_executions()` + `record_ticker_heartbeat()` were added by d9dd05b69d9b (#61791, "truthful execution ledger", 2026-07-09). The reporter's build (e440bf35) also carried #107485's `completed_occurrence()` in the due scan, which opens the same ledger on every tick — the first traceback in their errors.log is that in-loop hit (caught); the restart then hit the SAME corrupt ledger from the pre-loop recovery scan, which nothing caught: thread dead, no heartbeat, no error marker. --- cron/scheduler_provider.py | 13 +- cron/scheduler_thread.py | 53 +++++++ gateway/run.py | 16 +- tests/cron/test_ticker_startup_survival.py | 170 ++++++++------------- 4 files changed, 134 insertions(+), 118 deletions(-) create mode 100644 cron/scheduler_thread.py diff --git a/cron/scheduler_provider.py b/cron/scheduler_provider.py index 9f8ac8ee55..8332631be7 100644 --- a/cron/scheduler_provider.py +++ b/cron/scheduler_provider.py @@ -543,16 +543,15 @@ class InProcessCronScheduler(CronScheduler): # See #87644. _cycle_exc: BaseException | None = None # Enumeration and gating run on the ticker thread; a raising gate callable must - # fail THIS cycle (logged, no heartbeats), not end the thread (#111010). + # fail THIS cycle (logged, no heartbeats, NO ticks), not end the thread (#111010). + # Publish the list only once the gate has filtered it: a partial assignment would + # tick the ungated set — the exact stand-down the Desktop gate exists for (#100489). cycle_homes: list = [] try: - cycle_homes = [ - _profile_entry(e) for e in _existing_profile_homes(profile_homes) - ] + enumerated = [_profile_entry(e) for e in _existing_profile_homes(profile_homes)] if profile_gate is not None: - cycle_homes = [ - (name, home) for name, home in cycle_homes if profile_gate(name, home) - ] + enumerated = [(name, home) for name, home in enumerated if profile_gate(name, home)] + cycle_homes = enumerated except BaseException as e: logger.error("Cron profile enumeration error: %s", e, exc_info=True) _tick_error = f"{type(e).__name__}: {e}" diff --git a/cron/scheduler_thread.py b/cron/scheduler_thread.py new file mode 100644 index 0000000000..2552842909 --- /dev/null +++ b/cron/scheduler_thread.py @@ -0,0 +1,53 @@ +"""Supervised daemon thread for the in-process cron ticker. + +The gateway (and the Desktop backend) run ``InProcessCronScheduler.start`` on a bare daemon thread. +Every guard inside ``start`` keeps the loop alive on a per-tick failure, but nothing outside it could +notice a thread that had already ended — the gateway kept serving while ``ticker_heartbeat`` froze and +no job fired again until a restart (#111010). The supervisor is the missing outer layer: the +housekeeping loop asks it once a cycle to respawn a ticker that died while shutdown was not requested. +""" + +from __future__ import annotations + +import logging +import threading +from typing import Any, Callable, Mapping + +logger = logging.getLogger(__name__) + + +class SupervisedTickerThread: + """``threading.Thread``-shaped handle whose ``restart_if_dead`` respawns a dead ticker.""" + + def __init__(self, target: Callable[..., Any], *, args: tuple = (), kwargs: Mapping[str, Any] | None = None, + stop_event: threading.Event, name: str = "cron-scheduler") -> None: + self._target, self._args, self._kwargs = target, args, dict(kwargs or {}) + self._stop_event, self._name = stop_event, name + self._thread = self._spawn() + self.restarts = 0 + + def _spawn(self) -> threading.Thread: + return threading.Thread(target=self._target, args=self._args, kwargs=self._kwargs, daemon=True, name=self._name) + + def start(self) -> None: + self._thread.start() + + def is_alive(self) -> bool: + return self._thread.is_alive() + + def join(self, timeout: float | None = None) -> None: + self._thread.join(timeout) + + def restart_if_dead(self) -> bool: + """Respawn the ticker when it ended without ``stop_event``; True when a restart happened.""" + if self._stop_event.is_set() or self._thread.is_alive(): + return False + self.restarts += 1 + logger.error( + "Cron ticker thread %r died without a stop request; restarting (restart #%d). Scheduled jobs " + "did not fire while it was down — check errors.log for the escaping exception.", + self._name, self.restarts, + ) + self._thread = self._spawn() + self._thread.start() + return True diff --git a/gateway/run.py b/gateway/run.py index e91aade0b4..8a1fecf6f2 100644 --- a/gateway/run.py +++ b/gateway/run.py @@ -4547,6 +4547,7 @@ def _drain_restart_safe_cron_deliveries(adapters, loop, runner=None) -> None: def _start_gateway_housekeeping( stop_event: threading.Event, adapters=None, loop=None, interval: int = 60, cron_provider=None, runner=None, + cron_thread=None, ): """Background thread for gateway-only periodic chores (NOT cron). Separate from the cron trigger so chores run under any ``CronScheduler`` provider (external scale-to-zero has no 60s loop). @@ -4564,6 +4565,10 @@ def _start_gateway_housekeeping( (60, "Paste sweep", _housekeeping_paste_sweep)] if cron_provider is not None: chores.append((5, "Misfire catch-up sweep", lambda: _housekeeping_misfire_catch_up(cron_provider, adapters, loop))) + if cron_thread is not None: + # The ticker's own guards keep its loop alive; this is the outer layer for a thread that has + # already ended (#111010). Runs every tick so the outage is bounded by one housekeeping interval. + chores.append((1, "Cron ticker supervisor", cron_thread.restart_if_dead)) chores += [ (60, "Curator tick", _housekeeping_curator), (60, "Sync pull tick", _housekeeping_skill_sync), @@ -5113,9 +5118,10 @@ def _start_gateway_start_cron_and_housekeeping(runner): if isinstance(cron_provider, InProcessCronScheduler): cron_start_kwargs["can_dispatch"] = lambda: not ( runner._draining or runner._external_drain_active) - cron_thread = threading.Thread( - target=cron_provider.start, args=(cron_stop,), kwargs=cron_start_kwargs, daemon=True, - name="cron-scheduler") + # Supervised: a ticker that dies without a stop request is respawned by housekeeping (#111010). + from cron.scheduler_thread import SupervisedTickerThread + cron_thread = SupervisedTickerThread( + cron_provider.start, args=(cron_stop,), kwargs=cron_start_kwargs, stop_event=cron_stop) cron_thread.start() # External providers fire over loopback HTTP to THIS process's api_server; if it never came up (usually @@ -5140,7 +5146,7 @@ def _start_gateway_start_cron_and_housekeeping(runner): housekeeping_thread = threading.Thread( target=_start_gateway_housekeeping, args=(cron_stop,), kwargs={"adapters": runner.adapters, "loop": asyncio.get_running_loop(), - "cron_provider": cron_provider, "runner": runner}, + "cron_provider": cron_provider, "runner": runner, "cron_thread": cron_thread}, daemon=True, name="gateway-housekeeping") housekeeping_thread.start() return cron_stop, cron_provider, cron_thread, housekeeping_thread @@ -5148,7 +5154,7 @@ def _start_gateway_start_cron_and_housekeeping(runner): async def _start_gateway_shutdown_tail( runner, _control_server, cron_stop: threading.Event, cron_provider, - cron_thread: threading.Thread, housekeeping_thread: threading.Thread, + cron_thread: Any, housekeeping_thread: threading.Thread, _planned_stop_watcher_stop: threading.Event, _planned_stop_watcher_thread: threading.Thread, _signal_initiated_shutdown: list) -> bool: """Post-``wait_for_shutdown`` teardown; returns the process exit verdict (True = exit 0).""" diff --git a/tests/cron/test_ticker_startup_survival.py b/tests/cron/test_ticker_startup_survival.py index a12361fe22..4884741228 100644 --- a/tests/cron/test_ticker_startup_survival.py +++ b/tests/cron/test_ticker_startup_survival.py @@ -1,8 +1,7 @@ -"""The ticker's own contract — "an exception must not silently kill the daemon thread" — -must hold for every step that runs on the cron-scheduler thread, not just the tick body: -startup recovery, per-cycle profile enumeration/gating, and the status-marker writes. -The gateway starts this thread without a supervisor, so any escaping exception ends cron -silently while the gateway keeps running (#111010). +"""The ticker contract — "an exception must not silently kill the cron thread" — holds for every +step that runs on the cron-scheduler thread, not only the tick body: startup recovery, the +multiplex per-cycle enumeration/gate, and the status-marker writes. And when a ticker HAS ended, +gateway housekeeping respawns it (#111010). """ import threading @@ -20,130 +19,89 @@ def _wait_until(predicate, timeout=10.0, interval=0.005): return predicate() -def test_ticker_survives_startup_recovery_crash(tmp_path): - """A store that blows up during startup recovery must not kill the ticker thread: - the gateway keeps running and cron silently stops firing forever (#111010).""" +def _run_ticker(provider, stop, **kwargs): + thread = threading.Thread(target=provider.start, args=(stop,), kwargs={"interval": 0, **kwargs}, + daemon=True, name="cron-scheduler") + thread.start() + return thread + + +def test_ticker_survives_a_corrupt_ledger_at_startup(tmp_path, monkeypatch): + """A real corrupt ``executions.db`` (sqlite ``file is not a database``) used to escape + ``start()`` before the guarded loop and end the thread: gateway up, no job ever fires.""" from cron.scheduler_provider import InProcessCronScheduler + monkeypatch.setenv("HERMES_HOME", str(tmp_path)) + (tmp_path / "cron").mkdir() + (tmp_path / "cron" / "executions.db").write_bytes(b"not a sqlite database" * 64) ticks = [] - - def fake_tick(*args, **kwargs): - ticks.append(1) - return 0 - - def broken_recover(self): - raise RuntimeError("sqlite3.DatabaseError: database disk image is malformed") - stop = threading.Event() - provider = InProcessCronScheduler() - with ( - patch("cron.scheduler.tick", side_effect=fake_tick), - patch("cron.jobs.record_ticker_heartbeat", lambda **kw: None), - patch("cron.jobs.record_ticker_error", lambda *a, **kw: None), - patch.object(InProcessCronScheduler, "recover_interrupted", broken_recover), - ): - thread = threading.Thread( - target=provider.start, - args=(stop,), - kwargs={"interval": 0}, - daemon=True, - name="cron-scheduler", - ) - thread.start() - _wait_until(lambda: len(ticks) >= 1) - alive_after_first_tick = thread.is_alive() + with patch("cron.scheduler.tick", side_effect=lambda **kw: ticks.append(1)): + thread = _run_ticker(InProcessCronScheduler(), stop) + _wait_until(lambda: len(ticks) >= 2) + alive = thread.is_alive() stop.set() thread.join(timeout=5) - - assert alive_after_first_tick, "ticker thread died during startup recovery" - assert ticks, "ticker never ticked after startup recovery crashed" - assert not thread.is_alive() + assert alive and len(ticks) >= 2, "ticker died on the pre-loop recovery scan" + assert (tmp_path / "cron" / "ticker_heartbeat").exists() -def test_multiplex_ticker_survives_profile_gate_crash(tmp_path): - """A raising profile_gate callable runs on the ticker thread outside the guarded - tick body; its failure must degrade to a failed cycle, not a dead thread.""" +def test_multiplex_ticker_survives_a_raising_gate_without_ticking_ungated(tmp_path, caplog): + """A ``profile_gate`` that raises (Desktop stand-down probe) must fail the CYCLE: thread alive, + error logged, and — critically — zero ticks, since an unfiltered home list would tick the very + profiles the gate stands down for (#100489).""" from cron.scheduler_provider import InProcessCronScheduler home = tmp_path / "default" (home / "cron").mkdir(parents=True) - ticks = [] + stop = threading.Event() - def fake_tick(*args, **kwargs): - ticks.append(1) - return 0 - - def broken_gate(name, home): + def broken_gate(name, h): raise RuntimeError("gate boom") - stop = threading.Event() - provider = InProcessCronScheduler() - with ( - patch("cron.scheduler.tick", side_effect=fake_tick), - patch("cron.jobs.record_ticker_heartbeat", lambda **kw: None), - ): - thread = threading.Thread( - target=provider.start, - args=(stop,), - kwargs={ - "interval": 0, - "profile_homes": [("default", home)], - "profile_gate": broken_gate, - }, - daemon=True, - name="cron-scheduler", - ) - thread.start() - time.sleep(0.5) - alive_after_failed_cycles = thread.is_alive() + with patch("cron.scheduler.tick", side_effect=lambda **kw: ticks.append(1)), caplog.at_level("ERROR"): + thread = _run_ticker(InProcessCronScheduler(), stop, + profile_homes=[("default", home)], profile_gate=broken_gate) + _wait_until(lambda: sum("gate boom" in r.getMessage() for r in caplog.records) >= 2) + alive = thread.is_alive() stop.set() thread.join(timeout=5) - - assert alive_after_failed_cycles, "ticker thread died when profile_gate raised" - assert not thread.is_alive() + assert alive, "ticker thread died when profile_gate raised" + assert ticks == [], "a raising gate must not tick the ungated home list" -def test_ticker_survives_marker_write_crash(tmp_path): - """A misbehaving marker write (SystemExit from the store layer) must not end the - ticker — the loop body already swallows exactly this class of failure from the - provider SDKs (#32612); the marker writes deserve the same protection.""" - from cron.scheduler_provider import InProcessCronScheduler +def test_housekeeping_restarts_a_dead_ticker(monkeypatch): + """The supervisor is the outer layer: a ticker that ended without a stop request is respawned + on the next housekeeping tick; one that ended BECAUSE of the stop request is not.""" + import gateway.run as gateway_run + from cron.scheduler_thread import SupervisedTickerThread - ticks = [] - heartbeat_calls = {"n": 0} + class _OneTick: + def __init__(self): + self.waited = False - def fake_tick(*args, **kwargs): - ticks.append(1) - return 0 + def is_set(self): + return self.waited - def flaky_heartbeat(**kwargs): - heartbeat_calls["n"] += 1 - if heartbeat_calls["n"] == 2: - raise SystemExit("misbehaving marker write") + def wait(self, timeout=None): + self.waited = True + return True + starts = [] stop = threading.Event() - provider = InProcessCronScheduler() - with ( - patch("cron.scheduler.tick", side_effect=fake_tick), - patch("cron.jobs.record_ticker_heartbeat", side_effect=flaky_heartbeat), - patch.object(InProcessCronScheduler, "recover_interrupted", lambda self: 0), - patch("cron.jobs.record_ticker_error", lambda *a, **kw: None), - ): - thread = threading.Thread( - target=provider.start, - args=(stop,), - kwargs={"interval": 0}, - daemon=True, - name="cron-scheduler", - ) - thread.start() - _wait_until(lambda: len(ticks) >= 2) - alive_past_crash = thread.is_alive() - stop.set() - thread.join(timeout=5) - assert heartbeat_calls["n"] >= 2, "flaky heartbeat write never fired" - assert alive_past_crash, "ticker thread died on a marker-write SystemExit" - assert len(ticks) >= 2, "ticker stopped ticking after the marker-write crash" - assert not thread.is_alive() + def dying_ticker(stop_event): + starts.append(threading.current_thread().name) + raise RuntimeError("escaped") + + ticker = SupervisedTickerThread(dying_ticker, args=(stop,), stop_event=stop) + with patch("threading.excepthook", lambda args: None): + ticker.start() + _wait_until(lambda: not ticker.is_alive()) + gateway_run._start_gateway_housekeeping(_OneTick(), interval=0, cron_thread=ticker) + _wait_until(lambda: len(starts) == 2 and not ticker.is_alive()) + assert len(starts) == 2 and ticker.restarts == 1 + stop.set() + gateway_run._start_gateway_housekeeping(_OneTick(), interval=0, cron_thread=ticker) + assert len(starts) == 2, "a stopped ticker must not be respawned"