返回博客
官小西

deer-flow PR #4833 源码走读:持久化任务通知如何做到不丢、不重、不打扰

你让 Agent 起一个要跑几小时的活,然后合上笔记本去吃饭。回来打开浏览器:任务还在跑吗?Agent 还记得这个任务吗?结果出来之后,谁告诉你?

这三问是所有长时 Agent 系统的死穴。deer-flow 的 issue #4652(2026-08-03 提出)把最疼的一针写得很直白:任务 ID 靠对话上下文携带,多轮对话之后上下文压缩会把任务 ID 压没,大模型开始编造不存在的任务 ID;想绕过去就持续轮询,又会触发工具循环调用。这两个失败模式——幻觉和风暴——没有一个靠提示词能修好。

2026-08-22 合并的 PR #4833(commit 5ffc2d3,作者 AnnaSuSu,68 个文件 +4293/-121)是这套问题的完整解法。它属于 epic #4652(MCP 官方 Tasks 扩展协议 SEP-2663 支持)的第三步:

PR 合并时间 规模 交付内容
#4665 2026-08-08 +1677/-10,29 文件 durable task runtime 地基:持久化、轮询、租约
#4690 2026-08-15 +3218/-100,47 文件 第一个具体 driver:submit/status 工具接通
#4833 2026-08-22 +4293/-121,68 文件 可靠通知 + 会话隔离取消 + 聊天 UI 面板

截至 2026-08-23,deer-flow 仓库 80,572 Star、11,065 Fork。我在五月写过它的整体架构,这篇接着往下钻一层:一个「通知」在 DeerFlow 2.0 里到底经历了什么。

架构总览:三个租约循环 + 一张表

先看全貌。整套机制没有引入任何新依赖(PR 描述里专门强调 no new dependencies),没有消息队列,没有 Redis——所有可靠性都压在 PostgreSQL 一张表的模式上:

PR #4833 总体架构

四个角色:mcp_tasks 表(唯一事实源)、McpTaskService.run_once() 里的三条工作队列(轮询/取消/通知)、Gateway 的幂等 run 启动器、前端聊天页的后台任务面板。三条队列共享同一个骨架:认领(FOR UPDATE SKIP LOCKED + 租约)→ 执行 → 释放或提交。任何一步崩溃,租约过期后行自动回到可认领状态——没有内存态需要恢复,这就是服务重启自愈的全部秘密。

存储层:outbox 不需要单独的表

教科书上的 outbox 模式通常是一张独立的 outbox 表 + 一个中继进程。这里的设计更激进:outbox 直接长在任务行上,迁移 0013_mcp_task_notificationsmcp_tasks 加了 15 列(12 列通知 + 3 列取消)和 3 个索引:

# persistence/mcp_tasks/model.py(节选)
event_fingerprint: Mapped[str | None]        # sha256(事件快照)
event_version: Mapped[int]                   # 已观测到的最新事件版本
notified_version: Mapped[int]                # 已成功投递的事件版本
dispatch_version / dispatch_attempt / dispatch_event
notification_run_id / notification_error / notification_attempt_count
next_notification_at / notification_lease_owner / notification_lease_expires_at

「有没有没送出去的通知」这个判断,退化成一次索引上的比较:event_version > notified_version。一个版本对,就是一个锁存器(latch)。

新事件的产生靠指纹去重。每次轮询到远程快照,只有当事件快照的 sha256 变化了才 bump 版本号:

# persistence/mcp_tasks/sql.py(节选)
def _record_event_if_changed(row, *, tracking_degraded, now) -> bool:
    event = _notification_event(row, tracking_degraded=tracking_degraded)
    if event is None:
        return False
    fingerprint = _event_fingerprint(event)
    if fingerprint == row.event_fingerprint:
        return False                      # 没变化:一个通知都不发
    row.event_fingerprint = fingerprint
    row.event_version = int(row.event_version or 0) + 1
    if row.notification_status not in _INFLIGHT_NOTIFICATION_STATUSES:
        row.notification_status = "pending"
        row.next_notification_at = now     # 死信通道被新事件复活
    return True

注意最后一处:如果通知已经进了死信,新事件会把它拉回 pending 重走全程——死信杀死的是某个快照,不是这个任务的通知通道。这个细节是「不丢」语义的关键:毒丸消息烧不掉后续真正重要的通知。

投递层:通知 = 一次幂等 Agent run

这是整个 PR 最反直觉、也最值得咀嚼的设计。通知的「投递动作」不是推送一条 WebSocket 消息,也不是发 toast——是往任务所属的原会话线程里,注入一条用户不可见的 user turn,启动一次完整的 Agent run,让 LLM 把这个事件转述给用户

# gateway/services.py(节选)
def _mcp_task_notification_prompt(event: dict[str, Any]) -> str:
    payload = frame_untrusted_text(json.dumps(event, sort_keys=True, ...))
    instruction = (
        "A durable background MCP task has an update that requires the user's "
        "attention. Explain the update clearly and concisely. Do not expose or "
        "ask for a remote task ID. ..."
    )
    return f"{instruction}\n\n{payload}"

