diff --git a/agent/codex_runtime.py b/agent/codex_runtime.py index fb136e7c68..4e078705b2 100644 --- a/agent/codex_runtime.py +++ b/agent/codex_runtime.py @@ -51,8 +51,15 @@ def _codex_request_failure_details(error: BaseException) -> tuple[int | None, st def _coerce_usage_int(value: Any) -> int: - with suppress(ValueError): - if not isinstance(value, bool) and isinstance(value, (int, float, str)): + if isinstance(value, bool): + return 0 + if isinstance(value, int): + return max(value, 0) + if isinstance(value, float): + return max(int(value), 0) + if isinstance(value, str): + # Only the str->int parse is guarded; a float NaN still raises like it always has. + with suppress(ValueError): return max(int(value), 0) return 0 diff --git a/agent/relay_runtime.py b/agent/relay_runtime.py index 9039e431c9..abc0b68bae 100644 --- a/agent/relay_runtime.py +++ b/agent/relay_runtime.py @@ -158,16 +158,20 @@ class RelaySession: def _load_segments_config() -> dict[str, Any]: """gateway.telemetry.session_segments; both defaults OFF => rotation never fires.""" - segments: dict[str, Any] = {} - with contextlib.suppress(Exception): # config absence must not crash + on_compaction = False + max_turns = 0 + try: from gateway.run import _load_gateway_config # late import telemetry = (_load_gateway_config().get("gateway") or {}).get("telemetry") or {} segments = telemetry.get("session_segments") or {} - try: - max_turns = max(0, int(segments.get("max_turns", 0) or 0)) - except (TypeError, ValueError): - max_turns = 0 - return {"on_compaction": bool(segments.get("on_compaction", False)), "max_turns": max_turns} + on_compaction = bool(segments.get("on_compaction", False)) + try: + max_turns = max(0, int(segments.get("max_turns", 0) or 0)) + except (TypeError, ValueError): + max_turns = 0 + except Exception: # noqa: BLE001 - config absence (or a malformed section) must not crash + pass + return {"on_compaction": on_compaction, "max_turns": max_turns} _SEGMENTS_CONFIG = _Lazy(_load_segments_config) # cached at first read @@ -1217,10 +1221,13 @@ def _resolve_plugin_awaitable(value: Any) -> Any: """Resolve Relay's async plugin API from synchronous host construction.""" if not inspect.isawaitable(value): return value - with contextlib.suppress(RuntimeError): + # Only the "no running loop" probe is guarded: a RuntimeError raised by the awaitable itself + # (re-raised from the daemon thread) must propagate, not fall through to a second asyncio.run. + try: asyncio.get_running_loop() - return _run_on_daemon_thread(lambda: asyncio.run(value), name="hermes-nemo-relay-plugin-lifecycle") - return asyncio.run(value) + except RuntimeError: + return asyncio.run(value) + return _run_on_daemon_thread(lambda: asyncio.run(value), name="hermes-nemo-relay-plugin-lifecycle") def _session_id(event: dict[str, Any]) -> str: diff --git a/gateway/shutdown_watchdog.py b/gateway/shutdown_watchdog.py index 9859c5cd95..66d1e4d118 100644 --- a/gateway/shutdown_watchdog.py +++ b/gateway/shutdown_watchdog.py @@ -29,6 +29,7 @@ from utils import atomic_json_write logger = logging.getLogger(__name__) # Extra leash beyond ``agent.restart_drain_timeout`` so a slow-but-progressing drain survives. +# Matches the issue #66892 suggested hardening. DEFAULT_SHUTDOWN_WATCHDOG_GRACE_S = 60.0 DEFAULT_HEARTBEAT_INTERVAL_S = 30.0 DEFAULT_LOOP_FLOOR_TIMER_INTERVAL_S = 5.0 @@ -36,6 +37,9 @@ DEFAULT_LOOP_WATCHDOG_INTERVAL_S = 30.0 DEFAULT_LOOP_WATCHDOG_TIMEOUT_S = 10.0 # 3 sustained misses (~90-120s of loop block) escalate; stays tight because the heartbeat write is # off-loop. Slow loops tune gateway.loop_watchdog_* in config.yaml. +# The false-positive class that motivated raising this (the watchdog's own on-loop heartbeat fsync stalling +# the loop it monitors) is fixed at the root by the off-loop heartbeat write + two-witness probe (#90502), +# so the default stays tight for genuine wedges. DEFAULT_LOOP_WATCHDOG_MAX_STRIKES = 3 _HEARTBEAT_RELATIVE = ("state", "gateway.heartbeat") _WATCHDOG_DUMP_RELATIVE = ("logs", "gateway-shutdown-watchdog.log") @@ -165,7 +169,12 @@ def get_loop_heartbeat_path(home: Optional[Path] = None) -> Path: def get_loop_tick_socket_path(home: Optional[Path] = None, pid: Optional[int] = None) -> Path: """``/state/gateway.loop-tick..sock`` — PID-suffixed so a stale node from a dead process is never mistaken for this gateway's witness. Served by the loop itself - (``_tick_socket_handler``), so an answer proves the loop dispatches; the heartbeat cannot.""" + (``_tick_socket_handler``), so an answer proves the loop dispatches; the heartbeat cannot. + + Served by the gateway loop itself (see ``_tick_socket_handler``): an answer is direct proof that the + loop is dispatching, which is exactly the property the heartbeat file lost when its write moved off-loop + (#90502). + """ pid = int(pid if pid is not None else os.getpid()) return _home(home) / "state" / f"gateway.loop-tick.{pid}.sock" @@ -264,6 +273,9 @@ def arm_shutdown_watchdog( # Mirror _exit_after_graceful_shutdown: release PID file + runtime lock BEFORE the log drain # (never strand locks), then drain the log queue so logger.critical lands before os._exit. with contextlib.suppress(Exception): + # Mirror _exit_after_graceful_shutdown: release PID file + runtime lock BEFORE the log drain + # (locks must never be stranded), then drain the async log queue so the logger.critical above + # actually reaches the file before os._exit bypasses atexit. (#66892) from gateway.status import remove_pid_file, release_gateway_runtime_lock remove_pid_file() release_gateway_runtime_lock()