fix(cancel): add database fallback for run cancellation

This commit is contained in:
caoqianming 2026-08-21 13:50:32 +08:00
parent 6c3b8e98b2
commit 0a010c7ac3
5 changed files with 124 additions and 8 deletions

View File

@ -6,6 +6,10 @@
> 开发中的用户文案可先写入 `## Unreleased`;该区不会被前端解析,正式发布时再替换为数字版本和日期。 > 开发中的用户文案可先写入 `## Unreleased`;该区不会被前端解析,正式发布时再替换为数字版本和日期。
> 工程口径的完整记录见 `PROGRESS.md` / git log。 > 工程口径的完整记录见 `PROGRESS.md` / git log。
## Unreleased
- 修复服务更新或多实例切换期间,点击“停止”后对话可能一直停留在“停止中”的问题。
## 0.67.0 — 2026-08-21 ## 0.67.0 — 2026-08-21
- Blender 从单一回转窑模板升级为受管的通用静态三维场景创作,可组合基础体、拉伸/旋转/扫掠、管道、修改器、材质、灯光、文字和最多八个自定义视图;回转窑作为可复用组件继续提供。旧版回转窑工程不再支持在线续作。 - Blender 从单一回转窑模板升级为受管的通用静态三维场景创作,可组合基础体、拉伸/旋转/扫掠、管道、修改器、材质、灯光、文字和最多八个自定义视图;回转窑作为可复用组件继续提供。旧版回转窑工程不再支持在线续作。

View File

@ -165,7 +165,7 @@ Eval 与生产 core 解耦,通过现有 `/v1` API 创建专用任务、监听
**无感部署(蓝绿,0.41)**:生产双 systemd 实例(`ZCBOT_INSTANCE=blue/green`)+ nginx upstream 切流(对外仍 8765),取代单实例 restart 的 503 窗口。双实例并存的互踩由三件事消解:`tasks.run_owner`(0020)让 reaper 只收自己色、sandbox 容器名/label 带实例色各管各的、微信长轮询 PG advisory lock 选主;定时任务 claim 本就 SKIP LOCKED 天然安全。部署编排在 `deploy/update_bluegreen.sh`(流程/兜底见 RUN)。 **无感部署(蓝绿,0.41)**:生产双 systemd 实例(`ZCBOT_INSTANCE=blue/green`)+ nginx upstream 切流(对外仍 8765),取代单实例 restart 的 503 窗口。双实例并存的互踩由三件事消解:`tasks.run_owner`(0020)让 reaper 只收自己色、sandbox 容器名/label 带实例色各管各的、微信长轮询 PG advisory lock 选主;定时任务 claim 本就 SKIP LOCKED 天然安全。部署编排在 `deploy/update_bluegreen.sh`(流程/兜底见 RUN)。
**broker 外置(Redis pub/sub,✅ 0.42)**:0.41 落地时 event/cancel broker 留在进程内,代价是切换窗口内"刷新看不到旧实例 run 直播 / 停止送达不到"两个边缘。2026-07-06 重评(用户量已非个位数)决定实施 —— ①部署时总有 in-flight run,窗口边缘从偶发变常态;②单进程逼近天花板后,最近的扩容手段是稳态双实例同时接流量,broker 外置是它的硬前提(否则 POST 与 SSE 落不同实例直播全瞎)。选 Redis 不选 PG LISTEN/NOTIFY(token 级 delta 全过主库太吵:NOTIFY 全局队列 + 8KB payload 限制,把实时路径耦到 PG)、不选 nginx sticky hash(upstream 变更时 hash 重排,恰在部署窗口失效,治标)。实现(`web/broker.py`,LocalRunBroker/RedisRunBroker 同接口鸭子类型):`ZCBOT_REDIS_URL` env 开关,不设即进程内(dev 零影响);event 走 `PUBLISH zcbot:ev:<tid>` + 单条 async pubsub reader 路由到本地订阅 queue(Redis 只管跨进程一跳,fan-out 最后一公里仍在进程内);done = SETEX key(60s)+ publish 双通道,订阅先挂 channel 再查 key 不漏;cancel = SETEX/GET/DEL key,loop 在 chunk 间 poll(localhost ~0.1ms)。容错纪律:启动 ping 不通 fail-fast(同 sandbox init),运行中失败 10s 节流 log + 降级丢帧,reader 1s 退避重连。Redis 不解决的:run 本体仍绑进程(进程死 run 死,drain/reaper 不变)、线程池上限(调 `ZCBOT_RUN_MAX_WORKERS` 的事)。 **broker 外置(Redis pub/sub,✅ 0.42)**:0.41 落地时 event/cancel broker 留在进程内,代价是切换窗口内"刷新看不到旧实例 run 直播 / 停止送达不到"两个边缘。2026-07-06 重评(用户量已非个位数)决定实施 —— ①部署时总有 in-flight run,窗口边缘从偶发变常态;②单进程逼近天花板后,最近的扩容手段是稳态双实例同时接流量,broker 外置是它的硬前提(否则 POST 与 SSE 落不同实例直播全瞎)。选 Redis 不选 PG LISTEN/NOTIFY(token 级 delta 全过主库太吵:NOTIFY 全局队列 + 8KB payload 限制,把实时路径耦到 PG)、不选 nginx sticky hash(upstream 变更时 hash 重排,恰在部署窗口失效,治标)。实现(`web/broker.py`,LocalRunBroker/RedisRunBroker 同接口鸭子类型):`ZCBOT_REDIS_URL` env 开关,不设即进程内(dev 零影响);event 走 `PUBLISH zcbot:ev:<tid>` + 单条 async pubsub reader 路由到本地订阅 queue(Redis 只管跨进程一跳,fan-out 最后一公里仍在进程内);done = SETEX key(60s)+ publish 双通道,订阅先挂 channel 再查 key 不漏;cancel = SETEX/GET/DEL key,loop 在 chunk 间 poll(localhost ~0.1ms)。取消另以共享 PG 中已经提交的 `tasks.run_status='cancelling'` 作可靠兜底:运行线程每秒低频读取一次,未配置 Redis 的蓝绿窗口或 Redis 运行中故障时仍可跨实例停止PG 不承载 token 直播事件。容错纪律:启动 ping 不通 fail-fast(同 sandbox init),运行中失败 10s 节流 log + 降级丢帧,reader 1s 退避重连。Redis 不解决的:run 本体仍绑进程(进程死 run 死,drain/reaper 不变)、线程池上限(调 `ZCBOT_RUN_MAX_WORKERS` 的事)。
### 7.1 心智模型:Task 一等公民 + Dir 文件副视图 ### 7.1 心智模型:Task 一等公民 + Dir 文件副视图

