Skip to content
Closed
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
13 changes: 13 additions & 0 deletions cli-config.yaml.example
Original file line number Diff line number Diff line change
Expand Up @@ -870,6 +870,19 @@ agent:
# window on /restart, and keep it well under systemd's TimeoutStopSec.
# restart_drain_timeout: 0

# Cron-only floor under the same drain (seconds). Default 30.
# restart_drain_timeout above is written for chat turns, which are cheap to
# interrupt: the user is told the gateway is restarting and the session
# resumes on their next message. A cron run has no such safety net — it is
# recorded in jobs.json as a permanent failure, nobody is waiting on it, and
# a recurring job simply skips to its next schedule. So in-flight cron work
# gets its own grace window instead of inheriting the 0 above.
# Clamped at runtime to the shutdown-watchdog leash (restart_drain_timeout
# + 60s) minus teardown headroom, so values past ~50s need a matching
# TimeoutStopSec bump to take effect. Set 0 to opt out and drain cron on
# restart_drain_timeout like before.
# cron_drain_timeout: 30

# Upper bound (seconds) a submitted prompt waits for the deferred agent
# build (MCP discovery, model metadata, skills scan) before failing with a
# visible error. The wait is patient — the message is delivered as soon as
Expand Down
76 changes: 76 additions & 0 deletions gateway/restart.py
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,24 @@
DEFAULT_CONFIG["agent"]["restart_after_turn_timeout"]
)

# Cron-only floor under the ``stop()`` drain. ``restart_drain_timeout``
# defaults to 0 because interrupting a *chat* turn is cheap and recoverable:
# the user is told the gateway is restarting and the session is pre-marked
# resume_pending. An interrupted *cron* run has neither property — nobody is
# waiting on it, it lands in jobs.json as a permanent failure, and a recurring
# job just waits for its next schedule — so a zero-second drain silently
# destroys work. See #82161.
DEFAULT_GATEWAY_CRON_DRAIN_TIMEOUT = float(
DEFAULT_CONFIG["agent"]["cron_drain_timeout"]
)

# Seconds of the shutdown watchdog leash held back for the work that still has
# to happen after the drain returns: interrupt agents, kill tool subprocesses,
# mark in-flight jobs interrupted, disconnect adapters. Waiting for cron past
# that point trades a job that is killed *and recorded* for one that is
# SIGKILLed mid-write and stays wedged at ``last_status=running`` forever.
CRON_DRAIN_CLEANUP_RESERVE_S = 10.0


