Files
hermes-agent/tools/daemon_pool.py
Teknium e83816a4d1 review-fix(comments): restore lost #NNNN rationale comments across non-test source (mechanical sweep, condensed, code unchanged)
For each issue anchor present in BASE 63279301bc non-test .py and absent on HEAD, the BASE comment/docstring block was re-attached at the HEAD location of the code it explained (matched by the distinctive code line / enclosing def). Sentences already covered by an existing HEAD comment were deduped; the issue number always survives. Insert-only: no code lines changed.
2026-09-03 09:44:26 -07:00

57 lines
2.5 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 workers are non-daemon AND registered in ``_threads_queues``, whose atexit
hook joins every worker even after ``shutdown(wait=False)`` — one wedged worker
(tool blocked on network I/O, hung provider, stuck subagent) blocks interpreter
exit forever. This variant spawns daemon workers and skips that registration.
Use it for best-effort/interruptible work that must never hold the process open;
NOT for work that must complete before exit (durable writes 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 only does
this from 3.14; on 3.11-3.13 a bare worker starts with an EMPTY Context and
drops profile secret scope / HERMES_HOME override — under the multiplexed
gateway a credential read then fails closed with ``UnscopedSecretError``.
Unconditional: on 3.14+ ``ctx.run`` re-applies the same context (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)
# Carry the active profile into the review thread so MEMORY.md / skill review writes land in the
# right profile (#54937).
t = threading.Thread(
name=thread_name, target=_worker, daemon=True,
args=(weakref.ref(self, weakref_cb), self._work_queue, self._initializer, self._initargs),
)
t.start()
self._threads.add(t)