refactor: bound shared thread pool shutdown

This commit is contained in:
jxxghp
2026-08-24 14:51:38 +08:00
parent ad45bfac5d
commit 8826219173
7 changed files with 229 additions and 16 deletions
+2 -2
View File
@@ -478,8 +478,8 @@ def pytest_sessionfinish(session, exitstatus):
from app.runtime.thread import ThreadHelper
helper = ThreadHelper.get_existing_instance()
if helper:
helper.shutdown()
if helper and helper.shutdown() is False:
raise RuntimeError("shared thread pool did not converge")
except Exception as err:
_report_session_cleanup_error(session, "thread helper", err)
+12
View File
@@ -969,6 +969,18 @@ def test_stop_modules_propagates_message_queue_nonconvergence(monkeypatch):
_assert_completed_once(dependency)
def test_stop_modules_propagates_shared_thread_pool_nonconvergence(monkeypatch):
"""共享线程池仍有活动 Future 时必须由模块服务关闭结果向上暴露。"""
dependencies = _patch_module_shutdown_dependencies(monkeypatch)
dependencies["thread"].return_value = False
converged = asyncio.run(modules_initializer.stop_modules())
assert converged is False
for dependency in dependencies.values():
_assert_completed_once(dependency)
def test_stop_modules_drains_web_agent_tasks_before_persistence(monkeypatch):
"""关闭时先收口 Web Agent,再关闭持久化准入和数据库任务。"""
order = []
+103
View File
@@ -0,0 +1,103 @@
import threading
import time
from concurrent.futures import ThreadPoolExecutor
import pytest
from app.runtime import thread as thread_module
from app.runtime.thread import ThreadHelper
def _new_thread_helper() -> ThreadHelper:
"""构造不进入全局 Singleton 的隔离线程池 owner。"""
helper = object.__new__(ThreadHelper)
helper.__init__()
return helper
@pytest.mark.parametrize("use_legacy_pool", [False, True])
def test_thread_helper_shutdown_is_bounded_and_retryable(use_legacy_pool):
"""宿主 submit 和旧 pool.submit 都必须进入同一个可重试关闭 owner。"""
helper = _new_thread_helper()
entered = threading.Event()
release = threading.Event()
def blocked_work() -> str:
"""模拟无法由线程池强制取消的同步任务。"""
entered.set()
release.wait()
return "done"
assert isinstance(helper.pool, ThreadPoolExecutor)
submit = helper.pool.submit if use_legacy_pool else helper.submit
future = submit(blocked_work)
try:
assert entered.wait(timeout=1)
started_at = time.monotonic()
assert helper.shutdown(timeout=0.01) is False
assert time.monotonic() - started_at < 1
assert not future.done()
with pytest.raises(RuntimeError):
helper.submit(lambda: None)
with pytest.raises(RuntimeError):
helper.pool.submit(lambda: None)
finally:
release.set()
assert future.result(timeout=1) == "done"
assert helper.shutdown(timeout=1) is True
def test_thread_helper_shutdown_preserves_queued_work():
"""关闭封口不得取消已接受的排队任务,保持历史完成语义。"""
executor = thread_module._OwnedThreadPoolExecutor(max_workers=1)
entered = threading.Event()
release = threading.Event()
def blocked_work() -> str:
"""占用唯一 worker,确保后一任务仍在队列中。"""
entered.set()
release.wait()
return "first"
first = executor.submit(blocked_work)
queued = executor.submit(lambda: "queued")
try:
assert entered.wait(timeout=1)
assert executor.shutdown_bounded(timeout=0.01) is False
assert not queued.cancelled()
finally:
release.set()
assert first.result(timeout=1) == "first"
assert queued.result(timeout=1) == "queued"
assert executor.shutdown_bounded(timeout=1) is True
def test_thread_helper_shutdown_waits_for_worker_after_future_completion():
"""Future 已完成但用户回调仍阻塞时不得误报 worker 已收敛。"""
executor = thread_module._OwnedThreadPoolExecutor(max_workers=1)
work_release = threading.Event()
callback_entered = threading.Event()
callback_release = threading.Event()
future = executor.submit(lambda: work_release.wait())
def blocked_done_callback(_future) -> None:
"""模拟 Future 终态之后仍占用 worker 的第三方完成回调。"""
callback_entered.set()
callback_release.wait()
future.add_done_callback(blocked_done_callback)
work_release.set()
try:
assert callback_entered.wait(timeout=1)
assert future.done()
started_at = time.monotonic()
assert executor.shutdown_bounded(timeout=0.01) is False
assert time.monotonic() - started_at < 1
finally:
callback_release.set()
assert executor.shutdown_bounded(timeout=1) is True