流式传输架构
数据分片协议、传输无关的 StreamWriter,以及把 Redis Stream(持久历史)与 pub/sub(实时推)合起来、靠单调 id 让历史与实时无缝衔接的可恢复 SSE
引擎只管往一个流里写,送达是外层的事
Agent 一旦跑起来,输出就源源不断——文本 token、工具结果、状态转换、通知、元数据。这些都得可靠地送到客户端,还要在 Web(SSE)和桌面(IPC)上表现一致;但引擎本身不该操心网络怎么抖、走哪种传输。它只往一个传输无关的流里写,其余的分四层各管一段:
- StreamWriter——传输无关的事件写入接口;
- 数据分片协议——在 AI SDK 的内容分片旁边,捎带 agent 专属的事件;
- 任务传输层——SSE(默认)或 WebSocket,两条路径承载同一套帧格式;
- 可恢复 SSE——基于 Redis 的断连续传,仅 SSE 路径。
StreamWriter 接口
后端 Agent 通过传输无关的 StreamWriter 接口写入事件。业务层注入具体实现:
export interface StreamWriter {
write(event: { type: `data-${string}`; data: unknown; transient?: boolean }): void;
}
| 平台 | 实现方式 |
|---|---|
| Web——SSE(默认) | 封装 createUIMessageStream({ execute }) 返回的 UIMessageStreamWriter |
| Web——WebSocket | 同一个 UIMessageStreamWriter;帧经 wsHub.send(connectionId, ...) 点对点回推 |
| Chat 广播 | wsHub.broadcast() 做房间级扇出(仅文本 delta,chat service 专用) |
| 桌面端 (IPC) | 封装 Electron IPC send |
事件分为 持久化(通过 context.write() 存储在消息历史中)和 瞬时(通过 context.writeTransient()
仅展示但不持久化)两类。状态转换和通知通常是瞬时的;工具结果和文本内容是持久化的。
数据分片协议
Vercel AI SDK 已经定义了文本和工具调用分片的标准流式协议。Zapvol 在它之上加一层自定义的 数据分片事件,在内容流旁边捎带 agent 特有的元数据:
| 事件 | 载荷 | 用途 |
|---|---|---|
agent-init | Agent 配置、沙箱元数据、模型信息 | 在首个 token 之前初始化客户端 Agent 状态 |
agent-state | AgentStateKey 枚举值 | 驱动 UI 进度指示器的状态转换 |
notification | 消息字符串 + 严重级别 | 将 Agent 生成的通知以 toast 消息形式展示 |
title-updated | 自动生成的标题字符串 | 无需单独 API 调用即可更新任务标题 |
tool-stream | stdout/stderr 行分片 | 实时流式输出工具执行结果(如构建日志) |
所有事件与内容复用同一个 SSE 流,确保单个连接承载所有 Agent 到客户端的通信。
状态转换事件
toUIMessageStream() 函数将 AI SDK 流分片映射为状态事件:start → generating,start-step → executing,finish →
completed,error → error,abort → aborted。每次转换都作为瞬时的 agent-state 事件广播。
消息元数据
messageMetadata 回调将计时和用量数据附加到流消息上:
| 字段 | 设置时机 | 用途 |
|---|---|---|
streamCreatedAt | 流开始时 | 延迟测量基准 |
firstTokenAt | 首个文本/推理/工具调用时 | 首 token 时间指标 |
completedAt | 流结束时 | 总执行时长 |
totalUsage | 流结束时 | Token 用量统计(输入 + 输出) |
finishReason | 流结束时 | Agent 停止的原因 |
任务传输层——SSE(默认)vs WebSocket(可选)
核心构建一套流(task-runner.ts 的 startTaskUiStream);task-orchestrator.ts 上两层薄包装为它封帧:
| 包装层 | 入口 | 客户端经由 | 线上格式 |
|---|---|---|---|
execute | POST /api/tasks/:id/messages | DefaultChatTransport 搭标准 fetch | SSE 帧 data: <UIMessageChunk JSON>\n\n |
executeWs | /ws 上收到的 task:stream:start | DefaultChatTransport 搭 createWsFetch() 合成 body | WS text 帧 { type: "task:stream:frame", taskId, frame: UIMessageChunk } |
关键性质——格式等价:两条路径上 toUIMessageStream() 产出的 UIMessageChunk JSON 完全相同。SSE 加 data:
前缀,WS 外包一层信封。agent engine、writer、useChat 消费端对传输层完全无感。
怎么选
| 因素 | SSE | WebSocket |
|---|---|---|
| 默认启用 | 是:除非设置 VITE_TASK_WS_TRANSPORT=1 | 主动开启 |
| 断连续传 | 是:run produce 进 Redis stream buffer,重连读回(见下节) | 否:重连走 DB 历史,在途帧丢失 |
| 连接模型 | 每次任务一个 HTTP 请求 | 一条长连接,和 chat、未来功能共享 |
| 桌面端 | 不适用——桌面走 IPC | 不适用——customFetch 会绕过 WS 分支 |
| 多进程 | Redis Streams 做 resume | wsHub 的 Redis Pub/Sub 用于 chat 广播;任务帧经 wsHub.send(connectionId) 点对点回推,不过 Pub/Sub |
SSE 仍是推荐默认,因为它支持续传。WebSocket 适合将来的双向能力(更廉价的 abort 往返、实时光标提示、与 chat/presence 共用 socket),以及 SSE 被代理拦截的环境。
Abort 流程(两条路径语义一致)
客户端 Stop 按钮始终走 POST /api/tasks/:id/abort。该 handler 调用 taskOrchestrator.abort,顺序执行两步:
taskService.abort(userId, taskId)——校验归属 + 置isActive=false。未授权调用抛错。abortManager.abort(taskId)——触发runAgentLoop持有的AbortSignal,LLM 在下一个 await 点被打断。
WS 的 task:stream:abort 走同一个 orchestrator 方法,所以两条路径 Stop 行为完全一致。(客户端侧,useChat 内部 signal
abort 只会断开本地消费端,不会向服务端发 abort——匹配 SSE 断 HTTP 时的行为。)
可恢复 SSE(Redis Stream + pub/sub)
run 从不直接把流推给客户端——它 produce 进一个 Redis 支持的 StreamBuffer(redis-stream-buffer.ts),每个读者(live 的 execute 响应、以及中途重连)读的都是这同一个 buffer。正是这一点,让一次持续数分钟的 run 能扛过断连而无需重跑。
持久历史 + 实时推
buffer 对每个 chunk 合起来用两种 Redis 机制:
XADD把 chunk 追加进这一轮专属的一条 Redis Stream——这是(重)连读者回填用的持久历史。PUBLISH在XADD的同一刻把同一个 chunk 推到一个频道上——实时投递,事件驱动(不轮询、不阻塞XREAD),从而保住模型天然的 token 节奏。
一条共享的订阅连接经一张 dispatch map 多路复用所有频道——所以没有每流一条专用连接,这正是实时投递走 pub/sub + XRANGE 而非阻塞 XREAD 循环的原因。
读者如何自协调“历史↔实时”
读者与写入器之间没有握手——生产者可能是一个在读者连上之后才启动的 worker——所以它靠 Redis 单调的 stream id 给自己定序:
- 先 SUBSCRIBE(await 它),并缓冲这期间的实时推——快照之后 publish 的东西一个都不会漏。
- XRANGE 回填历史;记下发出的最后一个 id。
- 冲刷缓冲里
id > lastId的推送(去重),然后转入实时。
因为 stream id 单调,这三段是 id 连续的——无缺口、无重叠、无需去重表。正确性靠两个先后:生产者 XADD 先于 PUBLISH(id 来自 XADD),读者 SUBSCRIBE 先于 XRANGE。客户端生成的单调 id(<baseMs>-<seq>)让一个 chunk 的 XADD 与 PUBLISH 共用一条 pipeline。这替代了旧的每监听器 Promise 链去重(lib/resumable-stream.ts),改由消费端的 id 连续性保证——但 pub/sub 并没有消失,它就是实时推的那一半,如今配上 Stream 做持久历史。
生命周期与内存
一条 stream = 一轮 assistant 输出 的 chunks——produce 会先 DEL 掉旧键,所以一条 stream 绝不跨轮。每次写入刷新一个 grace TTL(600 秒),让长跑的一轮活着、也让崩溃的生产者(没写终止 entry)在此窗口内自然过期。终止 entry(经 XADD 写 end 字段)把键降到更短的保留窗口(300 秒)供迟到重连。若读者在生产者出现之前就已订阅,它在 120 秒初始等待上限后放弃。
| 操作 | Redis | 用途 |
|---|---|---|
| 追加 chunk | XADD + PUBLISH | 持久历史 entry + 实时推,共用一条 pipeline |
| 实时投递 | 频道上 SUBSCRIBE | 事件驱动推送,不轮询 |
| 回填 | 对 stream XRANGE | 为(重)连读者重放历史 |
| 完成 | 终止 end entry(XADD) | 向读者标记流结束 |
| 清理 | grace TTL(600 秒)→ 保留(300 秒) | 长跑的一轮保活;落定后进入重连窗口 |