refactor: enforce managed event dispatch

This commit is contained in:
jxxghp
2026-08-24 06:29:56 +08:00
parent 9a0d23857d
commit 46c5eec54e
4 changed files with 22 additions and 35 deletions
+9 -30
View File
@@ -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:
-2
View File
@@ -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,
@@ -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` 仍是不对外公开的宿主内部算法类。
### 总体判断
当前架构总体合理,已经从跨层混合的遗留单体收敛为**边界清晰的模块化单体**:
+2 -2
View File
@@ -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, {})