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.
This commit is contained in:
@@ -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():
|
||||
|
||||
@@ -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)
|
||||
|
||||
Reference in New Issue
Block a user