数据截至 (上游 commit 2822885e57e7)
传输与线协议:connection adapters 与 AG-UI RunAgentInput
30 秒导读: ChatClient 只懂两件事——「把消息发出去」和「订阅回流的 chunk」。 至于这些字节到底是走
fetch的 SSE、还是XMLHttpRequest的 NDJSON、还是一个 RPC 流,它一概不管。 本章讲的就是这层「传输适配器」:它把千奇百怪的网络实现,统一收拢成subscribe/send两个方法, 出站时按 AG-UI 的RunAgentInput打包请求体,入站时把服务器的字节流解析回一颗颗StreamChunk。
本章聚焦两个文件:
packages/ai-client/src/connection-adapters.ts—— 适配器的类型、归一化、内置工厂、健壮性守卫。packages/ai/src/utilities/ag-ui-wire.ts—— 出站序列化uiMessagesToWire(把UIMessage.parts翻成 AG-UI 线格式)。
ChatClient 那侧的「订阅循环怎么消费 chunk、状态机怎么转」属于第 1 章,这里不重复。
1. 这是什么(零基础也能懂)
一句话定义: connection adapter 是 ChatClient 和「网络」之间的可插拔转接头——它规定「一次对话请求怎么发、 流式回包怎么收」,但把「用什么协议、什么 HTTP 客户端」的选择权交给你。
它解决什么问题。 一个 headless 聊天客户端要能跑在各种环境里:
- 浏览器里走标准
fetch+ Server-Sent Events; - React Native 或老环境里
fetch流不可用,得退回XMLHttpRequest; - TanStack Start 的 server function 直接返回一个
AsyncIterable,根本没有 HTTP 层; - Cap'n Web 之类的 RPC,流是通过 RPC 通道回来的。
如果 ChatClient 把 fetch 写死,上面每一种都得改核心。适配器就是那道解耦缝:核心只依赖一个抽象接口,
具体传输作为「一个对象」传进来。
用起来什么样。 使用者几乎只写一行——挑一个内置工厂,交给 ChatClient:
// 示意,非 源码:三种传输,ChatClient 用法完全一致
import { fetchServerSentEvents, xhrHttpStream, stream } from '@tanstack/ai-client'
// 浏览器 SSE
const conn = fetchServerSentEvents('/api/chat')
// React Native 等无 fetch-stream 环境,退回 XHR + 换行分隔 JSON
const conn2 = xhrHttpStream('/api/chat')
// TanStack Start server function 直接给一个异步可迭代
const conn3 = stream((messages, data) => myServerFn({ messages, data }))
const client = new ChatClient({ connection: conn }) // 核心不变
一句话直觉: 把 connection adapter 想成电源转换头——插座(网络)千奇百怪,你的设备(ChatClient)
只认一种插头(subscribe/send)。转换头负责两件事:把你的电压「转出去」(打包请求),把回来的电流「转进来」(解析流)。
本节不出现底层解析细节。记住一件事:ChatClient 只认 subscribe/send,其余都是 适配器的活。
2. 顶层全景(它大概怎么转)
2.1 两种适配器「形状」
适配器有两种写法,是一个联合类型 ConnectionAdapter(connection-adapters.ts:965):
| 形状 | 你要实现的方法 | 心智模型 | 谁在用 |
|---|---|---|---|
ConnectConnectionAdapter | connect(messages,…) → AsyncIterable<StreamChunk> | 「一发一收」:调一次,拿一条流,for await 读完即止 | 全部内置工厂 |
SubscribeConnectionAdapter | subscribe() → AsyncIterable + send() → Promise | 「先订阅,后投递」:订阅是长期的口子,send 只管把请求推出去,chunk 从订阅口回来 | 需要长连接/多路复用的自定义传输(如共享 WebSocket) |
类型定义分别在 ConnectConnectionAdapter(connection-adapters.ts:814)和
SubscribeConnectionAdapter(connection-adapters.ts:929)。两者还各有一组可选的持久化钩子
(joinRun / hydrate,见 §3.1 末尾):有了它们,ChatClient 就能在整页刷新后 rejoin 在飞的 run。
为什么要两种? connect 简单直白,适合「一次请 求一条流」的 HTTP/SSE;subscribe/send 把「收」和「发」拆开,
适合那种「一个长连接服务多次发送」的场景。但 ChatClient 内部只想面对一种——于是有了归一化(§3)。
2.2 一次 send 的全景数据流
下面这张图是本章主干:从 ChatClient 调 send,到 chunk 回流被订阅口读到,中间适配器做了什么。
从上往下读,左侧是出站(打包+发送),右侧是入站(解析+回流)。
ChatClient.send(messages, body, signal, runContext)
│ (第1章:runContext 带 threadId/runId/clientTools/forwardedProps)
▼
┌───────────────── normalizeConnectionAdapter 归一化后的 send ──────────────────┐
│ │
│ ①出站打包 ⑤入站回流 │
│ buildRunAgentInputBody ──┐ ┌── push(chunk) 进队列 │
│ · uiMessagesToWire │ │ (activeBuffer / activeWaiters)│
│ · RunAgentInput 镜像字段 │ │ │ │
│ ▼ │ ▼ │
│ ②发请求(工厂内) ④解析字节 → StreamChunk │
│ fetch / XHR / 直连流 ─────► linesToSSEEvents │
│ │ · 跳过 :注释/event:/retry: │
│ ▼ · 记下 id: 偏移(断线重连用) │
│ ③服务器回字节流 (SSE/NDJSON) · [DONE] → 合成 RUN_FINISHED │
│ · JSON.parse 每行 │
│ │
└──────────────────────────────────────────────────────────────────────────────┘
│
▼
ChatClient.subscribe() 的 for-await 逐个拿到 chunk(第1章:喂给 StreamProcessor)
怎么读:①②③是「发出去」,④⑤是「收回来」。中间的 activeBuffer/activeWaiters 队列(§3.2)是把
「connect 的一条流」桥接到「subscribe 的长期订阅口」的关键。
2.3 部件一句话职责
| 部件 | 干什么 | 在哪 |
|---|---|---|
ConnectionAdapter 联合类型 | 定义两种适配器形状 | connection-adapters.ts:965 |
normalizeConnectionAdapter | 把任意适配器统一成 subscribe/send | connection-adapters.ts:975 |
buildRunAgentInputBody | 把消息+上下文打包成 AG-UI RunAgentInput 请求体 | connection-adapters.ts:1193 |
uiMessagesToWire | 把 UIMessage.parts 序列化成 AG-UI 线消息 | ag-ui-wire.ts:62 |
responseToSSEChunks / readStreamLines | 把 SSE/NDJSON 字节流解析回 StreamChunk | connection-adapters.ts:627 / :309 |
| 内置工厂 | 各种传输的现成实现(含 WebSocket) | connection-adapters.ts(§4) |
StreamTruncatedError / requireSyntheticId / abortableIterable | 健壮性守卫 | connection-adapters.ts:51 / :245 / :2513 |
3. 核心原理:把 connect 包成 subscribe/send
3.1 归一化要解决的小问题
ChatClient 构造时,不管你给的是 connect 型还是 subscribe/send 型,它都想拿到统一的 subscribe/send。
这就是 normalizeConnectionAdapter(connection-adapters.ts:975)的活。ChatClient 在
chat-client.ts:499 用它包一次,之后只调 this.connection.subscribe(...) / .send(...)。
它先做互斥校验——一个适配器只能是一种形状,同时给 connect 和 subscribe/send 直接抛错
(connection-adapters.ts:986-990)。然后分三条路:
normalizeConnectionAdapter(connection)
│
┌───────────────────────┼───────────────────────┐
▼ ▼ ▼
同时有 connect 和 有 subscribe+send 只有 connect
subscribe/send → 原样绑定返回 → 用 activeBuffer/
→ 抛错(形状冲突) (它已是目标形状) activeWaiters 队列包一层
原生就是 subscribe/send 的适配器直接 .bind 透传(connection-adapters.ts:992-1005)——
若还带可选的 joinRun/hydrate(rejoin 在飞 run / 服务器侧水合),也一并 bind 出去;
connect 型的包装同样只在底层真实现了 joinRun/hydrate 时才把它们暴露上去
(connection-adapters.ts:1137-1158)。难点全在第三条——
只有 connect 的适配器,怎么假装成 subscribe/send?
3.2 思路:一个 activeBuffer / activeWaiters 异步队列
connect 是「调用即得一条流」;subscribe 是「先开一个长期订阅口,chunk 陆续来」。要把前者变后者,
需要一个中间队列来解耦生产(send 里读 connect 流)和消费(subscribe 的 for-await)。
这个队列就两个数组(connection-adapters.ts:1014-1015):
activeBuffer—— chunk 来了但当前没人等,先囤在这。activeWaiters—— 有人await下一个 chunk 但还没货,把它的 resolve 回调囤在这。
push(connection-adapters.ts:1017)是生产端:有等待者就直接喂给它,否则塞进 buffer。
这是经典的「异步队列 / 单生产者单消费者」写法——永远只有 buffer 和 waiters 之一非空。
send() 侧(生产):push(chunk) subscribe() 侧(消费):next()
──────────────────────────── ────────────────────────────
有等待者? buffer 有货?
是 → 直接 resolve 那个 waiter 是 → 立刻 yield
否 → chunk 压进 activeBuffer 否 → 造一个 Promise,把 resolve
压进 activeWaiters,挂起等 push
3.3 精华细节:所有权转移给「最新订阅者」
一个坑:如果先后开了两次 subscribe(比如 reload 触发重订阅),老订阅口不能继续偷 chunk。
代码用一招「所有权转移」解决(connection-adapters.ts:1030-1036):
subscribe(abortSignal) {
// 把当前 buffer 整个搬给「我」这个最新订阅者,再把全局指针指向我的队列
const myBuffer = activeBuffer.splice(0) // splice(0):清空原数组并接管其内容
const myWaiters = []
activeBuffer = myBuffer // 之后 push 只会喂到「我」这
activeWaiters = myWaiters
return (async function* () { /* 从 myBuffer / myWaiters 逐个取 */ })()
}
splice(0) 把旧 buffer 掏空并接管,再把模块级的 activeBuffer/activeWaiters 重新指向本次订阅的私有队列。
自此 push 只会命中最新订阅者——旧订阅口自然断供。注释原话:把所有权转移给最新订阅者,让唯一一个
活跃的 subscribe() 收 chunk。
abort 怎么退出? 消费端挂起时,会给 abortSignal 挂一个 onAbort,一旦 abort 就 resolve(null)
(connection-adapters.ts:1045-1052);生成器看到 null 不 yield,循环条件 !abortSignal?.aborted 转假,干净退出。
3.4 真实实现:send 里读 connect 流并补终止事件
send(connection-adapters.ts:1058)才是真正调用底层 connect 的地方:它 for await 读 connect 的流,
每颗 chunk 都 push 进队列(connection-adapters.ts:1079),同时记录最近看到的 threadId/runId,
并监听是否出现过终止事件(RUN_FINISHED/RUN_ERROR,:1076-1078)。
关键在收尾:如果 connect 的流干净结束但从没发过终止事件,send 会合成一个 RUN_FINISHED
补上(connection-adapters.ts:1087-1106),让请求作用域的消费者能正常收尾;若 connect 抛错且没终止过,
则 合成 RUN_ERROR(:1107-1133)。合成事件复用调用方的 threadId/runId,好让 ChatClient 的
activeRunIds 追踪对得上(第1章)。这些 id 的兜底由 requireSyntheticId 守卫,见 §7.2。
4. 内置工厂:差异与选型
connect 型的实现全是现成工厂。它们只在三个维度上不同:HTTP 客户端(fetch/XHR/无)、
线格式(SSE / 换行 JSON / 直传对象)、以及是否解析 data:。
| 工厂 | 传输载体 | 线格式 | 是否解析 SSE(data:/[DONE]) | 典型场景 | 定义 |
|---|---|---|---|---|---|
fetchServerSentEvents | fetch 流 | SSE | 是 | 浏览器默认首选 | connection-adapters.ts:1259 |
fetchHttpStream | fetch 流 | 换行分隔 JSON(NDJSON) | 否(直接 JSON.parse 每行) | 服务器不发 SSE、发裸 JSON 流 | connection-adapters.ts:1430 |
xhrServerSentEvents | XMLHttpRequest | SSE | 是 | 无 fetch-stream 的环境(RN 等)要 SSE | connection-adapters.ts:1813 |
xhrHttpStream | XMLHttpRequest | 换行分隔 JSON | 否 | 无 fetch-stream 且服务器发裸 JSON | connection-adapters.ts:1898 |
webSocket | WebSocket | 帧内 JSON | 不适用 | 长连接 / 多端共享会话 | connection-adapters.ts:2099 |
stream | 无(直传 AsyncIterable) | 已是 StreamChunk | 不适用 | TanStack Start server function | connection-adapters.ts:2438 |
rpcStream | 无(RPC 通道) | 已是 StreamChunk | 不适用 | Cap'n Web 等 RPC 流 | connection-adapters.ts:2568 |
fetcherToConnectionAdapter | 用户给的 ChatFetcher | Response(按 SSE 解)或 AsyncIterable | 视返回值 | 把 fetcher 选项接进同一套管线(内部) | connection-adapters.ts:2468 |
选型速记:
- 能用
fetch流 + 服务器发 SSE →fetchServerSentEvents(最常见)。 - 目标环境
fetch流不可用(getResponseStreamReader会抛UnsupportedResponseStreamError,response-stream.ts:22)→ 换xhr*。 - 服务器发裸 JSON 行而非
data:包裹 → 选*HttpStream。 - 根本没有 HTTP(server function / RPC)→
stream/rpcStream,直接把AsyncIterable<StreamChunk>透传。
fetch 系与 xhr 系的共性: 两者 connect 里都做同样三步——解析 URL/options(支持函数式惰性求值)、
buildRunAgentInputBody 打包(§5)、发请求后把流交给解析器。差别只在「谁发字节」和「谁读字节」。
fetch 系读流用 readStreamLines(§6.1),XHR 系因为拿不到真正的 ReadableStream,改用
readXhrLines(connection-adapters.ts:1579)——靠 onprogress 增量读 responseText、按 offset 切新行。
stream / rpcStream 极简: 它们的 connect 就一句 yield* streamFactory(...)
(connection-adapters.ts:2450 / :2580),消息原样透传(保留 parts),转换成 ModelMessage 是服务端 chat() 的事。
fetcherToConnectionAdapter(内部桥): ChatClient 支持 fetcher 选项(一个直接发请求的函数,
types.ts:311 的 ChatFetcher)。这个工厂把 fetcher 包成 connect(chat-client.ts:151 处调用),
让 fetcher 走和其它适配器完全相同的 subscribe/send 管线。它强制要求 ChatClient 一定会传的
abortSignal 和 runContext,缺了直接抛错(connection-adapters.ts:2473-2482);fetcher 返回
Response 就按 SSE 解,返回 AsyncIterable 就用 abortableIterable 包一层可中断迭代(§7.3)。
断线续传(delivery durability)。 fetch/xhr 系的 SSE 与 NDJSON 工厂现在都内建
resumableStream 重连引擎(connection-adapters.ts:698):服务器若给事件打了 id: 偏移,
连接掉了会自动带 Last-Event-ID 重连、从最后偏移续播并去重;没打标签就退化成单次普通请求,
行为与从前完全一致。配套的两次新错误:DurableStreamIncompleteError(:72,续传无果)和
StreamReconnectLimitError(:87,连续无进展重连超上限,默认 5 次)。
5. 出站:把请求打包成 AG-UI RunAgentInput
5.1 请求体 buildRunAgentInputBody
服务器端遵循 AG-UI 协议,期望收到一个 RunAgentInput 形状的 JSON。buildRunAgentInputBody
(connection-adapters.ts:1193)就负责拼这个体:
| 字段 | 来源 | 说明 |
|---|---|---|
threadId / runId | runContext,缺则 generateRunId(...) 兜底 | 会话/本次运行的相关 id |
parentRunId | runContext.parentRunId(有才加) | 续跑/派生运行的父 id |
resume | runContext.resume(有才加) | interrupt 续跑的结算批(第 6 章) |
messages | uiMessagesToWire(messages)(§5.2) | 序列化后的线消息 |
tools | runContext.clientTools ?? [] | 客户端声明的工具(名/描述/JSON Schema) |
state / context | 空 {} / [] | AG-UI 结构占位 |
forwardedProps | 三层合并(见下) | 用户透传数据 |
data | { ...forwardedProps } | 旧字段名的镜像,向后兼容 |
优先级:later-spread-wins。 forwardedProps 由三个来源展开合并,后展开者覆盖前者
(connection-adapters.ts:1199-1206):
{ ...options.body, // ①静态适配器 body(构造时配置,最低优先)
...(runContext?.forwardedProps ?? {}), // ②本次运行的 forwardedProps
...data } // ③per-message 的 data(运行时,最高优先)
一句话:运行时值压过静态配置。ChatClient 那侧也遵循同样次序把 body 和 forwardedProps 合成
mergedBody 再传下来(chat-client.ts:2199-2247)。data 字段是给只认老字段名的消费者的镜像——
同一份数据挂两个名字,新老服务器都能读。
5.2 uiMessagesToWire:parts → AG-UI 线消息
uiMessagesToWire(ag-ui-wire.ts:62)把 TanStack 的 UIMessage(以 parts 数组为权威)翻成 AG-UI 线格式。
核心思想:每条锚点消息(system/user/assistant)原样带上 parts,再额外补 AG-UI 的镜像字段
(content、toolCalls),好让服务器端的 AG-UI Zod 解析通过。
不同 role 的处理:
| role | 产物 | 关键逻辑 |
|---|---|---|
system | 锚点 + content(纯文本) | collectText(parts),ag-ui-wire.ts:73-87 |
user | 锚点 + content(纯文本或多模态数组) | 有图/音/视/文档才用数组,否则纯文本,collectUserContent(:248) |
assistant | reasoning fan-out → 锚点 → tool fan-out | 见下,ag-ui-wire.ts:105-146 |
assistant 的「扇出」(fan-out)是精华。 一条 assistant 消息的 parts 里可能混着思考、文本、工具调用、工具结果。
除了把它们塞进锚点消息,还会额外把 thinking 和 tool-result 拆成独立的线消息,给严格的 AG-UI 服务器消费:
assistant.parts = [thinking, text, tool-call, tool-result]
│
├─► ① { role:'reasoning', id, content } ← thinking 扇出(锚点之前)
│ deriveReasoningId(msg.id, part) ag-ui-wire.ts:308-310
│
├─► ② 锚点 assistant 消息:{ ...msg, content?, toolCalls? }
│ content = collectText(parts) 仅 text !== '' 才带
│ toolCalls = collectToolCalls(parts) tool-call → function 镜像
│ ag-ui-wire.ts:290-306
│
└─► ③ { role:'tool', id, toolCallId, content } ← tool-result 扇出(锚点之后)
deriveToolMessageId(toolCallId) ag-ui-wire.ts:312-314
structured-output 回灌为 assistant content 的理由,以及 raw !== '' 守卫。
collectText(ag-ui-wire.ts:222)在拼 assistant 文本时,除了 text part,还会把已完成的
structured-output part 的原始 JSON(p.raw)当作文本喂回去(:237-243)。为什么?
因为多轮对话里,模型需要看见自己上一轮吐的那段结构化输出才能连贯——把它作为 assistant content 回灌就是干这个。
但有两个守卫:
- 只回灌
status === 'complete'的——流式中/出错的 part 会是残缺 JSON 片段,喂回去只会让模型犯迷糊。 p.raw !== ''守卫(ag-ui-wire.ts:240)——上游completeStructuredOutputPart会尽力填raw(调用方 → 已有 buffer →JSON.stringify(data)),但当data不可序列化(BigInt、循环引用)时, 这个兜底可能留下空串。守卫在此拦住:宁可不回灌,也不发一个''让模型看到「空的 assistant 轮次」。
6. 入站:把字节流解析回 StreamChunk
6.1 逐行读取与截断检测
readStreamLines(connection-adapters.ts:309)是最底层:从 reader 读字节、用 TextDecoder 增量解码、
按 \n 切行,把最后一个不完整的行留在 buffer 里(:331-332)等下次拼。
它的精华是截断检测。流正常结束时 buffer 应为空;若结束时 buffer 里还残留非空内容,
说明连接是在某一行中途被切断的(服务器崩溃、TCP 断、代理超时),于是抛 StreamTruncatedError
(connection-adapters.ts:351-358)。但有个例外:如果是消费者自己 abort 的(用户点了停止),
中途断行是预期的,不算 bug——所以 !abortSignal?.aborted 时才抛。结束前先 decoder.decode()
排空一次(:345-349),让多字节字符被截断的字节也显形为残缺尾部。
6.2 linesToSSEEvents:SSE 语义解析
responseToSSEChunks(connection-adapters.ts:627)是 fetch 系 SSE 的解析入口,真正的语义解析在
linesToSSEEvents(:419)——它在 readStreamLines 之上加 SSE 语义,fetch 与 XHR 两个 SSE 适配器共用:
跳过噪声行(connection-adapters.ts:442-448)——SSE 里有若干不是数据的行:
| 前缀 | 是什么 | 怎么处理 |
|---|---|---|
: | 注释行 | 跳过——代理/CDN 注入的 keepalive 心跳 |
event: | 事件名 | 跳过——本实现用 payload 里的 type 而非 SSE event 字段 |
id: | 事件偏移 | 不再跳过:记为 delivery-durability 的偏移(:428-436),断线重连时当 Last-Event-ID 用 |
retry: | 重连间隔 | 跳过——客户端不据此重连 |
数据行交给 parseSseDataLine(sse-utils.ts:1)剥掉 data: 前缀(容忍前导空格),也接受裸 JSON 行。
[DONE] 合成 RUN_FINISHED(connection-adapters.ts:450-466)——很多 SSE 服务器以 data: [DONE] 收尾,
但那不是一颗合法 StreamChunk。这里把它翻译成一个合成的 RUN_FINISHED,并复用最近见过的
threadId/runId/model(边解析边记在 lastThreadId 等,:468-475),让消费者收到一个带真实关联 id 的干净终止事件。
JSON 解析失败会直接抛,由消费者上报为错误。
6.3 chunkRunIds:给无 runId 的内容事件补 runId
一个协议现实:RUN_STARTED/RUN_FINISHED/RUN_ERROR 自带 runId,但内容事件
(TEXT_MESSAGE_CONTENT、TOOL_CALL_* 等)不带。可运行作用域的消费者(如「流式期间清空要抑制」的逻辑)
又需要知道每颗 chunk 属于哪次运行。
解法是一个 WeakMap,chunkRunIds(connection-adapters.ts:24-30)。在 §3.4 的 connect 包装里,
push(chunk, runContext?.runId)(:1079)会把调用方的 runId 盖章到这颗 chunk 上(push 内部
chunkRunIds.set,:1017-1020)。读取时统一走 getChunkRunId(connection-adapters.ts:37-44):优先用
WeakMap 里盖的请求方 runId——provider 可能在事件上盖自己的 run id,得让客户端的 run 身份赢——
没有才回退到 chunk 自带的 runId。用 WeakMap 是因为它以 chunk 对象为键、不阻止回收,
天然随 chunk 生命周期消失。
7. 健壮性守卫(三个小而关键的设计)
7.1 StreamTruncatedError:别把半截流当成功
StreamTruncatedError(connection-adapters.ts:51)在 §6.1 已讲:流以非空未终结 buffer 收尾就抛它。
意义在于让 ChatClient 转入 error 状态,而不是把一段被截断的流静默当成成功呈现给用户。
XHR 系在 finish(connection-adapters.ts:1622-1633)里也做同样判断。
7.2 requireSyntheticId:缺 id 就报错,绝不编造
当 §3.4 的 send 要合成 RUN_FINISHED/RUN_ERROR 时,得填 threadId/runId。这些本应由 ChatClient 的
runContext 提供(chat-client.ts:2234-2249)。requireSyntheticId(connection-adapters.ts:245)在
「上游流里也没见过、runContext 也没给」时直接抛错,而不是随手造一个假 id ——因为「此处缺 id」意味着有人
绕过了 ChatClient 的契约接线,应当暴露问题而非用假数据掩盖。这与文档创作里的「零编造」是同一种诚实。
(注:[DONE] 合成走另一条路,统一在 linesToSSEEvents 里用 ?? '' 兜底,:455-456,
因为那条路总有最近见过的上游 id 或调用方传入的 fallbackIds。)
7.3 abortableIterable:让不听话的迭代器也能被打断
fetcherToConnectionAdapter 若拿到一个 AsyncIterable,用 abortableIterable(connection-adapters.ts:2513)包一层。
问题:一个无视 signal 的生成器,会让 for await 一直挂到它自然结束。这个包装用
Promise.race([iterator.next(), abortPromise])(:2531)让「下一颗 chunk」和「abort 事件」赛跑——
abort 先到就 return,并在 finally 里调 iterator.return?.() 通知上游收尾。把「协作式取消」补齐成「强制可中断」。
8. 巧妙之处(可带走的技术)
- 所有权转移做多订阅安全(
connection-adapters.ts:1033):splice(0)+ 重指指针,让最新subscribe独占供给, 旧订阅口自动断供——不用引用计数、不用锁,一行搞定单消费者语义。 - 合成终止事件填补协议缝隙(
:1087/:450):无论底层流「干净结束不发终止」还是「[DONE]哨兵」, 都翻译成带真实 id 的RUN_FINISHED,让上层状态机永远能收尾。 - abort 感知的截断检测(
:356):同样是「buffer 非空收尾」,用户 abort 就当预期、否则当 bug—— 用一个信号位把「正常停止」和「异常截断」精确分开。 raw !== ''守卫防脏回灌(ag-ui-wire.ts:240):structured-output 序列化失败时宁可不发,也不让模型看到空轮次。- WeakMap 补 runId(
:24-30):给无 runId 的内容事件挂关联信息,又不干扰 GC。
9. 边界与局限
- SSE 层的
event:字段被无视(:442-448):本实现只认 payload 里的type。若你的服务器靠 SSE 的event:名区分事件,得改成把类型放进 JSON。id:则相反——它被当作 delivery-durability 偏移记下来 驱动断线重连(:428-436),服务器一旦打id:就进入了可续传模式。 - XHR 系不是真流,是轮询
responseText(readXhrLines,:1579):靠onprogress增量切片,内存里会累积 整个responseText(offset只移动读取位,不释放已读文本)。长响应下比fetch流更占内存——fetch系是首选, XHR 是兜底。 fetch流依赖运行时能力:Response.body.getReader/TextDecoder缺失会抛UnsupportedResponseStreamError(response-stream.ts:6),此时必须换xhr*或自定义传输。reasoningid 只求「够用」不求唯一:hashContent(ag-ui-wire.ts:316)是廉价确定性哈希, 注释明说容忍碰撞,因为 reasoning id 只对 AG-UI 外部消费者有意义,本项目自己的去重按toolCallId。
10. 横向对比(章内导航)
本章是「传输层」,和同组其它章的分工:
| 想了解 | 去哪 |
|---|---|
| chunk 到手后状态机怎么转、订阅循环怎么写 | 01-chat-client.md |
一颗颗 chunk 怎么拼成 UIMessage(与 §5 的反向过程) | 03-stream-processor.md |
useChat 怎么把 class 桥进框架 | 04-framework-bindings.md |
| 客户端工具、审批、续跑、多端生成态 | 06-tools-persistence-live.md |
一句话记忆: 出站 uiMessagesToWire 把 parts 摊平成线消息;入站 linesToSSEEvents 把字节拼回 chunk;
StreamProcessor 再把 chunk 拼回 parts——三者互为逆。
11. 代码地图(导航索引)
| 主题 | 文件路径 | 符号名 |
|---|---|---|
| 两种适配器形状 + 联合类型 | packages/ai-client/src/connection-adapters.ts:814, 929, 965 | ConnectConnectionAdapter / SubscribeConnectionAdapter / ConnectionAdapter |
| 归一化(核心) | packages/ai-client/src/connection-adapters.ts:975 | normalizeConnectionAdapter |
| connect 包装队列 | packages/ai-client/src/connection-adapters.ts:1014-1160 | push(闭包)/ activeBuffer / activeWaiters |
| 请求体打包 | packages/ai-client/src/connection-adapters.ts:1193 | buildRunAgentInputBody |
| 内置工厂 | packages/ai-client/src/connection-adapters.ts | fetchServerSentEvents / fetchHttpStream / xhrServerSentEvents / xhrHttpStream / webSocket / stream / rpcStream / fetcherToConnectionAdapter |
| 断线续传引擎 | packages/ai-client/src/connection-adapters.ts:698 | resumableStream / ReconnectTracker |
| SSE/行解析 | packages/ai-client/src/connection-adapters.ts:309, 419, 627, 1579 | readStreamLines / linesToSSEEvents / responseToSSEChunks / readXhrLines |
| SSE data 剥离 | packages/ai-client/src/sse-utils.ts | parseSseDataLine |
| 流能力探测 | packages/ai-client/src/response-stream.ts | getResponseStreamReader / createResponseStreamTextDecoder / UnsupportedResponseStreamError |
| 健壮性守卫 | packages/ai-client/src/connection-adapters.ts:51, 245, 2513 | StreamTruncatedError / requireSyntheticId / abortableIterable |
| runId 关联 | packages/ai-client/src/connection-adapters.ts:24-44 | chunkRunIds / getChunkRunId |
| 出站序列化(核心) | packages/ai/src/utilities/ag-ui-wire.ts | uiMessagesToWire |
| 文本/结构化回灌 | packages/ai/src/utilities/ag-ui-wire.ts | collectText(raw !== '' 守卫) |
| 多模态/工具/思考扇出 | packages/ai/src/utilities/ag-ui-wire.ts | collectUserContent / collectToolCalls / deriveReasoningId / deriveToolMessageId |
| 消费方(第1章) | packages/ai-client/src/chat-client.ts | resolveTransport / runContext 构造 / this.connection.send |