DeerFlow 端到端流式输出架构解析:LangGraph 事件如何同时送达 HTTP SSE 与嵌入式 Python 客户端
流式输出是任何长时运行 Agent 产品体验的命脉:浏览器、IM 渠道要看到 token 级打字机效果,Jupyter / 脚本 / 测试需要拿到逐步到达的原生对象。DeerFlow 为此维护了两条并行且刻意不复用的流式管道——Gateway(async + HTTP SSE + JSON)与
DeerFlowClient(sync + in-process + 原生 LangChain 对象)。本文完整梳理这两条路径的分工与契约、stream_mode三层语义、有界历史与gap恢复、客户端三个 id 去重集合的微妙不变式,以及防止"漏订messages模式"这类回归的测试策略。
阅读完本文,你将掌握:DeerFlow 为什么需要两条流式路径而非让客户端复用 Gateway、values / messages / custom 三种流模式各自的发射时机与事件契约、messages-tuple 与 messages 两种命名在协议层间的翻译关系、SSE 断线重连与 gap 恢复边界,以及嵌入式客户端三套 set[str] 各自守护的不变式。
TL;DR:先建立整体心智模型
- DeerFlow 存在两条并行的流式路径,服务两种截然不同的消费者模型:
- Gateway 路径:
async def/ HTTP SSE / JSON 序列化,服务浏览器前端与飞书 / Slack / Telegram 等 IM 渠道; - DeerFlowClient 路径:sync generator / in-process / 原生 LangChain 对象,服务 Jupyter、集成脚本与测试。 两条路径无法合并——消费者模型根本不同(详见 为什么有两条流式路径)。
- Gateway 路径:
- 两条路径都从
create_agent()工厂出发,核心订阅 LangGraph 的stream_mode=["values", "messages", "custom"]。其中values是节点级 state 快照,messages是 LLM token 级 delta,custom是显式StreamWriter事件;DeerFlow 内置 custom 事件同时通过 callback dispatch 暴露为astream_events(version="v2")的on_custom_event。 - 嵌入式客户端为每次
stream()调用维护三个set[str]:seen_ids/streamed_ids/counted_usage_ids。三者看似相似,实际各管理一个独立不变式,不能合并(详见 三个 id set 为什么不能合并)。
为什么有两条流式路径
两条路径服务的消费者模型根本不同,其差异可用一张对照表概括:
| 维度 | Gateway 路径 | DeerFlowClient 路径 |
|---|---|---|
| 入口 | FastAPI /runs/stream endpoint |
DeerFlowClient.stream(message) |
| 触发层 | worker.py run_agent |
client.py DeerFlowClient.stream |
| 执行模型 | async def + agent.astream() |
sync generator + agent.stream() |
| 事件传输 | StreamBridge(asyncio Queue)+ sse_consumer |
直接 yield |
| 序列化 | serialize(chunk) → 纯 JSON dict,匹配 LangGraph Platform wire 格式 |
StreamEvent.data,携带原生 LangChain 对象 |
| 消费者 | 前端 useStream React hook、飞书/Slack/Telegram channel、LangGraph SDK 客户端 |
Jupyter notebook、集成测试、内部 Python 脚本 |
| 生命周期管理 | RunManager:run_id 跟踪、disconnect 语义、multitask 策略、heartbeat |
每次 stream() 生成一个轻量 run_id 供 runtime context / tracing / per-run middleware 使用;函数返回即结束 |
| 断连恢复 | Last-Event-ID SSE 重连 |
无需要 |
两条路径的存在是 DRY 的刻意妥协:Gateway 的全部基础设施(async + Queue + JSON + RunManager)都是为了跨网络边界把事件送给 HTTP 消费者。当生产者(agent)和消费者(Python 调用栈)处于同一进程时,整套机制都是纯开销。
为什么不能让 DeerFlowClient 复用 Gateway
设计上曾经考虑过三种复用方案,均被否决:
- 让
client.stream()变成async def client.astream()——breaking change。Jupyter notebook 和同步脚本需要硬塞用不上的async for/asyncio.run(),DeerFlowClient"把 agent 当普通函数调用"的核心卖点直接消失。 - 在
client.stream()内部启动独立事件循环线程,用StreamBridge在 sync/async 之间做桥接——会引入线程池、队列、信号量。为了"消除重复"却把复杂度带进来,属于典型的 "wrong abstraction":开销高于复用收益。 - 让
run_agent自己兼容 sync mode——给 Gateway 增加一条用不到的死分支,污染 worker 的职责焦点。
因此两条路径的事件处理逻辑相似但不共享,这是刻意设计,不是疏忽。
源码佐证:这段取舍在 client.py
_stream_turn的 docstring 中有同样明确的阐述——run_agent是async def且使用agent.astream(),而 client 是 sync generator 使用agent.stream(),两者桥接需要每调用启动一个 event loop + thread;同时 Gateway 事件需要serialize()做 JSON/SSE wire 传输,in-process 调用方拿到的是可直接使用的StreamEvent结构。两条路径订阅哪些 LangGraph stream mode 的同步依赖行为断言而非共享常量来保证。
LangGraph stream_mode 三层语义
LangGraph 的 agent.stream(stream_mode=[...]) 是多路复用接口:一次订阅多个 mode,每个 mode 都是一个独立的事件源。
flowchart LR
classDef values fill:#B8C5D1,stroke:#5A6B7A,color:#2C3E50
classDef messages fill:#C9B8A8,stroke:#7A6B5A,color:#2C3E50
classDef custom fill:#B5C4B1,stroke:#5A7A5A,color:#2C3E50
subgraph LG["LangGraph agent graph"]
direction TB
Node1["node: LLM call"]
Node2["node: tool call"]
Node3["node: reducer"]
end
LG -->|"每个节点完成后"| V["values: 完整 state 快照"]
Node1 -->|"LLM 每产生一个 token"| M["messages: (AIMessageChunk, meta)"]
Node1 -->|"emit_custom_event()"| E["DeerFlow custom event helper"]
E -->|"StreamWriter.write()"| C["custom: 任意 dict"]
E -->|"dispatch_custom_event()"| A["astream_events(v2): on_custom_event"]
class V values
class M messages
class C custom
class A custom
三个核心 mode 的语义对照:
| Mode | 发射时机 | Payload | 粒度 |
|---|---|---|---|
values |
每个 graph 节点完成后 | 完整 state dict(title、messages、artifacts) | 节点级 |
messages |
LLM 每次 yield 一个 chunk;tool 节点完成时 | (AIMessageChunk | ToolMessage, metadata_dict) |
token 级 |
custom |
用户代码显式调用 StreamWriter.write() |
任意 dict | 应用定义 |
on_custom_event |
用户代码调用 dispatch_custom_event();通过 astream_events(version="v2") 消费 |
name + 任意 data |
应用定义 |
关键事实:这些 mode 不是"详细程度递增的梯度",而是互相独立的平行事件源。消费者必须显式订阅自己需要的那一个,缺订任何一个(尤其是 messages)都会静默丢失一整类事件——这正是历史回归的温床(见 为什么这个设计容易出 bug)。
DeerFlow 对 custom 事件的双通道约束
DeerFlow 自身产生的事件必须通过 custom_events.py 中定义的同步 / 异步 helper 发送。其约束为:
- 每个内置 payload 必须携带非空字符串
type;缺少合法type的 payload 只进入customstream,不会出现在astream_events。 - helper 的执行顺序是先写入
customstream,再做 best-effort callback dispatch:callback 名称取 payload 的type,data保留完整 payload。 - 这样原生 Gateway / Web UI /
DeerFlowClient的 custom 事件保持不变,astream_events消费者也能观察到同一事件。 - callback dispatch 的普通异常只记录 debug 日志,不允许打断原有 writer 链路;writer 自身的异常语义保持不变。
从源码看,emit_custom_event / aemit_custom_event 先调用 writer(payload),再经 _event_name 校验后走 dispatch_custom_event(event_name, payload)(异步版走 adispatch_custom_event);type 缺失时只记一条 debug 日志并跳过 callback。GraphBubbleUp 会被特殊透传,其余异常全部吞掉记为 debug——这是"writer 为主、callback 为 optional 观察通道"这一设计意图的直接实现。
两套命名的由来
同一件事在三个协议层有完全不同的名字:
Application HTTP / SSE LangGraph Graph
┌──────────────┐ ┌──────────────┐ ┌──────────────┐
│ frontend │ │ LangGraph │ │ agent.astream│
│ useStream │──"messages- │ Platform SDK │──"messages"──│ graph.astream│
│ Feishu IM │ tuple"──────│ HTTP wire │ │ │
└──────────────┘ └──────────────┘ └──────────────┘
- Graph 层(
agent.stream/agent.astream):LangGraph Python 直接 API,mode 叫"messages"。 - Platform SDK 层(
langgraph-sdkHTTP client):跨进程 HTTP 契约,mode 叫"messages-tuple"。 - Gateway worker 负责在两个命名体系之间显式翻译。
后果:DeerFlowClient.stream() 直接调用 agent.stream()(Graph 层),必须传 "messages";channels/manager.py 通过 langgraph-sdk 走 HTTP SDK,所以传 "messages-tuple"。这两个字符串不能互相替代,也不能被抽成"一个共享常量"——它们是不同协议层的 type alias,共享只会让某一层说不是它母语的话。
源码佐证:翻译逻辑现在集中在 runtime/stream_modes.py。
RunStreamMode定义了 DeerFlow runtime 边界完整支持的公开 mode 集合(values、messages-tuple、updates、debug、tasks、checkpoints、custom);normalize_stream_modes()负责归一化与白名单校验(不支持的 mode 直接抛UnsupportedStreamModeError);而真正的"翻译"由to_langgraph_stream_modes()完成——mapped = ["messages" if mode == "messages-tuple" else mode for mode in modes],随后去重并返回给agent.astream。Gateway 侧的 worker.pyrun_agent通过to_langgraph_stream_modes消费这份翻译,并在内部根据是否包含"messages-tuple"决定是否启用_LargeFileToolChunkBatcher(大文件工具 chunk 分批)等逻辑。
Gateway 路径:async + HTTP SSE
Gateway 路径是完整的多组件协作流水线,参与方与调用时序如下:
sequenceDiagram
participant Client as HTTP Client
participant API as FastAPI<br/>thread_runs.py
participant Svc as services.py<br/>start_run
participant Worker as worker.py<br/>run_agent (async)
participant Bridge as StreamBridge<br/>(asyncio.Queue)
participant Agent as LangGraph<br/>agent.astream
participant SSE as sse_consumer
Client->>API: POST /runs/stream
API->>Svc: start_run(body)
Svc->>Bridge: create bridge
Svc->>Worker: asyncio.create_task(run_agent(...))
Svc-->>API: StreamingResponse(sse_consumer)
API-->>Client: event-stream opens
par worker (producer)
Worker->>Agent: astream(stream_mode=lg_modes)
loop 每个 chunk
Agent-->>Worker: (mode, chunk)
Worker->>Bridge: publish(run_id, event, serialize(chunk))
end
Worker->>Bridge: publish_end(run_id)
and sse_consumer (consumer)
SSE->>Bridge: subscribe(run_id)
loop 每个 event
Bridge-->>SSE: StreamEvent
SSE-->>Client: "event: <name>\ndata: <json>\n\n"
end
end
各关键组件及其职责:
- worker.py
run_agent:在asyncio.Task中运行agent.astream(),把每个(mode, chunk)经serialize(chunk, mode=mode)转成 JSON 后调用bridge.publish(),结束时调用bridge.publish_end()。 runtime/stream_bridge:抽象 Queue 层的完整实现位于 runtime/stream_bridge/(含base.py抽象基类、memory.py与redis.py两种 backend)。publish/subscribe解耦生产者和消费者,支持Last-Event-ID重连、心跳、多订阅者 fan-out。从 base.py 可以看到StreamGap与END_SENTINEL(StreamEvent(id="", event="__end__", data=None))是 stream 结束 / 出现断档的两种显式信号。- services.py
sse_consumer:从 bridge 订阅事件,并配合format_sse()格式化成标准 SSE wire 帧(event: <name>\ndata: <json>\n\n,event_id可选)。 - runtime/serialization.py
serialize:mode-aware 序列化;在messagesmode 下serialize_messages_tuple把(chunk, metadata)转成[chunk.model_dump(), metadata],在valuesmode 下额外剥离__pregel_*内部键并丢弃 hide_from_ui 消息中的 base64 data: 图片块。
StreamBridge 的存在价值:当生产者(run_agent 任务)和消费者(HTTP 连接)运行在不同的 asyncio task 里时,需要一个可以跨 task 传递事件的中介。Queue 同时承担断连重连的缓冲区和多订阅者的 fan-out。
有界历史与 gap 恢复
Last-Event-ID 只保证在保留窗口内完整重放,并不代表无限历史。契约可以总结为:
- 有效游标仍在窗口内时:bridge 从该 ID 的下一条事件正常恢复。
- 有效游标已被
queue_maxsize裁剪,或在线消费者慢到落后水位线时:Gateway 会在任何部分重放之前,先发送一帧gap,而不是从当前最早事件静默地部分重放:
event: gap
data: {"code":"stream_replay_gap","run_id":"...","requested_event_id":"...","earliest_available_event_id":"...","latest_available_event_id":"...","recovery":"reload_durable_state"}
关于 gap 帧需要明确几个语义:
gap帧没有 SSEid:,其后也没有正常的end;当前订阅随即关闭。- 它是"恢复边界"而不是"客户端断开",因此不会触发
on_disconnect=cancel。 - 当缓冲区无任何保留事件时,
earliest_available_event_id与latest_available_event_id均为null。 - 客户端必须丢弃不再可信的瞬时状态,重新读取 thread checkpoint 与持久化的 run-event / message history,然后以
latest_available_event_id为游标跟随新事件(若缓冲区为空即latest_available_event_id为null,则不带游标重新加入流)。DeerFlow Web UI 会自动执行这一整套流程,且最多连续恢复五次。
Memory 与 Redis 两种 backend 在这一契约上的细微差异(来自原设计文档的补充说明):
- Redis:对无游标、空 stream 上已建立的阻塞等待遵循相同契约——第一次
XREAD唤醒的数据在交付前仍是 provisional baseline,bridge 会用下一次事务快照确认其尾 ID 仍在保留窗口;若生产者已裁剪该基线,订阅直接返回requested_event_id: null的gap。该检查有明显的性能代价:每轮订阅需要一个包含XRANGE、XREVRANGE、非阻塞XREAD的事务快照,空闲时还需要单独的阻塞XREAD来唤醒。 - "有效但已淘汰" vs malformed cursor 是不同策略:该契约只要求前者产生
gap。Redis 对 malformed ID 仍从 live tail 等待;Memory 对 malformed ID 及序号不低于水位线的未知 ID 仍采用"最早保留事件"策略。对序号已低于水位线的数字格式 foreign ID,Memory 无法再校验已淘汰的 timestamp,因此保守返回gap,优先保证客户端不会把不完整重放误认为完整。 - 孤儿 run 恢复:Memory 与 Redis 都只保留
queue_maxsize条数据事件。Redis backend 在每次publish()/publish_end()刷新 retained stream key 的 TTL;启动恢复与基于 worker lease 的周期恢复共用 Gateway stream terminalization 路径——RunManager先把 orphan run 持久化为error并写入显式stop_reason=orphan_recovered,随后 Gateway 发布END_SENTINEL并安排 stream cleanup。这里隐含一个判断:store-only SSE 与/waitconsumer 不能把普通 durable terminal status 当成流已完成,否则可能跳过延迟发布的 error 等尾部事件;只有orphan_recovered信号能在 heartbeat 时触发 END fallback——因为此时 producer 已被确认失联。TTL 是 Redis 内存和故障安全网,而不是正常的 subscriber 终止机制。
StreamBridge 可配置项
Bridge 的关键参数集中在 config/stream_bridge_config.py,默认值与边界如下:
| 参数 | 默认值 | 说明 |
|---|---|---|
type |
memory |
backend 类型:memory 为进程内事件日志(仅单进程可用);redis 用于多 worker Docker 部署(Redis Streams)。 |
redis_url |
None |
Redis URL;省略时按 DEER_FLOW_STREAM_BRIDGE_REDIS_URL → REDIS_URL → redis://localhost:6379/0 的顺序回退。 |
queue_maxsize |
256 |
每个 run 保留的最大事件数(memory 队列长度 / redis stream MAXLEN),最小为 1。 |
heartbeat_interval_seconds |
15.0 |
空闲时两次 heartbeat 之间的秒数(0 < v ≤ 86400,且拒绝布尔值),统一作用于 SSE 客户端、非流式 wait 请求与内部 stream 订阅者;显式传给 subscribe() 的值可覆盖单次订阅。 |
max_connections |
None |
Redis 连接池上限。每个活跃 SSE 客户端会持有一条阻塞在 XREAD ... BLOCK(最长 heartbeat_interval_seconds)的连接,数百并发客户端即数百连接。仅 redis bridge 生效。 |
stream_ttl_seconds |
86400 |
Redis stream key 滚动 TTL,每次 publish/publish_end 刷新;为 0 时关闭。是故障兜底而非正常终止机制。 |
recovered_stream_cleanup_delay_seconds |
60.0 |
孤儿 run 被恢复并发布 END marker 后,删除 stream key 前等待的秒数,给重连 SSE 客户端留出排空尾部信号的时间。仅 redis bridge 生效。 |
周期扫描、逐行状态写入和 Gateway callback 作为一个受监督的 single-flight 后台 task 执行:慢任务不会堆积,也不会阻塞唯一的 lease heartbeat;shutdown 优先收敛活跃 run 再处理恢复 task;尚未执行的延迟 stream cleanup 会改为立即删除。只有"runtime yield 前、无并发请求"的启动恢复会把最新受影响 thread 标记为 error,周期恢复不做非原子的 thread 投影。
DeerFlowClient 路径:sync + in-process
sequenceDiagram
participant User as Python caller
participant Client as DeerFlowClient.stream
participant Agent as LangGraph<br/>agent.stream (sync)
User->>Client: for event in client.stream("hi"):
Client->>Agent: stream(stream_mode=["values","messages","custom"])
loop 每个 chunk
Agent-->>Client: (mode, chunk)
Client->>Client: 分发 mode<br/>构建 StreamEvent
Client-->>User: yield StreamEvent
end
Client-->>User: yield StreamEvent(type="end")
对比之下,sync 路径每个环节的移动部件都显著更少:
- 没有
RunManager:一次stream()调用对应一次生命周期,只生成轻量run_id供 runtime context、tracing 与 per-run middleware 使用;函数返回即结束。从源码看,stream()还会基于run_id绑定 trace context——每次next()步进前后bind_trace_id/reset_trace_id成对执行,从不跨越yield,避免把 id 泄漏进调用方 Context 或触发跨 Context 的 GC 清理异常。 - 没有
StreamBridge:直接yield,生产与消费发生在同一个 Python 调用栈,不需要跨 task 中介。 - 没有 JSON 序列化:
StreamEvent.data直接携带原生 LangChain 值——AIMessage.content、usage_metadata的UsageMetadataTypedDict,以及非None的ToolMessage.artifact。messages-tuple工具结果与values快照里的工具消息都会保留 artifact;没有 artifact 的工具消息维持原有字段形状。Jupyter 用户拿到的是真正的类型,而不是经过网络序列化后的匿名值。 - 没有 asyncio:调用者可以直接写
for event in ...,不必写async for。
在 client.py 中,事件类型以 StreamEvent 数据类承载,type 为 "values" | "messages-tuple" | "custom" | "end",与 LangGraph SSE 协议保持一致,因此消费者可以在 HTTP 流与嵌入式模式间切换而无需改写事件处理逻辑。_stream_turn 的 docstring 明确列出各事件 payload 的 shape:
- type="values" data={"title": str|None, "messages": [...], "artifacts": [...]}
- type="custom" data={...}
- type="messages-tuple" data={"type": "ai", "content": <delta>, "id": str}
- type="messages-tuple" data={"type": "ai", "content": <delta>, "id": str, "usage_metadata": {...}}
- type="messages-tuple" data={"type": "ai", "content": "", "id": str, "tool_calls": [...]}
- type="messages-tuple" data={"type": "ai", "content": "", "id": str, "additional_kwargs": {...}}
- type="messages-tuple" data={"type": "tool", "content": str, "name": str, "tool_call_id": str, "id": str}
# Tool results also include "artifact" when the source ToolMessage has a non-None artifact.
- type="end" data={"usage": {"input_tokens": int, "output_tokens": int, "total_tokens": int}}
一个值得注意的边界:多轮对话需要 checkpointer。构造 DeerFlowClient 时传入 checkpointer 后,thread_id 才能保留跨轮上下文;否则每次 stream() / chat() 都是无状态调用,thread_id 仅用于文件隔离(uploads / artifacts)。另外,系统提示词(含日期、memory 与 skills 上下文)在 agent 首次创建时生成并按配置键缓存,长驻进程中若外部修改了 memory / skills 需调用 reset_agent() 强制重建。
消费语义:delta vs cumulative
LangGraph messages mode 给出的是 delta:每个 AIMessageChunk.content 只包含这一次新 yield 的 token,不是从头累计的完整文本。这与 LangChain 的 fs2 Stream 风格一致——上游发增量,下游负责累加:
- Gateway 路径:前端
useStreamReact hook 自己维护累加器。 - DeerFlowClient 路径:
chat()方法替调用者做累加。
DeerFlowClient.chat() 的 O(n) 累加器
client.py chat() 的实现是"逐 id 收集 delta 列表、最后一次性 join":
chunks: dict[str, list[str]] = {}
last_id: str = ""
for event in self.stream(message, thread_id=thread_id, **kwargs):
if event.type == "messages-tuple" and event.data.get("type") == "ai":
msg_id = event.data.get("id") or ""
delta = event.data.get("content", "")
if delta:
chunks.setdefault(msg_id, []).append(delta)
last_id = msg_id
return "".join(chunks.get(last_id, ()))
为什么不用 buffers[id] = buffers.get(id, "") + delta? CPython 的字符串 in-place concat 优化仅在 refcount=1 且 LHS 是 local name 时才生效;这里字符串存放在 dict 中又被 reassign,优化失效,每次都是 O(n) 拷贝 → 总体 O(n²)。实测对 50 KB / 5000 chunk 的回复,纯拷贝开销就要 100–300ms。用 list + "".join() 是真正的 O(n)。
需要补充的一点细节:chat() 返回的是最后一条完整 AI 消息累加出的文本——中间 AI 消息(如 planner 草稿)会被丢弃。需要逐 delta 拿每一条消息时,应直接使用 stream()。
三个 id set 为什么不能合并
DeerFlowClient.stream() 在一次调用生命周期内维护三个 set[str](源码位于 client.py _stream_turn):
seen_ids: set[str] = set() # values 路径内部 dedup
streamed_ids: set[str] = set() # messages → values 跨模式 dedup
counted_usage_ids: set[str] = set() # usage_metadata 幂等计数
乍看像是"三份几乎一样的东西",实际每个集合守护不同的不变式:
| Set | 负责的不变式 | 被谁填充 | 被谁查询 |
|---|---|---|---|
seen_ids |
连续两个 values 快照里的同一条 message 只生成一个 messages-tuple 事件 |
values 分支每处理一条消息就加入 | values 分支处理下一条消息前检查 |
streamed_ids |
一条消息若已通过 messages 模式 token 级流过,values 快照到达时不要再合成一次完整 messages-tuple |
messages 分支每发一个 AI/tool 事件就加入 | values 分支看到消息时检查 |
counted_usage_ids |
同一 usage_metadata 在 messages 末尾 chunk 与 values 快照的 final AIMessage 中各出现一份,累计总量只算一次 |
_account_usage() 每次接受 usage 就加入 |
_account_usage() 每次调用时检查 |
为什么不能只用一个 set
关键观察:同一个 message id 在这三个 set 里的加入时机完全不同:
sequenceDiagram
participant M as messages mode
participant V as values mode
participant SS as streamed_ids
participant SU as counted_usage_ids
participant SE as seen_ids
Note over M: 第一个 AI text chunk 到达
M->>SS: add(msg_id)
Note over M: 最后一个 chunk 带 usage
M->>SU: add(msg_id)
Note over V: snapshot 到达,包含同一条 AI message
V->>SE: add(msg_id)
V->>SS: 查询 → 已存在,跳过文本合成
V->>SU: 查询 → 已存在,不重复计数
seen_ids永远在 values 快照到达时加入,所以它是 "values 已处理" 的标记。一条只出现在 messages 流里的消息(罕见但可能),在seen_ids里永远不存在。streamed_ids在 messages 流的第一个有效事件时加入。一条只通过 values 快照到达的非 AI 消息(HumanMessage、被 truncate 的 tool 消息),在streamed_ids里永远不存在。counted_usage_ids只在看到非空usage_metadata时加入。一条完全没有 usage 的消息(tool message、错误消息)永远不会进入。
关于集合包含关系:counted_usage_ids ⊆ (streamed_ids ∪ seen_ids) 大致成立,但不是严格子集——一条消息可能在 messages 模式流完 text 后、但在最后一个带 usage 的 chunk 之前,就被 values snapshot 赶上:此时它已在 streamed_ids,却还不在 counted_usage_ids。若把它们合并成一个 dict-of-flags,这个微妙的时序依赖会从类型系统里消失,退化成注释里的一句话。三个独立的 set 则把不变式显式化了:每个 set 名都对应一个可以口头回答的问题。
从源码还可以看到
messages-tuple分支中带 usage 的"元数据型 follow-up"约定:若某 AI 消息的additional_kwargs在此前的 chunk 事件中从未发送过,values分支会补发一条 content 为空的 AI 事件仅携带增量additional_kwargs;客户端应按 message id 合并并忽略其文本渲染。_unsent_additional_kwargs用sent_additional_kwargs_by_id按 id 追踪已发送键,只发送真正的新增键值。
端到端:一次真实对话的事件时序
假设调用 client.stream("Count from 1 to 15"),LLM 给出 "one\ntwo\n...\nfifteen"(88 字符),tokenizer 把它拆成约 35 个 BPE chunk。事件到达序列的精简版如下:
sequenceDiagram
participant U as User
participant C as DeerFlowClient
participant A as LangGraph<br/>agent.stream
U->>C: stream("Count ... 15")
C->>A: stream(mode=["values","messages","custom"])
A-->>C: ("values", {messages: [HumanMessage]})
C-->>U: StreamEvent(type="values", ...)
Note over A,C: LLM 开始 yield token
loop 35 次,约 476ms
A-->>C: ("messages", (AIMessageChunk(content="ele"), meta))
C->>C: streamed_ids.add(ai-1)
C-->>U: StreamEvent(type="messages-tuple",<br/>data={type:ai, content:"ele", id:ai-1})
end
Note over A: LLM finish_reason=stop,最后一个 chunk 带 usage
A-->>C: ("messages", (AIMessageChunk(content="", usage_metadata={...}), meta))
C->>C: counted_usage_ids.add(ai-1)<br/>(无文本,不 yield)
A-->>C: ("values", {messages: [..., AIMessage(complete)]})
C->>C: ai-1 in streamed_ids → 跳过合成
C->>C: 捕获 usage (已在 counted_usage_ids,no-op)
C-->>U: StreamEvent(type="values", ...)
C-->>U: StreamEvent(type="end", data={usage:{...}})
四个关键观察:
- 用户看到 35 个 messages-tuple 事件,跨越约 476ms,每个事件携带一个 token delta 和同一个
id=ai-1。 - 最后那个
values快照里的AIMessage不会再触发一个完整的messages-tuple事件——因为ai-1 in streamed_ids跳过了合成。 end事件里的usage正好等于那一份 cumulative usage,不是它的两倍——counted_usage_ids在 messages 末尾 chunk 上已经吸收了 usage,values 分支对同一 id 的重复访问是 no-op。- 消费者拿到的
content是增量:"ele"只包含 3 个字符,不是"one\ntwo\n...ele"。要得到完整文本必须按id累加——chat()已经帮你做了。
为什么这个设计容易出 bug,以及测试策略
本文档的直接起因是 DeerFlow 历史上的一个真实回归(issue #1969):DeerFlowClient.stream() 原本只订阅 ["values", "custom"],漏掉了 "messages"。结果是 client.stream("hello") 退化为一次性返回,视觉上和 chat() 没区别——流式能力静默消失而 CI 依然全绿。
Custom 事件还有一条独立的回归边界:get_stream_writer() 产生的 chunk 不会自动成为 astream_events(version="v2") 的 on_custom_event,而 callback dispatch 也不会自动进入 stream_mode="custom"。因此测试必须使用真实最小 LangGraph 同时锁定两种 API,断言每个消费者各收到一次且 payload 相同;仅 mock 任一函数无法证明协议互操作。
这类 bug 有三个结构性原因:
- 多协议层命名:
messages/messages-tuple/ HTTP SSEmessages是同一概念的三个名字。在其中一层出错,不会在另外两层报错。 - 多消费者模型:Gateway 与 DeerFlowClient 是两套独立实现,没有单一的 "订阅哪些 mode" 的 single source of truth。前者订阅对了,不代表后者也订阅对了。
- mock 测试绕开了真实路径:老测试用
agent.stream.return_value = iter([dict_chunk, ...])喂 values 形状的 dict 来模拟 state 快照。这类输入永远不会进入messagesmode 分支,所以即便stream_mode少一个元素,CI 依然全绿。
防御手段:行为断言 + 真实 chunk shape mock
真正的防线是显式断言 "messages" mode 被订阅 + 用真实 chunk shape mock。回归测试位于 backend/tests/test_client.py 的 test_messages_mode_emits_token_deltas:
# backend/tests/test_client.py::test_messages_mode_emits_token_deltas
agent.stream.return_value = iter([
("messages", (AIMessageChunk(content="Hel", id="ai-1"), {})),
("messages", (AIMessageChunk(content="lo ", id="ai-1"), {})),
("messages", (AIMessageChunk(content="world!", id="ai-1"), {})),
("values", {"messages": [HumanMessage(...), AIMessage(content="Hello world!", id="ai-1")]}),
])
# ...
assert [e.data["content"] for e in ai_text_events] == ["Hel", "lo ", "world!"]
assert len(ai_text_events) == 3 # values snapshot must NOT re-synthesize
assert "messages" in agent.stream.call_args.kwargs["stream_mode"]
为什么这比"抽一个共享常量"更有效? 共享常量只能保证"用它的人写对字符串",但新增消费者的人可能根本不知道常量在哪。行为断言强制任何改动都要穿过实际执行路径——把订阅改回 ["values", "custom"] 会立刻让 assert "messages" in ... 失败。
活体信号:BPE 子词边界
回归的最终验证是让真实 LLM 数 1–15,然后确认输出里能看到 tokenizer 的子词切分:
[5.460s] 'ele' / 'ven' eleven 被拆成两个 token
[5.508s] 'tw' / 'elve' twelve 拆两个
[5.568s] 'th' / 'irteen' thirteen 拆两个
[5.623s] 'four'/ 'teen' fourteen 拆两个
[5.677s] 'f' / 'if' / 'teen' fifteen 拆三个
子词切分是 tokenizer 的外部事实,无法伪造。能看到它,就说明数据流逐 chunk 地穿过了整条管道,没有被任何中间层缓冲成整段。这种"活体信号"在流式系统里是比单元测试置信度更高的证据。
相关源码定位
| 关心什么 | 看这里 |
|---|---|
| DeerFlowClient 嵌入式流 | client.py DeerFlowClient.stream |
| Embedded ToolMessage artifact 序列化 | client.py _tool_message_event / _serialize_message |
chat() 的 delta 累加器 |
client.py DeerFlowClient.chat |
| Gateway async 流 | worker.py run_agent |
| Stream mode 公开集合与命名翻译 | runtime/stream_modes.py |
| StreamBridge 抽象与两种 backend | runtime/stream_bridge/ |
| Bridge 心跳 / 队列 / Redis TTL 配置 | config/stream_bridge_config.py |
| HTTP SSE 帧输出 | services.py sse_consumer / format_sse |
| 序列化到 wire 格式 | runtime/serialization.py |
| DeerFlow custom 事件双通道 helper | utils/custom_events.py |
| 飞书渠道的增量卡片更新 | channels/manager.py _handle_streaming_chat |
| Channels 自带的 delta/cumulative 防御性累加 | channels/manager.py _merge_stream_text |
| Frontend 支持的 mode 集合与 chat 流裁剪 | frontend/src/core/api/stream-mode.ts |
前端 useStream 消费入口 |
frontend/src/core/threads/hooks.ts |
| 核心回归测试 | backend/tests/test_client.py test_messages_mode_emits_token_deltas |
补充一句前端的契约细节:前端在 stream-mode.ts 中维护与后端一致的
SUPPORTED_RUN_STREAM_MODES集合(含values/messages-tuple/custom等),并对不支持的 mode 抛错或告警;forceChatRunStreamOptions会主动剔除values,避免 SDK 惰性消息追踪顺带请求values、在 graph 每一步后重传整段 thread state——这是"消费者必须只订阅自己需要的接口"这一原则在前端侧的对应实现。
atomcodeClaude Code 的开源替代方案。连接任意大模型,编辑代码,运行命令,自动验证 — 全自动执行。用 Rust 构建,极致性能。 | An open-source alternative to Claude Code. Connect any LLM, edit code, run commands, and verify changes — autonomously. Built in Rust for speed. Get StartedRust0624
Hy4-previewHy4 preview 是由腾讯混元团队研发的新一代混合专家(MoE)旗舰模型。模型总参数量 770B,每个 token 激活 49B,主干共包含78层,第一层采用标准 FFN,其余 77 层均为 MoE 结构,每层包含 256 个路由专家与 1 个共享专家,每个 token 激活 top-8 路由专家及共享专家。主干之外原生内置 1 层 MTP(总参数量 10B,激活 0.7B)以支持投机解码。Python00
GLM-5.3GLM-5.3 与 GLM-5.2 使用相同的基座模型——所有提升均来自后训练。与 GLM-5.2 相比,它在复杂编程和长程任务上的表现显著提升。Jinja00
GLM-5.3-FlashGLM-5.3-Flash (320B-A18B),是GLM-5系列的首个原生多模态模型。320B总参数,能力超过GLM-5.2Jinja00
Spark-X2.5-4BSpark-X2.5-4B 旨在让强大的 AI 更实用、更高效、更易获得。在广泛日常任务中表现强劲,涵盖对话、写作、翻译、推理、编码、工具调用以及智能体工作流,并在同等规模的开源模型中取得领先成绩。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00
Spark-X2.5-1.7BSpark-X2.5-1.7B 旨在让强大的 AI 更加实用、高效且易于获取。这些模型在广泛的日常任务中表现出色,涵盖对话、写作、翻译、推理、编程、工具调用和智能体工作流,并在同等规模的开源模型中取得领先结果。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00