Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
138 changes: 119 additions & 19 deletions cron/scheduler_provider.py
Original file line number Diff line number Diff line change
Expand Up @@ -165,30 +165,130 @@ def name(self) -> str:

def start(self, stop_event, *, adapters=None, loop=None, interval=60):
import logging
import threading
import time
from cron.scheduler import tick as cron_tick
from cron.jobs import record_ticker_heartbeat

logger = logging.getLogger("cron.scheduler_provider")
logger.info("In-process cron scheduler started (interval=%ds)", interval)

# Heartbeat once before the first sleep so `hermes cron status` sees a
# live ticker immediately after startup, not only after the first tick.
record_ticker_heartbeat()
while not stop_event.is_set():
ok = False
self._record_or_log_heartbeat(logger, success=False)

# Independent watchdog: bump the heartbeat on a fixed cadence that
# does NOT depend on cron_tick() returning. The main loop only
# records heartbeats after each tick, so a single tick that blocks
# (e.g. on .tick.lock contention, slow jobs.json I/O, or a long-
# running synchronous code path) leaves the previous heartbeat to
# age until the next iteration finally returns — easily minutes,
# which cron status then reports as "STALLED" (#60703). The watchdog
# runs at half the tick interval, so the heartbeat refreshes at
# least twice per cadence as long as THIS thread is runnable,
# independently of whether cron_tick is making progress. The
# watchdog marks success=False; the main loop's success-marker
# bumps remain the source of truth for "did a tick actually fire".
# Cleans up via stop_event, not atexit — gateway shutdown sets it,
# and we never want an atexit hook racing interpreter finalisation
# while the thread is mid-write.
watchdog_stop = threading.Event()

