流式传输架构

数据分片协议、传输无关的 StreamWriter,以及把 Redis Stream(持久历史)与 pub/sub(实时推)合起来、靠单调 id 让历史与实时无缝衔接的可恢复 SSE

引擎只管往一个流里写,送达是外层的事

Agent 一旦跑起来,输出就源源不断——文本 token、工具结果、状态转换、通知、元数据。这些都得可靠地送到客户端,还要在 Web(SSE)和桌面(IPC)上表现一致;但引擎本身不该操心网络怎么抖、走哪种传输。它只往一个传输无关的流里写,其余的分四层各管一段:

  1. StreamWriter——传输无关的事件写入接口;
  2. 数据分片协议——在 AI SDK 的内容分片旁边,捎带 agent 专属的事件;
  3. 任务传输层——SSE(默认)或 WebSocket,两条路径承载同一套帧格式;
  4. 可恢复 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-initAgent 配置、沙箱元数据、模型信息在首个 token 之前初始化客户端 Agent 状态
agent-stateAgentStateKey 枚举值驱动 UI 进度指示器的状态转换
notification消息字符串 + 严重级别将 Agent 生成的通知以 toast 消息形式展示
title-updated自动生成的标题字符串无需单独 API 调用即可更新任务标题
tool-streamstdout/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(可选)

任务流传输层 — SSE vs WebSocket 一个流核心 (buildUiStream),两层传输包装,相同的 UIMessageChunk 帧 SSE(默认) WebSocket(VITE_TASK_WS_TRANSPORT=1) ① 发送 ② 传输 ③ 服务端入口 ④ 服务端核心 ⑤ 线上帧 ⑥ 客户端解析 useChat + DefaultChatTransport(默认 fetch) fetch("POST /api/tasks/:id/messages") routes/tasks.ts → taskOrchestrator.execute buildUiStream → createUIMessageStreamResponse data: <UIMessageChunk JSON>\n\n 浏览器 SSE 解析器(DefaultChatTransport 内) useChat + DefaultChatTransport(createWsFetch) wsClient.send("task:stream:start", body) ws-service → taskOrchestrator.executeWs buildUiStream → reader.read() 循环 {type:"task:stream:frame", taskId, frame: ...} createWsFetch 重编码 → 同一 SSE 解析器 唯一 差异 汇合 — useChat 看到的 UIMessageChunk 完全一致 状态更新、渲染、hooks 全相同 —— 这条线以上代码对传输层零感知 续传:仅 SSE 路径(Redis 支持)。Abort:两条路径都调 taskOrchestrator.abort(共享 HTTP /abort)。

核心构建一套流task-runner.tsstartTaskUiStream);task-orchestrator.ts 上两层薄包装为它封帧:

包装层入口客户端经由线上格式
executePOST /api/tasks/:id/messagesDefaultChatTransport 搭标准 fetchSSE 帧 data: <UIMessageChunk JSON>\n\n
executeWs/ws 上收到的 task:stream:startDefaultChatTransportcreateWsFetch() 合成 bodyWS text 帧 { type: "task:stream:frame", taskId, frame: UIMessageChunk }

关键性质——格式等价:两条路径上 toUIMessageStream() 产出的 UIMessageChunk JSON 完全相同。SSE 加 data: 前缀,WS 外包一层信封。agent engine、writer、useChat 消费端对传输层完全无感。

怎么选

因素SSEWebSocket
默认启用是:除非设置 VITE_TASK_WS_TRANSPORT=1主动开启
断连续传是:run produce 进 Redis stream buffer,重连读回(见下节)否:重连走 DB 历史,在途帧丢失
连接模型每次任务一个 HTTP 请求一条长连接,和 chat、未来功能共享
桌面端不适用——桌面走 IPC不适用——customFetch 会绕过 WS 分支
多进程Redis Streams 做 resumewsHub 的 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,顺序执行两步:

  1. taskService.abort(userId, taskId)——校验归属 + 置 isActive=false。未授权调用抛错。
  2. 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 支持的 StreamBufferredis-stream-buffer.ts),每个读者(live 的 execute 响应、以及中途重连)读的都是这同一个 buffer。正是这一点,让一次持续数分钟的 run 能扛过断连而无需重跑。