idempotency_key = f"mcp-task:{task_id}:{dispatch_version}:{dispatch_attempt}"
record = await start_run(body, thread_id, request,
                         idempotency_key=idempotency_key,
                         require_existing_thread=True)

为什么值得这么做?因为「通知」的真正受众不是浏览器标签页,是那段对话的上下文。一条系统 toast 只能让此刻在线的用户看到;而一次注入原线程的 run,让 Agent 自己知道任务结束了——下一次用户问「刚才那个报告写完了吗」,Agent 不是查库,是真的知道。通知的终点从「屏幕」变成了「记忆」,这直接消灭了 #4652 描述的任务 ID 幻觉:ID 不再需要模型背诵,它在数据库里。

代价当然是每次通知烧一次 LLM 调用,所以幂等性是生死线。这里有两层:

  1. 状态机层:dispatched 状态记下 notification_run_id,后续轮询该 run 的终态来决定 delivered 还是 retry;
  2. runs 表层:迁移给 runs 加了 idempotency_key 唯一索引——同一 (task_id, version, attempt) 的重复启动(比如 worker 崩溃前没来得及标记)会被数据库兜底,返回同一条 run。

状态机:一次投递的完整旅程

notification_status 有七个态,每条边都是一次带租约守卫的条件更新(WHERE lease_owner = 本人 AND lease 未过期),崩溃只会让租约过期回到可认领:

notification_status 状态机

三个出口值得单独说:

线程忙(409)→ 弃旧投新。 通知 run 用 multitask_strategy="reject" 启动,线程正忙就吃 409 → ConflictError → 清空已存的事件快照、回到 pending,下次认领时用最新状态重建快照。注意这条路不计入失败次数——线程忙不是事件的错。对状态类通知来说这是正确的合并语义:中间态没有投递价值,投最新就够。

线程没了(404)→ 永久死信。 require_existing_thread=True 让通知 run 拒绝复活已被删除的会话(LookupError → 404 → PermanentNotificationError),直接进 dead_letter,一次重试都不浪费。这个区分很见功力:会重试的失败和不会重试的失败,从异常类型上就分开了errors.py 里专门定义了 PermanentNotificationError)。

投 5 次都没成 → 死信。 退避公式 min(poll_interval × 2^min(failures,16), max_backoff)_MAX_NOTIFICATION_ATTEMPTS = 5。有界重试是防通知风暴的最后一道闸——没有它,一个持续失败的投递目标会把重试队列变成 DoS 源头。

时序:正常路径 + 异常出口

把上面拼成一条时间线(含三个异常出口的位置):

投递时序图

取消:会话隔离 + 围栏

取消走的是和通知对称的独立队列(claim_cancel_requests),但有两个针对竞态的细节:

# request_cancel:取消请求会"围栏"掉在途轮询结果
row.cancel_requested_at = requested_at
row.lease_owner = None          # 在途轮询租约立即作废
row.lease_expires_at = None
# 重复取消必须保留已有取消租约,防止并发远程取消

用户侧的取消入口有两个:前端面板的停止按钮,和给 Agent 的两个内置工具 list_background_tasks / cancel_background_task。后者的存在意味着你可以直接对聊天框说「把刚才那个任务停了」,Agent 用自然语言匹配任务名即可——远程 MCP 任务 ID 全程不暴露给任何一方

安全边界:不信任任何外部字符串

远程 MCP 服务器返回的内容是不可信输入,PR 在三处独立设防:

  1. 进模型上下文前:事件 JSON 被 frame_untrusted_text() 包裹在不可信边界标记里,<system> 之类的伪造标签会被 HTML 转义(这是 DeerFlow 原有的提示注入防御中间件,本 PR 把它改名为公共原语复用);
  2. 进模型状态前:worker 把任务列表投影进 graph input 时,task_nameneutralize_untrusted_tags(),且只投影 4 个白名单字段、限 20 条;
  3. 出 API 前_public_task() 构造的对外形状只有 7 个白名单字段——测试里专门埋了 "must-not-leak"remote_task_iddriver_data.secret,断言它们不出现在任何输出里。

白名单是唯一正确的方向:黑名单永远列不全 JSON 里能藏什么。

前端:一张自适应轮询的面板

聊天页加了一个 Sheet 面板(active/recent 两段、展开看详情、一键取消),数据源是 react-query:

// frontend/src/core/background-tasks/hooks.ts(节选)
refetchInterval: (query) =>
  query.state.data?.some(isActiveBackgroundTask) ? 3000 : 15000,
refetchIntervalInBackground: false,

