review-fix(suppress-audit): langfuse/run_agent_cache/openviking — restore BASE exception semantics

This commit is contained in:
Teknium
2026-09-03 10:03:17 -07:00
parent 4749300508
commit fd333a67cd
3 changed files with 22 additions and 10 deletions

View File

@@ -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

View File

@@ -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)

View File

@@ -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