跳到主要内容

流式协议: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 + accumulatingToolArgsevent-accumulator.ts:61
事件类型定义AG-UI 标准事件(@ag-ui/core)+ Tambo 自定义事件packages/client/src/types/event.ts
applyJsonPatchprops_delta 里的 RFC 6902 patch 应用到当前 propspackages/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

主线走一遍(高层)

  1. 后端流出一串 AG-UI 事件(RUN_STARTED → 若干 TEXT_MESSAGE_* / TOOL_CALL_* / tambo.component.*RUN_FINISHED)。
  2. tambo-stream.ts 的处理循环逐个接住,给 TOOL_CALL_ARGS 事件预解析部分 JSON,然后 dispatch 一个 EVENT action。
  3. streamReducerevent.type 分派到对应 handler,不可变地返回一份新的 StreamState
  4. 每次返回的新状态里,那条正在生成的消息 / 组件都更「完整」了一点。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):

  1. 前半段 switch (action.type):处理 INIT_THREAD / SET_CURRENT_THREAD / LOAD_THREAD_MESSAGES管理性 action;EVENTbreak 落到后半段。
  2. 后半段 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/statetambo.component.props_deltaJSON Patch 应用到当前对象handleComponentDelta
工具参数TOOL_CALL_ARGSJSON 字符串 + 部分解析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_deltastate_delta 共用这一个 handler,只靠 field 参数区分,是「Rule of Three 之前不硬抽象」但这里两者确实同构的合理复用。

③ 工具参数:攒 JSON 字符串,末尾权威解析

工具调用的参数(如 {"city":"Tokyo","units":"metric"})也是一段段 delta 流下来的半截 JSON 字符串handleToolCallArgs 做两件事(event-accumulator.ts:1028):

  1. 把 delta 拼进 accumulatingToolArgs[toolCallId](纯字符串累加)。
  2. 尝试乐观地部分解析当前这段半截 JSON,好让 UI 边流边显示参数。

TOOL_CALL_END,才用标准 JSON.parse完整串做一次权威解析,并清掉暂存(event-accumulator.ts:1125)。「流式期间尽量解析、结束时权威解析」这个两段式,是下一节两个技巧的舞台。

3.4 三个易读性技巧

这三个技巧都服务同一个目标:流式过程中,让 UI 尽早、且尽量正确地显示出来

① partial-json:半截 JSON 也能解析

标准 JSON.parse('{"city":"Tok') 会直接抛错。但流式时我们手上就是这种半截串。第三方库 partial-jsonparse 能容忍不完整,尽力返回 {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:342toolTracker.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:356throttledStreamable.schedule,流段结束 flush)。用 gen(代际计数器)让 flush/重排后的旧回调失效,避免过期计时器乱开火(keyed-throttle.ts:34)。

注意分工:节流的是流式工具的重复执行,不是 reducer 本身——reducer 对每个事件都照常累积,保证状态逐帧完整;被节流的只是「拿这些参数去跑工具」这个副作用。


4. 一个事件如何走完全程(端到端追踪)

以「用户问天气,模型回一句话 + 一个 WeatherCard 组件」为例,串起前面所有机制。下面是事件序列与状态变化的对照:

#事件reducer 动作TamboThread 变化
1RUN_STARTEDhandleRunStartedthread.status = "streaming",记 runId
2TEXT_MESSAGE_STARThandleTextMessageStart新增一条 assistant 消息
3TEXT_MESSAGE_CONTENT×NhandleTextMessageContent该消息 text 块逐字变长
4TEXT_MESSAGE_ENDhandleTextMessageEndstreaming.messageId
5tambo.component.starthandleComponentStart加一个 component 块,streamingState:"started"
6tambo.component.props_delta×NhandleComponentDeltaprops 逐个 patch,streamingState:"streaming"
7tambo.component.endhandleComponentEndstreamingState:"done"
8RUN_FINISHEDhandleRunFinishedstatus = "idle",记 lastCompletedRunId

每一步 reducer 都返回一份新的完整状态,tambo-stream.ts 把它 snapshot 出来推给订阅者。于是用户看到的是:文字一个字一个字冒出来,卡片先出现骨架、数值再一个个填上——全部由这条事件流驱动。

streamingState 这个字段是连接本章与第 5 章的关键接口:"started" → "streaming" → "done" 三态,让渲染层知道该显示骨架、显示流动中的 props、还是标记完成(第 5 章据此渲染)。