有活跃任务 3 秒一刷,没有就降到 15 秒,标签页不可见时停刷。搭配 Agent 侧通知 run 的「结果进了对话流」,用户即使不看面板也不会错过终态。PR 描述的验证量也值得一提:后端 make test 11,653 通过,前端 994 个测试 + 生产构建的 E2E。

这套设计 vs 常规做法

维度 WebSocket/SSE 直推 客户端纯轮询 deer-flow:行上 outbox + Agent run
服务重启后 连接断,事件丢 不丢但延迟高 租约过期自动续投
重复投递 靠内存去重,弱 天然幂等 版本对 + 幂等 key 双保险
通知到达「谁」 浏览器标签页 浏览器标签页 会话上下文本身(Agent 知道)
毒丸消息 无限重试或丢 无重试 5 次后死信,新事件仍可投
中间态合并 409 自动弃旧投新
每条通知成本 一次推送 N 次空轮询 一次 LLM 调用
外部依赖 需要连接层 无(纯 PG)

我的判断:如果你的通知受众是「正在看屏幕的人」,直推 + 重连补发就够了;如果受众是「要继续对话的 Agent」,这套才对题。它贵(每条通知一次 LLM 调用),但买回来的是通知与上下文的合一——这个取舍在 Agent 产品里会越来越主流。

设计评分

维度 理由
崩溃安全性 9/10 全链路租约 + 幂等 key,无内存态;扣一分给 run 终态判定的轮询粒度
语义正确性 9/10 版本对判定、指纹去重、死信不杀通道、409 合并,四个细节都在点上
安全 9/10 三处独立白名单/转义 + 泄漏测试断言;边界想得很全
性能/成本 6/10 每条通知一次 LLM 调用;任务投影每 run 常驻上下文(限 20 条)
可迁移性 8/10 零新依赖、纯 PG 模式,任何有 PG 的系统都能抄;但「LLM 当渲染器」依赖 Agent 产品形态

不足与保留意见

诚实列几条我不会照单全收的地方:

  1. input_required 只展示不回传。任务进入「需要用户输入」状态时,通知会告诉你它在问什么,但这个 MVP 没法把答案送回去(PR 明说是未来 driver 能力)。交互闭环缺一半。
  2. 死信的用户体验是哑的。5 次失败后通知进 dead_letter,只有打开面板才能发现。对「结果完成」这种高价值事件,值得加一条管理员可见的告警路径。
  3. 投递延迟受轮询粒度限制。dispatched 状态下 run 终态靠下轮 run_once 才确认,延迟 = 轮询间隔。对分钟级任务无所谓,对秒级敏感场景要调参。
  4. 每 run 常驻的任务投影background_tasks(限 20 条)进每个 run 的 graph input,是持续 paid context——任务少的线程也在付这份钱。按需注入会更省。

趋势位置

这套东西不是孤例,它是 harness engineering 里「可靠性基础设施」一层的具体化:Agent 的工具调用从同步 RPC 走向分钟/小时级的持久任务,harness 就必须补上任务底座——提交、恢复、通知、取消,一个都不能靠模型自觉。再往上看一层,loop engineering 讨论的自主编排之所以敢让 Agent 自己跑长循环,前提正是这种 durable substrate 存在:循环可以断,状态不会丢。MCP 的 SEP-2663 Tasks 扩展把「任务」变成协议级概念之后,谁家 harness 能把持久任务做得最不声不响,谁的长时 Agent 就最接近可用。

结论:抄什么

如果你维护着任何有分钟级工作流的产品(视频生成、批量分析、定时任务),这个 PR 值得整段抄的是四件事,按优先级:

  1. 两个版本号event_version / notified_version)——「有没有没送出去的事件」用一次列比较回答,比任何 outbox 表都便宜;
  2. 指纹去重——状态没变就一个通知不发,用户感知的「不打扰」全靠它;
  3. 有界重试 + 死信不杀通道——5 次上限防风暴,新事件复活防丢失;
  4. 会重试的失败和不会重试的失败用不同异常类型——404 直接死信、409 不计失败,这个分类比任何重试参数都值钱。

至于「通知 = Agent run」,先看你的通知受众是人还是 Agent:是人,用直推;是 Agent,这套是迄今为止公开实现里最完整的一份。

参考资料

  1. AnnaSuSu — feat(mcp): complete durable task notifications and chat UI (PR #4833)(2026-08-22 合并,commit 5ffc2d3
  2. deer-flow 维护者 — [feat] 增加对MCP官方Tasks扩展协议(SEP-2663)的支持 (Issue #4652)(2026-08-03 提出)
  3. feat(mcp): add durable task runtime foundation (PR #4665)(2026-08-08 合并)
  4. feat(mcp): add ordinary durable task driver (PR #4690)(2026-08-15 合并)
  5. bytedance — deer-flow 仓库(截至 2026-08-23:80,572 Star / 11,065 Fork)
  6. modelcontextprotocol — Tasks Extension (SEP-2663) 提案
  7. 本站 — DeerFlow 2.0 深度解析(2026-05-07)