diff --git a/gateway/run.py b/gateway/run.py index 62fba594b4..399a653948 100644 --- a/gateway/run.py +++ b/gateway/run.py @@ -1302,6 +1302,19 @@ _AUTO_CONTINUE_FRESHNESS_SECS_DEFAULT = 60 * 60 # Override via ``config.yaml`` ``agent.gateway_startup_restore_drain_timeout``. _STARTUP_RESTORE_DRAIN_TIMEOUT_SECS_DEFAULT = 30.0 +# Default bound for the boot-time turn-machinery warm-up (#99373). On a +# fresh boot with no resume_pending sessions ``_finish_startup_restore`` +# used to open the inbound gate almost immediately, while the agent-side +# turn machinery (the run_agent/model_tools import graph, the tool-registry +# check_fn probes, the prompt builder) was still completely cold. A message +# arriving in that window was served with a skeleton system prompt: no +# context tier, no tool schemas (~1.7K tokens instead of ~14.6K). The +# warm-up runs BEFORE the gate opens so the first inbound turn starts with +# initialized machinery; the bound keeps a wedged init from making the +# gateway permanently unavailable. Override via ``config.yaml`` +# ``agent.gateway_startup_warmup_timeout`` (non-positive disables warm-up). +_STARTUP_WARMUP_TIMEOUT_SECS_DEFAULT = 20.0 + def _coerce_gateway_timestamp(value: Any) -> Optional[float]: """Best-effort conversion of stored gateway timestamps to epoch seconds. @@ -1388,6 +1401,60 @@ def _startup_restore_drain_timeout_secs() -> float: return float(_STARTUP_RESTORE_DRAIN_TIMEOUT_SECS_DEFAULT) +def _startup_warmup_timeout_secs() -> float: + """Max seconds the boot warm-up may hold the inbound gate shut (#99373). + + ``GatewayRunner._warm_turn_prerequisites`` initializes the agent-side + turn machinery BEFORE ``_finish_startup_restore`` opens the inbound + gate, so a message arriving seconds after boot can no longer be served + with a skeleton system prompt (no context tier, no tool schemas). The + warm-up is bounded so a wedged import or probe can never make the + gateway permanently unavailable — on timeout the gate opens anyway and + the warm-up finishes in the background. + + Reads ``HERMES_STARTUP_WARMUP_TIMEOUT`` (bridged from ``config.yaml`` + ``agent.gateway_startup_warmup_timeout`` at gateway startup, same + pattern as the other ``agent.*`` knobs). Non-positive disables the + warm-up entirely (restores the historical lazy-init behaviour). + """ + raw = os.environ.get("HERMES_STARTUP_WARMUP_TIMEOUT") + if raw is None or raw == "": + return float(_STARTUP_WARMUP_TIMEOUT_SECS_DEFAULT) + try: + return float(raw) + except (TypeError, ValueError): + return float(_STARTUP_WARMUP_TIMEOUT_SECS_DEFAULT) + + +def _warm_turn_machinery_sync() -> int: + """Synchronously initialize the turn prerequisites a first turn needs. + + Runs on an executor thread from ``_warm_turn_prerequisites``. Covers + exactly the lazy init observed inside skeleton turns (#99373): + + * the ``run_agent`` heavy import graph (the gateway imports it lazily + inside per-request handlers, so nothing else pulls it in at boot); + * ``model_tools.get_tool_definitions`` — materializes tool schemas and + primes the tool-registry ``check_fn`` TTL cache so availability + probes don't run (and fail cold) inside the user's first turn; + * the context-file tier (AGENTS.md / SOUL.md discovery + read). + + Returns the number of tool schemas materialized (logged for + diagnosability). + """ + import run_agent # noqa: F401 # heavy import graph, cached in sys.modules + import model_tools + + tool_defs = model_tools.get_tool_definitions(quiet_mode=True) + try: + from agent.prompt_builder import build_context_files_prompt + + build_context_files_prompt() + except Exception: + logger.debug("context-file warm-up failed (non-fatal)", exc_info=True) + return len(tool_defs) + + def _as_thread_info(info: Any) -> Optional[Tuple[str, str]]: """*info* as a (thread_id, initial_name) pair, or None if it isn't one. @@ -2760,6 +2827,10 @@ if _config_path.exists(): os.environ["HERMES_STARTUP_RESTORE_DRAIN_TIMEOUT"] = str( _agent_cfg["gateway_startup_restore_drain_timeout"] ) + if "gateway_startup_warmup_timeout" in _agent_cfg: + os.environ["HERMES_STARTUP_WARMUP_TIMEOUT"] = str( + _agent_cfg["gateway_startup_warmup_timeout"] + ) # config-authoritative knobs for the session-search index; same # bridge semantics as the agent settings above. _sessions_cfg = _cfg.get("sessions", {}) @@ -7228,6 +7299,7 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew _profile_failed_platforms: Optional[Dict[str, Dict[Platform, asyncio.Task]]] = None _systemd_watchdog: Optional[Any] = None _startup_restore_in_progress: bool = False + _startup_warmup_task: Optional[asyncio.Task] = None # ------------------------------------------------------------------ # Legacy per-session dict adapters. All per-session state lives in @@ -12400,6 +12472,78 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew drained += 1 return drained + def _start_startup_warmup(self) -> None: + """Kick off the boot turn-machinery warm-up in the background (#99373). + + Called from ``start()`` right after the startup-restore gate closes, + so the warm-up overlaps the (slow, network-bound) platform connects + instead of adding boot latency. ``_finish_startup_restore`` awaits + it (bounded) before opening the inbound gate. + """ + timeout = _startup_warmup_timeout_secs() + if timeout <= 0: + self._startup_warmup_task = None + return + self._startup_warmup_task = asyncio.ensure_future( + self._warm_turn_prerequisites() + ) + + async def _warm_turn_prerequisites(self) -> None: + """Initialize turn machinery off-loop before the gate opens (#99373). + + Runs ``_warm_turn_machinery_sync`` (run_agent import graph, tool + schemas + check_fn probe cache, context-file tier) on an executor + thread so the event loop — platform heartbeats, connects — stays + responsive. Never raises: a failed warm-up degrades to the + historical lazy init, it must not block startup. + """ + try: + loop = asyncio.get_running_loop() + t0 = time.monotonic() + tool_count = await loop.run_in_executor(None, _warm_turn_machinery_sync) + logger.info( + "Turn machinery warmed in %.1fs (%d tool schema(s) materialized)", + time.monotonic() - t0, + tool_count, + ) + except Exception: + logger.warning( + "Turn-machinery warm-up failed; first inbound turn will " + "initialize lazily", + exc_info=True, + ) + + async def _await_startup_warmup(self) -> None: + """Bounded wait for the boot warm-up before the inbound gate opens. + + On timeout the gate opens anyway (availability outranks prompt + completeness for a WEDGED init — same principle as the bounded + restore-drain wait above) and the warm-up continues in the + background; a late failure is still logged. + """ + task = getattr(self, "_startup_warmup_task", None) + if task is None or task.done(): + return + timeout = _startup_warmup_timeout_secs() + if timeout <= 0: + return + done, pending = await asyncio.wait({task}, timeout=timeout) + if pending: + logger.warning( + "Turn-machinery warm-up still running after %.0fs; opening " + "inbound gate anyway — the first turn may see lazily " + "initialized machinery (#99373). Warm-up continues in the " + "background.", + timeout, + ) + task.add_done_callback( + lambda t: GatewayRunner._log_late_background_failure( + t, + "boot turn-machinery warm-up failed after gate release", + level=logging.DEBUG, + ) + ) + async def _finish_startup_restore(self) -> None: """Wait (BOUNDED) for startup auto-resume, then release + drain inbound. @@ -12453,6 +12597,9 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew exc_info=(type(exc), exc, exc.__traceback__), ) self._startup_restore_tasks = [] + # Warm the turn machinery BEFORE the queue drains: replayed (and + # fresh) inbound turns must not build skeleton prompts (#99373). + await self._await_startup_warmup() drained = await self._drain_startup_restore_queue() self._startup_restore_in_progress = False if drained: @@ -13573,6 +13720,14 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew self._startup_restore_in_progress = True self._startup_restore_queue = [] self._startup_restore_tasks = [] + # Fresh-boot readiness (#99373): with no resume_pending sessions the + # gate above opens almost immediately, while the agent-side turn + # machinery (run_agent import graph, tool schemas, check_fn probes, + # context tier) is still cold — a message in that window was served + # with a skeleton system prompt. Start warming NOW so the work + # overlaps the network-bound platform connects below; + # _finish_startup_restore awaits it (bounded) before opening the gate. + self._start_startup_warmup() connected_count = 0 enabled_platform_count = 0 diff --git a/hermes_cli/config_defaults.py b/hermes_cli/config_defaults.py index 32dc955554..d403047f6e 100644 --- a/hermes_cli/config_defaults.py +++ b/hermes_cli/config_defaults.py @@ -327,6 +327,17 @@ DEFAULT_CONFIG = { # synchronously before the gate runs. Set to 0 to disable the bound # (historical "wait forever" behaviour). "gateway_startup_restore_drain_timeout": 30, + # Max seconds the boot turn-machinery warm-up (#99373) may hold the + # gateway's inbound gate shut. On a fresh boot the gateway warms the + # agent-side turn prerequisites (run_agent import graph, tool schemas + # + availability probes, context-file tier) BEFORE accepting inbound + # messages, so a message seconds after boot is no longer served with + # a skeleton system prompt (missing context files / tool schemas). + # On timeout the gate opens anyway and warm-up finishes in the + # background — a wedged init can't make the gateway permanently + # unavailable. Set to 0 to disable the warm-up (historical + # lazy-init behaviour). + "gateway_startup_warmup_timeout": 20, # Stale-stream ceiling for local providers (Ollama, oMLX, llama-cpp) in # seconds. When the base stale timeout is at its default (180s) and a # local endpoint is detected, this finite ceiling replaces the former diff --git a/tests/gateway/test_restart_resume_pending.py b/tests/gateway/test_restart_resume_pending.py index 2e99c2605b..fe147ae576 100644 --- a/tests/gateway/test_restart_resume_pending.py +++ b/tests/gateway/test_restart_resume_pending.py @@ -765,6 +765,98 @@ async def test_startup_restore_waits_for_resume_before_draining_inbound(): assert runner._startup_restore_in_progress is False +# --------------------------------------------------------------------------- +# Fresh-boot turn-machinery warm-up gate (#99373) +# --------------------------------------------------------------------------- + + +@pytest.mark.asyncio +async def test_fresh_boot_gate_stays_closed_until_warmup_completes(monkeypatch): + """#99373 regression: on a fresh boot (no resume_pending sessions) the + inbound gate must NOT open while the turn-machinery warm-up is still + running — a message in that window used to be served with a skeleton + system prompt (no context tier, no tool schemas).""" + runner, adapter = make_restart_runner() + runner._startup_restore_in_progress = True + runner._startup_restore_queue = [] + runner._startup_restore_tasks = [] # fresh boot: nothing to resume + + monkeypatch.setenv("HERMES_STARTUP_WARMUP_TIMEOUT", "5") + + warmup_done = asyncio.Event() + runner._startup_warmup_task = asyncio.create_task(warmup_done.wait()) + + handled: list[str] = [] + + async def fake_handle_message(event: MessageEvent) -> None: + handled.append(event.text) + + adapter.handle_message = fake_handle_message + + source = make_restart_source(chat_id="fresh-boot-chat") + inbound = MessageEvent( + text="early-bird", message_type=MessageType.TEXT, source=source + ) + # Inbound during the warm-up window queues instead of dispatching. + assert await runner._handle_message(inbound) is None + assert runner._startup_restore_queue == [inbound] + + finish_task = asyncio.create_task(runner._finish_startup_restore()) + for _ in range(5): + await asyncio.sleep(0) + # Warm-up still running -> gate still closed, nothing dispatched. + assert not finish_task.done() + assert runner._startup_restore_in_progress is True + assert handled == [] + + warmup_done.set() + await asyncio.wait_for(finish_task, timeout=5) + + # Gate opened only after warm-up; the queued message replayed. + assert runner._startup_restore_in_progress is False + assert handled == ["early-bird"] + assert runner._startup_restore_queue == [] + + +@pytest.mark.asyncio +async def test_wedged_warmup_cannot_hold_gate_shut_past_timeout(monkeypatch): + """Availability bound (#99373 / #98473 premise): a wedged warm-up must + not make the gateway permanently unavailable — the gate opens after the + bounded wait and the warm-up continues in the background.""" + runner, _adapter = make_restart_runner() + runner._startup_restore_in_progress = True + runner._startup_restore_queue = [] + runner._startup_restore_tasks = [] + + monkeypatch.setenv("HERMES_STARTUP_WARMUP_TIMEOUT", "0.1") + + never = asyncio.Event() + wedged = asyncio.create_task(never.wait()) + runner._startup_warmup_task = wedged + + await asyncio.wait_for(runner._finish_startup_restore(), timeout=5) + + assert runner._startup_restore_in_progress is False + assert not wedged.done() # warm-up not cancelled, continues in background + wedged.cancel() + + +@pytest.mark.asyncio +async def test_warmup_disabled_by_nonpositive_timeout(monkeypatch): + """gateway_startup_warmup_timeout <= 0 restores historical lazy init.""" + runner, _adapter = make_restart_runner() + runner._startup_restore_in_progress = True + runner._startup_restore_queue = [] + runner._startup_restore_tasks = [] + + monkeypatch.setenv("HERMES_STARTUP_WARMUP_TIMEOUT", "0") + runner._start_startup_warmup() + assert runner._startup_warmup_task is None + + await asyncio.wait_for(runner._finish_startup_restore(), timeout=5) + assert runner._startup_restore_in_progress is False + + # --------------------------------------------------------------------------- # Shutdown banner wording # ---------------------------------------------------------------------------