5. 巧妙之处(可借鉴)

  • 纯函数 reducer + 不可变更新:整个累积逻辑没有一处原地改数组/对象,全是 slice + 展开(updateMessageAtIndex / updateContentAtIndex,event-accumulator.ts:295/314)。好测、无并发副作用、天然兼容 useSyncExternalStore

  • 穷尽性检查兜底,不留静默 fallback:两个 switch 都用 const _x: never = event 收尾(event-accumulator.ts:650/1288)。上游 AG-UI 加了新事件类型而这里没处理,编译期就红,而不是运行时静默丢事件。

  • 两段式解析:流式期用 partial-json 尽早显示,结束时用标准 JSON.parse 权威定稿(event-accumulator.ts:1047 vs 1143)。既要「早」又要「准」,分开两处满足。

  • 预解析结果随 action 传递,避免重复劳动:处理循环解析一次 partial JSON,reducer 复用(action.parsedToolArgs,event-accumulator.ts:1044)。同一份半截 JSON 不解析两遍。

  • 代际计数器(gen)让过期计时器失效:keyed-throttle 里每次 flush/重排 gen++,旧 setTimeout 回调对不上代际就直接 return(keyed-throttle.ts:67)。这是处理「节流 + 可取消」的干净写法。

  • 乐观 UI 的占位线程迁移:先塞占位线程秒显消息,拿到真 id 再迁移清空(event-accumulator.ts:523)。


6. 边界与局限

  • 部分事件类型「已知但不处理」:STATE_SNAPSHOT / STATE_DELTA / MESSAGES_SNAPSHOT / STEP_STARTED 等一批 AG-UI 事件目前是 console.warn忽略(event-accumulator.ts:629)。是刻意的「未来阶段再支持」,不是 bug。

  • findContentById 是 O(n·m):按 id 找 content 块要从后往前扫所有消息 × 每条的 content(event-accumulator.ts:340)。代码里留了 TODO——高频流式 + 长会话下可能想加一个 contentId → 位置 的索引。目前够用。

  • 严格的 fail-fast:找不到消息 / content 块、TEXT_MESSAGE_END 的 messageId 对不上、工具参数 JSON.parse 失败,都直接 throw(如 event-accumulator.ts:924/1145)。这是设计选择:宁可炸出来,不要静默错乱的状态。上游事件序列若不合约,会在这里显式失败。

  • reasoning 是瞬时的:thinking 事件累积进 message.reasoning[],但不落库——从 API 重新加载的历史消息不会带 reasoning(types/message.ts:159)。


7. 代码地图(导航索引)

主题文件路径符号名
累积器主入口packages/client/src/utils/event-accumulator.tsstreamReducer
每线程状态容器packages/client/src/utils/event-accumulator.tsThreadState
顶层多线程状态packages/client/src/utils/event-accumulator.tsStreamState
初始状态packages/client/src/utils/event-accumulator.tscreateInitialState / createInitialThreadState
文字累积packages/client/src/utils/event-accumulator.tshandleTextMessageContent
工具参数累积packages/client/src/utils/event-accumulator.tshandleToolCallArgs / handleToolCallEnd
props/state 增量packages/client/src/utils/event-accumulator.tshandleComponentDelta
自定义事件分派packages/client/src/utils/event-accumulator.tshandleCustomEvent
占位线程迁移packages/client/src/utils/event-accumulator.tsPLACEHOLDER_THREAD_ID
Tambo 自定义事件类型packages/client/src/types/event.tsTamboCustomEvent / asTamboCustomEvent
线程/流式状态类型packages/client/src/types/thread.tsTamboThread / StreamingState
消息/内容块类型packages/client/src/types/message.tsTamboThreadMessage / Content / ComponentStreamingState
JSON Patch 封装packages/client/src/utils/json-patch.tsapplyJsonPatch
反严格化packages/client/src/utils/unstrictify.tsunstrictifyToolCallParamsFromSchema
按 key 节流packages/client/src/utils/keyed-throttle.tscreateKeyedThrottle
处理循环(dispatch 源头)packages/client/src/tambo-stream.tsTamboStream
事件流透传packages/client/src/utils/stream-handler.tshandleEventStream

上一章 02 决策循环:后端如何挑组件并流式吐 props · 下一章 04 一次消息的一生:HTTP 传输 + 客户端工具循环 · 返回 Tambo 全景与阅读地图