mirror of
https://github.com/jxxghp/MoviePilot.git
synced 2026-08-29 03:56:43 +08:00
fix: retain agent streaming flush owners
This commit is contained in:
@@ -58,6 +58,7 @@ class StreamingHandler:
|
||||
# 流式输出相关状态
|
||||
self._streaming_enabled = False
|
||||
self._flush_task: Optional[asyncio.Task] = None
|
||||
self._streaming_lifecycle_lock = asyncio.Lock()
|
||||
# 当前消息的发送信息(用于编辑消息)
|
||||
self._message_response: Optional[MessageResponse] = None
|
||||
# 已发送给用户的文本(用于追踪增量)
|
||||
@@ -162,6 +163,28 @@ class StreamingHandler:
|
||||
original_message_id: Optional[str] = None,
|
||||
original_chat_id: Optional[str] = None,
|
||||
title: str = "",
|
||||
):
|
||||
"""串行启动流式输出,禁止新一轮覆盖尚未结束的刷新 owner。"""
|
||||
async with self._streaming_lifecycle_lock:
|
||||
await self._start_streaming(
|
||||
channel=channel,
|
||||
source=source,
|
||||
user_id=user_id,
|
||||
username=username,
|
||||
original_message_id=original_message_id,
|
||||
original_chat_id=original_chat_id,
|
||||
title=title,
|
||||
)
|
||||
|
||||
async def _start_streaming(
|
||||
self,
|
||||
channel: Optional[str] = None,
|
||||
source: Optional[str] = None,
|
||||
user_id: Optional[str] = None,
|
||||
username: Optional[str] = None,
|
||||
original_message_id: Optional[str] = None,
|
||||
original_chat_id: Optional[str] = None,
|
||||
title: str = "",
|
||||
):
|
||||
"""
|
||||
启动流式输出。
|
||||
@@ -175,6 +198,10 @@ class StreamingHandler:
|
||||
:param original_message_id: 原始消息ID(如果是回复消息)
|
||||
:param original_chat_id: 原始聊天ID(如果是回复消息)
|
||||
"""
|
||||
if self._flush_task is not None:
|
||||
self._streaming_enabled = False
|
||||
await self._cancel_flush_task()
|
||||
|
||||
self._channel = channel
|
||||
self._source = source
|
||||
self._user_id = user_id
|
||||
@@ -208,6 +235,11 @@ class StreamingHandler:
|
||||
logger.debug("流式输出已启动")
|
||||
|
||||
async def stop_streaming(self) -> Tuple[bool, str]:
|
||||
"""串行停止流式输出,并等待本轮刷新与最终消息收口。"""
|
||||
async with self._streaming_lifecycle_lock:
|
||||
return await self._stop_streaming()
|
||||
|
||||
async def _stop_streaming(self) -> Tuple[bool, str]:
|
||||
"""
|
||||
停止流式输出。执行最后一次刷新确保所有内容都已发送。
|
||||
:return: (all_sent, final_text)
|
||||
|
||||
@@ -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 断线时的上游任务清理;阶段 40 已补齐异步防抖取消的终态所有权;阶段 41 已统一优雅重启兜底线程的唯一所有权;阶段 42 已补齐 Telegram typing 的多实例隔离和终态 owner;阶段 43 已统一 Discord typing 的异步 owner 和 shutdown 收尾;阶段 44 已清除 WebAgent 测试临时事件循环提前关闭产生的 CI 红注解;阶段 45 已统一影视与字幕搜索的请求级逐页任务编排;阶段 46 已收口启动性能门禁的托管 runner 假失败与诊断输出。
|
||||
> 实施进度:阶段 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 已补齐异步防抖取消的终态所有权;阶段 41 已统一优雅重启兜底线程的唯一所有权;阶段 42 已补齐 Telegram typing 的多实例隔离和终态 owner;阶段 43 已统一 Discord typing 的异步 owner 和 shutdown 收尾;阶段 44 已清除 WebAgent 测试临时事件循环提前关闭产生的 CI 红注解;阶段 45 已统一影视与字幕搜索的请求级逐页任务编排;阶段 46 已收口启动性能门禁的托管 runner 假失败与诊断输出;阶段 47 已补齐 Agent 渠道流式刷新任务的重入 owner。
|
||||
|
||||
## 当前复核结论(2026-08-24)
|
||||
|
||||
@@ -495,6 +495,17 @@
|
||||
样本。脚本 CLI 的只读/显式写入边界和启动基线文件均未改变;本阶段不涉及运行时 API、SDK/Compat、
|
||||
插件 ABI 或插件仓。
|
||||
|
||||
### 长期整改阶段 47:Agent 渠道流式刷新重入 owner 收口(2026-08-24)
|
||||
|
||||
- `StreamingHandler.start_streaming()` 原先会直接覆盖已有 `_flush_task`;旧任务若仍在线程池消息发送或编辑
|
||||
中,会继续读取已被新一轮改写的 channel/source/message 状态,且后续 `stop_streaming()` 只能等待新句柄。
|
||||
- start/stop 现在共享实例级异步生命周期锁。重复 start 会先关闭接收并等待旧 flush owner 真实结束,再
|
||||
发布新一轮消息上下文和新 owner;并发 stop 的最终刷新、消息 finalize 与状态重置完成前,新 start 也不能
|
||||
进入。定时刷新本身仍保持原有 0.3 秒节奏和自然完成语义,不在同步发送中途伪装成已取消。
|
||||
- 回归测试覆盖旧 owner 阻塞时重复 start 不覆盖上下文,以及 stop 最终刷新阻塞时新 start 必须等待。
|
||||
`StreamingHandler` 类路径、公开 `start_streaming()` / `stop_streaming()` 参数与返回值、渠道能力、消息格式、
|
||||
SDK/Compat 和 V1/V2/V3 插件合同均未改变;本阶段未修改插件仓。
|
||||
|
||||
### 总体判断
|
||||
|
||||
当前架构总体合理,已经从跨层混合的遗留单体收敛为**边界清晰的模块化单体**:
|
||||
|
||||
@@ -0,0 +1,101 @@
|
||||
"""Agent 渠道流式输出的 start/stop owner 生命周期测试。"""
|
||||
|
||||
import asyncio
|
||||
|
||||
import pytest
|
||||
|
||||
from app.agent.callback import StreamingHandler
|
||||
from app.schemas.types import NotificationChannel
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_repeated_streaming_start_retains_previous_flush_owner() -> None:
|
||||
"""重复启动必须先等待旧 flush owner 结束,再发布新一轮上下文。"""
|
||||
handler = StreamingHandler()
|
||||
handler._can_stream = lambda: True
|
||||
handler._source = "old-source"
|
||||
handler._streaming_enabled = True
|
||||
old_started = asyncio.Event()
|
||||
release_old = asyncio.Event()
|
||||
observed_sources: list[str | None] = []
|
||||
|
||||
async def old_flush_owner() -> None:
|
||||
"""在退出前读取 handler 上下文,暴露被新一轮提前覆盖的风险。"""
|
||||
old_started.set()
|
||||
await release_old.wait()
|
||||
observed_sources.append(handler._source)
|
||||
|
||||
old_task = asyncio.create_task(old_flush_owner(), name="old.flush")
|
||||
handler._flush_task = old_task
|
||||
await old_started.wait()
|
||||
|
||||
restart = asyncio.create_task(
|
||||
handler.start_streaming(
|
||||
channel=NotificationChannel.Feishu.value,
|
||||
source="new-source",
|
||||
),
|
||||
name="agent.streaming.restart",
|
||||
)
|
||||
await asyncio.sleep(0)
|
||||
|
||||
assert restart.done() is False
|
||||
assert handler._source == "old-source"
|
||||
assert handler._flush_task is old_task
|
||||
|
||||
release_old.set()
|
||||
await asyncio.wait_for(restart, timeout=1)
|
||||
new_task = handler._flush_task
|
||||
|
||||
assert old_task.done()
|
||||
assert observed_sources == ["old-source"]
|
||||
assert new_task is not None
|
||||
assert new_task is not old_task
|
||||
assert handler._source == "new-source"
|
||||
|
||||
await asyncio.wait_for(handler.stop_streaming(), timeout=1)
|
||||
assert new_task.done()
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_streaming_start_waits_for_concurrent_stop_final_flush() -> None:
|
||||
"""停止阶段的最终刷新未完成时,新启动不得改写消息上下文。"""
|
||||
handler = StreamingHandler()
|
||||
handler._can_stream = lambda: True
|
||||
handler._source = "stopping-source"
|
||||
handler._streaming_enabled = True
|
||||
flush_started = asyncio.Event()
|
||||
release_flush = asyncio.Event()
|
||||
|
||||
async def blocking_flush() -> None:
|
||||
"""把最终刷新停留在生命周期锁内,验证新启动等待。"""
|
||||
flush_started.set()
|
||||
await release_flush.wait()
|
||||
|
||||
handler._flush = blocking_flush
|
||||
stop_task = asyncio.create_task(
|
||||
handler.stop_streaming(),
|
||||
name="agent.streaming.stop",
|
||||
)
|
||||
await flush_started.wait()
|
||||
start_task = asyncio.create_task(
|
||||
handler.start_streaming(
|
||||
channel=NotificationChannel.Feishu.value,
|
||||
source="next-source",
|
||||
),
|
||||
name="agent.streaming.next",
|
||||
)
|
||||
await asyncio.sleep(0)
|
||||
|
||||
assert start_task.done() is False
|
||||
assert handler._source == "stopping-source"
|
||||
|
||||
release_flush.set()
|
||||
await asyncio.wait_for(stop_task, timeout=1)
|
||||
await asyncio.wait_for(start_task, timeout=1)
|
||||
next_flush_task = handler._flush_task
|
||||
|
||||
assert handler._source == "next-source"
|
||||
assert next_flush_task is not None
|
||||
|
||||
await asyncio.wait_for(handler.stop_streaming(), timeout=1)
|
||||
assert next_flush_task.done()
|
||||
Reference in New Issue
Block a user