def is_gateway_supervisor_process(
environ: Mapping[str, str] | None = None,
Expand Down Expand Up @@ -92,6 +110,64 @@ def parse_restart_after_turn_timeout(raw: object) -> float:
return max(0.0, value)


def parse_cron_drain_timeout(raw: object) -> float:
"""Parse the cron-only drain floor, falling back to the shared default.

``0`` is a deliberate opt-out — cron work is then interrupted on the same
budget as chat work, the pre-#82161 behaviour — and must not fall through
to the default, unlike empty/missing input.
"""
if raw is None:
return DEFAULT_GATEWAY_CRON_DRAIN_TIMEOUT
if isinstance(raw, str) and not raw.strip():
return DEFAULT_GATEWAY_CRON_DRAIN_TIMEOUT
try:
value = float(raw)
except (TypeError, ValueError):
return DEFAULT_GATEWAY_CRON_DRAIN_TIMEOUT
return max(0.0, value)


def resolve_cron_drain_budget(
drain_timeout: float,
cron_drain_timeout: float,
*,
watchdog_delay: float,
elapsed: float = 0.0,
cleanup_reserve_s: float = CRON_DRAIN_CLEANUP_RESERVE_S,
) -> float:
"""Seconds the shutdown drain may spend waiting on in-flight cron work.

The configured floor is clamped to what this process can actually honour.
The shutdown watchdog hard-exits at ``watchdog_delay`` and the service
manager's ``TimeoutStopSec`` is sized from the same drain timeout, so
waiting past that leash (minus ``cleanup_reserve_s`` for the teardown that
follows the drain) would swap a cleanly-interrupted job for a SIGKILL that
leaves it wedged mid-run — strictly worse than the bug being fixed.

Never returns less than ``drain_timeout``: the cron floor only ever
extends the wait, so an operator who deliberately configured a long
``restart_drain_timeout`` keeps it.
"""

def _seconds(value: object, fallback: float = 0.0) -> float:
try:
return max(float(value), 0.0) # type: ignore[arg-type]
except (TypeError, ValueError):
return fallback

drain = _seconds(drain_timeout)
floor = _seconds(cron_drain_timeout)
if floor <= 0.0:
return drain
ceiling = (
_seconds(watchdog_delay)
- _seconds(elapsed)
- _seconds(cleanup_reserve_s, CRON_DRAIN_CLEANUP_RESERVE_S)
)
return max(drain, min(floor, ceiling))


def resolve_restart_exit_wait_budget(
drain_timeout: float,
after_turn_timeout: float,
Expand Down
104 changes: 89 additions & 15 deletions gateway/run.py
Original file line number Diff line number Diff line change
Expand Up @@ -2215,6 +2215,8 @@ def _platform_has_bot_credential(platform: "Platform", platform_config: "Platfor
)
if "restart_drain_timeout" in _agent_cfg:
os.environ["HERMES_RESTART_DRAIN_TIMEOUT"] = str(_agent_cfg["restart_drain_timeout"])
if "cron_drain_timeout" in _agent_cfg:
os.environ["HERMES_CRON_DRAIN_TIMEOUT"] = str(_agent_cfg["cron_drain_timeout"])
if "gateway_auto_continue_freshness" in _agent_cfg:
os.environ["HERMES_AUTO_CONTINUE_FRESHNESS"] = str(
_agent_cfg["gateway_auto_continue_freshness"]
Expand Down Expand Up @@ -2433,12 +2435,15 @@ def _platform_has_bot_credential(platform: "Platform", platform_config: "Platfor
start_loop_liveness_watchdog,
)
from gateway.restart import (
DEFAULT_GATEWAY_CRON_DRAIN_TIMEOUT,
DEFAULT_GATEWAY_RESTART_AFTER_TURN_TIMEOUT,
DEFAULT_GATEWAY_RESTART_DRAIN_TIMEOUT,
GATEWAY_FATAL_CONFIG_EXIT_CODE,
GATEWAY_SERVICE_RESTART_EXIT_CODE,
parse_cron_drain_timeout,
parse_restart_after_turn_timeout,
parse_restart_drain_timeout,
resolve_cron_drain_budget,
)


Expand Down Expand Up @@ -5841,6 +5846,7 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew
_busy_text_mode: str = "interrupt"
_restart_drain_timeout: float = DEFAULT_GATEWAY_RESTART_DRAIN_TIMEOUT
_restart_after_turn_timeout: float = DEFAULT_GATEWAY_RESTART_AFTER_TURN_TIMEOUT
_cron_drain_timeout: float = DEFAULT_GATEWAY_CRON_DRAIN_TIMEOUT
_exit_code: Optional[int] = None
_draining: bool = False
_external_drain_active: bool = False
Expand Down Expand Up @@ -5978,6 +5984,7 @@ def __init__(self, config: Optional[GatewayConfig] = None):
self._busy_text_mode = self._load_busy_text_mode()
self._restart_drain_timeout = self._load_restart_drain_timeout()
self._restart_after_turn_timeout = self._load_restart_after_turn_timeout()
self._cron_drain_timeout = self._load_cron_drain_timeout()
self._provider_routing = self._load_provider_routing()
self._fallback_model = self._load_fallback_model()

Expand Down Expand Up @@ -8513,6 +8520,29 @@ def _load_restart_after_turn_timeout() -> float:
)
return value

@staticmethod
def _load_cron_drain_timeout() -> float:
"""Load the cron-only floor under the stop()/drain wait (#82161)."""
env_raw = os.getenv("HERMES_CRON_DRAIN_TIMEOUT")
if env_raw is not None and str(env_raw).strip() != "":
raw: object = env_raw
else:
cfg = _load_gateway_runtime_config()
raw = cfg_get(cfg, "agent", "cron_drain_timeout", default=None)
value = parse_cron_drain_timeout(raw)
# Warn only when the user supplied a non-empty value that failed to
# parse (parser falls back to the default). ``0`` is valid.
if raw is not None and str(raw).strip() != "":
try:
float(raw)
except (TypeError, ValueError):
logger.warning(
"Invalid cron_drain_timeout '%s', using default %.0fs",
raw,
DEFAULT_GATEWAY_CRON_DRAIN_TIMEOUT,
)
return value

@staticmethod
def _load_background_notifications_mode() -> str:
"""Load background process notification mode from config or env var.
Expand Down Expand Up @@ -9341,7 +9371,9 @@ async def _handle_active_session_busy_message(self, event: MessageEvent, session

return True

async def _drain_active_agents(self, timeout: float) -> tuple[Dict[str, Any], bool]:
async def _drain_active_agents(
self, timeout: float, cron_timeout: Optional[float] = None
) -> tuple[Dict[str, Any], bool]:
snapshot = self._snapshot_running_agents()
last_active_count = self._running_agent_count()
last_cron_count = self._active_cron_job_count()
Expand Down Expand Up @@ -9378,18 +9410,32 @@ def _maybe_update_status(force: bool = False) -> None:
return snapshot, False

_maybe_update_status(force=True)
if timeout <= 0:
return snapshot, True

deadline = asyncio.get_running_loop().time() + timeout
while (
(
len(self._running_agents)
or self._active_cron_job_count()
or self._active_api_run_count()
)
and asyncio.get_running_loop().time() < deadline
):
# Cron work drains on its own deadline. ``timeout``
# (``restart_drain_timeout``) defaults to 0 because interrupting a
# chat turn is announced and resumable; a cron run killed mid-flight
# is recorded in jobs.json as a permanent failure nobody is waiting
# on. Sharing one budget meant the default config could report
# ``timed_out=True`` after 0.00s with a cron job in flight and kill
# it — the drain never even entered this loop (#82161).
loop = asyncio.get_running_loop()
started = loop.time()
deadline = started + timeout
cron_deadline = started + (timeout if cron_timeout is None else cron_timeout)

def _still_draining() -> bool:
now = loop.time()
if (
len(self._running_agents) or self._active_api_run_count()
) and now < deadline:
return True
return bool(self._active_cron_job_count()) and now < cron_deadline

# Both budgets at 0 leave this loop unentered, which is the legacy
# "interrupt immediately" behaviour — expressed as an expired
# deadline rather than a special case, so the timed_out value below
# is always computed from real state instead of asserted up front.
while _still_draining():
_maybe_update_status()
await asyncio.sleep(0.1)
timed_out = (
Expand Down Expand Up @@ -13058,15 +13104,43 @@ def _phase_elapsed() -> float:

_cron_at_start = self._active_cron_job_count()
_api_at_start = self._active_api_run_count()
# In-flight cron work gets its own floor, clamped to the watchdog
# leash we're already running under so the extra wait can never
# cost us the post-drain cleanup window (#82161).
# getattr-guard: shutdown-path tests drive _stop_impl_body from
# bare doubles that aren't GatewayRunner instances, so they don't
# pick up the class-level default.
_cron_drain_cfg = getattr(
self, "_cron_drain_timeout", DEFAULT_GATEWAY_CRON_DRAIN_TIMEOUT
)
_cron_timeout = resolve_cron_drain_budget(
timeout,
_cron_drain_cfg,
watchdog_delay=resolve_shutdown_watchdog_delay(timeout),
elapsed=_phase_elapsed(),
)
if _cron_at_start and _cron_timeout > timeout:
logger.info(
"Shutdown drain: %d in-flight cron job(s) — waiting up to "
"%.0fs for them (cron_drain_timeout=%.0fs, "
"restart_drain_timeout=%.0fs)",
_cron_at_start,
_cron_timeout,
_cron_drain_cfg,
timeout,
)
_drain_started_at = time.monotonic()
active_agents, timed_out = await self._drain_active_agents(timeout)
active_agents, timed_out = await self._drain_active_agents(
timeout, _cron_timeout
)
_drain_elapsed = time.monotonic() - _drain_started_at
logger.info(
"Shutdown phase: drain done at +%.2fs (drain took %.2fs, "
"timed_out=%s, active_at_start=%d, active_now=%d, "
"cron_at_start=%d, cron_now=%d, "
"api_at_start=%d, api_now=%d)",
_phase_elapsed(),
time.monotonic() - _drain_started_at,
_drain_elapsed,
timed_out,
len(active_agents),
self._running_agent_count(),
Expand Down Expand Up @@ -13095,7 +13169,7 @@ def _phase_elapsed() -> float:
"Gateway drain timed out after %.1fs with %d active agent(s), "
"%d in-flight cron job(s), and %d api_server run(s); "
"interrupting remaining work.",
timeout,
_drain_elapsed,
self._running_agent_count(),
self._active_cron_job_count(),
self._active_api_run_count(),
Expand Down
9 changes: 9 additions & 0 deletions hermes_cli/config_defaults.py
Original file line number Diff line number Diff line change
Expand Up @@ -73,6 +73,15 @@
# (/restart, SIGUSR1), prefer restart_after_turn_timeout below so
# active turns finish *before* stop() begins (#77184).
"restart_drain_timeout": 0,
# Cron-only floor under the stop()/drain wait (seconds). A chat turn
# interrupted by a restart is announced to the user and resumed on
# their next message; an interrupted cron run is written to jobs.json
# as a permanent failure that nobody is waiting on, so it must not
# inherit restart_drain_timeout's 0 (#82161). Clamped at runtime to
# the shutdown-watchdog leash minus teardown headroom, so raising it
# past ~50s has no effect unless TimeoutStopSec is raised too.
# 0 = opt out (cron drains on restart_drain_timeout, legacy).
"cron_drain_timeout": 30,
# In-band restart wait for active turns to finish before stop()
# (seconds). /restart and SIGUSR1 refuse new work, then wait up to
# this cap for in-flight agents/cron/api runs to complete naturally
Expand Down
2 changes: 2 additions & 0 deletions tests/gateway/restart_test_helpers.py
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@
from gateway.config import GatewayConfig, Platform, PlatformConfig
from gateway.platforms.base import BasePlatformAdapter, SendResult
from gateway.restart import (
DEFAULT_GATEWAY_CRON_DRAIN_TIMEOUT,
DEFAULT_GATEWAY_RESTART_AFTER_TURN_TIMEOUT,
DEFAULT_GATEWAY_RESTART_DRAIN_TIMEOUT,
)
Expand Down Expand Up @@ -77,6 +78,7 @@ def make_restart_runner(
runner._restart_command_source = None
runner._restart_drain_timeout = DEFAULT_GATEWAY_RESTART_DRAIN_TIMEOUT
runner._restart_after_turn_timeout = DEFAULT_GATEWAY_RESTART_AFTER_TURN_TIMEOUT
runner._cron_drain_timeout = DEFAULT_GATEWAY_CRON_DRAIN_TIMEOUT
runner._stop_task = None
runner._busy_input_mode = "interrupt"
runner._update_prompt_pending = {}
Expand Down
1 change: 1 addition & 0 deletions tests/gateway/test_cron_active_work_drain.py
Original file line number Diff line number Diff line change
Expand Up @@ -85,6 +85,7 @@ async def test_in_flight_cron_job_marked_interrupted_on_forced_kill(self, monkey

runner, adapter = make_restart_runner()
runner._restart_drain_timeout = 0.01 # force the timeout path
runner._cron_drain_timeout = 0.01 # ...past the cron floor too (#82161)
adapter.disconnect = _make_async_noop()

sched._running_job_ids.add("job-1")
Expand Down
Loading
Loading