fix(cron): isolate multiplex profile failures per profile (#74878)
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>
This commit is contained in:
@@ -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}"
|
||||
|
||||
@@ -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)}
|
||||
|
||||
Reference in New Issue
Block a user