diff --git a/app/api/endpoints/agent.py b/app/api/endpoints/agent.py index f3b86cf7d..e2e6bcbb4 100644 --- a/app/api/endpoints/agent.py +++ b/app/api/endpoints/agent.py @@ -45,7 +45,6 @@ from app.application.agent import ( from app.chain.message import MessageChain from app.application.commands import get_command, get_commands from app.runtime.config import global_vars -from app.runtime.events import Event, EventManager from app.api.principal import ApiPrincipal from app.api.dependencies.agent import ( get_agent_chat_persistence, @@ -62,9 +61,12 @@ from app.application.messaging.chat import ( from app.application.security.user import get_configured_user_id_lookup from app.application.configuration import get_api_runtime_config_snapshot from app.application.messaging.agent import ( + attach_web_agent_message_queue, attach_web_agent_edit_queue, create_web_agent_background_task, + detach_web_agent_message_queue, detach_web_agent_edit_queue, + is_web_agent_message_for_user, ) from app.application.messaging.agent import agent_interaction_manager from app.application.messaging.agent import ( @@ -94,9 +96,6 @@ WEB_AGENT_STREAM_COALESCE_MAX_CHARS = 256 WEB_AGENT_STREAM_HEARTBEAT_SECONDS = 15.0 WEB_AGENT_STREAM_QUEUE_MAX_SIZE = 64 _WEB_AGENT_FILE_REGISTRY: dict[str, dict[str, Any]] = {} -_WEB_AGENT_MESSAGE_QUEUES: dict[str, list[Queue[_SchemaMessage]]] = {} -_WEB_AGENT_MESSAGE_LOCK = Lock() -_WEB_AGENT_MESSAGE_LISTENER_REGISTERED = False class _WebAgentEventPublisher: @@ -1400,149 +1399,6 @@ def _has_web_agent_traditional_interaction(user_id: str) -> bool: return has_pending_interaction(user_id) -def _extract_web_agent_message_from_event_data( - data: dict, -) -> Optional[_SchemaMessage]: - """ - 从 NoticeMessage 事件数据中提取 WebAgent 通知。 - - :param data: NoticeMessage 事件数据,兼容扁平字段和 message 包装格式 - :return: WebAgent 通知,不属于 WebAgent 或数据无效时返回 None - """ - if not isinstance(data, dict): - return None - - try: - message = data.get("message") - if isinstance(message, _SchemaMessage): - message = message - elif isinstance(message, dict): - message_data = copy.deepcopy(message) - message_data.pop("type", None) - message = _SchemaMessage(**message_data) - else: - message_data = copy.deepcopy(data) - message_data.pop("type", None) - message_data.pop("current_time", None) - message = _SchemaMessage(**message_data) - except Exception as err: - logger.debug(f"解析WebAgent通知事件失败: {err}") - return None - - channel = message.channel - channel_value = channel.value if isinstance(channel, NotificationChannel) else channel - if channel_value != NotificationChannel.WebAgent.value: - return None - return message - - -def _is_web_agent_message_for_user( - message: _SchemaMessage, - user_id: str, -) -> bool: - """ - 判断 NoticeMessage 事件是否属于当前 WebAgent 用户。 - - :param message: NoticeMessage 中的通知消息 - :param user_id: 当前登录用户 ID - :return: 可被本次 WebAgent 请求消费时返回 True - """ - try: - target_user = message.userid - return target_user is None or str(target_user) == str(user_id) - except Exception: - return False - - -def _get_web_agent_message_user_id(message: _SchemaMessage) -> Optional[str]: - """ - 从 NoticeMessage 事件中解析 WebAgent 目标用户。 - - :param message: NoticeMessage 中的通知消息 - :return: 用户 ID 字符串,事件不属于 WebAgent 时返回 None - """ - try: - channel = message.channel - channel_value = channel.value if isinstance(channel, NotificationChannel) else channel - if channel_value != NotificationChannel.WebAgent.value: - return None - user_id = message.userid - return str(user_id) if user_id is not None else None - except Exception: - return None - - -def _dispatch_web_agent_message_event(event: Event) -> None: - """ - 将 WebAgent NoticeMessage 分发给正在等待的请求队列。 - - :param event: NoticeMessage 广播事件 - """ - data = event.event_data if isinstance(event.event_data, dict) else {} - message = _extract_web_agent_message_from_event_data(data) - if not message: - return - with _WEB_AGENT_MESSAGE_LOCK: - user_id = _get_web_agent_message_user_id(message) - if user_id is None: - queues = [ - message_queue - for user_queues in _WEB_AGENT_MESSAGE_QUEUES.values() - for message_queue in user_queues - ] - else: - queues = list(_WEB_AGENT_MESSAGE_QUEUES.get(user_id) or []) - for message_queue in queues: - message_queue.put(message) - - -def _ensure_web_agent_message_listener() -> None: - """ - 确保 WebAgent NoticeMessage 全局监听器已注册。 - """ - global _WEB_AGENT_MESSAGE_LISTENER_REGISTERED - if _WEB_AGENT_MESSAGE_LISTENER_REGISTERED: - return - with _WEB_AGENT_MESSAGE_LOCK: - if _WEB_AGENT_MESSAGE_LISTENER_REGISTERED: - return - EventManager().add_event_listener( - EventType.NoticeMessage, - _dispatch_web_agent_message_event, - ) - _WEB_AGENT_MESSAGE_LISTENER_REGISTERED = True - - -def _attach_web_agent_message_queue(user_id: str, message_queue: Queue[_SchemaMessage]) -> None: - """ - 为当前 WebAgent 请求挂载通知收集队列。 - - :param user_id: 当前用户 ID - :param message_queue: 用于接收通知事件的队列 - """ - _ensure_web_agent_message_listener() - with _WEB_AGENT_MESSAGE_LOCK: - _WEB_AGENT_MESSAGE_QUEUES.setdefault(str(user_id), []).append(message_queue) - - -def _detach_web_agent_message_queue(user_id: str, message_queue: Queue[_SchemaMessage]) -> None: - """ - 移除当前 WebAgent 请求的通知收集队列。 - - :param user_id: 当前用户 ID - :param message_queue: 需要移除的队列 - """ - with _WEB_AGENT_MESSAGE_LOCK: - queues = _WEB_AGENT_MESSAGE_QUEUES.get(str(user_id)) - if not queues: - return - _WEB_AGENT_MESSAGE_QUEUES[str(user_id)] = [ - item for item in queues if item is not message_queue - ] - if not _WEB_AGENT_MESSAGE_QUEUES[str(user_id)]: - _WEB_AGENT_MESSAGE_QUEUES.pop(str(user_id), None) - - def _build_web_agent_command_items() -> list[dict]: """ 读取当前可用斜杠命令并转换为前端建议列表。 @@ -1629,7 +1485,7 @@ async def _collect_web_agent_traditional_events( edit_queue: Queue[dict] = Queue() user_id = str(current_user.id) - _attach_web_agent_message_queue(user_id, message_queue) + attach_web_agent_message_queue(user_id, message_queue) attach_web_agent_edit_queue(user_id, edit_queue) try: await run_in_threadpool( @@ -1668,13 +1524,13 @@ async def _collect_web_agent_traditional_events( break continue - if not _is_web_agent_message_for_user(message, user_id): + if not is_web_agent_message_for_user(message, user_id): continue events.extend(await _build_web_agent_message_events_async(message)) idle_deadline = time.monotonic() + WEB_AGENT_TRADITIONAL_IDLE_TIMEOUT_SECONDS return events finally: - _detach_web_agent_message_queue(user_id, message_queue) + detach_web_agent_message_queue(user_id, message_queue) detach_web_agent_edit_queue(user_id, edit_queue) diff --git a/app/application/messaging/agent.py b/app/application/messaging/agent.py index e6eb67e4b..3b816edae 100644 --- a/app/application/messaging/agent.py +++ b/app/application/messaging/agent.py @@ -1,4 +1,5 @@ import asyncio +import copy import uuid from dataclasses import dataclass, field from datetime import datetime, timedelta @@ -6,6 +7,8 @@ from queue import Queue from threading import Lock from typing import Awaitable, Callable, Dict, Iterable, List, Optional, Tuple, Union +from app.runtime.log import logger +from app.schemas.message import Message from app.schemas.types import NotificationChannel from app.runtime.tasks import get_task_registry @@ -177,6 +180,8 @@ agent_interaction_manager = AgentInteractionManager() _WEB_AGENT_EDIT_QUEUES: dict[str, list[Queue[dict]]] = {} _WEB_AGENT_EDIT_LOCK = Lock() +_WEB_AGENT_MESSAGE_QUEUES: dict[str, list[Queue[Message]]] = {} +_WEB_AGENT_MESSAGE_LOCK = Lock() _ChannelAdminResolver = Callable[[Optional[dict]], Iterable[Union[str, int]]] _CHANNEL_ADMIN_RESOLVERS: dict[str, _ChannelAdminResolver] = {} _WEB_AGENT_BACKGROUND_TASKS: set[asyncio.Task[object]] = set() @@ -371,6 +376,108 @@ def build_web_agent_message_update_event( } +def extract_web_agent_message_from_event_data(data: dict) -> Optional[Message]: + """ + 从 NoticeMessage 事件数据中提取 WebAgent 通知。 + + :param data: NoticeMessage 事件数据,兼容扁平字段和 message 包装格式 + :return: WebAgent 通知,不属于 WebAgent 或数据无效时返回 None + """ + if not isinstance(data, dict): + return None + + try: + message = data.get("message") + if isinstance(message, Message): + message = message + elif isinstance(message, dict): + message_data = copy.deepcopy(message) + message_data.pop("type", None) + message = Message(**message_data) + else: + message_data = copy.deepcopy(data) + message_data.pop("type", None) + message_data.pop("current_time", None) + message = Message(**message_data) + except Exception as err: + logger.debug(f"解析WebAgent通知事件失败: {err}") + return None + + channel = message.channel + channel_value = channel.value if isinstance(channel, NotificationChannel) else channel + if channel_value != NotificationChannel.WebAgent.value: + return None + return message + + +def is_web_agent_message_for_user(message: Message, user_id: str) -> bool: + """ + 判断 NoticeMessage 事件是否属于当前 WebAgent 用户。 + + :param message: NoticeMessage 中的通知消息 + :param user_id: 当前登录用户 ID + :return: 可被本次 WebAgent 请求消费时返回 True + """ + try: + target_user = message.userid + return target_user is None or str(target_user) == str(user_id) + except Exception: + return False + + +def _get_web_agent_message_user_id(message: Message) -> Optional[str]: + """返回 WebAgent 通知的目标用户 ID,无目标时返回 None。""" + try: + channel = message.channel + channel_value = channel.value if isinstance(channel, NotificationChannel) else channel + if channel_value != NotificationChannel.WebAgent.value: + return None + user_id = message.userid + return str(user_id) if user_id is not None else None + except Exception: + return None + + +def dispatch_web_agent_message_event(event: object) -> None: + """将 WebAgent NoticeMessage 分发给正在等待的请求队列。""" + event_data = getattr(event, "event_data", None) + data = event_data if isinstance(event_data, dict) else {} + message = extract_web_agent_message_from_event_data(data) + if not message: + return + with _WEB_AGENT_MESSAGE_LOCK: + user_id = _get_web_agent_message_user_id(message) + if user_id is None: + queues = [ + message_queue + for user_queues in _WEB_AGENT_MESSAGE_QUEUES.values() + for message_queue in user_queues + ] + else: + queues = list(_WEB_AGENT_MESSAGE_QUEUES.get(user_id) or []) + for message_queue in queues: + message_queue.put(message) + + +def attach_web_agent_message_queue(user_id: str, message_queue: Queue[Message]) -> None: + """为当前 WebAgent 请求挂载通知收集队列。""" + with _WEB_AGENT_MESSAGE_LOCK: + _WEB_AGENT_MESSAGE_QUEUES.setdefault(str(user_id), []).append(message_queue) + + +def detach_web_agent_message_queue(user_id: str, message_queue: Queue[Message]) -> None: + """移除当前 WebAgent 请求的通知收集队列。""" + with _WEB_AGENT_MESSAGE_LOCK: + queues = _WEB_AGENT_MESSAGE_QUEUES.get(str(user_id)) + if not queues: + return + _WEB_AGENT_MESSAGE_QUEUES[str(user_id)] = [ + item for item in queues if item is not message_queue + ] + if not _WEB_AGENT_MESSAGE_QUEUES[str(user_id)]: + _WEB_AGENT_MESSAGE_QUEUES.pop(str(user_id), None) + + def attach_web_agent_edit_queue(user_id: str, edit_queue: Queue[dict]) -> None: """ 为当前 WebAgent 请求挂载原消息编辑事件队列。 diff --git a/app/startup/initializers/modules.py b/app/startup/initializers/modules.py index e69aa4680..bad01bf15 100644 --- a/app/startup/initializers/modules.py +++ b/app/startup/initializers/modules.py @@ -71,6 +71,7 @@ from app.application.messaging.chat import ( get_configured_agent_chat_persistence, ) from app.application.messaging.agent import ( + dispatch_web_agent_message_event, shutdown_web_agent_background_tasks, wait_web_agent_background_tasks, ) @@ -855,6 +856,11 @@ async def init_modules() -> HostRuntime: user_auth() # 事件错误通知由启动组合层接入消息服务。 EventManager().set_error_notifier(notify_event_error) + # WebAgent 事件监听由组合根统一装配,HTTP 请求只管理自己的队列。 + EventManager().add_event_listener( + EventType.NoticeMessage, + dispatch_web_agent_message_event, + ) # 宿主类处理器在启动层显式登记,事件总线不再兜底 owner_class()。 configure_host_event_handler_resolver() # 加载模块 diff --git a/docs/refactor/backend-architecture-next-stage.md b/docs/refactor/backend-architecture-next-stage.md index 60f3dd05a..5d309653e 100644 --- a/docs/refactor/backend-architecture-next-stage.md +++ b/docs/refactor/backend-architecture-next-stage.md @@ -6,7 +6,7 @@ > 审计范围:宿主后端;排除 `app/plugins/**` 运行时插件副本 > 规范优先级:`AGENTS.md` 与 `docs/rules/` 高于本文 > 相关文档:`docs/architecture-overview.md`、`docs/refactor/backend-architecture-governance.md`、`docs/refactor/backend-module-refactor-compatibility.md` -> 实施进度:阶段 0~6 的宿主架构能力已完成收口;API/Application 公共复杂度基线已清零,启动组合根的 SystemConfigOper 构造点已由 14 降至 1;API 进程内后台任务已完成首批统一登记,插件仓适配和 Outbox 外围扩展仍按风险切片推进。Model/Base 查询与写装饰器、legacy 隐式会话外壳均已清零,插件 SDK 也不再导出宿主 Model。2026-08-23 的长期整改阶段 0 已恢复宿主、启动性能、官方插件和 SDK 契约门禁的可信基线;阶段 1a 已补齐 TaskRegistry owner 零债务门禁和诚实的关停超时语义;阶段 1b1 已收口整理 worker、pending 回放、失败通知、进程内 AI 重试、插件监控与事件投递的生命周期所有权;2026-08-24 的阶段 2 已将 212 个已观察宿主模块方法的 legacy aggregation 清零,并补齐可执行 fanout 与下载器文件 DTO 边界;阶段 3 已将消息交互和远程命令的订阅删除统一到 Application/UoW/outbox,宿主不再调用裸线程统计入口;阶段 4 已统一七种消息渠道的宿主回环与后台执行边界;阶段 5 已补齐事件窗口聚合任务的生命周期所有权;阶段 6 已统一插件文件操作的取消完成语义;阶段 7 已统一插件协程补偿的终态等待;阶段 8 已统一宿主同步函数的异步线程池入口;阶段 9 已统一工作流运行时的宿主获取路径;阶段 10 已统一模块、插件与调度运行时的显式 getter 调用;阶段 11 已清除系统配置 getter 的 Oper 形别名;阶段 12 已完成工作流域的显式 Chain 数据端口迁移;阶段 13 已收口用户、交互与消息链的数据端口;阶段 14 已收口音乐订阅数据端口;阶段 15 已收口站点数据端口;阶段 16 已收口媒体服务器数据端口;阶段 17 已收口下载数据端口;阶段 18 已收口主订阅数据端口;阶段 19 已收口整理数据端口;阶段 20 已收口 Agent 数据端口;阶段 21 已收口监控历史端口;阶段 22 已统一服务配置应用边界;阶段 23 已补齐媒体服务器 API 遗留的类形配置读取路径;阶段 24 已清除 Scheduler 内部无 owner 的协程提交双轨;阶段 25 已补齐 TaskRegistry 跨线程 owner 并迁移整理 AI 接管;阶段 26 已统一 Agent 会话清理提交;阶段 27 已统一历史 AI 进度 owner;阶段 28 已托管旧插件订阅统计线程;阶段 29 已统一 Emby 系条目转换并清零重复代码白名单;阶段 30 已收口插件市场请求级子任务;阶段 31 已托管搜索 AI 推荐任务;阶段 32 已清除事件调度器绕过生命周期 owner 的投递回退;阶段 33 已统一宿主 Agent 运行时的获取路径;阶段 34 已统一 durable-required 事件与 Outbox topic 事实源;阶段 35 已统一 LLM provider 管理 API 的运行时解析路径;阶段 36 已统一 WebAgent 音频能力访问边界;阶段 37 已统一插件输入事件发布路径。 +> 实施进度:阶段 0~6 的宿主架构能力已完成收口;API/Application 公共复杂度基线已清零,启动组合根的 SystemConfigOper 构造点已由 14 降至 1;API 进程内后台任务已完成首批统一登记,插件仓适配和 Outbox 外围扩展仍按风险切片推进。Model/Base 查询与写装饰器、legacy 隐式会话外壳均已清零,插件 SDK 也不再导出宿主 Model。2026-08-23 的长期整改阶段 0 已恢复宿主、启动性能、官方插件和 SDK 契约门禁的可信基线;阶段 1a 已补齐 TaskRegistry owner 零债务门禁和诚实的关停超时语义;阶段 1b1 已收口整理 worker、pending 回放、失败通知、进程内 AI 重试、插件监控与事件投递的生命周期所有权;2026-08-24 的阶段 2 已将 212 个已观察宿主模块方法的 legacy aggregation 清零,并补齐可执行 fanout 与下载器文件 DTO 边界;阶段 3 已将消息交互和远程命令的订阅删除统一到 Application/UoW/outbox,宿主不再调用裸线程统计入口;阶段 4 已统一七种消息渠道的宿主回环与后台执行边界;阶段 5 已补齐事件窗口聚合任务的生命周期所有权;阶段 6 已统一插件文件操作的取消完成语义;阶段 7 已统一插件协程补偿的终态等待;阶段 8 已统一宿主同步函数的异步线程池入口;阶段 9 已统一工作流运行时的宿主获取路径;阶段 10 已统一模块、插件与调度运行时的显式 getter 调用;阶段 11 已清除系统配置 getter 的 Oper 形别名;阶段 12 已完成工作流域的显式 Chain 数据端口迁移;阶段 13 已收口用户、交互与消息链的数据端口;阶段 14 已收口音乐订阅数据端口;阶段 15 已收口站点数据端口;阶段 16 已收口媒体服务器数据端口;阶段 17 已收口下载数据端口;阶段 18 已收口主订阅数据端口;阶段 19 已收口整理数据端口;阶段 20 已收口 Agent 数据端口;阶段 21 已收口监控历史端口;阶段 22 已统一服务配置应用边界;阶段 23 已补齐媒体服务器 API 遗留的类形配置读取路径;阶段 24 已清除 Scheduler 内部无 owner 的协程提交双轨;阶段 25 已补齐 TaskRegistry 跨线程 owner 并迁移整理 AI 接管;阶段 26 已统一 Agent 会话清理提交;阶段 27 已统一历史 AI 进度 owner;阶段 28 已托管旧插件订阅统计线程;阶段 29 已统一 Emby 系条目转换并清零重复代码白名单;阶段 30 已收口插件市场请求级子任务;阶段 31 已托管搜索 AI 推荐任务;阶段 32 已清除事件调度器绕过生命周期 owner 的投递回退;阶段 33 已统一宿主 Agent 运行时的获取路径;阶段 34 已统一 durable-required 事件与 Outbox topic 事实源;阶段 35 已统一 LLM provider 管理 API 的运行时解析路径;阶段 36 已统一 WebAgent 音频能力访问边界;阶段 37 已统一插件输入事件发布路径;阶段 38 已统一 WebAgent 通知事件监听与队列边界。 ## 当前复核结论(2026-08-24) @@ -398,13 +398,25 @@ - 宿主依赖图模块仍为 `806`,内部边由 `6547` 降至 `6545`;事件 producer 合同仍为 `78`,12 组禁止边 与唯一隔离 TMDB SCC 均未变化。 +### 长期整改阶段 38:WebAgent 通知事件监听边界收口(2026-08-24) + +- WebAgent 原消息编辑队列已经属于 `app.application.messaging.agent`,但 NoticeMessage 的解析、通知队列 + 和进程级监听仍位于 HTTP endpoint,并由首个请求惰性注册,形成同一请求的两套队列归属和注册时机。 +- 通知解析、按用户路由及队列挂载/释放现与编辑队列统一归入 messaging Application;startup 组合根在 + 事件消费启动前登记唯一 NoticeMessage listener。API 只持有当前请求队列,不再构造 EventManager 或 + 注册进程监听;通知 schema、用户广播语义、SSE 事件顺序和传统命令等待窗口均未变化。 +- 架构门禁拒绝 HTTP endpoint 新增 `register`/`add_event_listener`,V1/V2/V3 插件事件类型和 payload + 保持不变;事件合同仍跟踪同一个 NoticeMessage consumer,只把 owner 从 API 更正为 startup。 +- 宿主依赖图模块仍为 `806`,内部边由 `6545` 调整为 `6546`:删除 API 到事件总线的隐式边,并由 + Application 显式持有消息 schema 与诊断日志依赖;12 组禁止边与唯一隔离 TMDB SCC 均未变化。 + ### 总体判断 当前架构总体合理,已经从跨层混合的遗留单体收敛为**边界清晰的模块化单体**: - 继续采用单进程控制面是正确选择,不建议现在拆成微服务;插件、调度器、工作流、事件和数据库共享进程内状态,拆分会放大部署、事务和兼容成本。 - `foundation/domain/runtime/adapters/application/chain/api/startup` 的职责方向基本成立;宿主架构基线、复杂度 ratchet、异步阻塞 ratchet 当前均通过。 -- 依赖图当前为 `806` 个 Python 模块、`6545` 条内部导入边;唯一非平凡 SCC 位于隔离的 TMDB 第三方移植包内部,不应为了指标归零重写。 +- 依赖图当前为 `806` 个 Python 模块、`6546` 条内部导入边;唯一非平凡 SCC 位于隔离的 TMDB 第三方移植包内部,不应为了指标归零重写。 - 当前主要风险已经从“目录和依赖失控”转移到运行时协议、后台副作用的可靠性和遗留兼容面。换言之,下一阶段重点应是**语义收口和可验证性**,而不是继续搬文件或机械拆大文件。 综合评价:架构方向可持续,生产可用性较高;可演进性仍处于中等水平。现阶段没有静态审计发现必须立即推倒重来的 P0 架构问题,但存在需要按 P1/P2 计划治理的真实债务。 diff --git a/tests/fixtures/architecture/dependency-baseline.json b/tests/fixtures/architecture/dependency-baseline.json index 6a8fcddda..b12c8f0f3 100644 --- a/tests/fixtures/architecture/dependency-baseline.json +++ b/tests/fixtures/architecture/dependency-baseline.json @@ -13,8 +13,8 @@ "runtime_to_db": [], "workflow_to_db": [] }, - "edge_count": 6545, - "edge_sha256": "e74daad54d04652df3946334297eac0dd9a4560ca87f84365b25dc9f553436da", + "edge_count": 6546, + "edge_sha256": "9799edea47a44d52ae4d5d4bf3e38614b1d1a255f850891b3cc1d8fb85b40735", "edges": [ "app -> app.runtime", "app -> app.runtime.compat", @@ -1674,7 +1674,6 @@ "app.api.endpoints.agent -> app.chain.message", "app.api.endpoints.agent -> app.runtime", "app.api.endpoints.agent -> app.runtime.config", - "app.api.endpoints.agent -> app.runtime.events", "app.api.endpoints.agent -> app.runtime.execution", "app.api.endpoints.agent -> app.runtime.localization", "app.api.endpoints.agent -> app.runtime.log", @@ -2602,8 +2601,10 @@ "app.application.mediaserver -> app.schemas.system", "app.application.mediaserver -> app.schemas.types", "app.application.messaging.agent -> app.runtime", + "app.application.messaging.agent -> app.runtime.log", "app.application.messaging.agent -> app.runtime.tasks", "app.application.messaging.agent -> app.schemas", + "app.application.messaging.agent -> app.schemas.message", "app.application.messaging.agent -> app.schemas.types", "app.application.messaging.chat -> app.application", "app.application.messaging.chat -> app.application.database", diff --git a/tests/fixtures/architecture/runtime-contract-baseline.json b/tests/fixtures/architecture/runtime-contract-baseline.json index b7e04c4f3..8a2f72ab8 100644 --- a/tests/fixtures/architecture/runtime-contract-baseline.json +++ b/tests/fixtures/architecture/runtime-contract-baseline.json @@ -2114,7 +2114,7 @@ "EventType.NoticeMessage": { "consumers": [ { - "caller": "app.api.endpoints.agent", + "caller": "app.startup.initializers.modules", "count": 1 } ], diff --git a/tests/test_architecture_dependencies.py b/tests/test_architecture_dependencies.py index 0e9dfb3ec..36b2941cc 100644 --- a/tests/test_architecture_dependencies.py +++ b/tests/test_architecture_dependencies.py @@ -1232,6 +1232,25 @@ def test_application_services_do_not_resolve_event_manager_singleton(): assert violations == {} +def test_http_endpoints_do_not_register_process_event_listeners(): + """HTTP 端点不得拥有进程级事件监听器,监听装配必须留在 startup。""" + violations: dict[str, set[str]] = {} + for module_name, path in _discover_modules().items(): + if not module_name.startswith("app.api"): + continue + tree = ast.parse(path.read_text(encoding="utf-8"), filename=str(path)) + calls = { + node.func.attr + for node in ast.walk(tree) + if isinstance(node, ast.Call) + and isinstance(node.func, ast.Attribute) + and node.func.attr in {"register", "add_event_listener"} + } + if calls: + violations[module_name] = calls + assert violations == {} + + def test_agent_tools_do_not_import_entrypoint_internals(): """Agent 工具不得穿透导入 HTTP 端点、调度器与命令注册表内部实现。 diff --git a/tests/test_web_agent_stream.py b/tests/test_web_agent_stream.py index 1c718b82e..81d941785 100644 --- a/tests/test_web_agent_stream.py +++ b/tests/test_web_agent_stream.py @@ -13,7 +13,6 @@ from app.agent.orchestrator import agent_manager from app.api.endpoints.agent import ( _WebAgentEventPublisher, _WEB_AGENT_FILE_REGISTRY, - _WEB_AGENT_MESSAGE_QUEUES, _apply_web_agent_display_event, _build_web_agent_input_attachments, _build_web_agent_message_events, @@ -23,8 +22,6 @@ from app.api.endpoints.agent import ( _build_web_agent_traditional_callback_payload, _build_web_agent_display_message_from_events, _collect_web_agent_traditional_events, - _dispatch_web_agent_message_event, - _extract_web_agent_message_from_event_data, _get_web_agent_type, _has_web_agent_traditional_interaction, _prepare_web_agent_audio_attachment_path_async, @@ -38,8 +35,15 @@ from app.runtime.events import Event from app.db.oper.agentchat import AgentChatOper from app.db.models.agentchat import AgentChat from app.application.messaging.chat import AgentChatService, configure_agent_chat_service -from app.application.messaging.agent import build_web_agent_message_update_event -from app.application.messaging.agent import AgentInteractionOption, agent_interaction_manager +from app.application.messaging.agent import ( + AgentInteractionOption, + agent_interaction_manager, + attach_web_agent_message_queue, + build_web_agent_message_update_event, + detach_web_agent_message_queue, + dispatch_web_agent_message_event, + extract_web_agent_message_from_event_data, +) from app.application.messaging.skill import skill_interaction_manager from app.chain.message import MessageChain from app.schemas.notification import ChannelCapability, ChannelCapabilityManager @@ -561,7 +565,7 @@ def test_extract_web_agent_message_supports_wrapped_message_event(): userid="1", ) - extracted = _extract_web_agent_message_from_event_data( + extracted = extract_web_agent_message_from_event_data( {"message": message, "current_time": "2026-06-26 09:18:38"} ) @@ -571,7 +575,7 @@ def test_extract_web_agent_message_supports_wrapped_message_event(): def test_dispatch_web_agent_message_event_accepts_wrapped_message_event(): """WebAgent 等待队列应接收 message 包装格式的 NoticeMessage 事件。""" notice_queue = Queue() - _WEB_AGENT_MESSAGE_QUEUES["1"] = [notice_queue] + attach_web_agent_message_queue("1", notice_queue) message = schemas.Message( channel=NotificationChannel.WebAgent, source="web-agent", @@ -580,14 +584,14 @@ def test_dispatch_web_agent_message_event_accepts_wrapped_message_event(): ) try: - _dispatch_web_agent_message_event( + dispatch_web_agent_message_event( Event( EventType.NoticeMessage, {"message": message, "current_time": "2026-06-26 09:18:38"}, ) ) finally: - _WEB_AGENT_MESSAGE_QUEUES.pop("1", None) + detach_web_agent_message_queue("1", notice_queue) assert notice_queue.get_nowait() == message