Electric Streams 与 Vercel AI SDK 集成:用 Durable Transport 让 useChat 生成可恢复、可共享
【免费下载链接】electricThe agent platform built on sync.项目地址: https://gitcode.com/GitHub_Trending/el/electric
本篇技术指南讲解如何在基于 Electric Streams(Durable Streams 协议的托管实现)的应用中,通过@durable-streams/aisdk-transport把 Vercel AI SDK 的useChat默认 Transport 替换为持久化 Transport,让一次聊天生成在页面刷新、网络抖动、重渲染等场景下依然存活,并能够跨标签页、跨设备、跨用户与跨 Agent 恢复和共享。读完本文,你将掌握客户端与服务端的完整接入方式、resume 流程的设计要点,以及这一集成所依赖的 Durable Streams 底层协议原理。
为什么聊天生成需要"可恢复"
大多数 AI 应用在连接出现任何问题时都会中断:网络不稳定、页面导航或一次重渲染打断了正在进行的长时间生成,导致已经流式输出的内容丢失,用户不得不重新提问。Vercel AI SDK 意识到了这一挑战,并在其 UI 层提供了 Transport 接口作为扩展点——接入方可以自定义客户端与服务端之间的通信方式,同时保留正常的useChat数据流。
Electric Streams 提供的 Durable Streams 正是解决这类问题的原生数据原语:它们是持久、可寻址、追加写入、可按 offset 重放的 HTTP 流,专门服务于 agent 循环与实时数据场景。将 AI SDK 的 Transport 层替换为基于 Durable Streams 的实现,即可获得三方面能力:
- 韧性(resilience):生成过程中的消息与 token 全部落到持久化流中,断网重连后无需重新开始;
- 可恢复(resumability):任何客户端都可以从任意 offset 重新订阅并续读,刷新页面后无缝衔接;
- 协作(collaboration):连接到同一会话的多个客户端订阅并写入同一条流,天然支持多标签页、多设备、多用户与多 Agent 的实时与异步协作。
这一集成基于 Durable Streams 协议概述中描述的核心概念,并遵循我们在 Durable Transports for your AI SDK 发布说明中定义的Durable Sessions 模式。
安装
在项目中使用 pnpm 安装 transport 包:
pnpm add @durable-streams/aisdk-transport该包同时提供客户端 Transport(createDurableChatTransport)与服务端响应包装(toDurableStreamResponse),因此一个依赖即可覆盖两端。若你的服务端还需要直连 Durable Streams 服务器,可参考 Quickstart 中通过 curl 创建、追加、读取与实时 tail 一条流的基本流程。
客户端:把默认 Transport 换成createDurableChatTransport
在客户端,保持useChat的调用方式不变,仅将默认 Transport 替换为createDurableChatTransport,并开启resume:
import { useChat } from "@ai-sdk/react" import { createDurableChatTransport } from "@durable-streams/aisdk-transport" const transport = createDurableChatTransport({ api: "/api/chat" }) const chat = useChat({ transport, resume: true })api指向你现有的聊天接口(如/api/chat),请求/响应协议与 AI SDK 默认 Transport 保持一致;resume: true让useChat在初始化时主动尝试恢复上一次未完成的生成。
从实现角度看,createDurableChatTransport遵循 AI SDK Transport 的同一套模型:客户端不再以一次性请求/响应的方式消费服务端返回的 SSE 流,而是从 Durable Stream 上订阅并消费 token 数据。这意味着useChat的消息状态机、messages渲染与发送流程几乎不用改动,接入成本集中在"换一个 transport 对象"。
服务端:用toDurableStreamResponse包装消息流
在服务端,把 AI SDK 的 UI 消息流用toDurableStreamResponse包装后返回:
import { toDurableStreamResponse } from "@durable-streams/aisdk-transport" return toDurableStreamResponse({ source: result.toUIMessageStream(), stream: { writeUrl: buildWriteStreamUrl(streamPath), readUrl: buildReadProxyUrl(request, streamPath), headers: DURABLE_STREAMS_WRITE_HEADERS, }, })各字段的作用:
source:AI SDK 生成结果(result.toUIMessageStream())产生的 UI 消息流,即 LLM 流式输出的 chunk 序列;stream.writeUrl:Durable Stream 的写入地址,服务端将 AI SDK chunk 逐条 POST 追加到该流;stream.readUrl:Durable Stream 的读取地址,通常是一个由你控制的代理端点(proxy),用于拼接上游读 URL、附加读鉴权头并转发offset、live等查询参数;headers:写入流时携带的服务端鉴权头(如DURABLE_STREAMS_WRITE_HEADERS)。
服务端完成两件事:把 AI SDK 的 chunk 写入 Durable Streams,同时在响应头Location和响应体{ streamUrl }中返回读取 URL。客户端拿到该 URL 后即可订阅流、按 offset 续读,从而实现"同一条流"上的恢复与多端共享。
Resume flow:让刷新安全的完整流程
要让一次生成能够跨刷新存活,需要在业务层配合以下四个步骤(对应文档中的 Resume flow 一节):
- 持久化进行中生成所属的流 id:在生成进行期间,把当前活跃的 stream id 与聊天(chat)关联并持久化(例如写入数据库或本地存储),作为恢复的锚点;
- 新增一个重连端点:例如
GET /api/chat/:id/stream,用于按 chat id 查询当前是否有进行中的生成; - 按状态返回正确响应:没有活跃生成时返回
204(无内容,客户端无需恢复);存在活跃生成时返回200,并携带Location头与{ streamUrl }响应体,指向可续读的流地址; - 在
useChat中开启resume: true:客户端初始化时先请求重连端点,拿到streamUrl便从流的末尾位置恢复订阅,无缝接续被中断的生成。
这套"先探活、再续读"的设计,本质上是 Durable Streams 协议中消费者模型的应用:客户端保存上一次读取返回的Stream-Next-Offset,重连时从该 offset 继续GET,即可不重放整段会话(详见 Durable Streams 协议概述中的"Offsets"与"Consumers"小节)。
工作原理:Durable Streams 协议如何支撑恢复与共享
要理解为什么"换一个 transport"就能获得韧性,需要回到底层协议。Durable Streams 的每条流都是一个URL 可寻址、追加写入、持久有序的字节序列:数据一旦写入某个位置就不会改变,新数据只能追加到末尾,位置由不透明的、可字典序比较的offset标识。协议定义了创建(PUT)、追加(POST)、读取(GET)、元数据(HEAD)、关闭与删除六种操作,读取支持offset=-1(从头重放)与offset=now(仅订阅未来数据)两种哨兵值。
对 AI 聊天场景而言,最关键的三个特性是:
- 持久性与重放:写入的数据被持久化存储(Electric Streams 服务器以 Rust 实现,每条流按线上字节原样落盘,catch-up 读就是一次字节区间读取,默认
wal模式下追加在写入日志确认后才应答,天然可恢复)。因此 token chunk 不会因为页面刷新而丢失; - live 模式:消费者追赶完历史数据后,可通过
?live=sse或?live=long-poll实时订阅新数据。SSE 每约 60 秒由服务器周期性关闭以配合 CDN 连接折叠,客户端用最后一次control事件中的streamNextOffset重连,这正好对应resume: true场景下的"断线自动续读"; - 幂等写入与消费:协议支持
Producer-Id/Producer-Epoch/Producer-Seq三头部的幂等生产者语义,重试不会产生重复数据;读取则通过Stream-Next-Offset、Stream-Up-To-Date、Stream-Closed头驱动一个简单而可靠的读循环。
多端共享则来自"同一条流"的拓扑:任何客户端订阅并写入同一条流,服务端把 AI SDK chunk 也写入同一条流,所有订阅者都会实时收到同一份数据,无论数据来自哪个用户或哪个 Agent。这正是 Durable Streams 协议概述所描述的"catch-up 重放 + 实时扇出"统一由同一原语支撑的体现。
进阶:从原始字节流到 Durable Sessions
原始 Durable Stream 处理的是字节/token 流;当应用需要把 AI token 流与结构化状态(如工具调用结果、用户在线状态、共享文档)复用同一基础设施时,可以在其上叠加State Protocol,形成分层协议栈,也就是 Durable State 文档定义的Durable Sessions模式:
- Durable Streams—— 可靠、可恢复的字节投递;
- State Protocol—— 流上的结构化 CRUD 操作(
insert/update/delete变更事件与 snapshot / reset 控制事件); - 应用层协议—— AI SDK Transport、presence、CRDT(如 Yjs)等。
在这种分层下,一个 Agent 把 token 流式写入会话的同时,结构化状态(工具结果、用户在场、共享文档)可以流经同一套基础设施,支持多用户与多 Agent 的实时与异步协作。对消息历史的处理,集成遵循inversion of control(控制反转)原则:可以由你选择从会话流中物化历史(对应姊妹篇 TanStack AI 集成中的materializeSnapshotFromDurableStream思路),也可以把消息物化到 Postgres 等数据库,集成不强制规定持久化方式。
示例与下一步
原文档指向的包 README 与chat-aisdk示例应用位于上游 durable-streams 仓库(本仓库之外),建议结合本仓库内的以下材料做完整落地:
- Durable Streams 协议概述:offset、live 模式、生命周期、CDN 缓存等协议概念;
- Quickstart:启动
durable-streams-server并用 curl 完成建流、写入、读取、实时 tail 的全流程; - CLI 文档:用
durable-stream命令行管理流,适合联调验证; - Rust 服务器实现说明:了解持久化(WAL)、零拷贝读取、分层存储与 OpenTelemetry 观测等服务器侧实现;
- TanStack AI 集成:同一 Durable Sessions 模式在 TanStack AI Connection Adapter 上的对应实现,可作为模式对照。
接入完成后,一次典型的体验是:用户在弱网或刷新后回到页面,useChat自动通过重连端点找回进行中的生成并续读 token 流;多个标签页或设备订阅同一条流,实时看到同一份生成内容;而 Agent 产生的写入也汇入同一条流,实现人机共写一个会话。
【免费下载链接】electricThe agent platform built on sync.项目地址: https://gitcode.com/GitHub_Trending/el/electric
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考