可恢复 SSE —— 基于 Redis Streams 的流恢复 单调递增的 entry ID + cursor 连续性保证严格 FIFO 顺序 —— 零间隙,零重复 写入端 (Writer) 任务编排器 (Task Orchestrator) Redis 每个 run 一条 Stream(XADD / XREAD) Reader → 客户端 XREAD 从 cursor 读 → SSE 阶段 1 —— 正常流式 XADD 数据块1 → id 1-0 XREAD → 数据块1(cursor=1-0) XADD 数据块2 → id 2-0 XREAD → 数据块2(cursor=2-0) XADD 数据块3 ... XREAD → 数据块3 ... 阶段 2 —— 网络中断 XADD 数据块4(缓冲进 stream) XADD 数据块5(缓冲进 stream) 已断连 阶段 3 —— 客户端从上次 cursor 重连 重连 · XREAD 从 cursor 3-0 回放 3-0..5-0(历史),再阻塞等待实时 XADD 数据块6(实时) XREAD → 数据块6(cursor=6-0) XADD end-entry(终止) XREAD → end-entry → DONE 顺序保证 单调 stream ID + 每个 reader 的 cursor —— 历史与实时 entry 从同一个 XREAD 循环读出。 客户端感知 恢复的流 = 未中断的流。零间隙,零重复。

持久历史 + 实时推

buffer 对每个 chunk 合起来用两种 Redis 机制:

  • XADD 把 chunk 追加进这一轮专属的一条 Redis Stream——这是(重)连读者回填用的持久历史
  • PUBLISHXADD 的同一刻把同一个 chunk 推到一个频道上——实时投递,事件驱动(不轮询、不阻塞 XREAD),从而保住模型天然的 token 节奏。

一条共享的订阅连接经一张 dispatch map 多路复用所有频道——所以没有每流一条专用连接,这正是实时投递走 pub/sub + XRANGE 而非阻塞 XREAD 循环的原因。

读者如何自协调“历史↔实时”

读者与写入器之间没有握手——生产者可能是一个在读者连上之后才启动的 worker——所以它靠 Redis 单调的 stream id 给自己定序:

  1. 先 SUBSCRIBE(await 它),并缓冲这期间的实时推——快照之后 publish 的东西一个都不会漏。
  2. XRANGE 回填历史;记下发出的最后一个 id。
  3. 冲刷缓冲里 id > lastId 的推送(去重),然后转入实时。

因为 stream id 单调,这三段是 id 连续的——无缺口、无重叠、无需去重表。正确性靠两个先后:生产者 XADD 先于 PUBLISH(id 来自 XADD),读者 SUBSCRIBE 先于 XRANGE。客户端生成的单调 id(<baseMs>-<seq>)让一个 chunk 的 XADDPUBLISH 共用一条 pipeline。这替代了旧的每监听器 Promise 链去重(lib/resumable-stream.ts),改由消费端的 id 连续性保证——但 pub/sub 并没有消失,它就是实时推的那一半,如今配上 Stream 做持久历史。

生命周期与内存

一条 stream = 一轮 assistant 输出 的 chunks——produce 会先 DEL 掉旧键,所以一条 stream 绝不跨轮。每次写入刷新一个 grace TTL(600 秒),让长跑的一轮活着、也让崩溃的生产者(没写终止 entry)在此窗口内自然过期。终止 entry(经 XADDend 字段)把键降到更短的保留窗口(300 秒)供迟到重连。若读者在生产者出现之前就已订阅,它在 120 秒初始等待上限后放弃。

操作Redis用途
追加 chunkXADD + PUBLISH持久历史 entry + 实时推,共用一条 pipeline
实时投递频道上 SUBSCRIBE事件驱动推送,不轮询
回填对 stream XRANGE为(重)连读者重放历史
完成终止 end entry(XADD向读者标记流结束
清理grace TTL(600 秒)→ 保留(300 秒)长跑的一轮保活;落定后进入重连窗口
这页有帮助吗?