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:230 |
connection(SubscribeConnectionAdapter) | 传输层:send() 推请求、subscribe() 吐 chunk 流 | chat-client.ts:188 |
persistor(ChatPersistor,可选) | 持久化 + 抑制"清空后迟到的 chunk" | chat-client.ts:173-179 |
devtoolsBridge + events | 把内部事件桥到 devtools 事件总线;events 是它的 emitter 别名 | chat-client.ts:198-201 |
callbacksRef 是第五个关键结构(chat-client.ts:149-168):一个装着所有对外回调的可变引用盒,
让框架适配层能"热替换"回调而不必重建整个 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:149-168,构造期填充 chat-client.ts:203-220)。内部永远读 callbacksRef.current.onX,
而 updateOptions 只改盒子里的字段(chat-client.ts:1491-1517)——client 实例不变,回调却能热替换。
processor 的 events 如何桥到对外回调 + devtools。 构造 StreamProcessor 时传入一组 events
回调(chat-client.ts:235-432),它们是底层事件 → 对外回调 / devtools 事件的转接板。举两个代表:
// 示意,非源码 —— 转接板的两种典型形态(真实见 chat-client.ts:235-273)
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:240-258):置streaming状态,并向 devtools 发messageAppended。onStreamEnd(chat-client.ts:259-264):触发对外onFinish+ 回ready+resolveProcessing()—— 这一步让streamResponse里那个挂起的await processingComplete继续(§3.4)。onError(chat-client.ts:265-267)转reportStreamError;onToolCall(chat-client.ts:343-403) 自动执行客户端工具;onApprovalRequest(chat-client.ts:404-424)发审批事件——后两者的业务流 详见第 6 章,本章只点出"它们在构造期就被接到了 processor 上"。
devtoolsBridge 的角色。 this.events 是 devtoolsBridge.events(chat-client.ts:201),即
桥自己安装的 emitter。注释(chat-client.ts:121-126)说清了设计:这个 emitter 会自动附上
run/thread 上下文、并在每次事件后自动 emit 一张快照,所以 chat-client 全程只写
this.events.X(...),和没有 devtools 时一模一样。没配 devtools 时,工厂回退到
createNoOpChatDevtoolsBridge(chat-client.ts:198-199),所有 emit 都短路成空操作。
3.2 状态模型:四态 + 两个"这次别重复报错"的守卫
要解决的小问题: UI 需要知道"现在到哪一步了",而流式过程里状态转换又快又密,还得防止 同一次流的错误被报两次。
四个状态(ChatClientState,types.ts:107),按一次正常对话的推进顺序:
| 状态 | 含义 | 谁把它置进去 |
|---|---|---|
ready | 空闲,可以发下一条 | 初始值 / onStreamEnd / 收尾 |
submitted | 已发出请求,还没开始收 token | streamResponse 开头 chat-client.ts:861 |
streaming | 正在流式接收 | onStreamStart → setStatus('streaming') chat-client.ts:241 |
error | 本轮出错 | reportStreamError chat-client.ts:627 |
三个 setter 是状态变更的唯一入口,每个都做两件事:写字段 + 广播:
setStatus(chat-client.ts:502-506):写this.status,调onStatusChange,再emitSnapshot()。setIsLoading(chat-client.ts:496-500):isLoading是请求级开关(本地这一次请求在不在跑)。setSessionGenerating(chat-client.ts:520-525):会话级开关,带幂等短路(值没变就直接 return), 反映"共享会话是否在生成"——见 §3.6。
错误去重:errorReportedGeneration。 reportStreamError(chat-client.ts:616-633)先算
alreadyReported = errorReportedGeneration === streamGeneration:即"这一代流是否已经报过错"。
它总是 setError(更新 error 字段),但只有没报过时才把对外 onError 回调打出去,并盖章
errorReportedGeneration = streamGeneration(chat-client.ts:629-631)。这防止"RUN_ERROR chunk"
和"catch 块的异常"对同一次失败双重触发 onError。注意它对状态的处理很讲究:只有当前还在
isLoading / submitted / streaming 时才转 error(chat-client.ts:622-628),以保住"请求级错误
语义"——即使 RUN_ERROR 在 loading 翻 false 之后才姗姗来迟。
3.3 mergedBody:三层 forwardedProps 的优先级组装
要解决的小问题: 发到线上的"附带参数"(temperature、model 等)可能来自三个地方,谁覆盖谁?
三个来源槽,构造期就分开存(chat-client.ts:109-112, 185-186):
| 槽 | 来源 | 优先级 |
|---|---|---|
bodyOption | 构造/updateOptions 的 body(已废弃) | 最低 |
forwardedPropsOption | 构造/updateOptions 的 forwardedProps(规范字段) | 中 |
pendingMessageBody | sendMessage(content, body) 的每条消息 body 参数 | 最高 |
为什么分成三个槽而不是一个对象? 注释(chat-client.ts:105-108)点破:这样
updateOptions({ forwardedProps }) 不会误抹掉之前设过的 body(反之亦然)。合并只在发送时发生:
// 真实源码 chat-client.ts:902-906 —— 后展开的赢
const mergedBody = {
...this.bodyOption, // ① 废弃 body
...this.forwardedPropsOption,// ② forwardedProps 覆盖 body
...this.pendingMessageBody, // ③ 本条消息 body 最高
}
合并完立刻清空 pendingMessageBody(chat-client.ts:909),保证它只作用于这一次发送。
mergedBody 随后既作为 connection.send 的参数,又被塞进 runContext.forwardedProps
(chat-client.ts:947)走 AG-UI 线协议(线上字段细节见第 2 章)。
3.4 主线端到端:streamResponse() 怎么把这些串起来
这是全类最重要的私有方法(chat-client.ts:849-1087)。逐段走一遍(每段标注真实行号)。
① 并发闸门 + 代际号。 开头 if (this.isLoading) return false(chat-client.ts:851-853)拦住并发流。
随即 const generation = ++this.streamGeneration(chat-client.ts:856)——给这次流发一个自增的"代号"。
这个 generation 是本地常量,后面所有"我还是当前流吗"的判断都拿它跟 this.streamGeneration 比。
② 捕获 signal。 建 abortController 后立刻把 signal 存成局部常量(chat-client.ts:864-868)。
注释说明缘由:并发的 stop() 或 sendMessage() 会重新赋值 this.abortController,若 connect()
到时候才去读 this.abortController.signal 就可能拿到过期或 null 的 signal——所以先"拍照"下来。
③ onResponse 后的 abort 早退。 await onResponse() 后检查 if (signal.aborted) return false
(chat-client.ts:889-891)。注释(chat-client.ts:884-888)讲了这个 bug 的形状:如果在 onResponse
的 await 期间流被取消,取消时跑的 resolveProcessing() 是个空操作(此刻还没挂 promise),那么
下面 await processingComplete 会永久死锁。所以必须在分配 waitForProcessing() 之前早退。
④ 组装 mergedBody(见 §3.3)并清 pendingMessageBody。
⑤ 准备 processor + 起订阅 + 挂 promise。
processor.prepareAssistantMessage()(chat-client.ts:921)清掉上一次流的残留状态,
ensureSubscription()(chat-client.ts:924)保证订阅循环在跑,然后
const processingComplete = this.waitForProcessing()(chat-client.ts:927)——挂起一个 promise,
它会在 onStreamEnd / RUN 终态时被 resolveProcessing() 兑现。
⑥ 发送 + 等待处理完成。
// 真实源码 chat-client.ts:964-967
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:971-973 —— 我这盘磁带还是当前磁带吗?
if (generation !== this.streamGeneration) {
return false // 已被 reload()/新 sendMessage 取代,新流接管了 processor 和 processingResolve
}
若被取代就直接退,不碰任何共享状态——旧流的收尾不污染新流。接着若 status === 'error'
(chat-client.ts:977-988)也不算成功完成,发 run:errored 后返回 false。
⑧ 等客户端工具 + 定稿。 if (pendingToolExecutions.size > 0) await Promise.all(...)
(chat-client.ts:991-993)等所有客户端工具执行完,再 processor.finalizeStream()(幂等,chat-client.ts:996),
置 streamCompletedSuccessfully = true。
⑨ catch:区分 abort 与真错。 AbortError 走"取消"分支发 run:cancelled 后 return false
(chat-client.ts:1000-1010);其它错误先做代际校验 if (generation === this.streamGeneration)
再 reportStreamError(chat-client.ts:1011-1022)——被取代的流即使抛错也不报。
⑩ finally:被代际号守护的清理。 整个 finally 包在 if (generation === this.streamGeneration)
里(chat-client.ts:1028):只有仍是当前流才清理 abortController/isLoading/各种 current 字段
(chat-client.ts:1029-1037)。注释点明:被取代的流(如 reload() 起的新流)绝不能覆盖新流的
abortController 或 isLoading。清理后 drainPostStreamActions()(排空流中排队的动作),再做续跑判断:
// 真实源码 chat-client.ts:1061-1082(节选逻辑)
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 里、且受代际守护。
3.5 promise 握手:waitForProcessing / resolveProcessing
要解决的小问题: 主线 streamResponse 是一条 async 函数,而"流处理完了"这个信号来自另一 条
循环(订阅循环里 processor 触发的 onStreamEnd)。怎么让主线"睡到信号来了再醒"?
答案是一对配套方法,用一个字段 processingResolve 当"门铃按钮":
// 真实源码 chat-client.ts:595-598 & 724-730
private resolveProcessing(): void {
this.processingResolve?.() // 按门铃(若挂着 promise)
this.processingResolve = null // 一次性,按完即焚
}
private waitForProcessing(): Promise<void> {
this.resolveProcessing() // 先兑现任何过期的 promise(防上一次残留)
return new Promise<void>((resolve) => { this.processingResolve = resolve })
}
握手时序: 主线 waitForProcessing() 把 resolve 存进 processingResolve → 主线 await 睡下 →
订阅循环 处理到终态,onStreamEnd(chat-client.ts:263)或 updateRunLifecycle(chat-client.ts:487-489)
调 resolveProcessing() → 门铃响,主线醒。waitForProcessing 开头先 resolve 一次是为了兑现
"上一次 aborted 请求残留的 promise",避免野指针。多个地方都会按门铃(onStreamEnd、RUN 终态、
cancelInFlightStream、订阅循环出错),保证主线永不永久挂起。
3.6 订阅循环:consumeSubscription / startSubscription / ensureSubscription
要解决的小问题: chunk 是从连接持续流出的,需要一条独立于单次请求的长循环去消费, 而且这条循环的生死要能被优雅地起停、重启、容错。
三个方法分工:
| 方法 | 职责 | 位置 |
|---|---|---|
ensureSubscription() | "确保在跑":没订阅就 subscribe();订阅了但循环已 abort 就 subscribe({restart:true}) | chat-client.ts:707-718 |
startSubscription() | 真正起循环:建 subscriptionAbortController,跑 consumeSubscription,挂 .catch/.finally | chat-client.ts:638-666 |
consumeSubscription(signal) | 核心 for await:逐 chunk 处理 | chat-client.ts:671-702 |
consumeSubscription 一圈做什么(chat-client.ts:671-702):
// 真实源码骨架 chat-client.ts:672-701(节选)
const stream = this.connection.subscribe(signal)
for await (const chunk of stream) {
if (signal.aborted) break
const shouldIgnore = this.persistor?.shouldIgnoreChunk(chunk) ?? false
if (shouldIgnore) { /* 只更新 run 生命周期,跳过 processor(清空后迟到的 chunk) */ continue }
this.callbacksRef.current.onChunk(chunk) // ① 对外 onChunk
this.devtoolsBridge.observeChunk(chunk) // ② 喂 devtools
this.processor.processChunk(chunk) // ③ 拼消息(第 3 章)
this.updateRunLifecycle(chunk) // ④ 维护 run 集合 / session 状态
await new Promise((r) => setTimeout(r, 0)) // ⑤ 让出事件循环,给 UI 喘息
}
第 ⑤ 步的 setTimeout(0) 是个故意的让步——把控制权还给事件循环,让 UI 有机会重绘,避免
密集 chunk 把主线程堵死。shouldIgnore 分支处理"清空对话后才迟到的 chunk"(持久化关切,详见
第 6 章),但即便忽略它仍会更新 run 生命周期,以免 activeRunIds
与真实状态发散(注释在 chat-client.ts:692-697)。
startSubscription 的容错。 .catch(chat-client.ts:643-652)对非 AbortError 的错误置连接为
error、复位 session、报错,并兜底 resolveProcessing 防主线挂死;.finally(chat-client.ts:653-665)
有一个过期循环守卫:if (this.subscriptionAbortController?.signal !== signal) return——若这条循环
已被重启取代,就不动新循环的状态。这和 streamResponse 的代际守卫是同一种"我还是当前的吗"哲学。
subscribe()(chat-client.ts:1093-1106)/unsubscribe()(chat-client.ts:1112-1120)是对外的订阅开关:
前者置 isSubscribed、连接 connecting、起循环;后者取消在飞的流 + 中止订阅 + 复位 session。
3.7 run 生命周期:updateRunLifecycle / activeRunIds 与会话级生成态
要解决的小问题: 一次对话里可能有多个 run(如工具调用触发的续跑),而且"会话在不在生成"
这个状态要能被所有订阅者(甚至跨标签页/设备)看到——它不等于本地这一次请求的 isLoading。
两个"loading"要分清:
| 概念 | 含义 | 字段 / getter |
|---|---|---|
isLoading | 请求级:本地这一次 streamResponse 在不在跑 | getIsLoading() chat-client.ts:1382 |
sessionGenerating | 会话级:共享会话是否在生成(由 RUN 事件派生,跨端可见) | getSessionGenerating() chat-client.ts:1413 |
updateRunLifecycle(chat-client.ts:461-490)是唯一维护这套状态的地方:
RUN_STARTED:把 runId 加进activeRunIds、通知 persistor、setSessionGenerating(true)(chat-client.ts:465-471)。RUN_FINISHED/RUN_ERROR:按 runId 从activeRunIds删除;若RUN_ERROR没带 runId,当作 会话级错误清空所有 run(chat-client.ts:473-485)。之后按activeRunIds.size > 0更新 session 生成态, 并(默认)resolveProcessing()兑现主线握手(chat-client.ts:486-489)。
注释(chat-client.ts:692-697)强调:把 run 生命周期收进这一个方法,是为了让"忽略 chunk 的路径"
和"正常路径"不会分叉。getSessionGenerating 的文档注释(chat-client.ts:1407-1412)也点明它反映的是
"对所有订阅者可见的共享生成活动"——这是多端/群聊场景的基础(见第 6 章)。
4. 对外 API 表(逐个方法在做什么)
按"发起 / 控制 / 回填 / 生命周期 / getter"分组。业务细节链到对应章,这里只讲它在状态机里干什么。
4.1 发起与控制
| 方法 | 干什么(状态机视角) | 位置 |
|---|---|---|
sendMessage(content, body?) | 挂载 devtools → 空串/loading 时早退 → 归一化输入 → 存 pendingMessageBody → processor.addUserMessage → streamResponse() | chat-client.ts:771-794 |
append(message) | 归一化为 UIMessage(跳过 system)→ 加进 messages;若正 loading 则排队到流后再 streamResponse,否则立即发 | chat-client.ts:813-843 |
reload() | 找到最后一条 user 消息 → 取消在飞流 → removeMessagesAfter → 重发 streamResponse | chat-client.ts:1125-1149 |
stop() | cancelInFlightStream({setReadyStatus});若本地有流则复位 session;发 stopped 事件 | chat-client.ts:1154-1161 |
clear() | (有 persistor 时)快照 + 取消在飞流 + 复位 → processor.clearMessages → persistor.remove → 清 error | chat-client.ts:1166-1187 |
cancelInFlightStream(chat-client.ts:600-614)是 stop/reload/clear/updateOptions 共用的取消原语:
abort 掉 abortController、(可选)中止订阅、resolveProcessing()、setIsLoading(false)、(可选)回 ready。
4.2 回填(工具 / 审批)
| 方法 | 干什么 | 位置 |
|---|---|---|
addToolResult(result) | 回填客户端工具结果 → 校验输出 → processor.addToolResult → loading 时排队 checkForContinuation,否则立即 | chat-client.ts:1192-1244 |
addToolApprovalResponse({id, approved}) | 按 approval id 反查 toolCallId → processor.addToolApprovalResponse → 同样的"排队 or 立即续跑" | chat-client.ts:1260-1298 |
两者都遵循同一套续跑触发模式:流在跑就 queuePostStreamAction,不在跑就直接 checkForContinuation。
续跑闸门 checkForContinuation(chat-client.ts:1326-1352)用 continuationPending/continuationSkipped
两个标志做去重与补跑——避免链式审批场景下重复续跑。完整业务见第 6 章。
4.3 生命周期与配置
| 方法 | 干什么 | 位置 |
|---|---|---|
subscribe({restart?}) | 起订阅循环(独立于请求生命周期) | chat-client.ts:1093-1106 |
unsubscribe() | 取消在飞流 + 中止订阅 + 复位 session + 断连 | chat-client.ts:1112-1120 |
updateOptions(options) | 热更新配置:换 connection/fetcher(会重订阅)、独立替换 body/forwardedProps/context/tools、热替回调 | chat-client.ts:1435-1518 |
dispose() | unsubscribe() + 销毁 devtoolsBridge + 复位 devtoolsMounted | chat-client.ts:1520-1524 |
setMessagesManually(messages) | 直接改 processor 的消息并发快照 | chat-client.ts:1427-1430 |
updateOptions 换 transport 时很讲究(chat-client.ts:1445-1470):记住 wasSubscribed、取消在飞流、
复位 session、重建 connection,只有原本订阅着才重新 subscribe()——保持订阅状态一致。
4.4 只读 getter
getMessages(:1375)、getIsLoading(:1382)、getStatus(:1389)、getIsSubscribed(:1396)、
getConnectionStatus(:1403)、getSessionGenerating(:1413)、getError(:1420)——都是纯读字段,
无副作用,供框架适配层做快照订阅(见第 4 章)。
订阅入口在哪? 本类没有名为
subscribe/unsubscribe的"事件监听"方法——对外订阅状态变化 靠的是构造/updateOptions传入的onXChange回调族(存进callbacksRef);subscribe()/unsubscribe()控制的是连接订阅循环,不是状态监听。框架层就是靠热替这些回调 + getter 快照来"订阅"的。
5. 巧妙之处(可借鉴的技术)
-
代际号(generation)防串流——一个自增
int就解决了"异步收尾污染新任务"这个经典难题。 凡是 await 之后要改共享状态的地方,先if (generation !== this.streamGeneration) return(chat-client.ts:856, 971, 1011, 1028)。订阅循环用同构的 signal 守卫(chat-client.ts:655)。 -
signal 先拍照,再用——
const signal = this.abortController.signal(chat-client.ts:868)把 signal 从可变字段"钉"成局部常量,躲开并发重赋值导致的"拿到 null/过期 signal"(chat-client.ts:865-867注释)。 -
promise 握手解耦"发起"与"消费"——
waitForProcessing/resolveProcessing(chat-client.ts:595-598, 724-730) 让 async 主线能精确地"睡到另一条循环发信号"。waitForProcessing开头先兑现过期 promise,杜绝挂死。 -
onResponse 后的死锁预防——
if (signal.aborted) return false(chat-client.ts:889-891)在挂 promise 之前早退,避免"取消时的 resolveProcessing 是空操作 → 之后 await 永久挂起"(chat-client.ts:884-888注释)。 -
三槽 body 分开存、发送时合并——
bodyOption/forwardedPropsOption/pendingMessageBody(chat-client.ts:109-112)让updateOptions能改一个槽不误伤另一个,合并优先级清晰(chat-client.ts:902-906)。 -
run 生命周期单点收口——
updateRunLifecycle(chat-client.ts:461-490)是唯一维护activeRunIds与sessionGenerating的地方,让"忽略路径"和"正常路径"不分叉(chat-client.ts:692-697注释)。 -
回调 ref 盒 + devtools 透明桥——
callbacksRef(chat-client.ts:149-168)让回调可热替而不重建 client;this.events是 devtools 桥装的 emitter(chat-client.ts:121-126),自动附上下文 + 自动发快照, 使 chat-client 代码对"有没有 devtools"完全无感。
6. 边界与局限(诚实)
-
同一时刻只跑一条本地流。
streamResponse开头if (this.isLoading) return false(chat-client.ts:851-853)。想在流中再发消息?append会排队到流后(chat-client.ts:835-840), 而sendMessage在 loading 时直接早退不排队(chat-client.ts:777)——这是刻意的不对称。 -
本章不覆盖三块内部: chunk→parts 的拼装全在
StreamProcessor(第 3 章);传输/线协议在 connection adapters(第 2 章);工具执行、审批链、续跑去重、持久化、多端生成态的完整业务在第 6 章。 本章只讲"引擎的状态机与生命周期骨架"。 -
续跑判断依赖 processor 的
finishReason与areAllToolsComplete(chat-client.ts:1064, 1369)。shouldAutoSend(chat-client.ts:1359-1370)要求"最后一条 assistant 消息里至少有一个 tool-call"—— 纯文本回复没有可续跑的东西,直接返回 false。 -
#421兜底(chat-client.ts:1076-1081):bareRUN_FINISHED{stop}时 processor 没有 assistant 消息可发onStreamEnd,于是 finally 里补一个setStatus('ready')。这类"事件没按预期到达"的兜底 说明真实流协议存在边角情况。
7. 横向对比
ChatClient 走的是类状态机 + 显式生命周期路线:所有时序(代际、握手、run 集合)都摊在一个类里
手写。它刻意不用响应式框架的信号系统——因为它要 headless、要能被 React/Vue/Svelte/Solid 共用,
框架适配只是薄薄一层订阅回调(见第 4 章)。这与"把状态机塞进框架 store"
的方案相反:代价是手写时序较重(本章讲的那些守卫),收益是框架无关 + 可测试 + 可多端。
8. 代码地图(导航索引)
| 主题 | 文件路径 | 符号名 |
|---|---|---|
| 类定义与全部私有字段 | packages/ai-client/src/chat-client.ts:92-168 | ChatClient / callbacksRef |
| 构造期装配(processor events 转接板) | packages/ai-client/src/chat-client.ts:170-436 | constructor |
| 传输解析(connection vs fetcher 互斥) | packages/ai-client/src/chat-client.ts:77-90 | resolveTransport |
| 主线:流式生命周期 | packages/ai-client/src/chat-client.ts:849-1087 | streamResponse |
| promise 握手 | packages/ai-client/src/chat-client.ts:595-598, 724-730 | resolveProcessing / waitForProcessing |
| 订阅循环 | packages/ai-client/src/chat-client.ts:638-718 | startSubscription / consumeSubscription / ensureSubscription |
| run 生命周期 / 会话生成态 | packages/ai-client/src/chat-client.ts:461-490 | updateRunLifecycle / activeRunIds |
| 状态 setter | packages/ai-client/src/chat-client.ts:496-537 | setStatus / setIsLoading / setSessionGenerating |
| 错误去重 | packages/ai-client/src/chat-client.ts:616-633 | reportStreamError / errorReportedGeneration |
| body 三槽合并 | packages/ai-client/src/chat-client.ts:902-909 | mergedBody |
| 取消原语 | packages/ai-client/src/chat-client.ts:600-614 | cancelInFlightStream |
| 续跑闸门 | packages/ai-client/src/chat-client.ts:1326-1370 | checkForContinuation / shouldAutoSend |
| 配置热更新 | packages/ai-client/src/chat-client.ts:1435-1518 | updateOptions |
| 状态枚举 | packages/ai-client/src/types.ts:107 | ChatClientState |
| devtools 桥(空操作回退) | packages/ai-client/src/devtools-noop.ts | NoOpChatDevtoolsBridge |
| run 事件上下文类型 | packages/ai-client/src/events.ts:8-16 | ChatClientRunEventContext |