diff --git a/CHANGELOG.md b/CHANGELOG.md index 00d0673..432151c 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -6,6 +6,10 @@ > 开发中的用户文案可先写入 `## Unreleased`;该区不会被前端解析,正式发布时再替换为数字版本和日期。 > 工程口径的完整记录见 `PROGRESS.md` / git log。 +## Unreleased + +- 修复服务更新或多实例切换期间,点击“停止”后对话可能一直停留在“停止中”的问题。 + ## 0.67.0 — 2026-08-21 - Blender 从单一回转窑模板升级为受管的通用静态三维场景创作,可组合基础体、拉伸/旋转/扫掠、管道、修改器、材质、灯光、文字和最多八个自定义视图;回转窑作为可复用组件继续提供。旧版回转窑工程不再支持在线续作。 diff --git a/DESIGN.md b/DESIGN.md index 42bea27..2853c24 100644 --- a/DESIGN.md +++ b/DESIGN.md @@ -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)。 -**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:` + 单条 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:` + 单条 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 文件副视图 diff --git a/RUN.md b/RUN.md index 50f1bf9..40b026d 100644 --- a/RUN.md +++ b/RUN.md @@ -368,7 +368,7 @@ $env:ZCBOT_EVAL_TOKEN = "" | `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` | 必填 | | `GET /v1/tasks/{id}/events` | SSE 流(`event: ` + `data: `);订阅 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 轮询用 | 必填 | | `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 文件保留 | 必填 | @@ -1004,7 +1004,7 @@ sudo xfs_quota -x -c "limit -p bhard=10g zcbot_" /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}/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 再清空 | -| 点 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 级,无需处理 | | `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 | diff --git a/tests/test_run_cancel.py b/tests/test_run_cancel.py new file mode 100644 index 0000000..5bc8412 --- /dev/null +++ b/tests/test_run_cancel.py @@ -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() diff --git a/web/runs.py b/web/runs.py index 824f6e4..ea6a1ad 100644 --- a/web/runs.py +++ b/web/runs.py @@ -7,9 +7,11 @@ web 路由与渠道回调都经它起 run;`run_channel_conversation` 是微信/ from __future__ import annotations import asyncio +import time from uuid import UUID from sqlalchemy import select, update +from sqlalchemy.exc import SQLAlchemyError from core.paths import from_db_path from core.storage import session_scope @@ -21,6 +23,39 @@ from .broker import broker 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( task_id: UUID, user_id: UUID, user_message: 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 跑。 user_id 必须从 JWT 那侧透传过来 —— 决定 memory_block 读哪个 per-user 子树。 - cancel_check 桥 broker.is_cancelled,loop 在 stream chunk 间 + 工具调用之间 poll; - cancel 延迟 ~ 单 chunk 间隔(100ms 级);seedance 轮询间也读这个 cancel_check 用于 - 用户停止按钮(必须在 build_agent 阶段就传进去,因为 SeedanceTool ctor 持有它, - 不能像以前那样 build_agent 返回后再赋 agent.cancel_check)。 + cancel_check 优先读 broker(约 100ms 响应),并每秒读一次共享 DB 的 cancelling + 状态兜底。后者保证未配置 Redis 的蓝绿切换窗口、或 broker 临时故障时,停止请求 + 仍能跨实例送达。seedance 轮询间也读同一个 cancel_check 用于用户停止按钮(必须在 + build_agent 阶段就传进去,因为 SeedanceTool ctor 持有它,不能像以前那样 + build_agent 返回后再赋 agent.cancel_check)。 `ok` 收尾回 `idle`;`cancelled`(用户停止)与 `error` 一样落持久终态 —— 前端据此 补「已停止」/ 错误卡(扛过收尾重渲),下次起新 run(post_message 写 running)覆盖清掉。 @@ -45,7 +81,7 @@ def run_agent_bg( 避免重复落库。CLI 等非 Web 入口保持旧行为。 """ 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: broker.emit(task_id, {"type": "run_start"}) agent, session, sid, task_state, task_dir = build_agent(