fix(gateway): never construct SessionDB on the event-loop thread
SessionDB.__init__ runs schema init, and a migration against a contended state.db blocks for seconds. The goal/heartbeat path reached it synchronously on the gateway's event-loop thread (GoalManager() -> load_goal -> _get_session_db -> SessionDB()), so a contended DB starved the loop-liveness watchdog, which hard-exited with code 75 and the supervisor restarted straight back into the same state - an unbounded crash loop reported from an enterprise fleet. _get_session_db now detects a running loop on the calling thread: on a cache miss it kicks a one-shot background bootstrap thread and returns None immediately (every caller already degrades gracefully on None); the cached instance serves all later calls. Worker threads construct inline as before, with a lock-guarded cache so a bootstrap race keeps one instance and closes the loser. The heartbeat module shares this boundary via the same _get_session_db.
This commit is contained in:
committed by
Teknium
parent
22f0f22298
commit
8e81e2aaae
+64
-1
@@ -29,12 +29,14 @@ Nothing in this module touches the agent's system prompt or toolset.
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
import hashlib
|
||||
import json
|
||||
import logging
|
||||
import os
|
||||
import re
|
||||
import subprocess
|
||||
import threading
|
||||
import time
|
||||
from dataclasses import dataclass, field, asdict
|
||||
from datetime import datetime, timezone
|
||||
@@ -661,6 +663,23 @@ def _meta_key(session_id: str) -> str:
|
||||
|
||||
|
||||
_DB_CACHE: Dict[str, Any] = {}
|
||||
_DB_BOOTSTRAP_LOCK = threading.Lock()
|
||||
_DB_BOOTSTRAP_INFLIGHT: set = set()
|
||||
|
||||
|
||||
def _bootstrap_session_db(home: str) -> None:
|
||||
"""Construct SessionDB off-loop and populate the cache (worker thread)."""
|
||||
try:
|
||||
from hermes_state import SessionDB
|
||||
|
||||
db = SessionDB()
|
||||
except Exception as exc: # pragma: no cover
|
||||
logger.debug("GoalManager: background SessionDB() raised (%s)", exc)
|
||||
db = None
|
||||
with _DB_BOOTSTRAP_LOCK:
|
||||
if db is not None and home not in _DB_CACHE:
|
||||
_DB_CACHE[home] = db
|
||||
_DB_BOOTSTRAP_INFLIGHT.discard(home)
|
||||
|
||||
|
||||
def _get_session_db() -> Optional[Any]:
|
||||
@@ -671,6 +690,15 @@ def _get_session_db() -> Optional[Any]:
|
||||
``hermes_home`` path so profile switches still pick up the right DB.
|
||||
Defensive against import/instantiation failures so tests and
|
||||
non-standard launchers can still use the GoalManager.
|
||||
|
||||
Never constructs SessionDB on an event-loop thread. ``SessionDB.__init__``
|
||||
runs schema init, and a migration against a contended state.db blocks for
|
||||
seconds — on the gateway's loop thread that starves the loop-liveness
|
||||
watchdog, which hard-exits the process (exit 75) and crash-loops the
|
||||
gateway (Coatue field report, 2026-08-14). On a cache miss with a running
|
||||
loop we kick a one-shot background bootstrap and return None; every
|
||||
caller already degrades gracefully on None, and a later call returns the
|
||||
cached instance.
|
||||
"""
|
||||
try:
|
||||
from hermes_constants import get_hermes_home
|
||||
@@ -684,12 +712,47 @@ def _get_session_db() -> Optional[Any]:
|
||||
cached = _DB_CACHE.get(home)
|
||||
if cached is not None:
|
||||
return cached
|
||||
|
||||
try:
|
||||
asyncio.get_running_loop()
|
||||
except RuntimeError:
|
||||
on_loop_thread = False
|
||||
else:
|
||||
on_loop_thread = True
|
||||
|
||||
if on_loop_thread:
|
||||
with _DB_BOOTSTRAP_LOCK:
|
||||
# Re-check under the lock: a bootstrap may have finished between
|
||||
# the unlocked read above and here.
|
||||
cached = _DB_CACHE.get(home)
|
||||
if cached is not None:
|
||||
return cached
|
||||
if home not in _DB_BOOTSTRAP_INFLIGHT:
|
||||
_DB_BOOTSTRAP_INFLIGHT.add(home)
|
||||
threading.Thread(
|
||||
target=_bootstrap_session_db,
|
||||
args=(home,),
|
||||
name="goals-sessiondb-bootstrap",
|
||||
daemon=True,
|
||||
).start()
|
||||
return None
|
||||
|
||||
try:
|
||||
db = SessionDB()
|
||||
except Exception as exc: # pragma: no cover
|
||||
logger.debug("GoalManager: SessionDB() raised (%s)", exc)
|
||||
return None
|
||||
_DB_CACHE[home] = db
|
||||
with _DB_BOOTSTRAP_LOCK:
|
||||
existing = _DB_CACHE.get(home)
|
||||
if existing is not None:
|
||||
# A concurrent bootstrap won the race; keep one instance and
|
||||
# close ours so connections don't leak.
|
||||
try:
|
||||
db.close()
|
||||
except Exception:
|
||||
pass
|
||||
return existing
|
||||
_DB_CACHE[home] = db
|
||||
return db
|
||||
|
||||
|
||||
|
||||
@@ -0,0 +1,133 @@
|
||||
"""SessionDB bootstrap must never run schema init on an event-loop thread.
|
||||
|
||||
Coatue field report (2026-08-14): after an unclean shutdown plus a
|
||||
double-instance startup race, the v25 schema migration inside
|
||||
``SessionDB.__init__`` blocked the gateway's event-loop thread (reached via
|
||||
``_post_turn_goal_continuation`` → ``GoalManager()`` → ``load_goal`` →
|
||||
``_get_session_db``). The loop-liveness watchdog missed 3 probes and killed
|
||||
the process with exit 75; the supervisor restarted into the same state —
|
||||
an unbounded crash loop.
|
||||
|
||||
Contract fixed here, at the shared boundary (goals.py ``_get_session_db``,
|
||||
which the heartbeat module reuses):
|
||||
|
||||
- Called on a thread with a RUNNING event loop and no cached DB, it must
|
||||
NOT construct SessionDB inline. It returns None immediately (callers
|
||||
already degrade gracefully on None) and populates the cache from a
|
||||
background thread.
|
||||
- Called on a plain worker thread, it constructs inline as before.
|
||||
- Once cached, every thread gets the cached instance.
|
||||
"""
|
||||
|
||||
import asyncio
|
||||
import threading
|
||||
import time
|
||||
|
||||
import pytest
|
||||
|
||||
import hermes_cli.goals as goals
|
||||
|
||||
|
||||
class _RecordingDB:
|
||||
"""Stands in for SessionDB; records which thread constructed it."""
|
||||
|
||||
constructed_on: list = []
|
||||
|
||||
def __init__(self):
|
||||
_RecordingDB.constructed_on.append(threading.get_ident())
|
||||
|
||||
def get_meta(self, key):
|
||||
return None
|
||||
|
||||
|
||||
@pytest.fixture(autouse=True)
|
||||
def _clean_cache(monkeypatch):
|
||||
_RecordingDB.constructed_on = []
|
||||
monkeypatch.setattr(goals, "_DB_CACHE", {})
|
||||
yield
|
||||
|
||||
|
||||
def _patch_sessiondb(monkeypatch):
|
||||
import hermes_state
|
||||
|
||||
monkeypatch.setattr(hermes_state, "SessionDB", _RecordingDB)
|
||||
|
||||
|
||||
def test_loop_thread_cache_miss_returns_none_without_constructing(monkeypatch):
|
||||
"""On the event-loop thread with a cold cache: no inline construction,
|
||||
immediate None, background population."""
|
||||
_patch_sessiondb(monkeypatch)
|
||||
loop_thread_id = None
|
||||
inline_result = "UNSET"
|
||||
|
||||
async def main():
|
||||
nonlocal loop_thread_id, inline_result
|
||||
loop_thread_id = threading.get_ident()
|
||||
inline_result = goals._get_session_db()
|
||||
|
||||
asyncio.run(main())
|
||||
|
||||
# The loop-thread call must not have blocked on construction.
|
||||
assert inline_result is None
|
||||
# Background thread eventually populates the cache, OFF the loop thread.
|
||||
deadline = time.monotonic() + 5
|
||||
while time.monotonic() < deadline and not goals._DB_CACHE:
|
||||
time.sleep(0.02)
|
||||
assert goals._DB_CACHE, "background bootstrap never populated the cache"
|
||||
assert _RecordingDB.constructed_on, "SessionDB never constructed"
|
||||
assert all(t != loop_thread_id for t in _RecordingDB.constructed_on), (
|
||||
"SessionDB was constructed on the event-loop thread"
|
||||
)
|
||||
|
||||
|
||||
def test_loop_thread_cache_hit_returns_cached_instance(monkeypatch):
|
||||
"""A warm cache is returned directly even on the loop thread."""
|
||||
_patch_sessiondb(monkeypatch)
|
||||
sentinel = _RecordingDB()
|
||||
from hermes_constants import get_hermes_home
|
||||
|
||||
goals._DB_CACHE[str(get_hermes_home())] = sentinel
|
||||
result = "UNSET"
|
||||
|
||||
async def main():
|
||||
nonlocal result
|
||||
result = goals._get_session_db()
|
||||
|
||||
asyncio.run(main())
|
||||
assert result is sentinel
|
||||
|
||||
|
||||
def test_worker_thread_constructs_inline(monkeypatch):
|
||||
"""No running loop on the calling thread → construct inline, return it."""
|
||||
_patch_sessiondb(monkeypatch)
|
||||
db = goals._get_session_db()
|
||||
assert db is not None
|
||||
assert isinstance(db, _RecordingDB)
|
||||
assert _RecordingDB.constructed_on == [threading.get_ident()]
|
||||
|
||||
|
||||
def test_slow_construction_does_not_block_the_loop(monkeypatch):
|
||||
"""A SessionDB whose init blocks (locked-DB migration) must not stall
|
||||
the event loop past the watchdog probe window."""
|
||||
import hermes_state
|
||||
|
||||
class _BlockingDB:
|
||||
def __init__(self):
|
||||
time.sleep(2.0) # simulated contended migration
|
||||
|
||||
def get_meta(self, key):
|
||||
return None
|
||||
|
||||
monkeypatch.setattr(hermes_state, "SessionDB", _BlockingDB)
|
||||
elapsed = None
|
||||
|
||||
async def main():
|
||||
nonlocal elapsed
|
||||
t0 = time.monotonic()
|
||||
goals._get_session_db()
|
||||
elapsed = time.monotonic() - t0
|
||||
|
||||
asyncio.run(main())
|
||||
assert elapsed is not None and elapsed < 0.5, (
|
||||
f"loop-thread call blocked for {elapsed:.2f}s — watchdog territory"
|
||||
)
|
||||
Reference in New Issue
Block a user