From a2fea79de6bbee87d4e7ddc0b2cdc8f4c79638b2 Mon Sep 17 00:00:00 2001 From: Cyber-Yichen <180837074+Cyber-Yichen@users.noreply.github.com> Date: Wed, 2 Sep 2026 03:34:41 -0700 Subject: [PATCH] fix(cron): isolate multiplex profile failures per profile (#74878) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit One profile's broken cron store no longer takes the whole multiplex ticker down with it: - startup recovery loop: a per-profile exception (e.g. an unreadable executions.db raising sqlite3.DatabaseError) was uncaught and killed the ticker thread before its first tick — no profile ever fired. - tick loop: only CronTickYielded was caught per profile; any other exception escaped to the cycle-wide handler, skipping every remaining profile that cycle and marking all of them failed. Both loops now catch per profile, record the failure into THAT profile's ticker_last_error (`hermes cron status`), and keep ticking the siblings. The existing CronTickYielded/_profile_errors semantics and the #87644 EMFILE reclaim/backoff are preserved (backoff is applied once per cycle from the worst per-profile failure). Salvaged from PR #70747 (@Cyber-Yichen); the recovery test's real sqlite3.OperationalError shape is from PR #74888 (@OYLFLMH). Same class also reported in PR #74952 (@webtecnica). Co-authored-by: OYLFLMH <95945448+OYLFLMH@users.noreply.github.com> Co-authored-by: webtecnica <75556242+webtecnica@users.noreply.github.com> --- cron/scheduler_provider.py | 31 +++++++- tests/cron/test_scheduler_provider.py | 106 ++++++++++++++++++++++++++ 2 files changed, 136 insertions(+), 1 deletion(-) diff --git a/cron/scheduler_provider.py b/cron/scheduler_provider.py index 492a21b522..b6d43d7545 100644 --- a/cron/scheduler_provider.py +++ b/cron/scheduler_provider.py @@ -681,7 +681,7 @@ class InProcessCronScheduler(CronScheduler): """ import logging from cron.scheduler import tick as cron_tick - from cron.scheduler import CronTickYielded + from cron.scheduler import CronTickYielded, _is_fd_exhaustion from cron.jobs import ( clear_ticker_error, record_ticker_error, @@ -701,6 +701,8 @@ class InProcessCronScheduler(CronScheduler): # A profile may have been deleted since this snapshot was taken; # never recreate a deleted home's cron workspace via the heartbeat # below (#47368). + # One profile's broken store (corrupt executions.db, unreadable + # cron dir) must not abort startup for every other profile (#74878). for entry in _existing_profile_homes(profile_homes): home = entry[1] if isinstance(entry, tuple) else entry home_token = set_hermes_home_override(str(home)) @@ -714,6 +716,13 @@ class InProcessCronScheduler(CronScheduler): home, ) record_ticker_heartbeat() + except BaseException as e: + logger.error( + "Cron startup recovery error for profile at %s: %s", + home, + e, + exc_info=True, + ) finally: reset_hermes_home_override(home_token) @@ -722,6 +731,9 @@ class InProcessCronScheduler(CronScheduler): ok = False _tick_error = None _profile_errors: dict[str, str] = {} + # Worst per-profile failure this cycle (fd exhaustion wins) so the + # #87644 backoff/reclaim is applied once per cycle, not per profile. + _cycle_exc: BaseException | None = None try: if can_dispatch is not None and not can_dispatch(): logger.debug("Cron dispatch paused while gateway drains existing work") @@ -761,9 +773,26 @@ class InProcessCronScheduler(CronScheduler): # only ticker in the same cycle. logger.info("Cron tick yielded for profile at %s: %s", home, e) _profile_errors[str(home)] = f"{type(e).__name__}: {e}" + except BaseException as e: + # Any other failure is THIS profile's failure + # (#74878): record it against this profile's + # status and keep ticking the remaining profiles. + # BaseException for the same reason as the + # single-profile loop (#32612). + logger.error( + "Cron tick error for profile at %s: %s", + home, + e, + exc_info=True, + ) + _profile_errors[str(home)] = f"{type(e).__name__}: {e}" + if _cycle_exc is None or _is_fd_exhaustion(e): + _cycle_exc = e finally: reset_hermes_home_override(home_token) ok = not _profile_errors + if _cycle_exc is not None: + consecutive_failures = _note_tick_failure(_cycle_exc, consecutive_failures) except BaseException as e: logger.error("Cron tick error: %s", e, exc_info=True) _tick_error = f"{type(e).__name__}: {e}" diff --git a/tests/cron/test_scheduler_provider.py b/tests/cron/test_scheduler_provider.py index 12ac73560c..6bf710e4ed 100644 --- a/tests/cron/test_scheduler_provider.py +++ b/tests/cron/test_scheduler_provider.py @@ -786,3 +786,109 @@ def test_multiplex_missing_secondary_does_not_fall_back_to_shared(tmp_path): assert default_ad is shared assert sec_ad is not shared assert not sec_ad + + +def test_multiplex_ticker_isolates_profile_failures(tmp_path): + """A failing profile's tick must not skip healthy siblings in the same + cycle, nor darken their status (#74878).""" + from cron.jobs import get_ticker_last_error, record_ticker_error, use_cron_store + from cron.scheduler_provider import InProcessCronScheduler + from hermes_constants import get_hermes_home + + failing_home = tmp_path / "failing" + healthy_home = tmp_path / "healthy" + for home in (failing_home, healthy_home): + (home / "cron").mkdir(parents=True) + with use_cron_store(home): + record_ticker_error("RuntimeError: stale failure") + + stop = threading.Event() + tick_homes: list[str] = [] + + def _tick(*args, **kwargs): + home = str(get_hermes_home()) + tick_homes.append(home) + if home == str(failing_home): + raise RuntimeError("profile-local failure") + stop.set() + return 0 + + provider = InProcessCronScheduler() + with patch("cron.scheduler.tick", side_effect=_tick): + thread = threading.Thread( + target=provider.start, + args=(stop,), + kwargs={ + "interval": 0, + "profile_homes": [("failing", failing_home), ("healthy", healthy_home)], + }, + daemon=True, + ) + thread.start() + thread.join(timeout=5) + stop.set() + thread.join(timeout=5) + + assert not thread.is_alive() + assert str(healthy_home) in tick_homes, "healthy sibling was skipped" + assert not (failing_home / "cron" / "ticker_last_success").exists() + assert (healthy_home / "cron" / "ticker_last_success").exists() + with use_cron_store(failing_home): + assert get_ticker_last_error() == "RuntimeError: profile-local failure" + with use_cron_store(healthy_home): + assert get_ticker_last_error() is None + + +def test_multiplex_recovery_isolates_profile_failures(tmp_path): + """A startup-recovery error in one profile's ledger must not kill the + ticker thread before it ever ticks (#74878).""" + import sqlite3 + + from cron.scheduler_provider import InProcessCronScheduler + from hermes_constants import get_hermes_home + + failing_home = tmp_path / "failing" + healthy_home = tmp_path / "healthy" + for home in (failing_home, healthy_home): + (home / "cron").mkdir(parents=True) + + stop = threading.Event() + recovery_homes: list[str] = [] + tick_homes: list[str] = [] + + def _recover(): + home = str(get_hermes_home()) + recovery_homes.append(home) + if home == str(failing_home): + raise sqlite3.OperationalError("unable to open database file") + return 0 + + def _tick(*args, **kwargs): + tick_homes.append(str(get_hermes_home())) + if len(tick_homes) >= 2: + stop.set() + return 0 + + provider = InProcessCronScheduler() + with ( + patch.object(provider, "recover_interrupted", side_effect=_recover), + patch("cron.scheduler.tick", side_effect=_tick), + ): + thread = threading.Thread( + target=provider.start, + args=(stop,), + kwargs={ + "interval": 0, + "profile_homes": [("failing", failing_home), ("healthy", healthy_home)], + }, + daemon=True, + ) + thread.start() + thread.join(timeout=5) + stop.set() + thread.join(timeout=5) + + assert not thread.is_alive() + assert recovery_homes == [str(failing_home), str(healthy_home)] + # The failing profile stays in rotation: its ledger may still hold jobs. + assert set(tick_homes) == {str(failing_home), str(healthy_home)}