refactor(cron): dedup EMFILE tick-failure handling; share the fd-exhaustion text matcher
/simplify-code findings on the salvage stack: - the classify+reclaim+counter block was pasted verbatim into both ticker loops (_start and _start_multiplex) along with duplicated function-local imports — extracted _note_tick_failure() next to _backoff_wait_seconds so both loops share one implementation. - hermes_cli/cron.py's EMFILE hint reimplemented the text half of _is_fd_exhaustion with a case-SENSITIVE variation (drift risk) — split _is_fd_exhaustion_text() out and use it from both. 11 EMFILE tests + 54 provider/ticker tests green; ruff clean.
This commit is contained in:
+7
-2
@@ -1355,6 +1355,12 @@ def _is_lock_contention_errno(err: OSError) -> bool:
|
||||
return False
|
||||
|
||||
|
||||
def _is_fd_exhaustion_text(text: str) -> bool:
|
||||
"""Text-level half of :func:`_is_fd_exhaustion` (shared with the CLI hint)."""
|
||||
lowered = text.lower()
|
||||
return "too many open files" in lowered or "emfile" in lowered
|
||||
|
||||
|
||||
def _is_fd_exhaustion(exc: BaseException) -> bool:
|
||||
"""Return True when *exc* indicates file-descriptor exhaustion.
|
||||
|
||||
@@ -1364,8 +1370,7 @@ def _is_fd_exhaustion(exc: BaseException) -> bool:
|
||||
"""
|
||||
if isinstance(exc, OSError) and exc.errno in (errno.EMFILE, errno.ENFILE):
|
||||
return True
|
||||
text = str(exc).lower()
|
||||
return "too many open files" in text or "emfile" in text
|
||||
return _is_fd_exhaustion_text(str(exc))
|
||||
|
||||
|
||||
def _reclaim_fds_best_effort() -> None:
|
||||
|
||||
+24
-25
@@ -46,6 +46,24 @@ def _backoff_wait_seconds(interval: float, consecutive_failures: int) -> float:
|
||||
)
|
||||
|
||||
|
||||
def _note_tick_failure(exc: BaseException, consecutive_failures: int) -> int:
|
||||
"""Classify one failed tick and return the updated failure counter.
|
||||
|
||||
Shared by both ticker loops (#87644): on fd exhaustion, attempt
|
||||
reclamation (gc.collect + raise the soft nofile limit) so the NEXT tick
|
||||
can succeed, and bump the counter so ``_backoff_wait_seconds`` backs off
|
||||
exponentially while the process has no chance of making progress. Any
|
||||
other failure resets the counter — backoff is reserved for the
|
||||
self-inflicted EMFILE storm, not transient errors.
|
||||
"""
|
||||
from cron.scheduler import _is_fd_exhaustion, _reclaim_fds_best_effort
|
||||
|
||||
if _is_fd_exhaustion(exc):
|
||||
_reclaim_fds_best_effort()
|
||||
return consecutive_failures + 1
|
||||
return 0
|
||||
|
||||
|
||||
class CronScheduler(ABC):
|
||||
"""Axis-B trigger provider. Decides WHEN a due cron job fires.
|
||||
|
||||
@@ -369,10 +387,6 @@ class InProcessCronScheduler(CronScheduler):
|
||||
):
|
||||
import logging
|
||||
from cron.scheduler import tick as cron_tick
|
||||
from cron.scheduler import (
|
||||
_is_fd_exhaustion,
|
||||
_reclaim_fds_best_effort,
|
||||
)
|
||||
from cron.jobs import (
|
||||
clear_ticker_error,
|
||||
record_ticker_error,
|
||||
@@ -446,16 +460,10 @@ class InProcessCronScheduler(CronScheduler):
|
||||
# uid went unnoticed for ~14h with the reason buried in the
|
||||
# gateway log (#68483).
|
||||
record_ticker_error(f"{type(e).__name__}: {e}")
|
||||
if _is_fd_exhaustion(e):
|
||||
# EMFILE: try to free descriptors (gc.collect + raise the
|
||||
# soft nofile limit) so the NEXT tick can succeed; back off
|
||||
# exponentially so the exhausted process stops hammering
|
||||
# the store while it has no chance of making progress
|
||||
# (#87644).
|
||||
_reclaim_fds_best_effort()
|
||||
consecutive_failures += 1
|
||||
else:
|
||||
consecutive_failures = 0
|
||||
# EMFILE: reclaim fds + back off exponentially so the
|
||||
# exhausted process stops hammering the store while it has no
|
||||
# chance of making progress (#87644).
|
||||
consecutive_failures = _note_tick_failure(e, consecutive_failures)
|
||||
# Record liveness every iteration; bump the success marker only on a
|
||||
# clean tick, so status can tell "alive but failing every tick" from
|
||||
# "actually firing jobs" (#32612, #32895).
|
||||
@@ -485,10 +493,6 @@ class InProcessCronScheduler(CronScheduler):
|
||||
"""
|
||||
import logging
|
||||
from cron.scheduler import tick as cron_tick
|
||||
from cron.scheduler import (
|
||||
_is_fd_exhaustion,
|
||||
_reclaim_fds_best_effort,
|
||||
)
|
||||
from cron.jobs import (
|
||||
clear_ticker_error,
|
||||
record_ticker_error,
|
||||
@@ -547,13 +551,8 @@ class InProcessCronScheduler(CronScheduler):
|
||||
except BaseException as e:
|
||||
logger.error("Cron tick error: %s", e, exc_info=True)
|
||||
_tick_error = f"{type(e).__name__}: {e}"
|
||||
if _is_fd_exhaustion(e):
|
||||
# EMFILE: attempt fd reclamation so the next tick can
|
||||
# succeed, and back off exponentially (#87644).
|
||||
_reclaim_fds_best_effort()
|
||||
consecutive_failures += 1
|
||||
else:
|
||||
consecutive_failures = 0
|
||||
# EMFILE: reclaim fds + exponential backoff (#87644).
|
||||
consecutive_failures = _note_tick_failure(e, consecutive_failures)
|
||||
else:
|
||||
_tick_error = None
|
||||
# Record per-profile heartbeat after each tick cycle.
|
||||
|
||||
+2
-1
@@ -281,6 +281,7 @@ def cron_status():
|
||||
get_ticker_success_age,
|
||||
TICKER_INTERVAL_SECONDS,
|
||||
)
|
||||
from cron.scheduler import _is_fd_exhaustion_text as _cron_is_fd_exhaustion_text
|
||||
|
||||
# Allow ~3 missed ticker iterations (+ a little slack) before declaring
|
||||
# trouble. Derived from the shared interval constant so this threshold
|
||||
@@ -323,7 +324,7 @@ def cron_status():
|
||||
"gateway user, and prefer `docker exec -u <uid>:<gid>`.",
|
||||
Colors.YELLOW,
|
||||
))
|
||||
elif "Too many open files" in last_error or "EMFILE" in last_error:
|
||||
elif _cron_is_fd_exhaustion_text(last_error):
|
||||
print(color(
|
||||
" Hint: the ticker hit file-descriptor exhaustion "
|
||||
"(EMFILE). The scheduler now retries with backoff and "
|
||||
|
||||
Reference in New Issue
Block a user