Merge pull request #85147 from kshitijk4poor/feat/unified-deadline-layer
feat(agent): unified deadline layer — bounded execution primitive + timeout resolver (#85125 Phase 1)
This commit is contained in:
@@ -0,0 +1,544 @@
|
||||
"""Unified deadline layer — one bounded-execution primitive, one timeout resolver.
|
||||
|
||||
Phase 1 of the architectural fix for the timeout/hang backlog
|
||||
(https://github.com/NousResearch/hermes-agent/issues/85125).
|
||||
|
||||
The tree currently carries at least six site-local deadline mechanisms, each
|
||||
built for one incident, none shared (tool_executor batch deadline, telegram
|
||||
``_await_with_thread_deadline``, gateway turn lease, reasoning stale floors,
|
||||
``human_wait_ceiling``, per-MCP-handler timeouts). Every new stall report
|
||||
grows that list by one. This module is the shared foundation the call sites
|
||||
migrate onto in later phases:
|
||||
|
||||
* :func:`resolve_timeout` — one config-first resolution path for timeout
|
||||
values (``timeouts:`` section in config.yaml > legacy env var > default),
|
||||
so new surfaces stop inventing ``HERMES_*_TIMEOUT`` env vars (".env is for
|
||||
secrets only") and hardcoded literals stop ignoring user config
|
||||
(#63302, #53161, #43272 class).
|
||||
|
||||
* :func:`clamp_timeout` — platform-safe clamping. Large user-supplied
|
||||
timeouts overflow ``time_t`` inside ``threading.Lock.acquire(timeout=...)``
|
||||
/ ``Thread.join(timeout=...)`` on macOS and kill whole tool batches
|
||||
(#83220). Clamping at the shared boundary fixes that class once, for
|
||||
every consumer.
|
||||
|
||||
* :func:`run_bounded_async` — a wall-clock deadline for awaitables that does
|
||||
NOT depend on event-loop timers. ``asyncio.wait_for`` schedules its expiry
|
||||
on the loop; when the loop thread itself is blocked in a synchronous call
|
||||
(family A of the #84047 stall triage), every asyncio-based timeout in the
|
||||
process is silently disabled. This helper drives the deadline from a
|
||||
daemon ``threading.Timer`` (generalizing the proven telegram-adapter
|
||||
primitive) and abandons cancellation-shielded tasks instead of waiting for
|
||||
cancellation to complete. The telegram adapter's private copy
|
||||
(``plugins/platforms/telegram/adapter.py:_await_with_thread_deadline``)
|
||||
migrates onto this in Phase 2 of #85125 — do not let the two drift in the
|
||||
meantime; fix bugs here first.
|
||||
|
||||
* :func:`run_bounded_sync` — the same contract for synchronous callables
|
||||
bounded from a synchronous context (daemon worker thread, abandoned on
|
||||
expiry).
|
||||
|
||||
* :func:`kill_process_tree` — portable whole-tree termination so
|
||||
kill-on-timeout stops orphaning descendants (#71148, #59549, #84967,
|
||||
#68139 class). Existing site-local tree-kills that migrate onto this in
|
||||
Phase 4 of #85125: ``gateway/status.py`` (taskkill wrapper + psutil
|
||||
snapshot/reap pair) and ``tools/code_execution_tool.py`` (psutil
|
||||
recursive children kill).
|
||||
|
||||
Design invariants:
|
||||
|
||||
* Exceptions raised by the bounded operation propagate unchanged — callers
|
||||
keep their existing error handling. Only the *timeout* outcome is
|
||||
reified (as :class:`BoundedResult`), because that is the outcome the
|
||||
call sites keep getting wrong.
|
||||
* A timeout produced by this layer is OUR deadline, not the provider's.
|
||||
Callers that feed errors into ``agent/error_classifier.py`` should
|
||||
classify :class:`DeadlineExpired` distinctly from transport timeouts
|
||||
(the #59549 / #80323 misattribution class).
|
||||
* ``None`` timeout means unbounded, and non-positive resolved values are
|
||||
normalized to ``None`` (matching the existing
|
||||
``HERMES_CONCURRENT_TOOL_TIMEOUT_S`` convention).
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
import faulthandler
|
||||
import logging
|
||||
import os
|
||||
import subprocess
|
||||
import sys
|
||||
import threading
|
||||
import time
|
||||
from dataclasses import dataclass
|
||||
from typing import Any, Awaitable, Callable, Optional
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
__all__ = [
|
||||
"MAX_SAFE_TIMEOUT_S",
|
||||
"BoundedResult",
|
||||
"DeadlineExpired",
|
||||
"clamp_timeout",
|
||||
"resolve_timeout",
|
||||
"run_bounded_async",
|
||||
"run_bounded_sync",
|
||||
"kill_process_tree",
|
||||
]
|
||||
|
||||
# Upper bound for any timeout handed to platform wait primitives.
|
||||
#
|
||||
# CPython converts ``threading.Lock.acquire(timeout=...)`` /
|
||||
# ``Thread.join(timeout=...)`` deadlines to an absolute timestamp; very large
|
||||
# relative timeouts overflow ``time_t`` on macOS and raise
|
||||
# ``OverflowError: timestamp out of range for platform time_t`` (#83220).
|
||||
# One year is semantically "unbounded" for every wait in this codebase while
|
||||
# staying far below any platform conversion limit.
|
||||
MAX_SAFE_TIMEOUT_S = 31_536_000.0 # 365 days
|
||||
|
||||
# Grace period after a deadline fires before concluding the event loop thread
|
||||
# is blocked in a synchronous call and dumping stacks (family A diagnostics).
|
||||
_LOOP_BLOCKED_DUMP_GRACE_S = 5.0
|
||||
|
||||
|
||||
class DeadlineExpired(TimeoutError):
|
||||
"""A deadline enforced by this layer expired.
|
||||
|
||||
Distinct from transport/provider timeout types on purpose: when this is
|
||||
raised (or a :class:`BoundedResult` reports ``timed_out``), the timeout
|
||||
was Hermes's own bound — error classification must not attribute it to
|
||||
the provider (#59549 / #80323 misattribution class).
|
||||
"""
|
||||
|
||||
def __init__(self, label: str, timeout_s: float):
|
||||
super().__init__(f"deadline expired after {timeout_s:.1f}s: {label}")
|
||||
self.label = label
|
||||
self.timeout_s = timeout_s
|
||||
|
||||
|
||||
@dataclass(frozen=True, kw_only=True)
|
||||
class BoundedResult:
|
||||
"""Outcome of a bounded operation.
|
||||
|
||||
``timed_out`` is the reified outcome; on completion ``value`` holds the
|
||||
operation's return value. Operation exceptions are never captured here —
|
||||
they propagate to the caller unchanged.
|
||||
"""
|
||||
|
||||
timed_out: bool
|
||||
value: Any
|
||||
elapsed_s: float
|
||||
timeout_s: Optional[float]
|
||||
label: str
|
||||
|
||||
def raise_if_timed_out(self) -> Any:
|
||||
"""Return ``value``, raising :class:`DeadlineExpired` on timeout."""
|
||||
if self.timed_out:
|
||||
raise DeadlineExpired(self.label, float(self.timeout_s or 0.0))
|
||||
return self.value
|
||||
|
||||
|
||||
def clamp_timeout(timeout: Optional[float]) -> Optional[float]:
|
||||
"""Normalize a timeout value for platform wait primitives.
|
||||
|
||||
* ``None`` stays ``None`` (unbounded).
|
||||
* Non-positive values become ``None`` (unbounded) — matching the existing
|
||||
``HERMES_CONCURRENT_TOOL_TIMEOUT_S`` "0 disables the bound" convention.
|
||||
* Values above :data:`MAX_SAFE_TIMEOUT_S` are capped so they can never
|
||||
overflow ``time_t`` inside ``Lock.acquire`` / ``Thread.join`` on macOS
|
||||
(#83220).
|
||||
* Non-numeric values are treated as unset (``None``) with a warning
|
||||
rather than crashing the call path they were meant to protect.
|
||||
"""
|
||||
if timeout is None:
|
||||
return None
|
||||
try:
|
||||
value = float(timeout)
|
||||
except (TypeError, ValueError):
|
||||
logger.warning("clamp_timeout: non-numeric timeout %r; treating as unbounded", timeout)
|
||||
return None
|
||||
if value != value: # NaN
|
||||
logger.warning("clamp_timeout: NaN timeout; treating as unbounded")
|
||||
return None
|
||||
if value <= 0:
|
||||
return None
|
||||
return min(value, MAX_SAFE_TIMEOUT_S)
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Timeout resolution: config.yaml ``timeouts:`` section > legacy env var >
|
||||
# registered default.
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
def _timeouts_section() -> dict:
|
||||
"""Read the ``timeouts:`` root section from config.yaml (read-only).
|
||||
|
||||
Isolated for testability and so a broken config read can never take down
|
||||
the call path the timeout was protecting.
|
||||
"""
|
||||
try:
|
||||
from hermes_cli.config import load_config_readonly
|
||||
|
||||
section = load_config_readonly().get("timeouts")
|
||||
return section if isinstance(section, dict) else {}
|
||||
except Exception:
|
||||
logger.debug("timeouts: config read failed; using defaults", exc_info=True)
|
||||
return {}
|
||||
|
||||
|
||||
def _lookup_dotted(section: dict, key: str) -> Any:
|
||||
"""Walk ``a.b.c`` through nested dicts; return None when absent."""
|
||||
node: Any = section
|
||||
for part in key.split("."):
|
||||
if not isinstance(node, dict) or part not in node:
|
||||
return None
|
||||
node = node[part]
|
||||
return node
|
||||
|
||||
|
||||
def resolve_timeout(
|
||||
key: str,
|
||||
*,
|
||||
default: Optional[float],
|
||||
env_var: Optional[str] = None,
|
||||
) -> Optional[float]:
|
||||
"""Resolve a timeout in seconds for a dotted config key.
|
||||
|
||||
Precedence (established by the ``providers.*.request_timeout_seconds``
|
||||
pattern — config wins over the legacy env var):
|
||||
|
||||
1. ``timeouts.<key>`` in config.yaml (dotted key walks nested maps, e.g.
|
||||
``tools.concurrent_batch`` reads ``timeouts: {tools: {concurrent_batch: ...}}``)
|
||||
2. ``env_var`` when set and non-empty (legacy bridge — internal mechanism
|
||||
and back-compat only; new surfaces must not grow new user-facing
|
||||
``HERMES_*`` timeout env vars)
|
||||
3. ``default``
|
||||
|
||||
The winning value is passed through :func:`clamp_timeout`, so ``0`` or a
|
||||
negative value means "unbounded" and oversized values are made
|
||||
platform-safe. Invalid (non-numeric) config/env values fall through to
|
||||
the next source with a warning instead of breaking the protected path.
|
||||
"""
|
||||
raw = _lookup_dotted(_timeouts_section(), key)
|
||||
if raw is not None:
|
||||
# Explicit float() (clamp_timeout would also convert) so that invalid
|
||||
# config values FALL THROUGH to the env var / default instead of
|
||||
# resolving as unbounded — do not "simplify" this away. bool is
|
||||
# rejected because YAML `true` would silently become a 1-second
|
||||
# deadline; NaN is rejected for the same fall-through reason.
|
||||
if not isinstance(raw, bool):
|
||||
try:
|
||||
value = float(raw)
|
||||
if value == value: # not NaN
|
||||
return clamp_timeout(value)
|
||||
except (TypeError, ValueError):
|
||||
pass
|
||||
logger.warning("timeouts.%s: invalid value %r in config.yaml; ignoring", key, raw)
|
||||
|
||||
if env_var:
|
||||
env_raw = os.getenv(env_var, "").strip()
|
||||
if env_raw:
|
||||
try:
|
||||
return clamp_timeout(float(env_raw))
|
||||
except ValueError:
|
||||
logger.warning("invalid %s=%r; ignoring", env_var, env_raw)
|
||||
|
||||
return clamp_timeout(default)
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Bounded execution — async flavor.
|
||||
#
|
||||
# Generalizes plugins/platforms/telegram/adapter.py:_await_with_thread_deadline
|
||||
# (the #63309 fix): the deadline is driven by a daemon threading.Timer so a
|
||||
# blocked event loop cannot disable it, and a second timer dumps all thread
|
||||
# stacks when the loop provably failed to process the expiry — the one piece
|
||||
# of information loop-blocked hangs otherwise never surface.
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
def _consume_abandoned(task: "asyncio.Future[Any]") -> None:
|
||||
"""Observe an abandoned task's outcome so it never logs 'never retrieved'."""
|
||||
try:
|
||||
if not task.cancelled():
|
||||
task.exception()
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
|
||||
async def _run_abandon_cleanup(on_abandon: Callable[[], Awaitable[Any]]) -> None:
|
||||
"""Run abandonment cleanup fully fire-and-forget (its failures swallowed)."""
|
||||
try:
|
||||
await on_abandon()
|
||||
except Exception:
|
||||
logger.debug("deadline abandon-cleanup failed", exc_info=True)
|
||||
|
||||
|
||||
def _dump_blocked_loop_diagnostics(label: str, timeout_s: float) -> None:
|
||||
logger.warning(
|
||||
"[deadline] %r deadline (%.0fs) expired but the event loop has not "
|
||||
"processed the expiry after a further %.0fs — the loop thread appears "
|
||||
"BLOCKED in a synchronous call, which is why no asyncio timeout can "
|
||||
"fire. Dumping all thread stacks to stderr to identify the blocking "
|
||||
"frame.",
|
||||
label,
|
||||
timeout_s,
|
||||
_LOOP_BLOCKED_DUMP_GRACE_S,
|
||||
)
|
||||
try:
|
||||
faulthandler.dump_traceback(all_threads=True)
|
||||
except Exception:
|
||||
logger.debug("faulthandler traceback dump failed", exc_info=True)
|
||||
|
||||
|
||||
async def run_bounded_async(
|
||||
awaitable: Awaitable[Any],
|
||||
timeout: Optional[float],
|
||||
*,
|
||||
label: str = "operation",
|
||||
on_abandon: Optional[Callable[[], Awaitable[Any]]] = None,
|
||||
dump_on_blocked_loop: bool = True,
|
||||
) -> BoundedResult:
|
||||
"""Await ``awaitable`` under a wall-clock deadline independent of loop timers.
|
||||
|
||||
On completion returns ``BoundedResult(timed_out=False, value=...)``;
|
||||
exceptions from the operation (including ``asyncio.CancelledError`` from a
|
||||
caller cancelling *us*) propagate unchanged.
|
||||
|
||||
On timeout the underlying task is cancelled and **abandoned** — we do not
|
||||
await cancellation completion, because cancellation-shielded scopes (anyio,
|
||||
httpcore init, MCP SDK teardown) are exactly the paths that wedge forever.
|
||||
``on_abandon`` (zero-arg callable returning an awaitable) is scheduled as
|
||||
detached best-effort cleanup for the half-built state the abandoned task
|
||||
may leave behind. Returns ``BoundedResult(timed_out=True, value=None)``.
|
||||
|
||||
``timeout=None`` (or a non-positive resolved value) awaits unbounded.
|
||||
"""
|
||||
timeout_s = clamp_timeout(timeout)
|
||||
start = time.monotonic()
|
||||
if timeout_s is None:
|
||||
value = await awaitable
|
||||
return BoundedResult(timed_out=False, value=value, elapsed_s=time.monotonic() - start, timeout_s=None, label=label)
|
||||
|
||||
task = asyncio.ensure_future(awaitable)
|
||||
loop = asyncio.get_running_loop()
|
||||
deadline: "asyncio.Future[None]" = loop.create_future()
|
||||
loop_processed_expiry = threading.Event()
|
||||
|
||||
def _mark_expired() -> None:
|
||||
loop_processed_expiry.set()
|
||||
if not deadline.done():
|
||||
deadline.set_result(None)
|
||||
|
||||
def _expire_from_thread() -> None:
|
||||
loop.call_soon_threadsafe(_mark_expired)
|
||||
|
||||
def _watchdog_check() -> None:
|
||||
if not loop_processed_expiry.is_set():
|
||||
_dump_blocked_loop_diagnostics(label, timeout_s)
|
||||
|
||||
timer = threading.Timer(timeout_s, _expire_from_thread)
|
||||
timer.daemon = True
|
||||
timer.start()
|
||||
watchdog: Optional[threading.Timer] = None
|
||||
if dump_on_blocked_loop:
|
||||
watchdog = threading.Timer(
|
||||
timeout_s + _LOOP_BLOCKED_DUMP_GRACE_S, _watchdog_check
|
||||
)
|
||||
watchdog.daemon = True
|
||||
watchdog.start()
|
||||
try:
|
||||
try:
|
||||
done, _ = await asyncio.wait(
|
||||
{task, deadline}, return_when=asyncio.FIRST_COMPLETED
|
||||
)
|
||||
except asyncio.CancelledError:
|
||||
# The CALLER cancelled us. Without this, `task` would keep running
|
||||
# unobserved (and later log "exception was never retrieved") —
|
||||
# a leak the telegram original also had. Cancel + abandon it, then
|
||||
# let the cancellation propagate.
|
||||
task.cancel()
|
||||
task.add_done_callback(_consume_abandoned)
|
||||
raise
|
||||
if task in done:
|
||||
if not deadline.done():
|
||||
deadline.cancel()
|
||||
value = await task
|
||||
return BoundedResult(timed_out=False, value=value, elapsed_s=time.monotonic() - start, timeout_s=timeout_s, label=label)
|
||||
|
||||
task.cancel()
|
||||
task.add_done_callback(_consume_abandoned)
|
||||
if on_abandon is not None:
|
||||
cleanup = asyncio.ensure_future(_run_abandon_cleanup(on_abandon))
|
||||
cleanup.add_done_callback(_consume_abandoned)
|
||||
logger.warning("[deadline] %r timed out after %.1fs; task abandoned", label, timeout_s)
|
||||
return BoundedResult(timed_out=True, value=None, elapsed_s=time.monotonic() - start, timeout_s=timeout_s, label=label)
|
||||
finally:
|
||||
timer.cancel()
|
||||
if watchdog is not None:
|
||||
watchdog.cancel()
|
||||
# cancel() cannot stop a Timer whose callback is already running;
|
||||
# setting the event closes that race so a completed await can never
|
||||
# be misreported as a blocked loop.
|
||||
loop_processed_expiry.set()
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Bounded execution — sync flavor.
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
def run_bounded_sync(
|
||||
fn: Callable[[], Any],
|
||||
timeout: Optional[float],
|
||||
*,
|
||||
label: str = "operation",
|
||||
on_timeout: Optional[Callable[[], None]] = None,
|
||||
) -> BoundedResult:
|
||||
"""Run ``fn`` in a daemon worker thread under a wall-clock deadline.
|
||||
|
||||
On completion returns its value (exceptions re-raised in the caller).
|
||||
On expiry the worker thread is **abandoned** (daemon, so it cannot block
|
||||
interpreter exit), ``on_timeout`` (if given) runs best-effort in the
|
||||
caller's thread — e.g. to mark a backend suspect or kill a subprocess —
|
||||
and ``BoundedResult(timed_out=True)`` is returned.
|
||||
|
||||
Intended for infrequent, seconds-scale blocking backend calls. Do NOT
|
||||
use per-item in hot loops: each call spawns a thread, and every timeout
|
||||
permanently leaks an abandoned daemon thread — a wedged backend called
|
||||
in a retry loop would accumulate them.
|
||||
|
||||
``timeout=None`` (or non-positive) blocks until ``fn`` returns.
|
||||
"""
|
||||
timeout_s = clamp_timeout(timeout)
|
||||
start = time.monotonic()
|
||||
if timeout_s is None:
|
||||
return BoundedResult(timed_out=False, value=fn(), elapsed_s=time.monotonic() - start, timeout_s=None, label=label)
|
||||
|
||||
box: dict[str, Any] = {}
|
||||
done = threading.Event()
|
||||
|
||||
def _worker() -> None:
|
||||
try:
|
||||
box["value"] = fn()
|
||||
except BaseException as exc: # re-raised in caller; must not vanish
|
||||
box["exc"] = exc
|
||||
finally:
|
||||
done.set()
|
||||
|
||||
thread = threading.Thread(
|
||||
target=_worker, name=f"deadline-{label}", daemon=True
|
||||
)
|
||||
thread.start()
|
||||
if not done.wait(timeout_s):
|
||||
logger.warning("[deadline] %r timed out after %.1fs; worker abandoned", label, timeout_s)
|
||||
if on_timeout is not None:
|
||||
try:
|
||||
on_timeout()
|
||||
except Exception:
|
||||
logger.debug("deadline on_timeout callback failed", exc_info=True)
|
||||
return BoundedResult(timed_out=True, value=None, elapsed_s=time.monotonic() - start, timeout_s=timeout_s, label=label)
|
||||
|
||||
if "exc" in box:
|
||||
raise box["exc"]
|
||||
return BoundedResult(timed_out=False, value=box.get("value"), elapsed_s=time.monotonic() - start, timeout_s=timeout_s, label=label)
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Whole-tree process termination.
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
def kill_process_tree(pid: int, *, sig: Optional[int] = None) -> bool:
|
||||
"""Terminate ``pid`` and all its descendants, portably.
|
||||
|
||||
Kill-on-timeout that signals only the direct child orphans process trees
|
||||
(cron scripts, in-container shells, browser daemons — #71148 class).
|
||||
|
||||
* Windows: ``taskkill /F /T`` terminates the tree (``sig`` ignored;
|
||||
Windows has no equivalent). Console-window flash is suppressed via
|
||||
``windows_hide_flags`` and the exit code is checked, so a dead or
|
||||
inaccessible PID reports ``False`` like the POSIX path.
|
||||
* POSIX: the descendant set is snapshotted via psutil (a hard
|
||||
dependency) BEFORE any signal — once the parent dies its children are
|
||||
reparented and can no longer be found by a parent walk. Then the
|
||||
process group is signalled when ``pid`` leads one (covers
|
||||
grandchildren in the same session in one syscall), and every
|
||||
snapshotted descendant is signalled individually — which also reaches
|
||||
descendants that created their OWN sessions (a child that called
|
||||
``setsid``, exactly what user shell commands do; see
|
||||
tools/environments/base.py). ``sig`` defaults to ``SIGKILL``.
|
||||
psutil's identity-aware ``Process`` (PID + create time) means a
|
||||
recycled PID is never signalled.
|
||||
|
||||
Returns True when the target (or any of its tree) was signalled, False
|
||||
when the process was already gone or every termination call failed.
|
||||
"""
|
||||
if sys.platform == "win32":
|
||||
try:
|
||||
from hermes_cli._subprocess_compat import windows_hide_flags
|
||||
|
||||
creationflags = windows_hide_flags()
|
||||
except Exception:
|
||||
creationflags = 0
|
||||
try:
|
||||
proc = subprocess.run(
|
||||
["taskkill", "/F", "/T", "/PID", str(pid)],
|
||||
capture_output=True,
|
||||
timeout=15,
|
||||
check=False,
|
||||
creationflags=creationflags,
|
||||
)
|
||||
# taskkill exits non-zero for not-found / access-denied; keep the
|
||||
# cross-platform contract (False = nothing was terminated).
|
||||
return proc.returncode == 0
|
||||
except Exception:
|
||||
logger.debug("kill_process_tree: taskkill failed for pid %s", pid, exc_info=True)
|
||||
return False
|
||||
|
||||
import signal as _signal
|
||||
|
||||
if sig is None:
|
||||
sig = _signal.SIGKILL
|
||||
|
||||
# Snapshot descendants while the parent is still alive — after it dies
|
||||
# they reparent to init/subreaper and a parent walk finds nothing.
|
||||
descendants: list = []
|
||||
try:
|
||||
import psutil
|
||||
|
||||
descendants = psutil.Process(int(pid)).children(recursive=True)
|
||||
except Exception:
|
||||
# Already gone, or psutil unavailable in a stripped env — the
|
||||
# group-signal below still covers same-session descendants.
|
||||
descendants = []
|
||||
|
||||
signalled = False
|
||||
try:
|
||||
# NOTE: getpgid→killpg has an inherent TOCTOU (pid could be reaped and
|
||||
# recycled between the calls). All existing killpg sites share it; the
|
||||
# psutil sweep below is identity-aware and does not.
|
||||
pgid = os.getpgid(pid)
|
||||
except (ProcessLookupError, PermissionError, OSError):
|
||||
pgid = None
|
||||
try:
|
||||
if pgid is not None and pgid == pid:
|
||||
# pid leads its own group: one syscall covers the whole group.
|
||||
# (The == check guards against signalling the caller's own group
|
||||
# when pid is not a leader.)
|
||||
os.killpg(pgid, sig) # windows-footgun: ok — POSIX-only branch (win32 returns above)
|
||||
else:
|
||||
os.kill(pid, sig)
|
||||
signalled = True
|
||||
except ProcessLookupError:
|
||||
pass
|
||||
except (PermissionError, OSError):
|
||||
logger.debug("kill_process_tree: signal failed for pid %s", pid, exc_info=True)
|
||||
|
||||
# Sweep the snapshot: reaches descendants outside the parent's group
|
||||
# (their own setsid sessions) and the non-group-leader case.
|
||||
for child in descendants:
|
||||
try:
|
||||
if child.is_running(): # identity-aware: recycled PIDs skipped
|
||||
child.send_signal(sig)
|
||||
signalled = True
|
||||
except Exception:
|
||||
continue
|
||||
return signalled
|
||||
+13
-15
@@ -158,21 +158,19 @@ def _parse_tool_arguments(raw_arguments: Any) -> tuple[dict, Optional[str]]:
|
||||
|
||||
|
||||
def _resolve_concurrent_tool_timeout() -> float | None:
|
||||
raw = os.getenv("HERMES_CONCURRENT_TOOL_TIMEOUT_S", "").strip()
|
||||
if not raw:
|
||||
return _DEFAULT_CONCURRENT_TOOL_TIMEOUT_S
|
||||
try:
|
||||
value = float(raw)
|
||||
except ValueError:
|
||||
logger.warning(
|
||||
"invalid HERMES_CONCURRENT_TOOL_TIMEOUT_S=%r; using %.0fs",
|
||||
raw,
|
||||
_DEFAULT_CONCURRENT_TOOL_TIMEOUT_S,
|
||||
)
|
||||
return _DEFAULT_CONCURRENT_TOOL_TIMEOUT_S
|
||||
if value <= 0:
|
||||
return None
|
||||
return value
|
||||
"""Resolve the per-batch concurrent tool deadline.
|
||||
|
||||
Delegates to the unified resolver (#85125): ``timeouts.tools.concurrent_batch``
|
||||
in config.yaml wins, the legacy ``HERMES_CONCURRENT_TOOL_TIMEOUT_S`` env var
|
||||
remains the back-compat bridge, and ``0``/negative still disables the bound.
|
||||
"""
|
||||
from agent.deadline import resolve_timeout
|
||||
|
||||
return resolve_timeout(
|
||||
"tools.concurrent_batch",
|
||||
default=_DEFAULT_CONCURRENT_TOOL_TIMEOUT_S,
|
||||
env_var="HERMES_CONCURRENT_TOOL_TIMEOUT_S",
|
||||
)
|
||||
|
||||
|
||||
def _flush_session_db_after_tool_progress(
|
||||
|
||||
@@ -166,6 +166,22 @@ model:
|
||||
# gpt-5.4:
|
||||
# stale_timeout_seconds: 1800 # Longer non-stream stale timeout for slow large-context turns
|
||||
|
||||
# =============================================================================
|
||||
# Unified Timeouts (operation deadlines)
|
||||
# =============================================================================
|
||||
# One place to override Hermes's internal operation deadlines (seconds).
|
||||
# Keys are dotted paths resolved by agent/deadline.py:resolve_timeout().
|
||||
# Precedence: this section > legacy HERMES_* env var (back-compat) > built-in
|
||||
# default. 0 or a negative value disables the bound (unbounded); very large
|
||||
# values are clamped to a platform-safe maximum automatically.
|
||||
#
|
||||
# Currently resolved keys (more paths migrate here over time — see issue #85125):
|
||||
#
|
||||
# timeouts:
|
||||
# tools:
|
||||
# concurrent_batch: 420 # Deadline for a parallel tool-call batch
|
||||
# # (legacy env: HERMES_CONCURRENT_TOOL_TIMEOUT_S)
|
||||
|
||||
# =============================================================================
|
||||
# OpenRouter Provider Routing (only applies when using OpenRouter)
|
||||
# =============================================================================
|
||||
|
||||
@@ -1930,6 +1930,7 @@ _EXTRA_KNOWN_ROOT_KEYS = {
|
||||
"require_mention", # top-level convenience form honored by the gateway (#3979)
|
||||
"unauthorized_dm_behavior", # top-level form read by gateway/config.py
|
||||
"signal", # Signal settings bridged to env vars by gateway/config.py
|
||||
"timeouts", # unified timeout resolution section (agent/deadline.py, #85125)
|
||||
}
|
||||
_KNOWN_ROOT_KEYS = frozenset(DEFAULT_CONFIG.keys()) | _EXTRA_KNOWN_ROOT_KEYS
|
||||
|
||||
|
||||
@@ -0,0 +1,478 @@
|
||||
"""Tests for agent/deadline.py — the unified deadline layer (#85125).
|
||||
|
||||
Covers:
|
||||
* clamp_timeout normalization (None / non-positive / oversized / NaN / junk)
|
||||
* resolve_timeout precedence: config.yaml ``timeouts:`` > legacy env var > default
|
||||
* run_bounded_sync: completion, exception propagation, timeout + on_timeout
|
||||
* run_bounded_async: completion, exception propagation, timeout + abandonment
|
||||
of cancellation-shielded tasks, on_abandon cleanup
|
||||
* kill_process_tree: descendants of a session-leader child die with it (POSIX)
|
||||
* backward-compat contract of tool_executor._resolve_concurrent_tool_timeout
|
||||
after its migration onto resolve_timeout
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
import os
|
||||
import signal
|
||||
import subprocess
|
||||
import sys
|
||||
import threading
|
||||
import time
|
||||
|
||||
import pytest
|
||||
|
||||
from agent.deadline import (
|
||||
MAX_SAFE_TIMEOUT_S,
|
||||
BoundedResult,
|
||||
DeadlineExpired,
|
||||
clamp_timeout,
|
||||
kill_process_tree,
|
||||
resolve_timeout,
|
||||
run_bounded_async,
|
||||
run_bounded_sync,
|
||||
)
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# clamp_timeout
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
class TestClampTimeout:
|
||||
def test_none_stays_none(self):
|
||||
assert clamp_timeout(None) is None
|
||||
|
||||
def test_zero_and_negative_mean_unbounded(self):
|
||||
assert clamp_timeout(0) is None
|
||||
assert clamp_timeout(-5) is None
|
||||
|
||||
def test_normal_value_passes_through(self):
|
||||
assert clamp_timeout(420.0) == 420.0
|
||||
|
||||
def test_oversized_value_clamped_to_platform_safe_max(self):
|
||||
# The #83220 class: >time_t deadlines crash Lock.acquire on macOS.
|
||||
assert clamp_timeout(10**18) == MAX_SAFE_TIMEOUT_S
|
||||
|
||||
def test_clamped_value_safe_for_threading_primitives(self):
|
||||
# Regression proof for #83220: the clamped value itself must be
|
||||
# accepted by the exact primitive that used to overflow. Acquiring an
|
||||
# uncontended lock returns immediately regardless of timeout, so
|
||||
# passing the full clamped value is safe and actually exercises the
|
||||
# time_t conversion.
|
||||
big = clamp_timeout(float(10**15))
|
||||
assert big is not None
|
||||
lock = threading.Lock()
|
||||
assert lock.acquire(timeout=big)
|
||||
lock.release()
|
||||
|
||||
def test_nan_and_junk_treated_as_unbounded(self):
|
||||
assert clamp_timeout(float("nan")) is None
|
||||
assert clamp_timeout("not-a-number") is None # type: ignore[arg-type]
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# resolve_timeout
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
class TestResolveTimeout:
|
||||
def test_default_wins_when_nothing_configured(self, monkeypatch):
|
||||
monkeypatch.setattr("agent.deadline._timeouts_section", lambda: {})
|
||||
monkeypatch.delenv("HERMES_TEST_DEADLINE_X", raising=False)
|
||||
assert resolve_timeout("a.b", default=42.0, env_var="HERMES_TEST_DEADLINE_X") == 42.0
|
||||
|
||||
def test_env_var_beats_default(self, monkeypatch):
|
||||
monkeypatch.setattr("agent.deadline._timeouts_section", lambda: {})
|
||||
monkeypatch.setenv("HERMES_TEST_DEADLINE_X", "17.5")
|
||||
assert resolve_timeout("a.b", default=42.0, env_var="HERMES_TEST_DEADLINE_X") == 17.5
|
||||
|
||||
def test_config_beats_env_var(self, monkeypatch):
|
||||
monkeypatch.setattr(
|
||||
"agent.deadline._timeouts_section", lambda: {"a": {"b": 99}}
|
||||
)
|
||||
monkeypatch.setenv("HERMES_TEST_DEADLINE_X", "17.5")
|
||||
assert resolve_timeout("a.b", default=42.0, env_var="HERMES_TEST_DEADLINE_X") == 99.0
|
||||
|
||||
def test_dotted_key_walks_nested_maps(self, monkeypatch):
|
||||
monkeypatch.setattr(
|
||||
"agent.deadline._timeouts_section",
|
||||
lambda: {"tools": {"concurrent_batch": 300}},
|
||||
)
|
||||
assert resolve_timeout("tools.concurrent_batch", default=420.0) == 300.0
|
||||
|
||||
def test_zero_config_value_means_unbounded(self, monkeypatch):
|
||||
monkeypatch.setattr("agent.deadline._timeouts_section", lambda: {"a": {"b": 0}})
|
||||
assert resolve_timeout("a.b", default=42.0) is None
|
||||
|
||||
def test_invalid_config_value_falls_through_to_env(self, monkeypatch):
|
||||
monkeypatch.setattr(
|
||||
"agent.deadline._timeouts_section", lambda: {"a": {"b": "soon"}}
|
||||
)
|
||||
monkeypatch.setenv("HERMES_TEST_DEADLINE_X", "17.5")
|
||||
assert resolve_timeout("a.b", default=42.0, env_var="HERMES_TEST_DEADLINE_X") == 17.5
|
||||
|
||||
def test_invalid_env_value_falls_through_to_default(self, monkeypatch):
|
||||
monkeypatch.setattr("agent.deadline._timeouts_section", lambda: {})
|
||||
monkeypatch.setenv("HERMES_TEST_DEADLINE_X", "banana")
|
||||
assert resolve_timeout("a.b", default=42.0, env_var="HERMES_TEST_DEADLINE_X") == 42.0
|
||||
|
||||
def test_bool_config_value_rejected(self, monkeypatch):
|
||||
# YAML `true` must not silently become a 1-second deadline.
|
||||
monkeypatch.setattr("agent.deadline._timeouts_section", lambda: {"a": {"b": True}})
|
||||
assert resolve_timeout("a.b", default=42.0) == 42.0
|
||||
|
||||
def test_nan_config_value_falls_through(self, monkeypatch):
|
||||
# NaN must fall through to the next source, not resolve as unbounded.
|
||||
monkeypatch.setattr(
|
||||
"agent.deadline._timeouts_section", lambda: {"a": {"b": float("nan")}}
|
||||
)
|
||||
assert resolve_timeout("a.b", default=42.0) == 42.0
|
||||
|
||||
def test_broken_config_read_never_breaks_the_protected_path(self, monkeypatch):
|
||||
# _timeouts_section swallows config-load failures internally; prove
|
||||
# the public contract by making the underlying loader raise.
|
||||
import agent.deadline as dl
|
||||
|
||||
def _boom():
|
||||
raise RuntimeError("config unreadable")
|
||||
|
||||
monkeypatch.setattr("hermes_cli.config.load_config_readonly", _boom)
|
||||
assert dl._timeouts_section() == {}
|
||||
assert resolve_timeout("a.b", default=5.0) == 5.0
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# run_bounded_sync
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
class TestRunBoundedSync:
|
||||
def test_completion_returns_value(self):
|
||||
result = run_bounded_sync(lambda: "ok", 5.0, label="t")
|
||||
assert result.timed_out is False
|
||||
assert result.value == "ok"
|
||||
assert result.raise_if_timed_out() == "ok"
|
||||
|
||||
def test_unbounded_when_timeout_none(self):
|
||||
result = run_bounded_sync(lambda: 7, None, label="t")
|
||||
assert result.timed_out is False and result.value == 7
|
||||
|
||||
def test_exception_propagates_unchanged(self):
|
||||
class Boom(RuntimeError):
|
||||
pass
|
||||
|
||||
with pytest.raises(Boom):
|
||||
run_bounded_sync(lambda: (_ for _ in ()).throw(Boom("x")), 5.0, label="t")
|
||||
|
||||
def test_timeout_abandons_worker_and_reports(self):
|
||||
release = threading.Event()
|
||||
|
||||
def _wedged():
|
||||
release.wait(30)
|
||||
return "late"
|
||||
|
||||
start = time.monotonic()
|
||||
result = run_bounded_sync(_wedged, 0.2, label="wedged")
|
||||
elapsed = time.monotonic() - start
|
||||
assert result.timed_out is True
|
||||
assert result.value is None
|
||||
assert elapsed < 5.0 # returned near the deadline, not after 30s
|
||||
with pytest.raises(DeadlineExpired) as exc_info:
|
||||
result.raise_if_timed_out()
|
||||
assert "wedged" in str(exc_info.value)
|
||||
release.set()
|
||||
|
||||
def test_on_timeout_callback_runs(self):
|
||||
release = threading.Event()
|
||||
fired = []
|
||||
result = run_bounded_sync(
|
||||
lambda: release.wait(30),
|
||||
0.1,
|
||||
label="t",
|
||||
on_timeout=lambda: fired.append(True),
|
||||
)
|
||||
assert result.timed_out and fired == [True]
|
||||
release.set()
|
||||
|
||||
def test_on_timeout_callback_failure_is_swallowed(self):
|
||||
release = threading.Event()
|
||||
result = run_bounded_sync(
|
||||
lambda: release.wait(30),
|
||||
0.1,
|
||||
label="t",
|
||||
on_timeout=lambda: (_ for _ in ()).throw(RuntimeError("cleanup boom")),
|
||||
)
|
||||
assert result.timed_out is True
|
||||
release.set()
|
||||
|
||||
def test_deadline_expired_is_a_timeout_error(self):
|
||||
# Error-classification contract: our deadline must be catchable as
|
||||
# TimeoutError but distinguishable by type from transport timeouts.
|
||||
assert issubclass(DeadlineExpired, TimeoutError)
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# run_bounded_async
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
class TestRunBoundedAsync:
|
||||
def test_completion_returns_value(self):
|
||||
async def scenario():
|
||||
async def op():
|
||||
return "ok"
|
||||
|
||||
return await run_bounded_async(op(), 5.0, label="t")
|
||||
|
||||
result = asyncio.run(scenario())
|
||||
assert result.timed_out is False and result.value == "ok"
|
||||
|
||||
def test_unbounded_when_timeout_none(self):
|
||||
async def scenario():
|
||||
async def op():
|
||||
return 7
|
||||
|
||||
return await run_bounded_async(op(), None, label="t")
|
||||
|
||||
result = asyncio.run(scenario())
|
||||
assert result.timed_out is False and result.value == 7
|
||||
|
||||
def test_exception_propagates_unchanged(self):
|
||||
class Boom(RuntimeError):
|
||||
pass
|
||||
|
||||
async def scenario():
|
||||
async def op():
|
||||
raise Boom("x")
|
||||
|
||||
await run_bounded_async(op(), 5.0, label="t")
|
||||
|
||||
with pytest.raises(Boom):
|
||||
asyncio.run(scenario())
|
||||
|
||||
def test_timeout_returns_promptly(self):
|
||||
async def scenario():
|
||||
async def op():
|
||||
await asyncio.sleep(30)
|
||||
|
||||
start = time.monotonic()
|
||||
result = await run_bounded_async(op(), 0.2, label="slow")
|
||||
return result, time.monotonic() - start
|
||||
|
||||
result, elapsed = asyncio.run(scenario())
|
||||
assert result.timed_out is True
|
||||
assert elapsed < 5.0
|
||||
|
||||
def test_timeout_abandons_cancellation_shielded_task(self):
|
||||
"""The family-A killer case: asyncio.wait_for cannot expire a shielded
|
||||
scope; the thread-timer deadline must return anyway."""
|
||||
|
||||
async def scenario():
|
||||
hung = asyncio.Event()
|
||||
|
||||
async def inner():
|
||||
await hung.wait()
|
||||
|
||||
async def shielded():
|
||||
# Shield swallows the cancellation run_bounded_async issues.
|
||||
await asyncio.shield(asyncio.ensure_future(inner()))
|
||||
|
||||
start = time.monotonic()
|
||||
result = await run_bounded_async(shielded(), 0.2, label="shielded")
|
||||
elapsed = time.monotonic() - start
|
||||
hung.set() # release the orphan so the loop can drain
|
||||
await asyncio.sleep(0)
|
||||
return result, elapsed
|
||||
|
||||
result, elapsed = asyncio.run(scenario())
|
||||
assert result.timed_out is True
|
||||
assert elapsed < 5.0
|
||||
|
||||
def test_on_abandon_cleanup_runs_detached(self):
|
||||
async def scenario():
|
||||
cleaned = asyncio.Event()
|
||||
|
||||
async def _cleanup():
|
||||
cleaned.set()
|
||||
|
||||
async def op():
|
||||
await asyncio.sleep(30)
|
||||
|
||||
result = await run_bounded_async(
|
||||
op(), 0.1, label="t", on_abandon=_cleanup
|
||||
)
|
||||
await asyncio.wait_for(cleaned.wait(), timeout=5.0)
|
||||
return result
|
||||
|
||||
result = asyncio.run(scenario())
|
||||
assert result.timed_out is True
|
||||
|
||||
def test_completed_op_never_reports_timeout(self):
|
||||
# Race guard: completion just under the deadline must report success.
|
||||
async def scenario():
|
||||
async def op():
|
||||
await asyncio.sleep(0.01)
|
||||
return "made it"
|
||||
|
||||
return await run_bounded_async(op(), 5.0, label="t")
|
||||
|
||||
result = asyncio.run(scenario())
|
||||
assert result.timed_out is False and result.value == "made it"
|
||||
|
||||
def test_external_cancellation_cancels_inner_task(self):
|
||||
# If the CALLER cancels run_bounded_async, the inner task must not be
|
||||
# leaked running unobserved.
|
||||
async def scenario():
|
||||
started = asyncio.Event()
|
||||
inner_cancelled = asyncio.Event()
|
||||
|
||||
async def op():
|
||||
started.set()
|
||||
try:
|
||||
await asyncio.sleep(30)
|
||||
except asyncio.CancelledError:
|
||||
inner_cancelled.set()
|
||||
raise
|
||||
|
||||
outer = asyncio.ensure_future(
|
||||
run_bounded_async(op(), 25.0, label="t")
|
||||
)
|
||||
await started.wait()
|
||||
outer.cancel()
|
||||
with pytest.raises(asyncio.CancelledError):
|
||||
await outer
|
||||
await asyncio.wait_for(inner_cancelled.wait(), timeout=5.0)
|
||||
return True
|
||||
|
||||
assert asyncio.run(scenario()) is True
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# kill_process_tree
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
@pytest.mark.skipif(sys.platform == "win32", reason="POSIX process-group semantics")
|
||||
class TestKillProcessTree:
|
||||
def test_kills_descendants_of_session_leader(self, tmp_path):
|
||||
"""A child spawned with start_new_session must die with its own child.
|
||||
|
||||
This is the orphan-tree class (#71148): killing only the direct child
|
||||
leaves grandchildren running.
|
||||
"""
|
||||
started = tmp_path / "grandchild_started"
|
||||
marker = tmp_path / "grandchild_alive"
|
||||
grandchild_py = tmp_path / "grandchild.py"
|
||||
grandchild_py.write_text(
|
||||
"import pathlib, time\n"
|
||||
f"pathlib.Path({str(started)!r}).write_text('x')\n"
|
||||
"time.sleep(10)\n"
|
||||
f"pathlib.Path({str(marker)!r}).write_text('x')\n"
|
||||
)
|
||||
parent_py = tmp_path / "parent.py"
|
||||
parent_py.write_text(
|
||||
"import subprocess, sys, time\n"
|
||||
f"subprocess.Popen([sys.executable, {str(grandchild_py)!r}])\n"
|
||||
"time.sleep(10)\n"
|
||||
)
|
||||
proc = subprocess.Popen(
|
||||
[sys.executable, str(parent_py)], start_new_session=True
|
||||
)
|
||||
deadline = time.monotonic() + 10
|
||||
while not started.exists() and time.monotonic() < deadline:
|
||||
time.sleep(0.05)
|
||||
assert started.exists(), "grandchild never spawned — test harness broken"
|
||||
assert kill_process_tree(proc.pid) is True
|
||||
proc.wait(timeout=5)
|
||||
# Grandchild must be dead too: marker never appears.
|
||||
time.sleep(1.5)
|
||||
assert not marker.exists()
|
||||
|
||||
def test_kills_descendant_in_its_own_session(self, tmp_path):
|
||||
"""A descendant that setsid'd out of the parent's group must die too.
|
||||
|
||||
killpg on the parent's group cannot reach it; the psutil descendant
|
||||
sweep must (tools/environments/base.py documents user commands doing
|
||||
exactly this).
|
||||
"""
|
||||
started = tmp_path / "setsid_grandchild_started"
|
||||
marker = tmp_path / "setsid_grandchild_alive"
|
||||
grandchild_py = tmp_path / "grandchild.py"
|
||||
grandchild_py.write_text(
|
||||
"import pathlib, time\n"
|
||||
f"pathlib.Path({str(started)!r}).write_text('x')\n"
|
||||
"time.sleep(10)\n"
|
||||
f"pathlib.Path({str(marker)!r}).write_text('x')\n"
|
||||
)
|
||||
parent_py = tmp_path / "parent.py"
|
||||
parent_py.write_text(
|
||||
"import subprocess, sys, time\n"
|
||||
# grandchild leaves the parent's session/group entirely
|
||||
f"subprocess.Popen([sys.executable, {str(grandchild_py)!r}], start_new_session=True)\n"
|
||||
"time.sleep(10)\n"
|
||||
)
|
||||
proc = subprocess.Popen(
|
||||
[sys.executable, str(parent_py)], start_new_session=True
|
||||
)
|
||||
deadline = time.monotonic() + 10
|
||||
while not started.exists() and time.monotonic() < deadline:
|
||||
time.sleep(0.05)
|
||||
assert started.exists(), "grandchild never spawned — test harness broken"
|
||||
assert kill_process_tree(proc.pid) is True
|
||||
proc.wait(timeout=5)
|
||||
time.sleep(1.5)
|
||||
assert not marker.exists()
|
||||
|
||||
def test_already_dead_pid_returns_false(self):
|
||||
proc = subprocess.Popen([sys.executable, "-c", "pass"], start_new_session=True)
|
||||
proc.wait(timeout=10) # reaped: PID is gone from the process table
|
||||
assert kill_process_tree(proc.pid) is False
|
||||
|
||||
def test_non_group_leader_falls_back_to_single_kill(self):
|
||||
# Child in OUR process group: killpg would signal the test runner.
|
||||
proc = subprocess.Popen([sys.executable, "-c", "import time; time.sleep(30)"])
|
||||
try:
|
||||
assert os.getpgid(proc.pid) != proc.pid # not a leader
|
||||
assert kill_process_tree(proc.pid, sig=signal.SIGTERM) is True
|
||||
proc.wait(timeout=5)
|
||||
finally:
|
||||
if proc.poll() is None:
|
||||
proc.kill()
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# tool_executor migration contract
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
class TestConcurrentToolTimeoutMigration:
|
||||
"""_resolve_concurrent_tool_timeout keeps its exact legacy contract."""
|
||||
|
||||
def _resolver(self):
|
||||
from agent import tool_executor
|
||||
|
||||
return tool_executor._resolve_concurrent_tool_timeout
|
||||
|
||||
def test_default_unchanged(self, monkeypatch):
|
||||
monkeypatch.setattr("agent.deadline._timeouts_section", lambda: {})
|
||||
monkeypatch.delenv("HERMES_CONCURRENT_TOOL_TIMEOUT_S", raising=False)
|
||||
assert self._resolver()() == 420.0
|
||||
|
||||
def test_env_var_still_works(self, monkeypatch):
|
||||
monkeypatch.setattr("agent.deadline._timeouts_section", lambda: {})
|
||||
monkeypatch.setenv("HERMES_CONCURRENT_TOOL_TIMEOUT_S", "60")
|
||||
assert self._resolver()() == 60.0
|
||||
|
||||
def test_env_zero_still_disables(self, monkeypatch):
|
||||
monkeypatch.setattr("agent.deadline._timeouts_section", lambda: {})
|
||||
monkeypatch.setenv("HERMES_CONCURRENT_TOOL_TIMEOUT_S", "0")
|
||||
assert self._resolver()() is None
|
||||
|
||||
def test_env_invalid_still_falls_back_to_default(self, monkeypatch):
|
||||
monkeypatch.setattr("agent.deadline._timeouts_section", lambda: {})
|
||||
monkeypatch.setenv("HERMES_CONCURRENT_TOOL_TIMEOUT_S", "junk")
|
||||
assert self._resolver()() == 420.0
|
||||
|
||||
def test_new_config_key_wins(self, monkeypatch):
|
||||
monkeypatch.setattr(
|
||||
"agent.deadline._timeouts_section",
|
||||
lambda: {"tools": {"concurrent_batch": 300}},
|
||||
)
|
||||
monkeypatch.setenv("HERMES_CONCURRENT_TOOL_TIMEOUT_S", "60")
|
||||
assert self._resolver()() == 300.0
|
||||
Reference in New Issue
Block a user