def _watchdog() -> None:
# Half-cadence refresh; cap to at least 5s and at most the
# tick interval itself so a misconfigured 1s interval can't
# spam the heartbeat file.
cadence = max(5, min(interval // 2, interval))

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This order does not cap the cadence at interval: with the new test's interval=1, it evaluates to 5 seconds. That contradicts the test's one-second expectation and makes its third beat deadline-edge. Apply the lower bound before the upper cap, then adjust the test to avoid timing the final expected beat at the timeout boundary.

while not watchdog_stop.is_set():
self._record_or_log_heartbeat(logger, success=False)
watchdog_stop.wait(cadence)

watchdog_thread = threading.Thread(
target=_watchdog,
name="cron-ticker-heartbeat",
daemon=True,
)
watchdog_thread.start()

try:
while not stop_event.is_set():
ok = False
# Tick-duration guard: if cron_tick overshoots the configured
# interval we surface it as an anomaly instead of letting the
# loop appear hung. The watchdog above keeps the liveness
# signal fresh either way, so an overshoot reports as a logged
# warning, not a missed heartbeat (#60703).
tick_started = time.monotonic()
try:
cron_tick(verbose=False, adapters=adapters, loop=loop, sync=False)
ok = True
except BaseException as e:
# Catch BaseException (not just Exception) so a SystemExit from
# a misbehaving provider SDK / agent retry path does not kill
# the ticker thread silently (#32612). KeyboardInterrupt is
# intentionally caught here too — gateway shutdown is driven by
# 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)
finally:
tick_elapsed = time.monotonic() - tick_started
if tick_elapsed > interval:
logger.warning(
"Cron tick took %.1fs (interval=%ds) — gateway "
"process is alive but the ticker is slow; jobs "
"are firing, just late. Check for long-running "
"synchronous code paths or filesystem stalls "
"(#60703).",
tick_elapsed,
interval,
)
# 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).
self._record_or_log_heartbeat(logger, success=ok)
stop_event.wait(interval)
finally:
# Stop the watchdog so its heartbeat doesn't keep refreshing
# after the main loop exits; otherwise a clean gateway
# shutdown would leave the heartbeat fresh while the ticker
# itself is gone, masking teardown races the next time
# `cron status` runs.
watchdog_stop.set()
# Bound the join so a slow filesystem write can't stall
# shutdown; the thread is daemon=True and will be reaped
# with the process regardless.
watchdog_thread.join(timeout=1.0)

# Defined as a regular method (NOT @staticmethod) so the in-loop
# `self._record_or_log_heartbeat(...)` calls work — the heartbeat
# thread, the pre-loop liveness beat, and the post-tick success beat
# all share this single logger handle carried on `self`. A pure
# staticmethod would force the caller to thread the logger through
# every site, which invites drift if a future caller forgets to
# pass one.
def _record_or_log_heartbeat(self, logger, success: bool) -> None:
"""Heartbeat write that surfaces failures instead of swallowing them.

``record_ticker_heartbeat`` is intentionally best-effort and
silent — a transient disk error must never break the tick loop.
But when the SILENT best-effort path becomes the proximate cause
of a stall report (\"no heartbeat for Ns\" → \"check the
filesystem\"), we want an error-level breadcrumb in the gateway
log. Keeps the lightweight contract for callers while giving the
built-in ticker a chance to leave forensic evidence (#60703).
"""
try:
from cron.jobs import record_ticker_heartbeat
record_ticker_heartbeat(success=success)

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

record_ticker_heartbeat() already catches write failures internally (cron/jobs.py:694-702), so filesystem failures never reach this wrapper's except and no error is logged. Expose success/failure from the lower-level helper or log at the atomic-write boundary, and cover a real forced write failure.

except Exception as exc:
try:
cron_tick(verbose=False, adapters=adapters, loop=loop, sync=False)
ok = True
except BaseException as e:
# Catch BaseException (not just Exception) so a SystemExit from
# a misbehaving provider SDK / agent retry path does not kill
# the ticker thread silently (#32612). KeyboardInterrupt is
# intentionally caught here too — gateway shutdown is driven by
# 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)
# 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).
record_ticker_heartbeat(success=ok)
stop_event.wait(interval)
logger.error(
"Ticker heartbeat write failed (success=%s): %s",
success, exc,
)
except Exception:
# Logger itself may be torn down during interpreter
# finalisation; swallow so the loop itself is never
# jeopardised.
pass
200 changes: 200 additions & 0 deletions tests/cron/test_scheduler_provider.py
Original file line number Diff line number Diff line change
Expand Up @@ -666,3 +666,203 @@ def _boom(provider, base_url):

monkeypatch.setattr(ct, "_validate_cron_base_url", _boom)
assert _guard_job_credential_exfil({"id": "j8", "provider": "anthropic"}) is None


# ── F2c: heartbeat liveness independent of cron_tick() returning (#60703) ───


def test_heartbeat_refreshes_while_tick_blocks(tmp_path):
"""If cron_tick() blocks without returning, the watchdog heartbeat
thread keeps bumping the liveness file at half the tick interval,
so `hermes cron status` does not falsely report STALLED (#60703).

Simulates a hung tick by replacing ``cron.scheduler.tick`` with a
callable that sleeps forever, and asserts that the heartbeat file is
touched at least twice within ~3x the tick interval — proving the
watchdog is independent of the main tick loop making progress.
"""
import time as _t
from cron.scheduler_provider import InProcessCronScheduler

hb_path = tmp_path / "ticker_heartbeat"
hb_path.write_text("0")
age_before = hb_path.read_text()

# Use a tiny interval so the test stays fast. With interval=1 the
# watchdog cadence is max(5, min(0, 1)) = 1s, so we expect ~5
# refreshes over ~5s.
interval = 1
stop = threading.Event()
prov = InProcessCronScheduler()

cursor = {"started": False}
wall_start = _t.monotonic()

def _hung_tick(*a, **k):
cursor["started"] = True
# Block indefinitely; the test will set stop_event so the
# main loop finally exits, but until then this never returns.
stop.wait(60) # outlasts the test
return 0

sweeps = []
import cron.scheduler_provider as _sp
real_record = _sp.InProcessCronScheduler._record_or_log_heartbeat

def _counting(self, logger, success=False):
sweeps.append(_t.monotonic() - wall_start)
return real_record(self, logger, success=success)

with patch("cron.scheduler.tick", side_effect=_hung_tick), \
patch.object(
_sp.InProcessCronScheduler,
"_record_or_log_heartbeat",
new=_counting,
):
t = threading.Thread(
target=prov.start,
args=(stop,),
kwargs={"interval": interval},
daemon=True,
)
t.start()

# Wait for the hung tick to actually start before we start
# counting watchdog hits — otherwise a pre-tick pre-loop liveness
# beat could mislead the count.
deadline = _t.monotonic() + 2.0
while not cursor["started"] and _t.monotonic() < deadline:
_t.sleep(0.01)

# Force at least 3 watchdog sweeps. Watchdog cadence is
# max(5, interval//2, capped at interval). With interval=1,
# that's 1s; over ~3.5s we expect 3+ updates while tick is hung.
target_count = 3
ok = _wait_until(lambda: len(sweeps) >= target_count, timeout=5.0)
assert ok, (
f"watchdog only emitted {len(sweeps)} heartbeat(s) in ~{_t.monotonic() - wall_start:.1f}s "
"while tick was hung — heartbeat is not independent of cron_tick() (#60703)"
)
# Each sweep should advance monotonically.
assert sweeps == sorted(sweeps)

stop.set()
t.join(timeout=2.0)

assert not t.is_alive(), "ticker thread did not exit after stop_event was set"


def test_watchdog_stops_when_main_loop_exits(tmp_path):
"""The watchdog heartbeat thread must terminate when the main
loop exits; otherwise a clean gateway shutdown would leave the
liveness file fresh while the ticker itself is gone, masking
teardown races the next time `cron status` runs (#60703)."""
from cron.scheduler_provider import InProcessCronScheduler
import cron.scheduler_provider as _sp

hb_path = tmp_path / "ticker_heartbeat"
hb_path.write_text("x")

stop = threading.Event()
prov = InProcessCronScheduler()

sweeps_before = []
sweeps_after = []
phase = {"stopped": False}

real_record = _sp.InProcessCronScheduler._record_or_log_heartbeat

def counting(self, logger, success=False):
(sweeps_before if not phase["stopped"] else sweeps_after).append(1)
return real_record(self, logger, success=success)

with patch("cron.scheduler.tick", side_effect=lambda *a, **k: 0), \
patch.object(
_sp.InProcessCronScheduler,
"_record_or_log_heartbeat",
new=counting,
):
t = threading.Thread(target=prov.start, args=(stop,), kwargs={"interval": 60}, daemon=True)
t.start()

# Let the main loop tick a few times (interval=60 → loop pauses
# 60s on stop_event.wait, but the FIRST tick + tick_duration
# check returns quickly).
time.sleep(0.3)
assert len(sweeps_before) >= 1, "watchdog never started"

# Now stop the loop. The mock flips phase so we can detect any
# post-exit watchdog calls — there should be at most ONE trailing
# call (the watchdog's stop_event.wait was already in flight).
phase["stopped"] = True
stop.set()
t.join(timeout=2.0)
time.sleep(0.3) # let any rogue watchdog fire one more sweep

assert not t.is_alive(), "ticker thread did not exit after stop_event was set"
# Sweeps after stop should be zero or near-zero. We allow a small
# window because the watchdog may have been mid-way through its
# wait() call; what we are guarding against is runaway heartbeats.
assert len(sweeps_after) <= 1, (
f"watchdog emitted {len(sweeps_after)} heartbeats after main loop "
"exited — should stop cleanly with stop_event (#60703)"
)


def test_slow_tick_logs_warning_but_keeps_loop_alive():
"""A cron_tick() that takes longer than the configured interval is
logged as an anomaly but does not stop the loop. Combined with the
watchdog above, the heartbeat continues to refresh, so users see a
visible warning instead of a misleading STALLED status (#60703).
"""
import logging
from cron.scheduler_provider import InProcessCronScheduler

stop = threading.Event()
prov = InProcessCronScheduler()

calls = []

def _slow_tick(*a, **k):
calls.append(time.monotonic())
# Outlast a 0.1s interval by enough for the duration guard to
# trip; short enough to keep the test fast.
time.sleep(0.25)
return 0

caplog_handler = []

class CapturingHandler(logging.Handler):
def emit(self, record):
caplog_handler.append(record)

capturing = CapturingHandler()
logger = logging.getLogger("cron.scheduler_provider")
logger.addHandler(capturing)
try:
with patch("cron.scheduler.tick", side_effect=_slow_tick), \
patch("cron.jobs.record_ticker_heartbeat"):
t = threading.Thread(
target=prov.start,
args=(stop,),
kwargs={"interval": 0.1},
daemon=True,
)
t.start()
# Two iterations should be enough to prove the loop survives
# a slow tick.
assert _wait_until(lambda: len(calls) >= 2, timeout=5.0), (
"ticker did not keep looping after slow tick"
)
stop.set()
t.join(timeout=2.0)
finally:
logger.removeHandler(capturing)

assert not t.is_alive(), "ticker thread died after slow tick"
warnings = [r for r in caplog_handler if "Cron tick took" in r.getMessage()]
assert warnings, (
"tick-duration guard did not log a warning for a tick that "
"overshot the interval (#60703)"
)

Loading