流式协议:AG-UI 事件语言与状态累积器
30 秒导读: 后端不是一次性把「最终答案」发给前端,而是把「界面正在发生的每一点变化」拆成一串小事件(文字多了一个字、某个组件的某个 prop 变了)边算边发。前端有一个纯函数累积器(
streamReducer),把这一串增量事件逐个 reduce成一份完整的、随时可渲染的线程状态TamboThread。本章讲透这条「事件 → 线程状态」的数据管线。
本章聚焦数据怎么从事件变成状态。不讲组件如何被渲染成活的 React 元素(见 第 5 章),也不讲客户端工具怎么被触发执行(见 第 4 章)。承接 第 2 章——后端「大脑」挑好组件、开始流式吐 props 之后,吐出来的就是本章要解码的事件流。
1. 这是什么(零基础 也能懂)
一句话定义
AG-UI 事件是后端和前端之间约定的一套「界面变化的语言」;状态累积器是前端把这套语言「听懂并记下来」的机器。
它解决什么问题
想象你在看 AI 一边思考一边回答。如果后端要等整段回答 + 整个图表都算完再一次性发给你,你会盯着空白屏幕等好几秒。流式(streaming)的做法是:算出一点就发一点。
但「发一点」带来一个新麻烦:前端收到的是一堆碎片——
- 「新开一条 assistant 消息」
- 「这条消息的文字后面加上『你』」
- 「……再加上『好』」
- 「开始画一个
WeatherCard组件」 - 「它的
props.temperature从无变成2」 - 「……变成
23」
前端必须把这一串碎片攒成一个完整的、能直接渲染的东西。这个「攒」的过程,就是累积(accumulation)。
一句话直觉/类比
把它想成看直播弹幕拼一幅画:后端不断发「在坐标 (x,y) 涂一笔什么颜色」的指令(事件),前端一笔一笔照着画(累积),画布(
TamboThread)就随之越来越完整。你不需要它发整幅画,只要发「下一笔」。
用起来什么样
对应用开发者,这条管线基本是透明的——你 for await 一个流,每个事件后都能拿到一份「当前完整快照」:
// 示意,非源码:消费流的两种方式
// 1) 异步迭代:每个事件后拿到一份线程快照
for await (const { event, snapshot } of stream) {
console.log(event.type); // 例如 "TEXT_MESSAGE_CONTENT"
render(snapshot); // snapshot 是"累积到此刻"的完整 TamboThread
}
// 2) 只要最终结果
const finalThread = await stream.thread;
snapshot 每次都是累积到当前的完整状态,而不是「这一个事件本身」。这正是累积器的价值:把增量事件流,变成一连串「随时可渲染的完整状态」。
2. 顶层全景(它大概怎么转)
这条管线怎么读
从左到右是一次运行的数据流向:后端 SSE 事件流,经过一点轻加工,喂进纯函数 reducer,reducer 把它并进 TamboThread,组 件订阅这份状态。
后端 (SSE) 客户端处理循环 (tambo-stream.ts) 纯函数累积器 UI
────────── ────────────────────────────────────── ───────────────── ──────
streamReducer(
一串 handleEventStream() ── 逐个 event ──► state, action) 订阅
AG-UI ───► ├─ toolTracker.handleEvent(e) │ ──► TamboThread
事件 ├─ 若 TOOL_CALL_ARGS: 预解析 partial-json ├─ 按 event.type (messages,
(增量) │ → parsedToolArgs │ 分派 handler streaming,
├─ dispatch({type:'EVENT', event, ...}) ───► │ accumulating
│ (parsedToolArgs / toolSchemas 随行) └─ 返回**新** state ToolArgs)
└─ 若 props/state 增量: keyed-throttle 节流
一句话:处理循环负责取事件 + 轻加工,reducer 负责把事件折叠进状态,两者严格分工。
部件一句话职责
| 部件 | 干什么 | 文件 |
|---|---|---|
handleEventStream | 把 SDK 的异步流原样透传为 AGUIEvent,顺带 debug 日志 | packages/client/src/utils/stream-handler.ts:47 |
| 处理循环 | for await 每个事件,预解析工具参数,dispatch 给 reducer,再节流触发流式工具 | packages/client/src/tambo-stream.ts:312 |
streamReducer | 核心:(state, action) => newState,把一个事件折叠进线程状态 | packages/client/src/utils/event-accumulator.ts:386 |
ThreadState | 单个线程的累积容器:thread + streaming + accumulatingToolArgs | event-accumulator.ts:61 |
| 事件类型定义 | AG-UI 标准事件(@ag-ui/core)+ Tambo 自定义事件 | packages/client/src/types/event.ts |
applyJsonPatch | 把 props_delta 里的 RFC 6902 patch 应用到当前 props | packages/client/src/utils/json-patch.ts:19 |
parsePartialJson | 把「还没收完」的半截 JSON 也尽量解析出来 | 第三方 partial-json(在 event-accumulator.ts:42 引入) |
unstrictifyToolCallParamsFromSchema | 把严格模式 schema 塞进来的 null 还原掉 | packages/client/src/utils/unstrictify.ts:160 |
createKeyedThrottle | 按 key 独立节流,控制流式工具的重复执行频率 | packages/client/src/utils/keyed-throttle.ts:54 |
主线走一遍(高层)
- 后端流出一串 AG-UI 事件(
RUN_STARTED→ 若干TEXT_MESSAGE_*/TOOL_CALL_*/tambo.component.*→RUN_FINISHED)。 tambo-stream.ts的处理循环逐个接住,给TOOL_CALL_ARGS事件预解析部分 JSON,然后dispatch一个EVENTaction。streamReducer按event.type分派到对应 handler,不可变地返回一份新的StreamState。- 每次返回的新状态里,那条正在生成的消息 / 组件都更「完整」了一点。UI 订阅这份状态,自然随之刷新。
3. 核心原理
本章的核心是四个机制,由浅入深:① 事件语言(有哪些事件)→ ② 累积器骨架(reducer 怎么组织)→ ③ 增量如何落地(文字追加、props patch、工具参数攒 JSON)→ ④ 三个易读性技巧(部分解析 / 反严格化 / 节流)。
3.1 事件语言:AG-UI 标准事件 + Tambo 扩展
它要解决的小问题
前后端要对「界面正在发生什么」达成同一套词汇。Tambo 复用了开源的 AG-UI(Agent-UI 协议,@ag-ui/core)标准事件,再加上自己的一组自定义事件来表达「生成式组件」这种 AG-UI 没有的概念。
两类事件
第一类:AG-UI 标准事件,reducer 直接 switch (event.type) 分派(event-accumulator.ts:563):
| 事件 | 含义 |
|---|---|
RUN_STARTED / RUN_FINISHED / RUN_ERROR | 一次运行的开始 / 正常结束 / 出错 |
TEXT_MESSAGE_START / _CONTENT / _END | 一条文字消息的开始 / 追加一段 delta / 结束 |
TOOL_CALL_START / _ARGS / _END / _RESULT | 一次工具调用的开始 / 参数增量 / 参数结束 / 结果 |
THINKING_TEXT_MESSAGE_START / _CONTENT / _END | 模型「思考」(reasoning)内容的开始 / 增量 / 结束 |
start/content/end 这种「三段式」是 AG-UI 的通用节奏:开一个坑 → 往坑里灌增量 → 封坑。
第二类:Tambo 自定义事件,统一裹在 AG-UI 的 CUSTOM 事件里,靠 name 区分(types/event.ts:101):
name | 含义 |
|---|---|
tambo.component.start | 开始渲染一个生成式组件(带 componentName / componentId) |
tambo.component.props_delta | 组件 props 的一次 JSON Patch 增量 |
tambo.component.state_delta | 组件 state 的一次 JSON Patch 增量 |
tambo.component.end | 该组件流式结束 |
tambo.run.awaiting_input | 运行暂停,等客户端执行工具(交给第 4 章) |
tambo.message.parent | 声明某消息是「在生成另一条消息期间」产生的(MCP sampling / elicitation) |
类型守卫:把泛泛的 CustomEvent 收窄成 Tambo 事件
AG-UI 的 CustomEvent 只知道有个 name: string,不知道 value 长啥样。asTamboCustomEvent 用一个 switch (event.name) 把它收窄成精确类型,遇到不认识的 name 返回 undefined:
// packages/client/src/types/event.ts:143 asTamboCustomEvent(节选)
switch (event.name) {
case "tambo.component.start":
return event as ComponentStartEvent;
case "tambo.component.props_delta":
return event as ComponentPropsDeltaEvent;
// …其余分支…
default:
return undefined; // 不是已知 Tambo 事件
}
reducer 的 handleCustomEvent 拿到收窄后的事件,再对 customEvent.name 做一次 switch 分派,并用 never 做穷尽性检查——将来加了新自定义事件却忘了处理,TypeScript 会在编译期报错(event-accumulator.ts:1288)。这是「fail-fast、不留静默兜底」原则的体现。
3.2 累积器骨架:streamReducer 与 ThreadState
它要解决的小问题
怎么把「一串事件」变成「一份状态」,而且要可预测、可测试、不出并发怪 bug?答案是借用 React useReducer 那套模型:一个纯函数 (state, action) => newState,永远返回新对象、绝不原地改。
状态长什么样
顶层状态 StreamState 是「多线程」的——一个 app 可以同时有多个会话:
StreamState
├─ currentThreadId: string ← 当前激活的线程
└─ threadMap: Record<threadId, ThreadState>
│
ThreadState (event-accumulator.ts:61)
├─ thread: TamboThread ← 会渲染的东西:messages[]、status…
├─ streaming: StreamingState ← 瞬时流式状态:runId、messageId、error…
├─ accumulatingToolArgs: Record<toolCallId, string>
│ ← 正在拼的工具参数 JSON 字符串
└─ lastCompletedRunId?: string
三个字段各司其职,值得记住它们的分工:
thread—— 要持久渲染的内容:消息列表、每条消息的 content 块、线程状态。streaming—— 瞬时的流式元信息:当前runId、正在写哪条messageId、reasoning 起始时间、错误。运行结束会被清成idle。accumulatingToolArgs—— 一个临时暂存区:工具参数是一段一段流下来的 JSON 字符串,得先攒成完整串才能JSON.parse;攒的过程就放这里,TOOL_CALL_END时清掉。
骨架:先处理非事件 action,再 switch 事件
streamReducer 分两段(event-accumulator.ts:386):
- 前半段
switch (action.type):处理INIT_THREAD/SET_CURRENT_THREAD/LOAD_THREAD_MESSAGES等管理性 action;EVENT走break落到后半段。 - 后半段
switch (event.type):真正的事件累积,每个事件类型对应一个handleXxx纯函数,返回新的ThreadState,最后统一并回threadMap。
// event-accumulator.ts:563 事件分派骨架(节选)
switch (event.type) {
case EventType.RUN_STARTED:
updatedThreadState = handleRunStarted(threadState, event);
break;
case EventType.TEXT_MESSAGE_CONTENT:
updatedThreadState = handleTextMessageContent(threadState, event);
break;
// …
default: {
const _exhaustiveCheck: never = event; // 穷尽性检查
throw new UnreachableCaseError(_exhaustiveCheck);
}
}
一个巧思:占位线程「迁移」
新会话在拿到服务器真实 threadId 之前,为了立刻把用户消息显示出来,会先塞进一个占位线程 PLACEHOLDER_THREAD_ID(event-accumulator.ts:214)。等 RUN_STARTED 带回真实 id,reducer 把占位线程里的消息迁移到真实线程,再清空占位线程(event-accumulator.ts:523)。这就是聊天框里「点发送,消息秒出现」的乐观 UI(optimistic UI)背后的机制。
3.3 增量如何落地(三种典型累积)
同样是「攒」,不同内容攒法不同。这里对比三种最典型的:
| 内容 | 事件 | 累积方式 | handler |
|---|---|---|---|
| 文字 | TEXT_MESSAGE_CONTENT | 字符串拼接(旧 + delta) | handleTextMessageContent |
| 组件 props/state | tambo.component.props_delta | JSON Patch 应用到当前对象 | handleComponentDelta |
| 工具参数 | TOOL_CALL_ARGS | 攒 JSON 字符串 + 部分解析 | handleToolCallArgs |
① 文字:最简单的拼接
找到当前消息的最后一个 content 块,如果是 text 就把 delta 接到尾巴上,否则新开一个 text 块:
// event-accumulator.ts:884 handleTextMessageContent(节选)
const updatedContent: Content[] = isTextBlock
? [...content.slice(0, -1), { ...lastContent, text: lastContent.text + event.delta }]
: [...content, { type: "text", text: event.delta }];
注意全程用展开 + slice 造新数组,不 push——这是整份文件的铁律:不可变更新。
② 组件 props:JSON Patch 增量
组件的 props 不是简单追加,而是「改某个字段」。后端发的是 RFC 6902 JSON Patch 操作数组(如 [{op:"replace", path:"/temperature", value:23}])。handleComponentDelta 找到组件块,把 patch 应用上去,并把 streamingState 标为 "streaming":
// event-accumulator.ts:1394 handleComponentDelta(节选)
const currentValue = field === "props"
? componentContent.props
: (componentContent.state ?? {});
const updatedValue = applyJsonPatch(currentValue, operations);
applyJsonPatch 是对 fast-json-patch 的薄封装,关键参数是 mutate = false——让它克隆而非原地改,并开 validate = true 校验,patch 失败时抛带上下文的错误(json-patch.ts:19)。props_delta 和 state_delta 共用这一个 handler,只靠 field 参数区分,是「Rule of Three 之前不硬抽象」但这里两者确实同构的合理复用。
③ 工具参数:攒 JSON 字符串,末尾权威解析
工具调用的参数(如 {"city":"Tokyo","units":"metric"})也是一段段 delta 流下来的半截 JSON 字符串。handleToolCallArgs 做两件事(event-accumulator.ts:1028):
- 把 delta 拼进
accumulatingToolArgs[toolCallId](纯字符串累加)。 - 尝试乐观地部分解析当前这段半截 JSON,好让 UI 边流边显示参数。
到 TOOL_CALL_END,才用标准 JSON.parse 对完整串做一次权威解析,并清掉暂存(event-accumulator.ts:1125)。「流式期间尽量解析、结束时权威解析」这个两段式,是下一节两个技巧的舞台。
3.4 三个易读性技巧
这三个技巧都服务同一个目标:流式过程中,让 UI 尽早、且尽量正确地显示出来。
① partial-json:半截 JSON 也能解析
标准 JSON.parse('{"city":"Tok') 会直接抛错。但流式时我们手上就是这种半截串。第三方库 partial-json 的 parse 能容忍不完整,尽力返回 {city:"Tok"}。reducer 只接受解析结果是「非数组的对象」,否则保持不变:
// event-accumulator.ts:1046 handleToolCallArgs(节选)
try {
const parsed: unknown = parsePartialJson(newAccumulatedJson);
if (typeof parsed === "object" && parsed !== null && !Array.isArray(parsed)) {
parsedInput = parsed as Record<string, unknown>;
}
} catch { /* 还解析不出来 — 保持 input 不变 */ }
一处小优化: 处理循环在
dispatch之前已经替TOOL_CALL_ARGS预解析过一次(tambo-stream.ts:342调toolTracker.parsePartialArgs),把结果作为action.parsedToolArgs随事件传进来。reducer 若收到就直接用,避免重复解析(event-accumulator.ts:1044)。
② unstrictify:还原严格模式塞进来的 null
这是个不显然但很实际的坑。当用 OpenAI 结构化输出(structured outputs) 严格模式时,所有可选参数会被改造成「必填 + 可为 null」,于是模型对没填的参数发 null。若直接把这些 null 交给组件,可能踩到「本该缺省却收到 null」的行为。
unstrictifyToolCallParamsFromSchema 拿原始 schema 对照,把「原本可选、且不可为 null」的参数上那些 strictness 造出来的 null 剥掉(有默认值则填默认值),同时保留 _tambo_* 服务端注入参数、丢弃 schema 里没有的幻觉键(unstrictify.ts:160):
// event-accumulator.ts:1086 流式期间也做反严格化(节选)
if (toolSchemas) {
const schema = toolSchemas.get(toolUseContent.name);
if (schema) {
parsedInput = unstrictifyToolCallParamsFromSchema(schema, parsedInput);
}
}
TOOL_CALL_ARGS(流式中)和 TOOL_CALL_END(结束时)都做这步,保证 reducer 任何时刻吐出的工具 参数都是「符合原始 schema 语义」的值,而不只是最后才对。
③ keyed-throttle:按 key 独立节流
流式工具(streamable tool,参数一边流一边就想执行的工具)有个矛盾:TOOL_CALL_ARGS 可能几十毫秒来一发,若每发都重新执行工具,太费。createKeyedThrottle 提供一个每个 key 独立的 leading + trailing 节流器:
- 某 key 第一次来:立刻执行(leading edge)。
- 冷却窗口内再来:只更新暂存值、打上「有尾巴」标记。
- 冷却结束若有尾巴:用最新值再执行一次(trailing edge),再续一个冷却窗口。
flush():强制把所有待执行的尾巴立即执行,清空计时器。
处理循环用它包住流式工具执行:每个 toolCallId 独立节流,「~每 100ms 一次」(tambo-stream.ts:356 调 throttledStreamable.schedule,流段结束 flush)。用 gen(代际计数器)让 flush/重排后的旧回调失效,避免过期计时器乱开火(keyed-throttle.ts:34)。
注意分工:节流的是流式工具的重复执行,不是 reducer 本身——reducer 对每个事件都照常累积,保证状态逐帧完整;被节流的只是「拿这些参数去跑工具」这个副作用。