From b48cc47a4c4212c70552c81338eff146ceb3abdd Mon Sep 17 00:00:00 2001 From: ethernet Date: Tue, 8 Sep 2026 17:04:44 -0400 Subject: [PATCH] fix(tools): use Python 3.14 daemon worker contexts Python 3.14 moves initializer state into a worker context. The old worker arguments fail on the first submission to a fresh executor. Use the current stdlib worker signature while preserving daemon shutdown and per-submission profile isolation. Cover worker initialization and context isolation on reused threads. Verified on Python 3.14.7: 41 focused tests pass through scripts/run_tests.sh. Ruff and diff checks pass. Packaged release verification remains pending. --- tests/tools/test_daemon_pool.py | 64 ++++++++++++++++++++++++--------- tools/daemon_pool.py | 18 +++++----- 2 files changed, 55 insertions(+), 27 deletions(-) diff --git a/tests/tools/test_daemon_pool.py b/tests/tools/test_daemon_pool.py index 9370afb46f..5b14b231d6 100644 --- a/tests/tools/test_daemon_pool.py +++ b/tests/tools/test_daemon_pool.py @@ -31,6 +31,34 @@ def test_workers_are_daemon_threads(): pool.shutdown(wait=True) +def test_initializer_runs_on_each_worker_before_tasks(): + local = threading.local() + ready = threading.Barrier(2, timeout=10) + initialized = [] + lock = threading.Lock() + marker = object() + + def initialize(value): + local.value = value + with lock: + initialized.append(threading.current_thread()) + + def task(): + ready.wait() + return local.value, threading.current_thread() + + with DaemonThreadPoolExecutor(max_workers=2, initializer=initialize, initargs=(marker,)) as pool: + futures = [pool.submit(task) for _ in range(2)] + results = [future.result(timeout=10) for future in futures] + + workers = {worker for _, worker in results} + assert len(workers) == 2 + assert set(initialized) == workers + assert len(initialized) == len(workers) + assert all(value is marker for value, _ in results) + assert all(worker.daemon and worker not in _threads_queues for worker in workers) + + def test_idle_worker_reuse(): pool = DaemonThreadPoolExecutor(max_workers=4) try: @@ -70,27 +98,29 @@ def test_wedged_worker_does_not_block_interpreter_exit(): def test_submit_propagates_caller_contextvars(): - """Pool workers inherit contextvars set in the submitting context. - - Stdlib ThreadPoolExecutor snapshots the caller's context with - ``copy_context()``; some bundled CPython runtime builds strip that, so - the daemon pool restores it explicitly. Without the fix this returns - the default because the worker runs in a bare context. - """ + """A reused worker must not mix profile scopes or leak task mutations.""" from contextvars import ContextVar var = ContextVar("daemon_pool_test_var", default="unset") - pool = DaemonThreadPoolExecutor(max_workers=1) - try: - token = var.set("hello") - try: - seen = pool.submit(var.get).result(timeout=10) - finally: - var.reset(token) - assert seen == "hello" - finally: - pool.shutdown(wait=True) + def read_and_mutate(): + seen = var.get() + var.set("worker mutation") + return seen, threading.current_thread() + + with DaemonThreadPoolExecutor(max_workers=1) as pool: + futures = [] + for profile in ("first", "second"): + token = var.set(profile) + try: + futures.append(pool.submit(read_and_mutate)) + finally: + var.reset(token) + results = [future.result(timeout=10) for future in futures] + assert [seen for seen, _ in results] == ["first", "second"] + assert results[0][1] is results[1][1] + assert pool.submit(var.get).result(timeout=10) == "unset" + assert var.get() == "unset" def _repo_root(): diff --git a/tools/daemon_pool.py b/tools/daemon_pool.py index 5682a8f6ae..f3acc7bdd4 100644 --- a/tools/daemon_pool.py +++ b/tools/daemon_pool.py @@ -24,11 +24,11 @@ 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).""" + """Keep each task in its caller's profile scope, even on a reused worker. + + Thread-start context cannot track later submissions from other profiles + (#54937). The stdlib worker context manages initialization, not contextvars. + """ ctx = copy_context() def _run_with_context(*call_args, **call_kwargs): @@ -36,8 +36,8 @@ class DaemonThreadPoolExecutor(ThreadPoolExecutor): 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. + # Match CPython 3.14 worker startup, but keep abandoned work out of the + # interpreter's atexit joins: daemon=True and no _threads_queues entry. if self._idle_semaphore.acquire(timeout=0): return @@ -46,11 +46,9 @@ class DaemonThreadPoolExecutor(ThreadPoolExecutor): 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), + args=(weakref.ref(self, weakref_cb), self._create_worker_context(), self._work_queue), ) t.start() self._threads.add(t)