数据截至 (上游 commit 2822885e57e7)
headless 核心:ChatClient 状态机与流式生命周期
30 秒导读:
ChatClient是 TanStack AI 前端的中央引擎——一个不认识 React/Vue/Svelte 的纯 TypeScript 类。你调sendMessage('你好'),它负责把这句话变成一次网络请求、把回来的 流式碎片拼成消息、在中途自动跑客户端工具 / 等审批 / 续跑,并把每一次状态变化广播给 订阅者(框架适配层)和 devtools。这一章讲透这台引擎本身的状态机与流式生命周期; 工具/审批/持久化的业务细节见第 6 章,传输线协议见 第 2 章,chunk→parts 的拼装见第 3 章。
1. 这是什么(零基础也能懂)
一句话定义: ChatClient 是一个框架无关的聊天状态机——它把"一轮对话从发出到结束"
的所有杂活(网络、流式解析、状态流转、打断、去重、续跑)封装成一个类,对外只暴露一小把方法。
解决什么问题 / 给谁用。 假设你在写一个聊天 UI:用户敲一句话,你要发请求、要一个 token 一个
token 地把回复画出来、中途模型可能要调工具、可能要你点"允许"、网络可能断、用户可能点"停止"或
"重发"。这些时序和边界情况极其容易写乱。ChatClient 把它们收进一个类里,让上层框架适配层
(useChat,见第 4 章)只需"订阅状态 + 转发几个方法调用"。
它对外能做什么(API 一览):
- 发消息 / 追加消息 / 重发 / 停止 / 清空(
sendMessageappendreloadstopclear) - 回填工具结果与审批回复(
addToolResultaddToolApprovalResponse) - 订阅生命周期、更新配置、销毁(
subscribe/unsubscribeupdateOptionsdispose) - 一组只读 getter(
getMessagesgetStatusgetIsLoadinggetErrorgetSessionGenerating…)
用起来什么样(最小真实用法):
// 示意,非源码 —— 演示 ChatClient 的对外形态
const client = new ChatClient({
fetcher: myFetcher, // 或 connection: 某个 adapter
onMessagesChange: (msgs) => render(msgs), // 每次消息变了就重绘
onStatusChange: (s) => setBadge(s), // ready/submitted/streaming/error
})
await client.sendMessage('帮我查一下今天的待办') // 发出并流式接收
client.stop() // 用户点了停止
一句话直觉: 把 ChatClient 想成一台磁带录音机——sendMessage 按下"录制",
它一边从网络收流(磁带走),一边把内容写进 UI;stop 是"停止键",reload 是"倒带重录"。
难点全在:同一时刻只能有一盘磁带在走,而用户随时会按停止 / 倒带 / 换带——引擎必须保证
旧磁带的收尾动作绝不污染新磁带。这就是本章的主角:代际号(generation)防串流。
本节不出现底层代码。目标:你现在知道"这是台干嘛的引擎"。
2. 顶层全景(它大概怎么转)
ChatClient 内部由四大协作件组成,外加一条贯穿 始终的主线方法 streamResponse()。
2.1 四大协作件
| 部件 | 干什么 | 装配位置(chat-client.ts) |
|---|---|---|
processor(StreamProcessor) | 把底层 chunk 拼成 UIMessage[],并回调各种"消息事件" | 构造函数 chat-client.ts:600 |
connection(SubscribeConnectionAdapter) | 传输层:send() 推请求、subscribe() 吐 chunk 流 | chat-client.ts:499 |
persistor(ChatPersistor,可选) | 持久化消息 + resume 快照;"清空后迟到的 chunk"由 clearedStreamTracker 抑制(chat-client.ts:322) | chat-client.ts:483-488 |
devtoolsBridge + events | 把内部事件桥到 devtools 事件总线;events 是它的 emitter 别名 | chat-client.ts:509-512 |
callbacksRef 是第五个关键结构(chat-client.ts:421-451):一个装着所有对外回调的可变引用盒,
让框架适配层能"热替换"回调而不必重建整个 client(见 §3.1。
2.2 一张图:一次 sendMessage 怎么流动
怎么读这张图:从上到下是时间顺序;左边是你调的方法,中间是引擎主线, 右边是两条并行的循环(订阅循环收 chunk、processor 回调驱动状态)。
你的代码 ChatClient 主线 两条并行循环
───────────── ───────────────────── ────────────────────
sendMessage(text)
│ 加 user 消息 processor.addUserMessage
└──────────────► streamResponse() ← 主线入口
│
① ++generation / 建 abortController
② setStatus('submitted') / isLoading=true
③ 组装 mergedBody(见 §3.3)
④ ensureSubscription() ───────────────► 订阅循环启动
⑤ processingComplete = waitForProcessing() (等一个 promise)
⑥ connection.send(...) ──推请求──► ┌── consumeSubscription
│ │ for await chunk:
⑦ await processingComplete │ processor.processChunk
│ ◄── resolveProcessing() ──┤ → 回调 onStreamStart/
│ 由 onStreamEnd 触发 │ onTextUpdate/onStreamEnd
⑧ 代际校验 / 状态校验 │ updateRunLifecycle
⑨ finally:清理 + 续跑判断 └── (setTimeout(0) 让出 UI)
│
◄── 状态/消息变化经 callbacksRef 广播给你
主线一句话: streamResponse() 发起请求 + 挂起一个 promise,而真正"流式拼消息"发生在
另一条订阅循环 consumeSubscription 里;两者靠 waitForProcessing/resolveProcessing 这对
promise 握手汇合。这就是它能"一边收流一边不阻塞、又能在收完后精确继续"的关键。
3. 核心原理(逐个机制,由浅入深)
3.1 构造期装配:对外回调装进一个可变的 ref 盒
要解决的小问题: React 每次渲染都会给你新的回调函数(新的 onFinish 闭包)。如果 client
直接把回调"焊死"在自己身上,那要么每次渲染重建 client(丢状态),要么回调永远是旧的(闭包过期)。
思路: 把所有对外回调塞进一个可变引用盒 callbacksRef = { current: {...} }
(chat-client.ts:421-451,构造期填充 chat-client.ts:514-535)。内部永远读 callbacksRef.current.onX,
而 updateOptions 只改盒子里的字段(chat-client.ts:2996-3036)——client 实例不变,回调却能热替换。
processor 的 events 如何桥到对外回调 + devtools。 构造 StreamProcessor 时传入一组 events
回调(chat-client.ts:605-809),它们是底层事件 → 对外回调 / devtools 事件的转接板。举两个代表:
// 示意,非源码 —— 转接板的两种典型形态(真实见 chat-client.ts:606-634)
events: {
onMessagesChange: (messages) => {
this.persistor?.notifyMessagesChanged(messages) // ① 通知持久化
this.callbacksRef.current.onMessagesChange(messages) // ② 广播给框架层
},
onStreamEnd: (message) => {
this.callbacksRef.current.onFinish(message) // 对外 onFinish
this.setStatus('ready') // 状态回到 ready
this.resolveProcessing() // 关键:兑现 promise 握手
},
}
onStreamStart(chat-client.ts:610-628):置streaming状态,并向 devtools 发messageAppended。onStreamEnd(chat-client.ts:629-634):触发对外onFinish+ 回ready+resolveProcessing()—— 这一步让streamResponse里那个挂起的await processingComplete继续(§3.4)。onError(chat-client.ts:635-637)转reportStreamError;onToolCall(chat-client.ts:713-773) 自动执行客户端工具;onApprovalRequest(chat-client.ts:774-794)发审批事件——后两者的业务流 详见第 6 章,本章只点出"它们在构造期就被接到了 processor 上"。
devtoolsBridge 的角色。 this.events 是 devtoolsBridge.events(chat-client.ts:512),即
桥自己安装的 emitter。注释(chat-client.ts:386-391)说清了设计:这个 emitter 会自动附上
run/thread 上下文、并在每次事件后自动 emit 一张快照,所以 chat-client 全程只写
this.events.X(...),和没有 devtools 时一模一样。没配 devtools 时,工厂回退到
createNoOpChatDevtoolsBridge(chat-client.ts:509-511),所有 emit 都短路成空操作。
3.2 状态模型:四态 + 两个"这次别重复报错"的守卫
要解决的小问题: UI 需要知道"现在到哪一步了",而流式过程里状态转换又快又密,还得防止 同一次流的错误被报两次。
四个状态(ChatClientState,types.ts:362),按一次正常对话的推进顺序:
| 状态 | 含义 | 谁把它置进去 |
|---|---|---|
ready | 空闲,可以发下一条 | 初始值 / onStreamEnd / 收尾 |
submitted | 已发出请求,还没开始收 token | streamResponse 开头 chat-client.ts:2158 |
streaming | 正在流式接收 | onStreamStart → setStatus('streaming') chat-client.ts:611 |
error | 本轮出错 | reportStreamError chat-client.ts:1566 |
三个 setter 是状态变更的唯一入口,每个都做两件事:写字段 + 广播:
setStatus(chat-client.ts:1402-1406):写this.status,调onStatusChange,再emitSnapshot()。setIsLoading(chat-client.ts:1396-1400):isLoading是请求级开关(本地这一次请求在不在跑)。setSessionGenerating(chat-client.ts:1420-1425):会话级开关,带幂等短路(值没变就直接 return), 反映"共享会话是否在生成"——见 §3.6。
错误去重:errorReportedGeneration。 reportStreamError(chat-client.ts:1566-1583)先算
alreadyReported = errorReportedGeneration === streamGeneration:即"这一代流是否已经报过错"。
它总是 setError(更新 error 字段),但只有没报过时才把对外 onError 回调打出去,并盖章
errorReportedGeneration = streamGeneration(chat-client.ts:1579-1582)。这防止"RUN_ERROR chunk"
和"catch 块的异常"对同一次失败双重触发 onError。注意它对状态的处理很讲究:只有当前还在
isLoading / submitted / streaming 时才转 error(chat-client.ts:1572-1578),以保住"请求级错误
语义"——即使 RUN_ERROR 在 loading 翻 false 之后才姗姗来迟。
3.3 mergedBody:三层 forwardedProps 的优先级组装
要解决的小问题: 发到线上的"附带参数"(temperature、model 等)可能来自三个地方,谁覆盖谁?
三个来源槽,构造期就分开存(chat-client.ts:349-352, 495-497):
| 槽 | 来源 | 优先级 |
|---|---|---|
bodyOption | 构造/updateOptions 的 body(已废弃) | 最低 |
forwardedPropsOption | 构造/updateOptions 的 forwardedProps(规范字段) | 中 |
pendingMessageBody | sendMessage(content, body) 的每条消息 body 参数 | 最高 |
为什么分成三个槽而不是一个对象? 注释(chat-client.ts:490-494)点破:这样
updateOptions({ forwardedProps }) 不会误抹掉之前设过的 body(反之亦然)。合并只在发送时发生:
// 真实源码 chat-client.ts:2199-2203 —— 后展开的赢
const mergedBody = {
...this.bodyOption, // ① 废弃 body
...this.forwardedPropsOption,// ② forwardedProps 覆盖 body
...this.pendingMessageBody, // ③ 本条消息 body 最高
}
合并完立刻清空 pendingMessageBody(chat-client.ts:2206),保证它只作用于这一次发送。
mergedBody 随后既作为 connection.send 的参数,又被塞进 runContext.forwardedProps
(chat-client.ts:2247)走 AG-UI 线协议(线上字段细节见第 2 章)。
3.4 主线端到端:streamResponse() 怎么把这些串起来
这是全类最重要的私有方法(chat-client.ts:2132-2413)。逐段走一遍(每段标注真实行号)。
① 并发闸门 + 代际号。 开头 if (this.isLoading) return false(chat-client.ts:2133-2136)拦住并发流。
随即 const generation = ++this.streamGeneration(chat-client.ts:2138-2139)——给这次流发一个自增的"代号"。
这个 generation 是本地常量,后面所有"我还是当前流吗"的判断都拿它跟 this.streamGeneration 比。
随后接管并清空待续跑的 resume 字段(pendingResumeThreadId 等,chat-client.ts:2140-2151)——
那是 interrupt 续跑的输入,业务见第 6 章。
② 捕获 signal。 建 abortController 后立刻把 signal 存成局部常量(chat-client.ts:2161-2165)。
注释说明缘由:并发的 stop() 或 sendMessage() 会重新赋值 this.abortController,若 connect()
到时候才去读 this.abortController.signal 就可能拿到过期或 null 的 signal——所以先"拍照"下来。
③ onResponse 后的 abort 早退。 await onResponse() 后检查 if (signal.aborted) return false
(chat-client.ts:2186-2188)。注释(chat-client.ts:2181-2185)讲了这个 bug 的形状:如果在 onResponse
的 await 期间流被取消,取消时跑的 resolveProcessing() 是个空操作(此刻还没挂 promise),那么
下面 await processingComplete 会永久死锁。所以必须在分配 waitForProcessing() 之前早退。
④ 组装 mergedBody(见 §3.3)并清 pendingMessageBody。
⑤ 准备 processor + 起订阅 + 挂 promise。
processor.prepareAssistantMessage()(chat-client.ts:2218)清掉上一次流的残留状态,
ensureSubscription()(chat-client.ts:2221)保证订阅循环在跑,然后
const processingComplete = this.waitForProcessing()(chat-client.ts:2224)——挂起一个 promise,
它会在 onStreamEnd / RUN 终态时被 resolveProcessing() 兑现。
⑥ 发送 + 等待处理完成。
// 真实源码 chat-client.ts:2265-2273
await this.connection.send(messages, mergedBody, signal, runContext)
// Wait for subscription loop to finish processing all chunks
await processingComplete
send() 把请求推给传输层(chunk 会从订阅循环那边冒出来,不在这里);await processingComplete
则挂起主线,直到订阅循环把这一轮的 chunk 全处理完并触发兑现。这是"发起"与"消费"解耦的接缝。
⑦ 代际防串流校验。 醒来后第一件事:
// 真实源码 chat-client.ts:2277-2279 —— 我这盘磁带还是当前磁带吗?
if (generation !== this.streamGeneration) {
return false // 已被 reload()/新 sendMessage 取代,新流接管了 processor 和 processingResolve
}
若被取代就直接退,不碰任何共享状态——旧流的收尾不污染新流。接着若 status === 'error'
(chat-client.ts:2281-2294)也不算成功完成,发 run:errored 后返回 false。
⑧ 等客户端工具 + 定稿。 if (pendingToolExecutions.size > 0) await Promise.all(...)
(chat-client.ts:2296-2299)等所有客户端工具执行完,再 processor.finalizeStream()(幂等,chat-client.ts:2302),
置 streamCompletedSuccessfully = true。
⑨ catch:区分 abort 与真错。 AbortError 走"取消"分支发 run:cancelled 后 return false
(chat-client.ts:2306-2316);其它错误先做代际校验 if (generation === this.streamGeneration)
再 reportStreamError(chat-client.ts:2317-2328)——被取代的流即使抛错也不报。
⑩ finally:被代际号守护的清理。 整个 finally 包在 if (generation === this.streamGeneration)
里(chat-client.ts:2334):只有仍是当前流才清理 abortController/isLoading/各种 current 字段
(chat-client.ts:2335-2343)。注释点明:被取代的流(如 reload() 起的新流)绝不能覆盖新流的
abortController 或 isLoading。清理后 drainPostStreamActions()(排空流中排队的动作),再做续跑判断:
// 真实源码 chat-client.ts:2367-2393(节选逻辑)
if (streamCompletedSuccessfully) {
const lastPart = messages.at(-1)?.parts.at(-1)
const { finishReason } = this.processor.getState()
if (lastPart?.type === 'tool-result' && finishReason !== 'stop' && this.shouldAutoSend()) {
await this.checkForContinuation() // 工具跑完且模型想继续 → 自动续跑
} else if (this.status !== 'ready') {
this.setStatus('ready') // 兜底:bare RUN_FINISHED{stop} 没触发 onStreamEnd 的 #421 场景
}
}
续跑只在满足三条件时发生: 最后一个 part 是 tool-result、finishReason !== 'stop'、
且 shouldAutoSend() 为真。续跑与工具的完整业务见第 6 章;本章只强调它
发生在 finally 里、且受代际守护。finally 的最后还处理消息队列:成功且不续跑时 drainQueue()
自动发排队的消息(chat-client.ts:2394-2400);失败/abort 时 flushQueue() 直接清空——不把排队的消息
发向一个大概率已坏的端点(chat-client.ts:2402-2408)。