diff --git a/app/runtime/event/dispatch.py b/app/runtime/event/dispatch.py index 2cd87e81f..4a2750eeb 100644 --- a/app/runtime/event/dispatch.py +++ b/app/runtime/event/dispatch.py @@ -2,7 +2,6 @@ from __future__ import annotations -import asyncio import inspect import time from collections.abc import Callable @@ -25,20 +24,17 @@ class EventDispatcher: *, registry: EventRegistry, binding_resolver: EventBindingResolver, - executor: Callable[[], Any], - event_loop: Callable[[], Any], event_factory: Callable[..., Any], error_handler: Callable[..., None], - async_handle_sink: Callable[[Any], bool] | None = None, - sync_handle_sink: ( - Callable[[Callable[..., Any], tuple[Any, ...]], bool] | None - ) = None, + async_handle_sink: Callable[[Any], bool], + sync_handle_sink: Callable[ + [Callable[..., Any], tuple[Any, ...]], + bool, + ], ) -> None: - """注入注册表、绑定器、执行器和错误策略回调。""" + """注入注册表、绑定器、生命周期提交器和错误策略回调。""" self._registry = registry self._binding_resolver = binding_resolver - self._executor = executor - self._event_loop = event_loop self._event_factory = event_factory self._error_handler = error_handler self._async_handle_sink = async_handle_sink @@ -126,28 +122,11 @@ class EventDispatcher: ) if inspect.iscoroutinefunction(handler): coroutine = self.safe_invoke_async(handler, isolated) - if self._async_handle_sink: - self._async_handle_sink(coroutine) - continue - try: - asyncio.run_coroutine_threadsafe(coroutine, self._event_loop()) - except RuntimeError: - coroutine.close() - logger.warning( - "事件 %s 的异步处理器无法投递,事件循环已停止", - event.event_type, - ) + self._async_handle_sink(coroutine) else: - if self._sync_handle_sink: - self._sync_handle_sink( - self.safe_invoke_sync, - (handler, isolated), - ) - continue - self._executor().submit( + self._sync_handle_sink( self.safe_invoke_sync, - handler, - isolated, + (handler, isolated), ) def safe_invoke_sync(self, handler: Callable, event: Any) -> None: diff --git a/app/runtime/events.py b/app/runtime/events.py index bee558470..fef91b56f 100644 --- a/app/runtime/events.py +++ b/app/runtime/events.py @@ -156,8 +156,6 @@ class EventManager(metaclass=Singleton): self.__dispatcher = EventDispatcher( registry=self.__registry, binding_resolver=self.__binding_resolver, - executor=lambda: self.__executor, - event_loop=lambda: global_vars.loop, event_factory=Event, error_handler=lambda **kwargs: self.__handle_event_error(**kwargs), async_handle_sink=self.__register_async_handle, diff --git a/docs/refactor/backend-architecture-next-stage.md b/docs/refactor/backend-architecture-next-stage.md index 7c6fc10fe..55dfbbe6c 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 推荐任务。 +> 实施进度:阶段 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 的投递回退。 ## 当前复核结论(2026-08-24) @@ -331,6 +331,16 @@ - 依赖基线仅新增 `app.chain.search -> app.runtime.tasks`,模块保持 `806`、内部边为 `6546`,12 组禁止边 与唯一隔离 TMDB SCC 均未变化。 +### 长期整改阶段 32:事件投递所有权单路径(2026-08-24) + +- `EventManager` 的正式组合早已向 `EventDispatcher` 注入同步和异步 handle sink,但调度器仍 + 允许不注入 sink,并直接调用线程池或 `run_coroutine_threadsafe()`;该回退无法进入事件总线 + 的 owner 句柄表,因而绕过 stop/drain 的取消、等待和超时诊断。 +- 调度器现在把两类 sink 收紧为必需的生命周期依赖,并删除自行提交的第二套实现;所有广播 + handler 都必须先被 `EventManager` 登记,停止状态下由同一个 sink 拒绝并关闭协程。 +- 事件类型、payload、优先级、同步/异步 handler 签名、插件监听注册和 SDK/Compat 映射均未修改; + `EventDispatcher` 仍是不对外公开的宿主内部算法类。 + ### 总体判断 当前架构总体合理,已经从跨层混合的遗留单体收敛为**边界清晰的模块化单体**: diff --git a/tests/test_correlation.py b/tests/test_correlation.py index ee211a6ff..a9a9f9412 100644 --- a/tests/test_correlation.py +++ b/tests/test_correlation.py @@ -112,10 +112,10 @@ def test_event_dispatch_restores_producer_correlation_id() -> None: dispatcher = EventDispatcher( registry=registry, binding_resolver=resolver, - executor=MagicMock(), - event_loop=MagicMock(), event_factory=Event, error_handler=MagicMock(), + async_handle_sink=MagicMock(), + sync_handle_sink=MagicMock(), ) with correlation_scope("producer-request"): event = Event(EventType.SystemError, {})