mirror of
https://github.com/jxxghp/MoviePilot.git
synced 2026-08-29 12:06:51 +08:00
fix: close search streams on disconnect
This commit is contained in:
@@ -256,6 +256,9 @@ async def _iter_batched_search_events(
|
||||
if next_event_task and not next_event_task.done():
|
||||
next_event_task.cancel()
|
||||
await asyncio.gather(next_event_task, return_exceptions=True)
|
||||
close_iterator = getattr(iterator, "aclose", None)
|
||||
if close_iterator is not None:
|
||||
await close_iterator()
|
||||
|
||||
if pending_append_event:
|
||||
yield pending_append_event
|
||||
@@ -274,10 +277,11 @@ async def _stream_search_events(request: Request, event_source: AsyncIterator[di
|
||||
last_event_type = "none"
|
||||
last_stage = "none"
|
||||
termination_reason = "source_exhausted"
|
||||
batched_events = _iter_batched_search_events(event_source)
|
||||
logger.info(f"渐进式搜索流已建立,搜索ID:{search_id},路径:{request_path}")
|
||||
try:
|
||||
has_sent_final_replace = False
|
||||
async for event in _iter_batched_search_events(event_source):
|
||||
async for event in batched_events:
|
||||
last_event_type = event.get("type") or "unknown"
|
||||
last_stage = event.get("stage") or last_stage
|
||||
if await request.is_disconnected():
|
||||
@@ -321,6 +325,7 @@ async def _stream_search_events(request: Request, event_source: AsyncIterator[di
|
||||
transmitted_bytes += len(payload.encode("utf-8"))
|
||||
yield payload
|
||||
finally:
|
||||
await batched_events.aclose()
|
||||
elapsed = time.monotonic() - started_at
|
||||
logger.info(
|
||||
f"渐进式搜索流结束,搜索ID:{search_id},路径:{request_path},"
|
||||
|
||||
@@ -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 通知事件监听与队列边界。
|
||||
> 实施进度:阶段 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 断线时的上游任务清理。
|
||||
|
||||
## 当前复核结论(2026-08-24)
|
||||
|
||||
@@ -410,6 +410,14 @@
|
||||
- 宿主依赖图模块仍为 `806`,内部边由 `6545` 调整为 `6546`:删除 API 到事件总线的隐式边,并由
|
||||
Application 显式持有消息 schema 与诊断日志依赖;12 组禁止边与唯一隔离 TMDB SCC 均未变化。
|
||||
|
||||
### 长期整改阶段 39:搜索 SSE 断线清理收口(2026-08-24)
|
||||
|
||||
- 搜索链的站点并发任务已由各异步生成器在 `finally` 中取消并等待,但 API 批处理和传输包装器在
|
||||
客户端断开时只退出循环,上游关闭依赖异步生成器回收时机;现在两层包装器都显式关闭可关闭的
|
||||
上游迭代器,断线返回前即可触发站点任务清理。
|
||||
- 回归测试锁定客户端在首个事件到达时断开后,上游 `finally` 已在流函数返回前执行;心跳、批量
|
||||
append/replace、字幕签名、SSE payload 与缓存头均未修改,V1/V2/V3 插件合同不受影响。
|
||||
|
||||
### 总体判断
|
||||
|
||||
当前架构总体合理,已经从跨层混合的遗留单体收敛为**边界清晰的模块化单体**:
|
||||
|
||||
@@ -69,6 +69,45 @@ def test_batched_search_events_emit_heartbeat_while_source_is_idle(monkeypatch):
|
||||
assert asyncio.run(_read_heartbeat()) == {"type": "heartbeat"}
|
||||
|
||||
|
||||
def test_search_stream_disconnect_closes_upstream_before_return():
|
||||
"""客户端断开时应在流返回前关闭上游,立即触发搜索任务清理。"""
|
||||
|
||||
async def _consume_disconnected_stream():
|
||||
"""消费一个会立即断开的搜索流,并返回上游关闭状态。"""
|
||||
source_closed = asyncio.Event()
|
||||
|
||||
async def _source():
|
||||
"""输出首个事件后等待包装器显式关闭。"""
|
||||
try:
|
||||
yield {"type": "progress", "stage": "searching", "items": []}
|
||||
finally:
|
||||
source_closed.set()
|
||||
|
||||
async def _disconnected():
|
||||
"""模拟客户端在首个业务事件到达时已经断开。"""
|
||||
return True
|
||||
|
||||
request = SimpleNamespace(
|
||||
url=SimpleNamespace(path="/api/v1/search/title/stream"),
|
||||
headers={},
|
||||
query_params={},
|
||||
is_disconnected=_disconnected,
|
||||
)
|
||||
payloads = [
|
||||
payload
|
||||
async for payload in search_endpoint._stream_search_events(
|
||||
request,
|
||||
_source(),
|
||||
)
|
||||
]
|
||||
return payloads, source_closed.is_set()
|
||||
|
||||
payloads, source_closed = asyncio.run(_consume_disconnected_stream())
|
||||
|
||||
assert payloads == []
|
||||
assert source_closed is True
|
||||
|
||||
|
||||
def test_search_stream_response_disables_proxy_buffering(monkeypatch):
|
||||
"""搜索 SSE 响应应显式禁用缓存和 Nginx 代理缓冲。"""
|
||||
|
||||
|
||||
Reference in New Issue
Block a user