传输与线协议: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:244):
| 形状 | 你要实现的方法 | 心智模型 | 谁在用 |
|---|---|---|---|
ConnectConnectionAdapter | connect(messages,…) → AsyncIterable<StreamChunk> | 「一发一收」:调一次,拿一条流,for await 读完即止 | 全部七个内置工厂 |
SubscribeConnectionAdapter | subscribe() → AsyncIterable + send() → Promise | 「先订阅,后投递」:订阅是长期的口子,send 只管把请求推出去,chunk 从订阅口回来 | 需要长连接/多路复用的自定义传输(如共享 WebSocket) |
类型定义分别在 ConnectConnectionAdapter(connection-adapters.ts:212)和
SubscribeConnectionAdapter(connection-adapters.ts:224)。
为什么要两种? 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 / 直连流 ─────► responseToSSEChunks │
│ │ · 跳过 :注释/event:/id:/retry: │
│ ▼ · [DONE] → 合成 RUN_FINISHED │
│ ③服务器回字节流 (SSE/NDJSON) · JSON.parse 每行 │
│ │
└──────────────────────────────────────────────────────────────────────────────┘
│
▼
ChatClient.subscribe() 的 for-await 逐个拿到 chunk(第1章:喂给 StreamProcessor)
怎么读:①②③是「发出去」,④⑤是「收回来」。中间的 activeBuffer/activeWaiters 队列(§3.2)是把
「connect 的一条流」桥接到「subscribe 的长期订阅口」的关键。
2.3 部件一句话职责
| 部件 | 干什么 | 在哪 |
|---|---|---|
ConnectionAdapter 联合类型 | 定义两种适配器形状 | connection-adapters.ts:244 |
normalizeConnectionAdapter | 把任意适配器统一成 subscribe/send | connection-adapters.ts:254 |
buildRunAgentInputBody | 把消息+上下文打包成 AG-UI RunAgentInput 请求体 | connection-adapters.ts:425 |
uiMessagesToWire | 把 UIMessage.parts 序列化成 AG-UI 线消息 | ag-ui-wire.ts:47 |
responseToSSEChunks / readStreamLines | 把 SSE/NDJSON 字节流解析回 StreamChunk | connection-adapters.ts:144 / :94 |
| 七个内置工厂 | 各种传输的现成实现 | connection-adapters.ts(§4) |
StreamTruncatedError / requireSyntheticId / abortableIterable | 健壮性守卫 | connection-adapters.ts:41 / :60 / :1002 |
3. 核心原理:把 connect 包成 subscribe/send
3.1 归一化要解决的小问题
ChatClient 构造时,不管你给的是 connect 型还是 subscribe/send 型,它都想拿到统一的 subscribe/send。
这就是 normalizeConnectionAdapter(connection-adapters.ts:254)的活。ChatClient 在
chat-client.ts:188 用它包一次,之后只调 this.connection.subscribe(...) / .send(...)。
它先做互斥校验——一个适配器只能是一种形状,同时给 connect 和 subscribe/send 直接抛错
(connection-adapters.ts:265)。然后分三条路:
normalizeConnectionAdapter(connection)
│
┌───────────────────────┼───────────────────────┐
▼ ▼ ▼
同时有 connect 和 有 subscribe+send 只有 connect
subscribe/send → 原样绑定返回 → 用 activeBuffer/
→ 抛错(形状冲突) (它已是目标形状) activeWaiters 队列包一层
原生就是 subscribe/send 的适配器直接 .bind 透传(connection-adapters.ts:271-276);难点全在第三条——
只有 connect 的适配器,怎么假装成 subscribe/send?
3.2 思路:一个 activeBuffer / activeWaiters 异步队列
connect 是「调用即得一条流」;subscribe 是「先开一个长期订阅口,chunk 陆续来」。要把前者变后者,
需要一个中间队列来解耦生产(send 里读 connect 流)和消费(subscribe 的 for-await)。
这个队列就两个数组(connection-adapters.ts:285-286):
activeBuffer—— chunk 来了但当前没人等,先囤在这。activeWaiters—— 有人await下一个 chunk 但还没货,把它的 resolve 回调囤在这。
push(connection-adapters.ts:288)是生产端:有等待者就直接喂给它,否则塞进 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:301-307):
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:316-323);生成器看到 null 不 yield,循环条件 !abortSignal?.aborted 转假,干净退出。
3.4 真实实现:send 里读 connect 流并补终止事件
send(connection-adapters.ts:329)才是真正调用底层 connect 的地方:它 for await 读 connect 的流,
每颗 chunk 都 push 进队列(connection-adapters.ts:350),同时记录最近看到的 threadId/runId,
并监听是否出现过终止事件(RUN_FINISHED/RUN_ERROR,:347)。
关键在收尾:如果 connect 的流干净结束但从没发过终止事件,send 会合成一个 RUN_FINISHED
补上(connection-adapters.ts:356-372),让请求作用域的消费者能正常收尾;若 connect 抛错且没终止过,
则合成 RUN_ERROR(:373-393)。合成事件复用调用方的 threadId/runId,好让 ChatClient 的
activeRunIds 追踪对得上(第1章)。这些 id 的兜底由 requireSyntheticId 守 卫,见 §6.2。
4. 七个内置工厂:差异与选型
connect 型的实现全是现成工厂。它们只在三个维度上不同:HTTP 客户端(fetch/XHR/无)、
线格式(SSE / 换行 JSON / 直传对象)、以及是否解析 data:。
| 工厂 | 传输载体 | 线格式 | 是否解析 SSE(data:/[DONE]) | 典型场景 | 定义 |
|---|---|---|---|---|---|
fetchServerSentEvents | fetch 流 | SSE | 是 | 浏览器默认首选 | connection-adapters.ts:490 |
fetchHttpStream | fetch 流 | 换行分隔 JSON(NDJSON) | 否(直接 JSON.parse 每行) | 服务器不发 SSE、发裸 JSON 流 | connection-adapters.ts:574 |
xhrServerSentEvents | XMLHttpRequest | SSE | 是 | 无 fetch-stream 的环境(RN 等)要 SSE | connection-adapters.ts:825 |
xhrHttpStream | XMLHttpRequest | 换行分隔 JSON | 否 | 无 fetch-stream 且服务器发裸 JSON | connection-adapters.ts:895 |
stream | 无(直传 AsyncIterable) | 已是 StreamChunk | 不适用 | TanStack Start server function | connection-adapters.ts:939 |
rpcStream | 无(RPC 通道) | 已是 StreamChunk | 不适用 | Cap'n Web 等 RPC 流 | connection-adapters.ts:1045 |
fetcherToConnectionAdapter | 用户给的 ChatFetcher | Response(按 SSE 解)或 AsyncIterable | 视返回值 | 把 fetcher 选项接进同一套管线(内部) | connection-adapters.ts:963 |
选型速记:
- 能用
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:663)——靠 onprogress 增量读 responseText、按 offset 切新行。
stream / rpcStream 极简: 它们的 connect 就一句 yield* streamFactory(...)
(connection-adapters.ts:947 / :1053),消息原样透传(保留 parts),转换成 ModelMessage 是服务端 chat() 的事。
fetcherToConnectionAdapter(内部桥): ChatClient 支持 fetcher 选项(一个直接发请求的函数,
types.ts:56 的 ChatFetcher)。这个工厂把 fetcher 包成 connect(chat-client.ts:88 处调用),
让 fetcher 走和其它适配器完全相同的 subscribe/send 管线。它强制要求 ChatClient 一定会传的
abortSignal 和 runContext,缺了直接抛错(connection-adapters.ts:968-977);fetcher 返回
Response 就按 SSE 解,返回 AsyncIterable 就用 abortableIterable 包一层可中断迭代(§6.3)。
5. 出站:把请求打包成 AG-UI RunAgentInput
5.1 请求体 buildRunAgentInputBody
服务器端遵循 AG-UI 协议,期望收到一个 RunAgentInput 形状的 JSON。buildRunAgentInputBody
(connection-adapters.ts:425)就负责拼这个体:
| 字段 | 来源 | 说明 |
|---|---|---|
threadId / runId | runContext,缺则 generateRunId(...) 兜底 |