Files
hermes-agent/tools/daemon_pool.py
T
MattMaximo 6f625f7381 fix(tools): propagate caller contextvars in DaemonThreadPoolExecutor.submit
Some bundled CPython runtime builds strip stdlib ThreadPoolExecutor's
copy_context() propagation, so work submitted to the daemon pool runs in a
bare context. Under the multiplexed gateway this dropped the profile
secret scope in pool workers: the context-compression timeout fence
resolved auxiliary provider keys (SURPLUS_API_KEY) with
UnscopedSecretError, silently degrading LLM compression to lossy
deterministic summaries and driving re-read loops in affected sessions.

Restore stdlib semantics in submit() by snapshotting the caller's context
and running the callable inside it (a no-op re-application on runtimes
that already propagate). Mirrors the gateway's
_run_in_executor_with_context pattern.

Tests: daemon pool worker sees caller contextvars; scoped get_secret works
in a daemon-pool worker under multiplex while scoped misses still fail
closed (no env leak).
2026-09-01 22:28:52 -07:00

96 lines
4.0 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""Shared daemon-thread ThreadPoolExecutor.
Stdlib ``ThreadPoolExecutor`` workers are non-daemon AND are registered in
``concurrent.futures.thread._threads_queues``, whose atexit hook
(``_python_exit``) joins every worker unconditionally — even after
``shutdown(wait=False)``. A single wedged worker (tool blocked on network
I/O, hung provider daemon, stuck subagent) therefore blocks interpreter
exit forever. This is the root cause of multi-minute CLI exits on long
sessions: every abandoned concurrent-tool batch leaves workers that the
exit hook insists on joining.
``DaemonThreadPoolExecutor`` spawns daemon workers and skips the
``_threads_queues`` registration, so:
- ``_python_exit`` never joins them, and
- the interpreter's non-daemon thread join at shutdown skips them.
Semantics are otherwise identical (initializer/initargs, work queue,
idle-thread reuse) and, since #95119, so is context propagation:
``submit`` snapshots the submitting context with ``copy_context()`` and
runs each work item inside it, matching stdlib ``ThreadPoolExecutor``.
That matters because some bundled CPython runtime builds omit stdlib's
context propagation entirely, which silently drops contextvar-based state
(profile secret scope, HERMES_HOME override) in pool workers — e.g. the
context-compression timeout fence resolved auxiliary provider keys with
``UnscopedSecretError`` under the multiplexed gateway. Use it for any
pool whose work is best-effort or
independently interruptible and must never hold the process open:
concurrent tool execution, background memory sync, catalog fan-out,
subagent timeout wrappers. Do NOT use it for work that must complete
before exit (durable writes) — those belong on foreground threads with
explicit bounded joins.
"""
from __future__ import annotations
import threading
import weakref
from concurrent.futures import ThreadPoolExecutor
from concurrent.futures.thread import _worker
from contextvars import copy_context
__all__ = ["DaemonThreadPoolExecutor"]
class DaemonThreadPoolExecutor(ThreadPoolExecutor):
"""ThreadPoolExecutor variant whose workers do not block process exit."""
def submit(self, fn, /, *args, **kwargs):
"""Submit a callable, propagating the caller's contextvars.
Stdlib ``ThreadPoolExecutor`` snapshots the submitting context with
``copy_context()`` and runs each work item inside it, so pool
workers inherit contextvar state such as the multiplexed profile
secret scope. Some bundled CPython runtime builds strip that
propagation from the stdlib executor (their ``_WorkItem.run`` calls
the callable directly), which broke auxiliary LLM key resolution
from the context-compression timeout fence with
``UnscopedSecretError``. Restore the stdlib behavior explicitly so
the daemon pool behaves identically on every runtime; on runtimes
that already propagate, the inner ``ctx.run`` re-applies the same
immutable context and is a no-op.
"""
ctx = copy_context()
def _run_with_context(*call_args, **call_kwargs):
return ctx.run(fn, *call_args, **call_kwargs)
return super().submit(_run_with_context, *args, **kwargs)
def _adjust_thread_count(self) -> None:
# Mirrors CPython's implementation (3.8–3.13) with two changes:
# daemon=True and no _threads_queues registration.
if self._idle_semaphore.acquire(timeout=0):
return
def weakref_cb(_, q=self._work_queue):
q.put(None)
num_threads = len(self._threads)
if num_threads < self._max_workers:
thread_name = "%s_%d" % (self._thread_name_prefix or self, num_threads)
t = threading.Thread(
name=thread_name,
target=_worker,
args=(
weakref.ref(self, weakref_cb),
self._work_queue,
self._initializer,
self._initargs,
),
daemon=True,
)
t.start()
self._threads.add(t)