refactor: unify message queue shutdown convergence

This commit is contained in:
jxxghp
2026-08-24 14:25:59 +08:00
parent 415335b215
commit ad45bfac5d
5 changed files with 123 additions and 13 deletions
+12
View File
@@ -957,6 +957,18 @@ def test_stop_modules_propagates_false_without_skipping_later_cleanup(monkeypatc
_assert_completed_once(dependency)
def test_stop_modules_propagates_message_queue_nonconvergence(monkeypatch):
"""消息队列线程未终止时必须由模块服务关闭结果向上暴露。"""
dependencies = _patch_module_shutdown_dependencies(monkeypatch)
dependencies["stop_message"].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 = []
+53 -2
View File
@@ -1,4 +1,6 @@
import threading
import time
from unittest.mock import MagicMock
from app.application.messaging.message import MessageQueueManager, TemplateHelper, stop_message
from app.foundation.singleton import SingletonClass
@@ -11,20 +13,69 @@ def test_message_queue_stop_wakes_idle_monitor(monkeypatch):
manager.__init__(check_interval=10)
started_at = time.monotonic()
manager.stop()
converged = manager.stop()
elapsed = time.monotonic() - started_at
assert converged is True
assert elapsed < 1
assert not manager.thread.is_alive()
def test_message_queue_stop_reports_blocked_callback_and_supports_retry(monkeypatch):
"""发送回调仍阻塞时应有限返回 False,并保留线程供后续重试。"""
entered = threading.Event()
release = threading.Event()
def blocked_send_callback(*_args, **_kwargs) -> None:
"""模拟无法由消息队列强制取消的同步渠道调用。"""
entered.set()
release.wait()
monkeypatch.setattr(MessageQueueManager, "init_config", lambda self: None)
manager = object.__new__(MessageQueueManager)
manager.__init__(send_callback=blocked_send_callback, check_interval=0)
manager.queue.put({"args": ("payload",), "kwargs": {}})
assert entered.wait(timeout=1)
try:
started_at = time.monotonic()
assert manager.stop(timeout=0.01) is False
assert time.monotonic() - started_at < 1
assert manager.thread.is_alive()
finally:
release.set()
assert manager.stop(timeout=1) is True
assert not manager.thread.is_alive()
def test_stop_message_does_not_initialize_absent_services(monkeypatch):
"""消息服务未初始化时,关闭入口不应为了清理而创建后台资源"""
monkeypatch.setattr(SingletonClass, "_instances", {})
assert MessageQueueManager.get_existing_instance() is None
assert TemplateHelper.get_existing_instance() is None
stop_message()
assert stop_message() is True
assert MessageQueueManager not in SingletonClass._instances
assert TemplateHelper not in SingletonClass._instances
def test_stop_message_aggregates_queue_failure_and_closes_template(monkeypatch):
"""消息队列未收敛时仍应关闭模板缓存并向生命周期返回 False。"""
queue_manager = MagicMock()
queue_manager.stop.return_value = False
template_helper = MagicMock()
monkeypatch.setattr(
SingletonClass,
"_instances",
{
MessageQueueManager: queue_manager,
TemplateHelper: template_helper,
},
)
assert stop_message(timeout=0.25) is False
queue_manager.stop.assert_called_once_with(timeout=0.25)
template_helper.close.assert_called_once_with()