fix(cron): stale ticker yields its tick to a fresh gateway

A long-lived process whose checkout was updated underneath it (hot git
pull, interrupted hermes update) serves mixed sys.modules. When such a
stale process races a fresh gateway for the cron tick lock and wins the
minute, every agent job it dispatches can die on ImportErrors whose real
cause is staleness — and the fresh gateway's ticker skips the same minute
as lock-loser, so the user's scheduled job fires broken or not at all.

tick() now checks, BEFORE acquiring the tick lock:

  skew detected (boot fingerprint != disk revision)
    AND this process does not own the gateway runtime lock
    AND that lock is held (a fresh gateway is alive)
      -> raise CronTickYielded, skipping the tick entirely

Each arm alone keeps the old behavior:
- skew + self-owned lock -> proceed (delivery-path stale-code hint stays
  the surface for gateway-owned dispatches)
- skew + no lock holder -> proceed (desktop-standalone users must not
  lose their only ticker to a silent yield)
- skew None (non-git install, no boot fingerprint, probe failure) ->
  proceed; yielding is a certainty claim, never a guess

The yield RAISES instead of returning 0 so the provider loops record it
via record_ticker_error and mark the heartbeat success=False — a yielded
tick must not look like a healthy one (hermes cron status shows why),
mirroring the EMFILE propagation contract (#87644). Yield logging is
throttled to once per skew episode. Self-healing: when the fresh gateway
dies, its lock releases and the stale ticker's next tick proceeds.

Multiplex loop: a yield for one profile no longer cancels sibling
profiles' ticks in the same cycle; only the yielding profile records an
unsuccessful beat.

gateway/status.py gains owns_gateway_runtime_lock() —
is_gateway_runtime_lock_active() is True for the lock's own owner too, so
a caller deciding whether to yield to a FRESH gateway needs the
in-process handle as the discriminator.
This commit is contained in:
Jay. (neocode24)
2026-08-28 07:44:12 +09:00
committed by Teknium
parent fd1d8271db
commit 9a7732b45f
6 changed files with 479 additions and 8 deletions

View File

@@ -0,0 +1 @@
neocode24

View File

@@ -206,6 +206,91 @@ def _detect_gateway_code_skew() -> tuple[str, str] | None:
return None
class CronTickYielded(RuntimeError):
"""A stale-code ticker yielded this tick to a fresh gateway.
Raised by ``tick()`` BEFORE the tick lock is acquired when the process is
provably running stale code (boot fingerprint ≠ disk), it does NOT own the
gateway runtime lock, and that lock is held by another (fresh) process.
Fresh code picks the job up within one tick interval, so the stale process
must stay out of the dispatch race entirely — including lock contention,
which would otherwise starve the fresh ticker on a busy minute.
Raised instead of returned so the provider loops
(``cron/scheduler_provider.py``) record it via ``record_ticker_error`` and
mark the heartbeat ``success=False``: a yielded tick is NOT a healthy tick
(``hermes cron status`` must not show green while jobs only fire from the
other process). Liveness stays visible — the loop keeps beating and keeps
yielding; if the fresh gateway dies, its lock releases and the stale
ticker's next tick proceeds normally (self-healing, no restart needed).
Skew detection returning ``None`` (non-git install, no boot fingerprint —
e.g. a one-shot CLI tick, or any probe failure) never yields: yield only
on certainty, fail open otherwise.
"""
def __init__(self, boot_rev: str, disk_rev: str) -> None:
self.boot_rev = boot_rev
self.disk_rev = disk_rev
super().__init__(
f"Cron tick yielded to a fresh gateway process (stale code: "
f"booted on {boot_rev}, disk is at {disk_rev})"
)
# Log the yield at most once per episode: a stale ticker that keeps yielding
# for hours must not spam the error log every interval. Reset when the
# condition clears (proceeds without yielding) or the skew changes.
_YIELD_LOG_INTERVAL_SECONDS = 3600.0
_last_yield_log: dict[str, object] = {}
def _should_yield_tick_to_fresh_gateway() -> tuple[str, str] | None:
"""Decide whether this tick must yield to a fresher gateway process.
Returns the ``(boot_rev, disk_rev)`` skew labels when ALL of: this process
has a boot fingerprint that differs from the checkout on disk (code
skew), it does not own the gateway runtime lock, and some other process
currently holds that lock — i.e. a fresh gateway is alive and will
dispatch due jobs itself. Returns ``None`` otherwise.
Every probe failure returns ``None``: the gateway-status import, the lock
probe, and skew detection are each individually fail-open. Yielding is a
certainty claim, never a guess.
"""
skew = _detect_gateway_code_skew()
if skew is None:
return None
try:
from gateway import status as _gateway_status
except Exception:
return None
try:
if _gateway_status.owns_gateway_runtime_lock():
return None
if not _gateway_status.is_gateway_runtime_lock_active():
return None
except Exception:
return None
return skew
def _log_tick_yield_once(reason: str) -> None:
"""Log the yield at error level once per episode (skew signature)."""
global _last_yield_log
now = time.monotonic()
last_reason = _last_yield_log.get("reason")
last_at = _last_yield_log.get("at", 0.0)
if last_reason != reason or (now - float(last_at)) >= _YIELD_LOG_INTERVAL_SECONDS:
logger.error(
"Cron tick yielded: this process is running stale code (%s) and a "
"fresher gateway owns the runtime lock — jobs will fire from that "
"process. Restart this one to reclaim its ticks.",
reason,
)
_last_yield_log = {"reason": reason, "at": now}
def _summarize_cron_failure_for_delivery(job: dict, error: str | None) -> str:
"""Return a compact one-line failure message for chat delivery.
@@ -7934,6 +8019,22 @@ def tick(
Returns:
Number of jobs executed (0 if another tick is already running)
"""
# Stale-code yield gate — BEFORE the lock race (#stale-tick-preemption).
# A long-lived process whose checkout was updated underneath it (hot
# ``git pull``, interrupted ``hermes update``) serves MIXED sys.modules:
# every agent job it dispatches can die on ImportErrors whose real cause
# is staleness. When this process is provably stale AND a fresher
# process holds the gateway runtime lock, that process's own ticker
# dispatches due jobs — this one must not even enter the lock race and
# preempt dispatch on a busy minute. When no fresh holder exists
# (desktop-standalone users), yielding would silently kill the user's
# only ticker, so the tick proceeds and job failures surface through the
# delivery path's stale-code hint instead.
_skew = _should_yield_tick_to_fresh_gateway()
if _skew is not None:
_log_tick_yield_once(f"boot={_skew[0]} disk={_skew[1]}")
raise CronTickYielded(_skew[0], _skew[1])
lock_dir, lock_file = _get_lock_paths()
_ensure_cron_dir(lock_dir)

View File

@@ -562,6 +562,7 @@ class InProcessCronScheduler(CronScheduler):
default_profile=None,
):
import logging
from cron.scheduler import CronTickYielded
from cron.scheduler import tick as cron_tick
from cron.jobs import (
clear_ticker_error,
@@ -630,7 +631,14 @@ class InProcessCronScheduler(CronScheduler):
# stop_event (set by the main thread's signal handler), not by
# an exception in this daemon thread, so swallowing it and
# re-checking stop_event keeps shutdown clean.
logger.error("Cron tick error: %s", e, exc_info=True)
if isinstance(e, CronTickYielded):
# Expected while this process is stale and a fresh gateway
# owns the runtime lock: not an error to debug, but it IS
# recorded below so status shows why ticks aren't firing
# from here. tick() already logged it once per episode.
logger.info("Cron tick yielded: %s", e)
else:
logger.error("Cron tick error: %s", e, exc_info=True)
# Persist the failure reason next to the heartbeat markers so
# `hermes cron status`/`list` (separate processes) can show
# WHY ticks fail, not just that the success marker is stale —
@@ -673,6 +681,7 @@ class InProcessCronScheduler(CronScheduler):
"""
import logging
from cron.scheduler import tick as cron_tick
from cron.scheduler import CronTickYielded
from cron.jobs import (
clear_ticker_error,
record_ticker_error,
@@ -712,6 +721,7 @@ class InProcessCronScheduler(CronScheduler):
while not stop_event.is_set():
ok = False
_tick_error = None
_profile_errors: dict[str, str] = {}
try:
if can_dispatch is not None and not can_dispatch():
logger.debug("Cron dispatch paused while gateway drains existing work")
@@ -742,28 +752,44 @@ class InProcessCronScheduler(CronScheduler):
sync=False,
can_dispatch=can_dispatch,
)
except CronTickYielded as e:
# This profile is served stale and a fresh
# gateway owns its runtime lock: record the yield
# for THIS profile's status only, and keep
# ticking the remaining profiles — one profile's
# fresh gateway must not cancel another profile's
# 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}"
finally:
reset_hermes_home_override(home_token)
ok = True
ok = not _profile_errors
except BaseException as e:
logger.error("Cron tick error: %s", e, exc_info=True)
_tick_error = f"{type(e).__name__}: {e}"
# 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.
# Record per-profile heartbeat after each tick cycle. Distinguish
# a COMPLETED cycle (``_tick_error`` unset) — where each profile's
# beat reflects its own outcome, so a yielding profile does not
# darken healthy siblings — from an aborted one (exception), where
# no profile completed and all beats are unsuccessful (#32612).
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))
try:
with use_cron_store(home):
record_ticker_heartbeat(success=ok)
_home_ok = (
_tick_error is None and str(home) not in _profile_errors
)
record_ticker_heartbeat(success=_home_ok)
# Surface the failure reason (or clear it) per profile
# so `hermes cron status` can show WHY ticks fail
# (#68483).
if ok:
if _home_ok:
clear_ticker_error()
elif str(home) in _profile_errors:
record_ticker_error(_profile_errors[str(home)])
elif _tick_error:
record_ticker_error(_tick_error)
finally:

View File

@@ -1000,6 +1000,20 @@ def release_gateway_runtime_lock() -> None:
_clear_running_pid_cache()
def owns_gateway_runtime_lock() -> bool:
"""Return True when THIS process holds the gateway runtime lock.
``is_gateway_runtime_lock_active`` answers "does *anyone* hold the lock?"
and deliberately returns True for the lock's own owner — a caller deciding
whether to yield to a *fresh* gateway (e.g. the cron tick's stale-code
gate) must distinguish self-ownership from another process's lock, and
re-probing the lock file cannot: acquiring a probe handle on a lock this
process already holds succeeds on POSIX. The in-process handle is the only
discriminator, so expose it as a tiny predicate.
"""
return _gateway_lock_handle is not None
def is_gateway_runtime_lock_active(lock_path: Optional[Path] = None) -> bool:
"""Return True when some process currently owns the gateway runtime lock."""
global _gateway_lock_handle

View File

@@ -249,9 +249,19 @@ def cron_list(show_all: bool = False):
def cron_tick():
"""Run due jobs once and exit."""
from cron.scheduler import tick
from cron.scheduler import CronTickYielded, tick
try:
tick(verbose=True)
except CronTickYielded as exc:
# Not expected on this surface (a one-shot CLI process has no boot
# fingerprint, so the yield gate is inert) — but if a future caller
# records one, report cleanly instead of a traceback.
print(color(f"✗ {exc}", Colors.YELLOW))
print(
" A fresher gateway process owns the runtime lock and will fire "
"due jobs; this stale process yielded its tick."
)
return 1
except OSError as exc:
# tick() now propagates real lock-acquisition failures (EMFILE,
# EACCES on open, ...) instead of swallowing them as contention

View File

@@ -0,0 +1,319 @@
"""Stale-code cron tick yield gate.
A long-lived process whose checkout was updated underneath it (hot git pull /
interrupted ``hermes update``) serves mixed ``sys.modules``; when such a
process races a fresh gateway for the cron tick lock and wins, every agent
job it dispatches can die on ImportErrors whose real cause is staleness.
The gate in ``cron.scheduler.tick`` (before lock acquisition) yields — raises
``CronTickYielded`` — only when ALL of:
1. this process is provably stale (boot fingerprint ≠ disk revision),
2. it does NOT own the gateway runtime lock, and
3. some other process currently holds that lock.
Any probe failure or missing fingerprint means "proceed" (fail-open). Yielded
ticks must surface as failed ticks (``record_ticker_error`` + heartbeat
``success=False``), never as healthy ones.
"""
from __future__ import annotations
import logging
import threading
import time
from unittest.mock import patch
import pytest
import cron.scheduler as scheduler_mod
SKEW = ("boot0123abcd", "disk4567efgh")
def _wait_until(predicate, timeout=10.0, interval=0.005):
"""Block until ``predicate()`` is truthy or ``timeout`` elapses."""
deadline = time.monotonic() + timeout
while time.monotonic() < deadline:
value = predicate()
if value:
return value
time.sleep(interval)
return predicate()
def _gate_mocks(monkeypatch, *, owns: bool, active: bool, skew=SKEW):
"""Point the yield gate's three probes at fixed answers."""
from gateway import status as gateway_status
monkeypatch.setattr(scheduler_mod, "_detect_gateway_code_skew", lambda: skew)
monkeypatch.setattr(gateway_status, "owns_gateway_runtime_lock", lambda: owns)
monkeypatch.setattr(
gateway_status, "is_gateway_runtime_lock_active", lambda lock_path=None: active
)
class TestTickYieldGate:
def test_skew_plus_foreign_gateway_yields_before_lock(self, monkeypatch, tmp_path):
"""(a) skew + a fresh foreign gateway holds the runtime lock → the
tick yields WITHOUT acquiring the tick lock or dispatching."""
_gate_mocks(monkeypatch, owns=False, active=True)
def _no_lock_allowed(*_a, **_kw):
raise AssertionError("tick lock must not be acquired after a yield")
monkeypatch.setattr(scheduler_mod, "_get_lock_paths", _no_lock_allowed)
with pytest.raises(scheduler_mod.CronTickYielded) as excinfo:
scheduler_mod.tick(verbose=False)
# The exception carries the forensics for the ticker-error record.
assert SKEW[0] in str(excinfo.value)
assert SKEW[1] in str(excinfo.value)
def test_skew_plus_self_owned_lock_proceeds(self, monkeypatch, tmp_path):
"""(b) skew + this process IS the gateway → proceed. The delivery
path's stale-code hint (part 3) stays the surface for this case."""
_gate_mocks(monkeypatch, owns=True, active=True)
# Reaching the lock paths proves the gate let the tick proceed.
called = []
real = scheduler_mod._get_lock_paths
def _spy():
called.append(1)
return real()
monkeypatch.setattr(scheduler_mod, "_get_lock_paths", _spy)
assert scheduler_mod.tick(verbose=False) == 0 # empty store: no jobs
assert called, "tick must proceed when this process owns the gateway lock"
def test_skew_plus_no_gateway_proceeds(self, monkeypatch, tmp_path):
"""(c) skew + no gateway alive (desktop-standalone) → proceed:
yielding would silently kill the user's only ticker."""
_gate_mocks(monkeypatch, owns=False, active=False)
assert scheduler_mod.tick(verbose=False) == 0
def test_no_fingerprint_proceeds(self, monkeypatch, tmp_path):
"""(d) skew is None (non-git install, no boot fingerprint, probe
failure) → always proceed."""
_gate_mocks(monkeypatch, owns=False, active=True, skew=None)
assert scheduler_mod.tick(verbose=False) == 0
def test_lock_probe_exception_proceeds(self, monkeypatch, tmp_path):
"""Fail-open: the lock probe raising must not yield or crash the
tick."""
from gateway import status as gateway_status
monkeypatch.setattr(scheduler_mod, "_detect_gateway_code_skew", lambda: SKEW)
def _boom(lock_path=None):
raise OSError("lock probe failed")
monkeypatch.setattr(gateway_status, "is_gateway_runtime_lock_active", _boom)
assert scheduler_mod.tick(verbose=False) == 0
def test_yield_logs_once_per_episode(self, monkeypatch, caplog):
"""The yield log is throttled: repeated yields with the same skew
signature log once, not once per tick interval."""
_gate_mocks(monkeypatch, owns=False, active=True)
with caplog.at_level(logging.ERROR, logger="cron.scheduler"):
scheduler_mod._log_tick_yield_once("boot=a disk=b")
scheduler_mod._log_tick_yield_once("boot=a disk=b")
scheduler_mod._log_tick_yield_once("boot=a disk=b")
yield_logs = [
r for r in caplog.records if "Cron tick yielded" in r.getMessage()
]
assert len(yield_logs) == 1
# A NEW skew signature (the checkout moved again) is a new episode.
with caplog.at_level(logging.ERROR, logger="cron.scheduler"):
scheduler_mod._log_tick_yield_once("boot=a disk=c")
yield_logs = [
r for r in caplog.records if "Cron tick yielded" in r.getMessage()
]
assert len(yield_logs) == 2
class TestYieldedTickIsAFailedTick:
"""(e) A yielded tick must not be recorded as a healthy tick."""
def test_provider_records_error_and_unsuccessful_heartbeat(self):
from cron.scheduler_provider import InProcessCronScheduler
beats: list[bool] = []
errors: list[str] = []
stop = threading.Event()
prov = InProcessCronScheduler()
with patch(
"cron.scheduler.tick",
side_effect=scheduler_mod.CronTickYielded(SKEW[0], SKEW[1]),
), patch(
"cron.jobs.record_ticker_heartbeat",
side_effect=lambda success=False: beats.append(success),
), patch(
"cron.jobs.record_ticker_error",
side_effect=lambda msg: errors.append(msg),
), patch("cron.jobs.clear_ticker_error") as clear:
t = threading.Thread(
target=prov.start, args=(stop,), kwargs={"interval": 0}, daemon=True
)
t.start()
assert _wait_until(lambda: len(beats) >= 3), "ticker did not keep beating"
stop.set()
t.join(timeout=5)
assert not t.is_alive(), "ticker must keep yielding, not die"
assert errors, "yield reason must be persisted for `hermes cron status`"
assert "CronTickYielded" in errors[0] or "yielded" in errors[0]
assert beats[-1] is False, "a yielded tick is not a successful tick"
clear.assert_not_called()
def test_ticker_takes_over_when_gateway_lock_releases(self):
"""Self-healing: once the fresh gateway dies (lock inactive), the
yielding ticker's next tick proceeds normally."""
from cron.scheduler_provider import InProcessCronScheduler
beats: list[bool] = []
state = {"foreign_gateway_alive": True}
stop = threading.Event()
prov = InProcessCronScheduler()
def _tick(*args, **kwargs):
# Gate logic is exercised for real; emulate the two probe arms.
if state["foreign_gateway_alive"]:
raise scheduler_mod.CronTickYielded(SKEW[0], SKEW[1])
return 0
with patch("cron.scheduler.tick", side_effect=_tick), patch(
"cron.jobs.record_ticker_heartbeat",
side_effect=lambda success=False: beats.append(success),
), patch("cron.jobs.record_ticker_error"), patch("cron.jobs.clear_ticker_error"):
t = threading.Thread(
target=prov.start, args=(stop,), kwargs={"interval": 0}, daemon=True
)
t.start()
assert _wait_until(lambda: len(beats) >= 2)
state["foreign_gateway_alive"] = False # gateway lock released
assert _wait_until(lambda: True in beats), "ticker never took over"
stop.set()
t.join(timeout=5)
assert not t.is_alive()
def test_multiplex_yield_in_one_profile_does_not_cancel_siblings(self):
"""Multiplex loop: a yield for profile A (its fresh gateway owns the
runtime lock) must not cancel profile B's tick in the same cycle —
B may have no other ticker. B's heartbeat stays healthy; A's records
the yield."""
from cron.scheduler_provider import InProcessCronScheduler
home_a = "/tmp/home-a"
home_b = "/tmp/home-b"
ticked: list[str] = []
per_home_beats: dict[str, list[bool]] = {home_a: [], home_b: []}
per_home_errors: dict[str, list[str]] = {home_a: [], home_b: []}
state = {"cycles": 0}
stop = threading.Event()
prov = InProcessCronScheduler()
def _tick(*args, **kwargs):
from hermes_constants import get_hermes_home
home = str(get_hermes_home())
ticked.append(home)
if home == home_a:
# Profile A: stale process, fresh foreign gateway → yield.
raise scheduler_mod.CronTickYielded(SKEW[0], SKEW[1])
# Profile B ticks fine.
return 0
def _beat(success=False):
from hermes_constants import get_hermes_home
per_home_beats[str(get_hermes_home())].append(success)
def _err(msg):
from hermes_constants import get_hermes_home
per_home_errors[str(get_hermes_home())].append(msg)
with patch("cron.scheduler.tick", side_effect=_tick), patch(
"cron.jobs.record_ticker_heartbeat", side_effect=_beat
), patch("cron.jobs.record_ticker_error", side_effect=_err), patch(
"cron.jobs.clear_ticker_error"
):
t = threading.Thread(
target=prov.start,
args=(stop,),
kwargs={
"interval": 0,
"profile_homes": [home_a, home_b],
},
daemon=True,
)
t.start()
assert _wait_until(lambda: len(ticked) >= 6), "multiplex loop did not run"
stop.set()
t.join(timeout=5)
assert not t.is_alive()
# Every cycle ticked BOTH homes — B was never cancelled by A's yield.
assert home_b in ticked and home_a in ticked
# A recorded unsuccessful beats + its yield reason; B stayed healthy.
# (The first beat per home is the pre-loop initial heartbeat, always
# unsuccessful — skip it when judging cycle outcomes.)
assert per_home_beats[home_a], "profile A must still beat (liveness)"
assert all(b is False for b in per_home_beats[home_a][1:])
assert per_home_errors[home_a], "profile A must record the yield reason"
assert per_home_beats[home_b][1:] and all(
b is True for b in per_home_beats[home_b][1:]
), "profile B must stay healthy"
assert not per_home_errors[home_b], "profile B must not inherit A's yield"
class TestGatewayLockOwnershipProbe:
"""The self-ownership accessor backing the gate, against the real lock."""
def test_ownership_follows_acquire_and_release(self, tmp_path, monkeypatch):
from gateway import status as gateway_status
monkeypatch.setenv("HERMES_HOME", str(tmp_path))
assert gateway_status.owns_gateway_runtime_lock() is False
assert gateway_status.acquire_gateway_runtime_lock() is True
try:
assert gateway_status.owns_gateway_runtime_lock() is True
# Self-ownership is distinguishable even though the shared
# liveness probe also reports active for the owner.
assert gateway_status.is_gateway_runtime_lock_active() is True
finally:
gateway_status.release_gateway_runtime_lock()
assert gateway_status.owns_gateway_runtime_lock() is False
def test_ownership_probe_is_profile_scoped_by_process_home(
self, tmp_path, monkeypatch
):
"""The gateway lock is a process-level identity file (#56986): under
a multiplex profile-home override the probe must still resolve the
launch home's lock, not the overridden profile's."""
from gateway import status as gateway_status
from hermes_constants import (
reset_hermes_home_override,
set_hermes_home_override,
)
monkeypatch.setenv("HERMES_HOME", str(tmp_path))
other_profile = tmp_path / "profiles" / "other"
other_profile.mkdir(parents=True)
token = set_hermes_home_override(str(other_profile))
try:
assert gateway_status.acquire_gateway_runtime_lock() is True
try:
assert gateway_status.owns_gateway_runtime_lock() is True
# The lock file lives in the launch home, not the override.
assert (tmp_path / "gateway.lock").exists()
assert not (other_profile / "gateway.lock").exists()
finally:
gateway_status.release_gateway_runtime_lock()
finally:
reset_hermes_home_override(token)