4
RUN.md
View File

@ -368,7 +368,7 @@ $env:ZCBOT_EVAL_TOKEN = "<dedicated-eval-user-jwt>"
| `GET /v1/tasks/{id}/messages` | LiteLLM payload 透传;另带 `artifact_refs`(助手产物)与 `attachment_refs`(用户附件)。两者均以 `null` 表示旧消息、`[]` 表示新消息明确为空、非空数组表示 task-relative 结构化引用 | 必填 | | `GET /v1/tasks/{id}/messages` | LiteLLM payload 透传;另带 `artifact_refs`(助手产物)与 `attachment_refs`(用户附件)。两者均以 `null` 表示旧消息、`[]` 表示新消息明确为空、非空数组表示 task-relative 结构化引用 | 必填 |
| `POST /v1/tasks/{id}/messages` | `{content, attachments?:[{path,kind,label?}], image_model?=""}` 发消息;`attachments` 路径由后端按当前 task working_dir 校验并规范化,允许纯附件消息;旧客户端省略该字段继续兼容。返 `{events_url,run_id}`,其中 `run_id` 是本轮 user message UUID**`run_status` 是 running/cancelling → 409**UI 应 disable send 直到 SSE `done` | 必填 | | `POST /v1/tasks/{id}/messages` | `{content, attachments?:[{path,kind,label?}], image_model?=""}` 发消息;`attachments` 路径由后端按当前 task working_dir 校验并规范化,允许纯附件消息;旧客户端省略该字段继续兼容。返 `{events_url,run_id}`,其中 `run_id` 是本轮 user message UUID**`run_status` 是 running/cancelling → 409**UI 应 disable send 直到 SSE `done` | 必填 |
| `GET /v1/tasks/{id}/events` | SSE 流(`event: <type>` + `data: <json>`);订阅 task 当前活动 | 必填 | | `GET /v1/tasks/{id}/events` | SSE 流(`event: <type>` + `data: <json>`);订阅 task 当前活动 | 必填 |
| `POST /v1/tasks/{id}/cancel` | 协作式 cancel;`run_status != running` → 409;LLM 走 streaming,chunk 间 poll cancel — 延迟 100ms 级,基本秒退 | 必填 | | `POST /v1/tasks/{id}/cancel` | 协作式 cancel;`run_status != running` → 409;LLM 走 streaming,优先从 broker 轮询取消(约 100ms并每秒读取共享 DB 状态兜底跨实例送达 | 必填 |
| `GET /v1/procs` | 当前用户全部后台进程(bg proc,§8.12;shell/run_python `background=true` 启动);纯文件系统读取,前端运行条 5s 轮询用 | 必填 | | `GET /v1/procs` | 当前用户全部后台进程(bg proc,§8.12;shell/run_python `background=true` 启动);纯文件系统读取,前端运行条 5s 轮询用 | 必填 |
| `POST /v1/tasks/{id}/procs/{proc_id}/kill` | 强制终止后台进程(host 杀进程树 / docker rm 容器);幂等 | 必填 | | `POST /v1/tasks/{id}/procs/{proc_id}/kill` | 强制终止后台进程(host 杀进程树 / docker rm 容器);幂等 | 必填 |
| `POST /v1/tasks/{id}/clear` | 清空当前 task 全部 messages + reset `tasks.tokens_prompt/completion/cost_cny` 三列累计 + `run_status='idle'`;`usage_events`(账单记账)**不动**,只 `message_id` 列变 NULL;run 活跃中(running/cancelling)→ 409(先 cancel);FS 文件保留 | 必填 | | `POST /v1/tasks/{id}/clear` | 清空当前 task 全部 messages + reset `tasks.tokens_prompt/completion/cost_cny` 三列累计 + `run_status='idle'`;`usage_events`(账单记账)**不动**,只 `message_id` 列变 NULL;run 活跃中(running/cancelling)→ 409(先 cancel);FS 文件保留 | 必填 |
@ -1004,7 +1004,7 @@ sudo xfs_quota -x -c "limit -p bhard=10g zcbot_<user_uuid>" /opt
| `POST /v1/tasks/{id}/messages` 返 409 `task already has an active run` | 上一条消息的 BG run 还没跑完;等流式 done 或点 stop / `POST .../cancel`;服务异常下 `run_status``running`/`cancelling`,启动 reaper 会清 | | `POST /v1/tasks/{id}/messages` 返 409 `task already has an active run` | 上一条消息的 BG run 还没跑完;等流式 done 或点 stop / `POST .../cancel`;服务异常下 `run_status``running`/`cancelling`,启动 reaper 会清 |
| `POST /v1/tasks/{id}/cancel` 返 409 `task not running` | `run_status` 不是 `running`(idle / cancelling / error 都不能 cancel);dev SPA 自动忽略不报错 | | `POST /v1/tasks/{id}/cancel` 返 409 `task not running` | `run_status` 不是 `running`(idle / cancelling / error 都不能 cancel);dev SPA 自动忽略不报错 |
| `POST /v1/tasks/{id}/clear` 返 409 `task has an active run` | 当前 run 还没跑完;先点停止 / `POST .../cancel` 等流式 done 再清空 | | `POST /v1/tasks/{id}/clear` 返 409 `task has an active run` | 当前 run 还没跑完;先点停止 / `POST .../cancel` 等流式 done 再清空 |
| 点 stop 后流式没立刻停 | streaming 改造后正常路径秒退;若仍卡可能是 ① httpx 连接 close 没立刻关(GC 时机)/ ② 模型 thinking 阶段长时间不吐 chunk,等下一个 chunk 到达才能 poll cancel(罕见) | | 点 stop 后流式没立刻停 | 正常单实例/Redis broker 路径约 100ms 响应;蓝绿切换或 broker 故障时由共享 DB 每秒轮询兜底。若数秒仍未停止,检查应用日志中的 DB/Redis 错误及任务是否卡在不可协作中断的外部工具调用 |
| `[startup] reaped N stale active run(s)` | 上次 web 进程未正常 finish 留下 N 个孤儿 run,启动 lifespan 自动标 error。info 级,无需处理 | | `[startup] reaped N stale active run(s)` | 上次 web 进程未正常 finish 留下 N 个孤儿 run,启动 lifespan 自动标 error。info 级,无需处理 |
| `seedream` tool 没出现在对话里 | `.env` 没设 `ARK_API_KEY`,build_agent 跳过注册。设了重启 web 即可;无需迁移、无需 DB 改动 | | `seedream` tool 没出现在对话里 | `.env` 没设 `ARK_API_KEY`,build_agent 跳过注册。设了重启 web 即可;无需迁移、无需 DB 改动 |
| 图像模型选了「GPT 生图」却仍走 seedream / 报错 | `.env` 没设 `UNIFYLLM_API_KEY`(gpt_image tool 未注册,静默 fallback 豆包);或服务器没代理出口(直连 unifyllm.ai TLS 失败)。跑 `scripts/diag_unifyllm.py` 验证连通后重启 web | | 图像模型选了「GPT 生图」却仍走 seedream / 报错 | `.env` 没设 `UNIFYLLM_API_KEY`(gpt_image tool 未注册,静默 fallback 豆包);或服务器没代理出口(直连 unifyllm.ai TLS 失败)。跑 `scripts/diag_unifyllm.py` 验证连通后重启 web |

76
tests/test_run_cancel.py Normal file
View File

@ -0,0 +1,76 @@
"""Web run cancellation: broker fast path and shared-DB fallback."""
from __future__ import annotations
import unittest
from contextlib import contextmanager
from types import SimpleNamespace
from unittest.mock import MagicMock, patch
from uuid import uuid4
from sqlalchemy.exc import SQLAlchemyError
class RunCancelCheckTests(unittest.TestCase):
def test_broker_signal_stops_without_db_poll(self) -> None:
from web import runs
tid = uuid4()
broker = MagicMock()
broker.is_cancelled.return_value = True
with (
patch.object(runs, "broker", broker),
patch.object(runs, "session_scope") as scope,
):
check = runs._RunCancelCheck(tid)
self.assertTrue(check())
scope.assert_not_called()
broker.is_cancelled.assert_called_once_with(tid)
def test_db_cancelling_falls_back_across_processes(self) -> None:
from web import runs
tid = uuid4()
broker = MagicMock()
broker.is_cancelled.return_value = False
result = MagicMock()
result.scalar_one_or_none.return_value = "cancelling"
session = SimpleNamespace(execute=MagicMock(return_value=result))
@contextmanager
def fake_scope():
yield session
with (
patch.object(runs, "broker", broker),
patch.object(runs, "session_scope", fake_scope),
patch.object(runs.time, "monotonic", side_effect=[10.0, 10.5, 11.0]),
):
check = runs._RunCancelCheck(tid)
self.assertFalse(check(), "一秒节流窗口内不查询 DB")
self.assertTrue(check(), "共享 DB 的 cancelling 应成为可靠兜底")
session.execute.assert_called_once()
def test_db_poll_failure_does_not_abort_run(self) -> None:
from web import runs
tid = uuid4()
broker = MagicMock()
broker.is_cancelled.return_value = False
@contextmanager
def broken_scope():
raise SQLAlchemyError("db unavailable")
yield
with (
patch.object(runs, "broker", broker),
patch.object(runs, "session_scope", broken_scope),
patch.object(runs.time, "monotonic", side_effect=[20.0, 21.0]),
):
self.assertFalse(runs._RunCancelCheck(tid)())
if __name__ == "__main__":
unittest.main()

View File

@ -7,9 +7,11 @@ web 路由与渠道回调都经它起 run;`run_channel_conversation` 是微信/
from __future__ import annotations from __future__ import annotations
import asyncio import asyncio
import time
from uuid import UUID from uuid import UUID
from sqlalchemy import select, update from sqlalchemy import select, update
from sqlalchemy.exc import SQLAlchemyError
from core.paths import from_db_path from core.paths import from_db_path
from core.storage import session_scope from core.storage import session_scope
@ -21,6 +23,39 @@ from .broker import broker
from .sinks import WebEventSink from .sinks import WebEventSink
class _RunCancelCheck:
"""Fast broker signal with a shared-DB fallback for cross-instance cancel.
Redis remains the low-latency path when configured. The database poll is
deliberately throttled: it covers blue/green overlap and broker outages
without turning the per-stream-chunk cancellation check into DB traffic.
"""
_DB_POLL_SECONDS = 1.0
def __init__(self, task_id: UUID) -> None:
self.task_id = task_id
self._next_db_poll = time.monotonic() + self._DB_POLL_SECONDS
def __call__(self) -> bool:
if broker.is_cancelled(self.task_id):
return True
now = time.monotonic()
if now < self._next_db_poll:
return False
self._next_db_poll = now + self._DB_POLL_SECONDS
try:
with session_scope() as s:
status = s.execute(
select(Task.run_status).where(Task.task_id == self.task_id)
).scalar_one_or_none()
return status == "cancelling"
except SQLAlchemyError:
# DB 短暂不可用不能打断正常生成broker 快路径仍会在后续检查中生效。
return False
def run_agent_bg( def run_agent_bg(
task_id: UUID, user_id: UUID, user_message: str, task_id: UUID, user_id: UUID, user_message: str,
image_variant: str = "", video_variant: str = "", image_variant: str = "", video_variant: str = "",
@ -31,10 +66,11 @@ def run_agent_bg(
sink 通过 broker.emit 桥事件回 asyncio loop;agent.run sync,所以在 to_thread sink 通过 broker.emit 桥事件回 asyncio loop;agent.run sync,所以在 to_thread
user_id 必须从 JWT 那侧透传过来 决定 memory_block 读哪个 per-user 子树 user_id 必须从 JWT 那侧透传过来 决定 memory_block 读哪个 per-user 子树
cancel_check broker.is_cancelled,loop stream chunk + 工具调用之间 poll; cancel_check 优先读 broker( 100ms 响应)并每秒读一次共享 DB cancelling
cancel 延迟 ~ chunk 间隔(100ms );seedance 轮询间也读这个 cancel_check 用于 状态兜底后者保证未配置 Redis 的蓝绿切换窗口 broker 临时故障时停止请求
用户停止按钮(必须在 build_agent 阶段就传进去,因为 SeedanceTool ctor 持有它, 仍能跨实例送达seedance 轮询间也读同一个 cancel_check 用于用户停止按钮(必须在
不能像以前那样 build_agent 返回后再赋 agent.cancel_check) build_agent 阶段就传进去,因为 SeedanceTool ctor 持有它,不能像以前那样
build_agent 返回后再赋 agent.cancel_check)
`ok` 收尾回 `idle`;`cancelled`(用户停止) `error` 一样落持久终态 前端据此 `ok` 收尾回 `idle`;`cancelled`(用户停止) `error` 一样落持久终态 前端据此
已停止/ 错误卡(扛过收尾重渲),下次起新 run(post_message running)覆盖清掉 已停止/ 错误卡(扛过收尾重渲),下次起新 run(post_message running)覆盖清掉
@ -45,7 +81,7 @@ def run_agent_bg(
避免重复落库CLI 等非 Web 入口保持旧行为 避免重复落库CLI 等非 Web 入口保持旧行为
""" """
from core.agent_builder import build_agent, sync_task_tokens from core.agent_builder import build_agent, sync_task_tokens
cancel_check = lambda tid=task_id: broker.is_cancelled(tid) cancel_check = _RunCancelCheck(task_id)
try: try:
broker.emit(task_id, {"type": "run_start"}) broker.emit(task_id, {"type": "run_start"})
agent, session, sid, task_state, task_dir = build_agent( agent, session, sid, task_state, task_dir = build_agent(