数据截至 (上游 commit 2822885e57e7)
chunk → parts 引擎:StreamProcessor 如何拼出 UIMessage
30 秒导读: 上一章的 connection adapter 把服务端的线协议解成一串 AG-UI chunk 事件(
TEXT_MESSAGE_CONTENT、TOOL_CALL_ARGS……)。这一章的StreamProcessor是把这串增量、乱序、会丢事件的流,一条条累加成 UI 真正拿去渲染的那个数据结构——UIMessage.parts。它是整个前端"看得见的消息"的唯一来源。
本章讲透 packages/ai/src/activities/chat/stream/ 这个目录:核心是 processor.ts 里约 2570 行的 StreamProcessor 类,配套是 message-updaters.ts(纯函数改 parts)、strategies.ts(节流)、json-parser.ts(容忍半截 JSON)、types.ts(内部状态类型)。
这些 parts 怎么被 hook 消费在第 4 章,怎么被渲染成 DOM 在第 5 章。本章只负责一件事:把 chunk 变成 parts。
1. 这是什么(零基础也能懂)
1.1 先认识产物:parts 化的消息
在很多聊天 SDK 里,一条 AI 消息就是一个字符串 content。TanStack AI 不是——它的一条消息是一个 parts 数组:文本、思考、工具调用、工具结果、结构化输出、图片……每一样都是数组里的一个 part。
UIMessage 的形状(packages/ai-client/src/types.ts:623 UIMessage,运行时同构定义在 packages/ai/src/types.ts:560):
UIMessage {
id: "msg_abc"
role: "assistant"
parts: [
{ type: "thinking", content: "我先查一下天气……" }
{ type: "text", content: "让我帮你查一下。" }
{ type: "tool-call", id: "call_1", name: "getWeather", state: "input-complete", arguments: '{"city":"上海"}' }
{ type: "tool-result", toolCallId: "call_1", content: "18°C 多云" }
{ type: "text", content: "上海现在 18 度,多云。" }
]
}
为什么要拆成 parts? 因为 UI 要对不同内容做不同的事:思考块折叠、工具调用画成卡片、文本做 markdown 渲染。一坨字符串做不到,parts 数组天然给了 UI 分类渲染的抓手。
parts 的类型定义(packages/ai/src/types.ts:504 MessagePart):
| part 类型 | 装什么 | 关键字段 |
|---|---|---|
text | 助手回复正文 | content |
thinking | 推理/思考内容 | content · stepId · signature |
tool-call | 一次工具调用 | id · name · arguments(JSON 串)· state · output · approval |
tool-result | 工具执行结果(喂回 LLM 用) | toolCallId · content · state |
structured-output | 流式结构化输出 | status · raw · partial · data |
ui-resource | MCP Apps 的 ui:// 组件 | resource · toolCallId · toolName |
image/audio/video/document | 多模态 | 各自的 source |
1.2 再认识难点:输入是"碎的、乱的、可能缺的"
StreamProcessor 拿到的不是这个漂亮结构,而是一串碎片事件。同样一条消息,在流里长这样:
TEXT_MESSAGE_START { messageId: "msg_abc" }
TEXT_MESSAGE_CONTENT { delta: "让我" }
TEXT_MESSAGE_CONTENT { delta: "帮你" }
TEXT_MESSAGE_CONTENT { delta: "查一下。" }
TOOL_CALL_START { toolCallId: "call_1", toolCallName: "getWeather" }
TOOL_CALL_ARGS { toolCallId: "call_1", delta: '{"ci' }
TOOL_CALL_ARGS { toolCallId: "call_1", delta: 'ty":"上海"}' }
TOOL_CALL_END { toolCallId: "call_1" }
RUN_FINISHED { runId: "run_1" }
难点有三层,后面每一节都在解决它们:
- 增量: 文本和工具参数都是一小段一小段来的,要累加。
- 半截: 工具参数
{"ci还不是合法 JSON,但 UI 想边流边预览,需要容错解析。 - 可能缺:
TOOL_CALL_END可能因为适配器 bug 或断流永远不来,状态机必须有兜底。
1.3 一句话直觉
把 StreamProcessor 想成一个专门给聊天流做的 reducer: state 是 UIMessage[],每个 chunk 是一个 action,processChunk 是 dispatch,每个 handler 就是一个 case,message-updaters.ts 里的纯函数是 reducer 本体。每次 state 变了就 onMessagesChange 通知外面重渲染。
2. 顶层全景(它大概怎么转)
2.1 一张图:从 chunk 到 UI
先看整条数据流。从上到下读:事件进来,分派到 handler,handler 调纯函数改 parts,改完广播出去。
connection adapter(第 2 章)吐出的 chunk 流
│
▼
┌──────────────────────────────────────────────┐
│ StreamProcessor.process(stream) │
│ for await (chunk of stream) processChunk() │
└──────────────────────────────────────────────┘
│
▼
┌──────────────────────────────────────────────┐
│ processChunk(): 一个大 switch(chunk.type) │ ← 中央分派
│ TEXT_* │ TOOL_CALL_* │ REASONING_* │ CUSTOM │
│ RUN_* │ STEP_* │ MESSAGES_SNAPSHOT │ ... │
└──────────────────────────────────────────────┘
│ 每个 case → 一个 handler
▼
┌────────────────────────────── ────────────────┐
│ handler:更新两处状态 │
│ (a) messageStates → 每条消息的流式草稿状态 │
│ (b) this.messages → 通过 message-updaters │
│ 纯函数拼出的 UIMessage[] │
└──────────────────────────────────────────────┘
│
▼
events.onMessagesChange([...messages]) ← 广播全量数组
│
▼
useChat(第 4 章) → UI 组件(第 5 章)渲染
怎么读这张图: 中间那个 switch 是心脏;它左手维护一份"内部草稿"(messageStates),右手把草稿投影成对外的 UIMessage[],每次投影完就整份广播出去。
2.2 五个文件各干什么
| 文件 | 职责 | 关键符号 |
|---|---|---|
processor.ts | 状态机主体:分派 chunk、维护流式状态、拼装 parts | StreamProcessor · processChunk |
message-updaters.ts | 一组纯函数,输入旧 messages + 一个改动,返回新 messages | updateTextPart · updateToolCallPart |
strategies.ts | 决定"文本累加到什么程度才广播一次",做 UI 节流 | ImmediateStrategy · PunctuationStrategy |
json-parser.ts | 用 partial-json 库解析半截的工具参数 | parsePartialJSON |
types.ts | 内部状态类型(不是对外的 UIMessage) | MessageStreamState · InternalToolCallState |
2.3 一条最短主线
以"助手回一句纯文本"为例,走一遍不进代码:
TEXT_MESSAGE_START到达 → 建一条空的 assistantUIMessage。- 每个
TEXT_MESSAGE_CONTENT→ 把delta累加进草稿,按策略择机把最新文本写进那条消息的textpart。 RUN_FINISHED→finalizeStream()收尾,冲刷未广播的文本,触发onStreamEnd。
真正复杂的是工具调用、思考、结构化输出这三条支线,下面逐个拆。
3. 核心原理
3.0 双份状态:草稿态 vs 投影态
进入细节前,先建立最重要的一个心智模型:StreamProcessor 同时维护两份状态,这是它所有逻辑的地基。
| 草稿态(内部) | 投影态(对外) | |
|---|---|---|
| 是什么 | messageStates: Map<msgId, MessageStreamState> | messages: UIMessage[] |
| 存什么 | 累加缓冲:currentSegmentText、toolCalls Map、thinkingSteps Map | UI 真正渲染的 parts |
| 谁读 | 只有 processor 自己 | 通过 onMessagesChange 给外部 |
| 定义 | packages/ai/src/activities/chat/stream/types.ts:58 MessageStreamState | packages/ai/src/types.ts:560 UIMessage |
字段声明在 processor.ts:182-217:messages、messageStates、activeMessageIds、toolCallToMessage、structuredMessageIds 等。
为什么要两份? 因为累加逻辑(比如"上一段文本是什么 、要不要另起新段")是脏活,不该塞进对外的 UIMessage。草稿态承担脏活,投影态保持干净、随时可渲染。后面每个机制都是"更新草稿 → 投影成 parts"这一个套路的变体。
3.1 中央分派:processChunk 的 switch
它要解决的小问题: 一个流里混着十几种事件类型,得把每种路由到对的处理逻辑。
思路: 一个大 switch(chunk.type),一 case 一 handler,没列到的类型故意忽略(如 STATE_SNAPSHOT、STATE_DELTA)。
真实实现在 processor.ts:538 processChunk。它先(可选)录制 chunk 供回放测试,然后分派。事件分几组:
| 事件组 | 成员 | handler |
|---|---|---|
| 文本 | TEXT_MESSAGE_START / _CONTENT / _END | handleTextMessage*Event |
| 工具 | TOOL_CALL_START / _ARGS / _END / _RESULT | handleToolCall*Event |
| 推理 | REASONING_MESSAGE_CONTENT、STEP_STARTED / _FINISHED | handleReasoning* / handleStep* |
| 运行 | RUN_STARTED / _FINISHED / _ERROR | handleRun*Event |
| 快照/自定义 | MESSAGES_SNAPSHOT、CUSTOM | handleMessagesSnapshotEvent · handleCustomEvent |
注意 REASONING_START / _END 等只是 break 掉不处理(processor.ts:626-631)——它们是纯边界信号,没有需要累加的内容。
3.2 逐消息状态机:懒创建与消息路由
它要解决的小问题: 一次流里可能有多条助手消息(自动续跑、多智能体),事件又常常不带 messageId。得知道"这个 chunk 属于哪条消息"。
关键设计一:懒创建。 prepareAssistantMessage()(processor.ts:302)不立刻建消息,只重置流式状态。真正的消息在第一个有内容的 chunk 到达时,由 ensureAssistantMessage() 懒创建。
为什么? 注释说得很直白(processor.ts:295-301):自动续跑有时不产生任何内容,提前建消息会让 UI 闪一条空消息。懒创建避免这个抖动。
关键设计二:三张路由表。 事件不带 messageId 时,靠这几张表把它落到对的消息:
| 表 | 键 → 值 | 作用 |
|---|---|---|
activeMessageIds | Set | 当前活跃的消息;getActiveAssistantMessageId() 从尾部找最近的助手消息 |
toolCallToMessage | toolCallId → msgId | TOOL_CALL_ARGS/END 不带 msgId,靠它反查 |
messageStates | msgId → 草稿态 | 每条消息的累加缓冲 |
ensureAssistantMessage()(processor.ts:743)是落地点,它按优先级找目标消息:① 传入的 preferredId → ② 活跃助手消息 → ③ 断线重连时从已存在的 this.messages 里"水合"出草稿态(把已有文本塞回 currentSegmentText,让后续 delta 接着拼,processor.ts:761-785)→ ④ 都没有就新建一条。第 ③ 步是重连/续跑不重复建消息的关键。
幂等清理: removeMessagesAfter()(processor.ts:454)在 reload/retry 时,会把上面几张路由表里指向已删消息的条目一并删掉——注释(processor.ts:456-464)解释了原因:否则续跑的流可能把 delta 落到已失效的表项上,污染新消息。
3.3 文本累加:delta 累加、段落切分、节流
它要解决的小问题: 文本是一段段来的,要累加;工具调用之后又可能有新一段文本,不能和前一段混在一起。
真实实现在 processor.ts:1275 handleTextMessageContentEvent。核心逻辑三步:
① 只认 delta。 文本增量只从 chunk.delta 累加(processor.ts:1316-1317)——早期的"content 兜底 + startsWith 累积值检测"已移除,线协议约定文本增量一律走 delta。
② 段落切分。 如果这段文本前面发生过工具调用(hasToolCallsSinceTextStart),且检测到是新段(isNewTextSegment,processor.ts:2111),就先把旧段广播出去,再清空缓冲开新段(processor.ts:1296-1313)。这样工具调用前后的两段文本会变成 parts 数组里两个独立的 text part。
③ 节流广播。 不是每个 delta 都广播,而是问一句策略:
// 示意,非源码 —— 节流的核心判断
const shouldEmit = chunkStrategy.shouldEmit(chunkPortion, currentSegmentText)
if (shouldEmit && currentSegmentText !== lastEmittedText) {
emitTextUpdateForMessage(messageId) // 这步才真正改 parts + 广播
}
对应 processor.ts:1324-1331。策略由 strategies.ts 提供:
| 策略 | 何时广播 | 用途 |
|---|---|---|
ImmediateStrategy(默认) | 每个 chunk 都广播 | 最实时 |
PunctuationStrategy | chunk 含 .,!?;:\n 时 | 按句子自然停顿 |
BatchStrategy | 每 N 个 chunk | 降低 UI 更新频率 |
WordBoundaryStrategy | chunk 以空白结尾 | 不切断单词 |
CompositeStrategy | 任一子策略同意(OR) | 组合 |
策略接口极简(strategies.ts + types.ts:38 ChunkStrategy):就一个 shouldEmit(chunk, accumulated) => boolean。
投影那一步(replace vs push): emitTextUpdateForMessage(processor.ts:2251)调 updateTextPart(message-updaters.ts:25),它有个关键判断——如果消息最后一个 part 是 text 就替换它(同段续写),否则 push 一个新 text part(工具调用后的新段):
// message-updaters.ts:38-44 —— 决定续写还是另起
if (lastPart && lastPart.type === 'text') {
parts[parts.length - 1] = { type: 'text', content } // 同段:替换
} else {
parts.push({ type: 'text', content }) // 新段:追加
}
这一行就是"parts 数组里为什么文本会分成好几段"的答案。
3.4 工具调用:三段生命周期与状态机
工具调用是本章最复杂的一条支线。它的 part 有一个 state 字段,在流里逐步推进。
状态机全貌(取值定义在 packages/ai-client/src/types.ts:342 ToolCallState):
TOOL_CALL_START TOOL_CALL_ARGS(首个 delta) TOOL_CALL_END
│ │ │
▼ ▼ ▼
awaiting-input ──────► input-streaming ──────── ──► input-complete
│
(若工具需要审批,走 CUSTOM 事件) │
approval-requested ─► approval-responded
│
(拿到结果 addToolResult / RESULT) ▼
complete / error
三个 handler 对应三段:
① START(processor.ts:1347): 在草稿态 toolCalls Map 里建一条 InternalToolCallState,同时 append 一个 tool-call part,状态 awaiting-input。这里做了两件防御:把 TOOL_CALL_START 上带到的 provider metadata 存进 part(比如 Gemini 的 thoughtSignature,下一轮原样回传,processor.ts:1369-1372);把 toolCallId → messageId 存进路由表供后续 ARGS/END 反查(processor.ts:1388)。注释强调 START 必须先于 ARGS 到,否则 ARGS 因查不到 ID 被静默丢弃(processor.ts:1340-1341)。
② ARGS(processor.ts:1422): 把 delta 拼进 arguments 字符串,首个非空 delta 把状态推到 input-streaming,然后每次都试着 partial-parse 一遍供 UI 预览:
// processor.ts:1446-1448 —— 边流边解析半截 JSON
existingToolCall.parsedArguments = this.jsonParser.parse(
existingToolCall.arguments, // 可能还是 '{"ci' 这种半截
)
解析器是 json-parser.ts:56 parsePartialJSON,底层用 partial-json 库;解析失败返回 undefined,不抛错(json-parser.ts:31-43) ——早期流数据太少解析失败是预期行为。
③ END(processor.ts:1480)—— 参数定稿 + canonical input 回填: 早期版本里 TOOL_CALL_END 有"带不带 result"的双重角色;现在结果一律走 TOOL_CALL_RESULT,END 只负责把参数定稿——handler 顶部注释(processor.ts:1468-1479)写明 "Tool output arrives on TOOL_CALL_RESULT, not on this event"。END 的新活是:若携带了解析好的 input(spec 字段或 metadata.tanstack.input),就在没见过 ARGS delta 时回填 arguments 串(Anthropic server_tool_use / web_search 这类一次性给全 input 的适配器,issue #839),并用 canonical 值覆盖渲染 part 的 input(processor.ts:1496-1531)。
结果到达:TOOL_CALL_RESULT(processor.ts:1553) 做"一份给 UI、一份给 LLM"的双写(processor.ts:1577-1600):updateToolCallWithOutput 更新 tool-call part 的 output 字段(给 UI 看),再 updateToolResultPart 建独立的 tool-result part(给下一轮 LLM 看)。这个分裂是理解工具结果的关键。
3.5 客户端工具与审批:CUSTOM 事件
它要解决的小问题: 有些工具要在浏览器端执行,有些要用户点确认才能跑。这些不是标准 AG-UI 事件,走 CUSTOM 通道。
handleCustomEvent(processor.ts:1954)按 chunk.name 分派:
name | 含义 | 动作 |
|---|---|---|
tool-input-available | 客户端工具待执行 | 触发 events.onToolCall,让宿主跑工具 |
approval-requested | 工具需审批 | 把 part 置 approval-requested,触发 onApprovalRequest |
structured-output.start / .complete | 结构化输出边界 | 见 3.7 |
ui-resource | MCP Apps 组件 | 在消息上挂一个 ui-resource part |
| 其它 | 用户自定义 | 转发给 events.onCustomEvent |
注意 handler 顶部注释(processor.ts:1940-1952):tool-input-available / approval-requested 如今是 legacy/replay 兼容入口——现行核心流用 RUN_FINISHED.outcome.type === 'interrupt' 表达"等用户动作",这两个 CUSTOM 事件不再是新发射的权威来源。
宿主处理完后,回调两个入口把结果写回消息:
addToolResult(toolCallId, output, error?)(processor.ts:339):客户端工具执行完,写output+ 建tool-resultpart。addToolApprovalResponse(approvalId, approved)(processor.ts:383):用户批/拒后,把 part 置approval-responded。
审批事件里的一个细节: RUN_FINISHED 之后 activeMessageIds 已清空,所以 approval-requested 解析目标消息时会回退到 toolCallToMessage 表(processor.ts:2029-2033)——因为那张表在 finalize 后仍保留。