mirror of
https://github.com/jxxghp/MoviePilot.git
synced 2026-08-29 03:56:43 +08:00
fix: drain cancelled debounce tasks
This commit is contained in:
+27
-7
@@ -184,9 +184,19 @@ class AsyncDebouncer(BaseDebouncer):
|
||||
"""
|
||||
super().__init__(*args, **kwargs)
|
||||
self.task: Optional[asyncio.Task] = None
|
||||
self._retired_tasks: set[asyncio.Task[Any]] = set()
|
||||
self.lock = asyncio.Lock()
|
||||
self.is_cooling_down = False
|
||||
|
||||
def _cancel_active_task(self) -> None:
|
||||
"""取消当前任务并保留 owner,直到任务进入真实终态。"""
|
||||
task = self.task
|
||||
if task is None or task.done():
|
||||
return
|
||||
task.cancel()
|
||||
self._retired_tasks.add(task)
|
||||
task.add_done_callback(self._retired_tasks.discard)
|
||||
|
||||
async def __call__(self, *args, **kwargs) -> None:
|
||||
"""
|
||||
异步调用防抖函数。
|
||||
@@ -205,8 +215,7 @@ class AsyncDebouncer(BaseDebouncer):
|
||||
self.log_info("前沿模式 (async): 立即执行协程。")
|
||||
await self.func(*args, **kwargs)
|
||||
|
||||
if self.task and not self.task.done():
|
||||
self.task.cancel()
|
||||
self._cancel_active_task()
|
||||
|
||||
self.is_cooling_down = True
|
||||
self.task = asyncio.create_task(self._end_cool_down())
|
||||
@@ -226,7 +235,7 @@ class AsyncDebouncer(BaseDebouncer):
|
||||
后沿模式的逻辑。
|
||||
"""
|
||||
if self.task and not self.task.done():
|
||||
self.task.cancel()
|
||||
self._cancel_active_task()
|
||||
self.log_debug("后沿模式 (async): 检测到新的调用,已取消旧任务。")
|
||||
|
||||
self.task = asyncio.create_task(self._delayed_execute(*args, **kwargs))
|
||||
@@ -246,12 +255,23 @@ class AsyncDebouncer(BaseDebouncer):
|
||||
|
||||
async def cancel(self) -> None:
|
||||
"""
|
||||
取消任何挂起的调用,并重置状态。
|
||||
取消并等待所有挂起调用进入终态,再重置状态。
|
||||
"""
|
||||
async with self.lock:
|
||||
if self.task and not self.task.done():
|
||||
self.task.cancel()
|
||||
self.task = None
|
||||
current_task = asyncio.current_task()
|
||||
tasks = {
|
||||
task
|
||||
for task in (*self._retired_tasks, self.task)
|
||||
if task is not None and task is not current_task
|
||||
}
|
||||
for task in tasks:
|
||||
if not task.done():
|
||||
task.cancel()
|
||||
if tasks:
|
||||
await asyncio.gather(*tasks, return_exceptions=True)
|
||||
self._retired_tasks.difference_update(tasks)
|
||||
self.task = None
|
||||
if tasks:
|
||||
self.log_info("异步防抖器被手动取消。")
|
||||
self.is_cooling_down = False
|
||||
|
||||
|
||||
@@ -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 已统一插件输入事件发布路径;阶段 38 已统一 WebAgent 通知事件监听与队列边界;阶段 39 已补齐搜索 SSE 断线时的上游任务清理。
|
||||
> 实施进度:阶段 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 通知事件监听与队列边界;阶段 39 已补齐搜索 SSE 断线时的上游任务清理;阶段 40 已补齐异步防抖取消的终态所有权。
|
||||
|
||||
## 当前复核结论(2026-08-24)
|
||||
|
||||
@@ -418,6 +418,14 @@
|
||||
- 回归测试锁定客户端在首个事件到达时断开后,上游 `finally` 已在流函数返回前执行;心跳、批量
|
||||
append/replace、字幕签名、SSE payload 与缓存头均未修改,V1/V2/V3 插件合同不受影响。
|
||||
|
||||
### 长期整改阶段 40:异步防抖取消终态收口(2026-08-24)
|
||||
|
||||
- `AsyncDebouncer` 在后续调用替换旧任务时只发出取消并覆盖句柄,公开 `cancel()` 也在任务真正退出前
|
||||
清空引用;现在被替换的任务保留为 retired owner,`cancel()` 会取消并等待活动任务和 retired 任务
|
||||
全部进入终态后再返回。
|
||||
- 回归测试覆盖旧任务在取消清理中阻塞、后续任务同时存在的场景。`debounce()` 装饰器、leading/trailing
|
||||
触发规则、同步 `Debouncer`、`app.utils.debounce` 兼容映射和 V1/V2/V3 插件导入路径均未修改。
|
||||
|
||||
### 总体判断
|
||||
|
||||
当前架构总体合理,已经从跨层混合的遗留单体收敛为**边界清晰的模块化单体**:
|
||||
|
||||
@@ -0,0 +1,59 @@
|
||||
"""同步与异步防抖生命周期测试。"""
|
||||
|
||||
import asyncio
|
||||
|
||||
import pytest
|
||||
|
||||
from app.runtime.debounce import AsyncDebouncer
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_async_debouncer_cancel_waits_for_replaced_task_cleanup() -> None:
|
||||
"""公开取消必须等待已被后续调用替换的任务完成取消清理。"""
|
||||
first_started = asyncio.Event()
|
||||
first_cancelled = asyncio.Event()
|
||||
allow_first_cleanup = asyncio.Event()
|
||||
|
||||
async def delayed_work(value: str) -> None:
|
||||
"""让首个调用在取消后停留,暴露 retired task 的终态时序。"""
|
||||
if value != "first":
|
||||
await asyncio.Event().wait()
|
||||
return
|
||||
first_started.set()
|
||||
try:
|
||||
await asyncio.Event().wait()
|
||||
except asyncio.CancelledError:
|
||||
first_cancelled.set()
|
||||
await allow_first_cleanup.wait()
|
||||
raise
|
||||
|
||||
debouncer = AsyncDebouncer(delayed_work, interval=0)
|
||||
await debouncer("first")
|
||||
await first_started.wait()
|
||||
|
||||
await debouncer("second")
|
||||
await first_cancelled.wait()
|
||||
cancel_task = asyncio.create_task(debouncer.cancel())
|
||||
await asyncio.sleep(0)
|
||||
|
||||
assert cancel_task.done() is False
|
||||
allow_first_cleanup.set()
|
||||
await asyncio.wait_for(cancel_task, timeout=1)
|
||||
assert debouncer.task is None
|
||||
assert debouncer._retired_tasks == set()
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_async_debouncer_cancel_is_idempotent_without_active_task() -> None:
|
||||
"""没有活动任务时重复取消应立即返回并保持冷却状态关闭。"""
|
||||
|
||||
async def no_op() -> None:
|
||||
"""提供不会被实际调度的异步函数。"""
|
||||
|
||||
debouncer = AsyncDebouncer(no_op, interval=1, leading=True)
|
||||
|
||||
await debouncer.cancel()
|
||||
await debouncer.cancel()
|
||||
|
||||
assert debouncer.task is None
|
||||
assert debouncer.is_cooling_down is False
|
||||
Reference in New Issue
Block a user