From 4575146acd58814cc5d5b5e2c73190609866327e Mon Sep 17 00:00:00 2001 From: jxxghp Date: Mon, 24 Aug 2026 05:25:43 +0800 Subject: [PATCH] refactor: require scheduler task ownership --- app/scheduler.py | 62 ++++++++----------- .../adr/0007-background-action-reliability.md | 3 +- .../backend-architecture-next-stage.md | 11 +++- tests/test_scheduler_lifecycle.py | 9 +++ 4 files changed, 48 insertions(+), 37 deletions(-) diff --git a/app/scheduler.py b/app/scheduler.py index 73ececbb6..160b00657 100644 --- a/app/scheduler.py +++ b/app/scheduler.py @@ -1333,7 +1333,7 @@ class Scheduler(ConfigReloadMixin, metaclass=SingletonClass): self, coro: Any, *, - job_id: Optional[str] = None, + job_id: str, generation: int = 0, on_unstarted_cancel: Optional[Callable[[], None]] = None, ) -> bool: @@ -1343,7 +1343,7 @@ class Scheduler(ConfigReloadMixin, metaclass=SingletonClass): - 仅调用方循环可用:在当前循环排队为独立任务 - 无运行中循环(测试/CLI):新建循环同步执行,确保进度不丢失 - 带有 job 标识的句柄由 Scheduler 自己持有,关闭时可以取消并等待。 + job 标识是所有权键;所有句柄都由 Scheduler 持有,关闭时可以取消并等待。 """ try: running_loop = asyncio.get_running_loop() @@ -1356,42 +1356,34 @@ class Scheduler(ConfigReloadMixin, metaclass=SingletonClass): and not target_loop.is_closed() ) if running_loop and (not target_loop_available or running_loop is target_loop): - if job_id is not None: - with self._lock: - if not self._accepts_handle(job_id, generation): - coro.close() - return False - handle = running_loop.create_task(coro) - registered = self._register_handle( - job_id=job_id, - generation=generation, - loop=running_loop, - handle=handle, - ) - if on_unstarted_cancel: - handle.add_done_callback( - lambda submitted: ( - on_unstarted_cancel() - if submitted.cancelled() - else None - ) - ) - return registered - else: - running_loop.create_task(coro) - return True - elif target_loop_available: - if job_id is not None: - return self._submit_cross_thread( - coro, - target_loop=target_loop, + with self._lock: + if not self._accepts_handle(job_id, generation): + coro.close() + return False + handle = running_loop.create_task(coro) + registered = self._register_handle( job_id=job_id, generation=generation, - on_unstarted_cancel=on_unstarted_cancel, + loop=running_loop, + handle=handle, ) - else: - asyncio.run_coroutine_threadsafe(coro, target_loop) - return True + if on_unstarted_cancel: + handle.add_done_callback( + lambda submitted: ( + on_unstarted_cancel() + if submitted.cancelled() + else None + ) + ) + return registered + elif target_loop_available: + return self._submit_cross_thread( + coro, + target_loop=target_loop, + job_id=job_id, + generation=generation, + on_unstarted_cancel=on_unstarted_cancel, + ) elif self._lifecycle_state in {"stopping", "stopped"}: coro.close() return False diff --git a/docs/adr/0007-background-action-reliability.md b/docs/adr/0007-background-action-reliability.md index 0bf772340..819fb1d17 100644 --- a/docs/adr/0007-background-action-reliability.md +++ b/docs/adr/0007-background-action-reliability.md @@ -93,7 +93,8 @@ Event Contract Registry 是 53 个事件的逐项机器清单。下表按相同 `module.imdb.cache_clear`;同步调用方式和无运行事件循环时的立即清理行为保持不变,宿主关停后不再 接受新的清理任务。 - Scheduler 的协程作业与异步进度收尾由 Scheduler 自有句柄表持有;同步 `start()` / `stop()` ABI 保持, - 生命周期关闭入口等待目标事件循环确认真实收尾,跨线程取消代理不作为任务完成凭据。 + 生命周期关闭入口等待目标事件循环确认真实收尾,跨线程取消代理不作为任务完成凭据。内部事件循环提交 + 必须携带 `job_id` owner,当前循环和跨线程路径均登记句柄,不保留 fire-and-forget 分支。 ### Plugin package mutations diff --git a/docs/refactor/backend-architecture-next-stage.md b/docs/refactor/backend-architecture-next-stage.md index d72b7cfb0..f91fb0d01 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 遗留的类形配置读取路径。 +> 实施进度:阶段 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 的协程提交双轨。 ## 当前复核结论(2026-08-24) @@ -249,6 +249,15 @@ - `MediaServerHelper` 继续服务于运行实例发现和插件 SDK 兼容,类路径、方法、配置 Schema、路由及响应 均未修改,V2/V3 插件无需迁移。 +### 长期整改阶段 24:Scheduler 内部协程所有权收口(2026-08-24) + +- `_submit_to_loop()` 的三个生产调用点都已携带 job owner,但私有签名仍允许省略 `job_id`,并为当前循环 + 与跨线程主循环各保留一条不登记句柄的 fire-and-forget 分支;现将 owner 收紧为必填并删除双轨。 +- 进度更新、同步作业收尾和协程作业继续按 generation 登记同一 Scheduler 句柄表;停止接收后拒绝提交, + shutdown 取消并等待真实终态,跨线程代理不被误认作完成。 +- `Scheduler.start()`、同步 `stop()`、作业定义、插件调度方法、进度 payload 和 SDK/Compat 均未修改, + V2/V3 插件行为保持兼容。 + ### 总体判断 当前架构总体合理,已经从跨层混合的遗留单体收敛为**边界清晰的模块化单体**: diff --git a/tests/test_scheduler_lifecycle.py b/tests/test_scheduler_lifecycle.py index 41f16a0ea..09f8e2909 100644 --- a/tests/test_scheduler_lifecycle.py +++ b/tests/test_scheduler_lifecycle.py @@ -2,6 +2,7 @@ import asyncio import gc +import inspect import threading import warnings @@ -69,6 +70,14 @@ def _scheduler(job_id: str, func) -> Scheduler: return scheduler +def test_internal_loop_submission_requires_job_owner() -> None: + """Scheduler 内部协程桥接不得再提供省略 job owner 的游离提交分支。""" + parameter = inspect.signature(Scheduler._submit_to_loop).parameters["job_id"] + + assert parameter.default is inspect.Parameter.empty + assert parameter.annotation is str + + @pytest.mark.anyio async def test_stop_async_cancels_and_awaits_scheduler_owned_job(monkeypatch) -> None: """关闭后已投递协程必须取消并完成收尾,不得遗留 owner 句柄。"""