From fd333a67cd52fd4044ddad2c9ef8845a9395cb66 Mon Sep 17 00:00:00 2001 From: Teknium <127238744+teknium1@users.noreply.github.com> Date: Thu, 3 Sep 2026 10:03:17 -0700 Subject: [PATCH] =?UTF-8?q?review-fix(suppress-audit):=20langfuse/run=5Fag?= =?UTF-8?q?ent=5Fcache/openviking=20=E2=80=94=20restore=20BASE=20exception?= =?UTF-8?q?=20semantics?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- gateway/run_agent_cache.py | 17 ++++++++++++----- plugins/memory/openviking/__init__.py | 12 ++++++++---- plugins/observability/langfuse/__init__.py | 3 ++- 3 files changed, 22 insertions(+), 10 deletions(-) diff --git a/gateway/run_agent_cache.py b/gateway/run_agent_cache.py index a7845d8eba..717076129d 100644 --- a/gateway/run_agent_cache.py +++ b/gateway/run_agent_cache.py @@ -596,13 +596,15 @@ class GatewayAgentCacheMixin: ) def _spawn_release_thread(self, target, args: tuple, name: str, *, inline_fallback: bool) -> None: - """Run a release on a daemon thread; ``inline_fallback`` runs it inline when no thread can start (interpreter shutdown).""" + """Run a release on a daemon thread. ``inline_fallback`` runs it inline (best-effort) when no + thread can start (interpreter shutdown); otherwise a spawn failure propagates, as on main.""" try: threading.Thread(target=target, args=args, daemon=True, name=name).start() except Exception: - if inline_fallback: - with suppress(Exception): - target(*args) + if not inline_fallback: + raise + with suppress(Exception): + target(*args) def _finalizable_unexpired_session_entry(self, key: str): """Session-store entry for ``key`` when the expiry watcher will still finalize it; None when @@ -760,7 +762,12 @@ class GatewayAgentCacheMixin: "Agent cache pressure: anon RSS %dMB over budget %dMB — evicting %d LRU session(s): %s", rss_mb, bounds.memory_high_mb, evicted_count, ", ".join(key for key, _ in plan), ) - self._spawn_release_thread(self._release_pressure_batch, (plan,), "agent-cache-pressure", inline_fallback=True) + try: + threading.Thread(target=self._release_pressure_batch, args=(plan,), daemon=True, + name="agent-cache-pressure").start() + except Exception: + # Thread spawn failed (interpreter shutdown): release inline, unguarded (as on main). + self._release_pressure_batch(plan) # _release_pressure_batch drains `plan` in place (so the trim runs with no lingering agent # refs) — len(plan) is 0 once the daemon thread finishes, hence the pre-captured count. return evicted_count diff --git a/plugins/memory/openviking/__init__.py b/plugins/memory/openviking/__init__.py index b1f0caf3f3..270c81f7d6 100644 --- a/plugins/memory/openviking/__init__.py +++ b/plugins/memory/openviking/__init__.py @@ -206,10 +206,14 @@ def _atexit_commit_sessions(): if provider is None: return _last_active_provider = None - with suppress(Exception): # best-effort at shutdown time - provider.on_session_end([]) - with suppress(Exception): - provider._release_run_lock() + try: + with suppress(Exception): # best-effort at shutdown time + provider.on_session_end([]) + finally: + # ``finally`` (as on main): the run lock is released even when on_session_end + # dies of a BaseException (KeyboardInterrupt during atexit). + with suppress(Exception): + provider._release_run_lock() atexit.register(_atexit_commit_sessions) diff --git a/plugins/observability/langfuse/__init__.py b/plugins/observability/langfuse/__init__.py index bf28c17293..a6edc25227 100644 --- a/plugins/observability/langfuse/__init__.py +++ b/plugins/observability/langfuse/__init__.py @@ -920,7 +920,8 @@ def on_session_finalize(*, session_id: str = "", reason: str = "", **_: Any) -> keys = [k for k in _TRACE_STATE if not session_id or k == session_id or any(f in k for f in fragments)] for key in keys: _finish_trace(key) - _flush(client) + with _failsafe("finalize flush"): + client.flush() # Shut down only at true process exit (not /new, /reset, session expiry: the # cached client must keep exporting). Doing it while modules are intact keeps