十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

agents/streams 深入解析:为 Cloudflare Lifecycle Object 构建可断点续传的持久化增量输出流

agents/streams 深入解析:为 Cloudflare Lifecycle Object 构建可断点续传的持久化增量输出流 agents/streams 深入解析为 Cloudflare Lifecycle Object 构建可断点续传的持久化增量输出流【免费下载链接】agentsBuild and deploy AI Agents on Cloudflare项目地址: https://gitcode.com/GitHub_Trending/agents1/agents本文基于仓库docs/agents/streams.md并结合packages/agents/src/streams/源码与测试展开讲解agents/streams这一实验性能力如何在Agent或普通 Durable Object 上获得“有序、持久、可回放、可续传”的增量输出流。读完你将掌握流的打开/追加/收尾/消费完整生命周期、基于单调游标的重放续读、SSE 一键服务、块式存储与事务化 cutover以及它与 Tasks、Chat 的组合方式并能在自己的 Lifecycle Object 上直接落地。实验性声明agents/streams导出的所有 API 在持久化输出面稳定之前可能随版本变化见 streams.md。使用时应锁定版本并留意 CHANGELOG。一、Streams 是什么为 Lifecycle Object 提供持久化增量输出agents/streams为 Lifecycle Object 增加持久的增量输出能力每个流stream对应一条有序的持久化 chunk 日志带一个单调递增的游标cursor消费端支持“先重放、后实时跟随replay-then-tail”的读取方式并具备终态terminal status。它的核心价值在于消费端断线重连从自己记录的下一个游标开始重放不多发、不少发生产端崩溃恢复生产者中途死亡时恰好留下它已持久化追加的 chunks重放的生产者可以从流自身游标继续天然具备“续传而非重做”的语义无需 Alarm该能力只依赖标准能力服务storage 与 events因此也可以运行在 facet 上。源码头注释印证了这一设计定位Streams拥有cf_agents_streams与cf_agents_stream_blocks两张表实现“有序、持久化的每流 chunk 日志 单调游标 重放/续读 终态”该终态同时充当 Tasks 能力的恢复证据见 packages/agents/src/streams/streams.ts。从结构上看它继承LifecycleCapability只消费 Lifecycle 提供的标准服务因此是一个“即插即用”的组合式能力与Scheduler、WebSockets等能力处于同一抽象层。二、安装与使用两种宿主接入方式2.1 普通 Durable Object能力组合import { DurableObject } from cloudflare:workers; import { Lifecycle } from agents/lifecycle; import { Streams } from agents/streams; export class ReportObject extends DurableObjectEnv { readonly streams new Streams(); readonly lifecycle Lifecycle.install(this).use(this.streams); }Lifecycle.install(this)会构造生命周期并安装运行时处理器fetch、alarm、WebSocket 系列.use(this.streams)将该能力注册进生命周期。能力按注册顺序执行启动阶段顺序运行各能力的onStart请求阶段以中间件方式分发——第一个返回Response的能力接管请求返回undefined则继续传递见 docs/agents/lifecycle.md。2.2 Agent在组合根构造器中安装对Agent而言把额外能力装进组合根是标准套路——在构造函数里调用this.lifecycle.use(...)export class ReportAgent extends AgentEnv { readonly streams new Streams(); constructor(ctx: AgentContext, env: Env) { super(ctx, env); this.lifecycle.use(this.streams); } }Agent本身已经构造了 lifecycle并安装 WebSockets 能力因此这里只需追加注册。2.3 可配置项new Streams(options)目前支持一个配置项见 streams.ts配置项类型默认值说明maxChunkBytesnumberDEFAULT_MAX_CHUNK_BYTES 1_048_5761 MiB单个序列化 chunk 的字节数上限超出抛StreamSerializationError值得说明的是测试里用new Streams({ maxChunkBytes: 1024 })构造过最小宿主见 packages/agents/src/tests/capabilities/streams.ts而 chat 的ResumableStream适配器则用1_900_000的更高上限容纳打包段见下文“Chat 场景”。streams包的全部导出见 packages/agents/src/streams/index.ts包括Streams、三个错误类StreamClosedError、StreamNotFoundError、StreamSerializationError、sseResponse以及全部类型。三、生产端open → append → close/errorconst stream await this.streams.open(reply:123, { metadata }); stream.append(chunk); // 同步持久化写入唤醒实时读者 stream.close(); // 或 stream.error(reason)3.1 chunk 与游标chunk 是JSON 值StreamJsonstring / number / boolean / null / 数组 / 纯 JSON 对象默认单 chunk 上限 1 MiB可用maxChunkBytes调整每次append分配下一个单调递增的序号——即该流的游标cursor。序号从 0 开始重放时保持稳定见 packages/agents/src/streams/types.tsappend是同步持久化写入返回本次 chunk 的seq并立即唤醒正在实时跟随的读者#wake机制。3.2 open() 的幂等语义open()对id 幂等见 streams.ts重开一个存活streaming的流返回定位在当前游标的 writer——这是“续传生产者”的基础重开一个已收尾settled的流抛StreamClosedError对同一 id 二次收尾no-op幂等保证恢复调用方安全。收尾本身是“流状态变为精确”的时刻#settleRow用一条受保护的UPDATE ... WHERE state streaming在同一同步块内写入最终游标、closed_at、chunk_count与error_message见 streams.ts。3.3 writer 句柄open()返回的StreamWriter暴露三个成员见 types.tsstreamId流 idcursor只读下一个将要分配的序号append(chunk): number持久化追加并唤醒读者close(options?)/error(reason?, options?)以completed/errored收尾二者均幂等commit回调仅在本次调用真正完成收尾时执行。3.4 追加护栏append fence源码中#append的顺序值得关注序列化先于护栏检查——因为JSON.stringify可能执行用户对象的toJSON()后者可能同步重入该流随后在同一同步块内完成“状态检查 → 读日志尾部 → 写块”中间不允许 await 或调用用户代码否则会重新引入丢失更新的竞态。Durable Object 单隔离、单同步块执行模型保证了状态检查与写入的原子性见 streams.ts。id 本身有约束必须是非空字符串且不超过 256 字符否则直接抛错#validateStreamId见 streams.ts。四、消费端重放再实时跟随for await (const chunk of this.streams.read(reply:123, { from, signal })) { // 先重放 from 之后已持久化的 chunks然后实时跟随新追加 // 流收尾时结束 } const status await this.streams.status(reply:123); // { state: streaming | completed | errored, cursor, ... } | nullread(streamId, { from, signal })from为包含起点序号默认 0先重放、后跟随signal中止会抛出其 reason见 types.ts。读取不依赖生产者存活状态read()内部实际上委托给readBatches()见 streams.ts因此语义完全一致status()返回StreamStatus | null包含state、cursor 持久化 chunk 总数 / 下一个序号、可选的tag/metadata/error以及createdAt/updatedAt/closedAt。对存活流游标和活跃时间从 chunk 日志尾部推导流的行内计数器在存活期“刻意不维护”对终态流则直接读取 settle 时写入的精确值见 streams.tslist({ state, tag, limit })按状态、标签过滤最新在前ORDER BY created_at DESC, stream_id DESC默认limit100state可传单个值或数组见 streams.tsdelete(streamId)删除一个已收尾的流及其 chunk 日志返回是否删除存活流必须先 settle否则抛错见 streams.ts。读取行为细节源码注释与实现共同确认对一个已errored的流read()仍会先把已持久化的 chunks 全部产出然后正常结束——终端结果通过status()查询。当流在读取中途被删除时读取立即返回#state为undefined即return。五、Tagid 之外的业务查找键Tag 是 id 的查找侧lookup sideopen(id, { tag })会给流打上一个被索引的应用键——例如请求 id、一个会话——并且刻意不要求唯一。一个产生连续多个流的操作一次重试的 turn、一次重新生成的回复会为每个流打上同一个 taglist({ tag, limit: 1 })就能找到最新的那个流结果最新在前tag 在创建时固定重开一个存活流并传入不同tag 会抛错——这是“配置冲突”而非“新流”见 streams.ts 与 types.ts。作者的建议原文档原意只用一个 id 就够的场景先用 id只有当“一个操作可能拥有多个流”时才需要 tag。源码层面cf_agents_streams表上建有idx_cf_agents_streams_tag (tag, created_at)索引支撑该查找。六、批量读取与 onUpToDate当消费者按写付费时——一次 SSE flush、一次 RPC 跳、一次历史追加——应当批量读取而不是逐 chunk 处理for await (const batch of this.streams.readBatches(reply:123, { from, batchSize: 50 // 重放阶段每个数组的上限默认 100 })) { flush(batch); // StreamChunk[] —— 一个积压只产生一次写而非每个 chunk 一次 }readBatches()的生命周期与read()完全一致重放 → 跟随 → 收尾结束差别只在粒度重放阶段每次产出最多batchSize个 chunk 的数组batchSize默认 100最小强制为 1见 streams.ts实时跟随阶段每次醒来后累积的所有新 chunk 合并为一个数组产出。readBatches还额外接受onUpToDate回调在读者首次到达持久化日志尾部时调用一次。注意“追平caught-up”与“结束ended”是两个概念一个存活流在实时跟随期间始终处于“up to date”状态——可用它来 flush 重放出的 UI或点亮“live”指示器见 types.ts。实现上短批次即“已到尾部”的证据回调同步执行后重新轮询避免“先唤醒后注册”导致丢事件见 streams.ts。七、与 Tasks 组合流即恢复证据Tasks 重放机制正是围绕“流”的契约设计的一个 task 步骤向它并不拥有的流追加数据因为生产者的循环从流自身的持久化游标开始所以中断后的重放就是一次续传——流就是恢复证据。readonly tasks new Tasks({ definitions: { generatev1: async (input: GenerateInput, step: TaskStep) { return step.do(stream, async () { const stream await this.streams.open(input.streamId); // 续传生产者从流自身游标开始所以中断后的重放不会重复任何 chunk。 for (let i stream.cursor; i input.total; i) { stream.append(await this.produce(i)); } stream.close(); }); } } });几个要点零依赖两个能力互不 import“composes through checkpointed cursors, never imports”见 streams.ts真实进程杀死下依然成立死前已追加的 chunks 恰是死后status()报告的内容——这一点由SIGKILL e2e 测试套件证明仓库中的 packages/agents/src/tests/capabilities/streams.ts 提供了TaskStreamComposeObject最小宿主它记录每次生产者重新进入时的entryCursors验证重放是“续传而非重做”从而证明“无 chunk 重复”的不变量。八、服务层一条调用完成整个 SSE 生命周期对 SSE 场景一个辅助函数即可服务流的整个生命周期import { sseResponse } from agents/streams; async onRequest(request: Request) { return sseResponse(this.streams, reply:123, { request }); }agents/streams会做以下全部事情见 packages/agents/src/streams/sse.ts原生续传每个 chunk 的序号写在 SSE 的id:字段上。浏览器EventSource重连时自动发送Last-Event-ID辅助函数从下一个 chunk 继续——零客户端代码实现游标持久化也支持?from查询参数。注意实现细节Last-Event-ID必须“先检查存在再解析”否则首次连接会因Number(null) 0而跳过第 0 个 chunk见 sse.ts重放从续传点重放已持久化 chunks控制事件到达日志尾部时发出up-to-date控制事件随后实时跟随新追加心跳跟随期间周期性发送: heartbeat注释帧默认每 30 秒heartbeatMs可配0 关闭避免空闲代理掐断连接终态帧流completed发doneerrored发error携带记录的 reason断连中止请求的signal在客户端断开时中止跟随流不存在时返回 404。SSEResponseOptions完整可配项见 sse.ts选项默认说明request无用于续传Last-Event-ID头 /?from与断连中止from由 request 推导显式指定起始序号覆盖 request 推导的续传点signal无额外中止信号与request.signal组合batchSize100每次写入的最大 chunk 数heartbeatMs30000心跳间隔毫秒0 禁用其他传输方式下read()/readBatches()就是可直接自行接管的异步迭代器。examples/next/streams 是端到端演示项目。九、存储设计rollover blocks 与事务化 cutover9.1 块式存储rollover blocksChunks 以滚动块的形式存储见 streams.ts每个流一行block不断通过UPDATE ... SET body body || , || chunk追加直到块体达到256 KBBLOCK_MAX_CHARS见 streams.ts下一次追加就 INSERT 新的一行每次追加无论如何都恰好一次计费行写入要么是长块体的 UPDATE要么是下一块的 INSERT——与“每 chunk 一行”的日志计费相同但上千个 chunk 的流只有屈指可数的几行因此删除它只需要几次写入而不是上千次重放时逐块解析块体是逗号连接的 JSON 文本[${body}]即可还原见 streams.ts块表cf_agents_stream_blocks使用WITHOUT ROWIDCloudflare 按“写行数”计费索引维护普通 rowid 表会为主键维护隐藏的唯一索引——每次 INSERT/DELETE 多计一行WITHOUT ROWID让主键即表本身chunk 追加精确计费一行见 streams.ts。而流的元数据表刻意保留 rowid 表rowid 作为同毫秒创建时间下“最新在前”排序的确定性 tiebreak且每个流只计费一次每次 turn 一行而非每 chunk。9.2 cutover一次事务完成“收尾 落库 清理”流是临时的一旦其内容变成了别的东西一条会话消息、一份报告它的行就成了死重。cutover在一个 SQLite 事务内完成结束流 → 运行你自己的同步写入 → 删除流的行stream.close({ commit: () sessionSync.upsert(message), // 仅允许同步写入 discard: true // 在同一事务内删除该流的行 });语义保证见 types.ts 与 streams.ts要么消息存在且流已删除要么commit抛错、settle 回滚、流仍然存活——绝无中间态也不留给 retention 清扫任何残留error(reason, { commit, discard })对失败的生产者同样适用commit绝不能 awaitSession 句柄的__DO_NOT_USE_WILL_BREAK__sync().upsert()是对应的同步消息写入它返回一个notify()用于在事务提交后派发变更推送事件与读者唤醒属于 SQLite 之外的副作用只在事务提交后触发回滚的 cutover 不会产生任何事件若流已是终态或被删除commit不会执行也不会删除任何东西幂等。原文档给出了一组真实 Durable Object 上的实测数据400-chunk 的 chat turn每次写 10 个 chunk旧日志方案写入 42 行、清扫又要 42 行块方案写入仍是 42 行但 cutover 只需3 行。仓库中的 packages/agents/src/e2e-tests/stream-cutover-crash.test.ts 在真实 Durable Object 死亡的各个关键点块追加提交前/后、滚动到下一块前后、会话持久化后清理前等验证每次重启要么找到“精确已提交的流前缀”要么找到“最终会话消息”——两者必居其一绝不两无也绝不应用两次。组合宿主见CutoverHarnessObjectStreams Sessions 同一 Lifecycle见 packages/agents/src/tests/capabilities/streams.ts。十、Chat 场景AIChatAgent 与 Think 的落地点AIChatAgent与Think把运行中的 turn 输出存放在这里ResumableStream来自agents/chat是 Streams 之上的薄适配层见 packages/agents/src/chat/resumable-stream.ts把约10 个线上 chunk打包进**一个存储段segment**以节省写入次数CHUNK_BUFFER_SIZE 10并提高 backing Streams 的maxChunkBytes到 1_900_000createChatStreams()见 resumable-stream.ts每个 turn 都以cutover结束助手消息、流的 settle、流行的删除在一个事务内提交——不再有 alarm 驱动的清扫崩溃遗留的流只有两种归宿仍处于streaming恢复时从它重建消息或在下一次流启动时被回收连同超过一小时未活跃的 in-flight 行一起ABANDONED_STREAM_RETENTION_MS 60 * 60 * 1000以“最后 chunk 活动”而非“打开时间”度量避免误杀长时活跃流见 resumable-stream.ts已有的cf_ai_chat_stream_*表会自动迁移到该能力上schema v1 → v2 采用惰性折叠某个流的旧逐-chunk 行在首次触碰该流时才折叠成块启动绝不一次性为整条旧日志买单见 streams.ts。打包模式值得任何高频生产者复制把已经持有的数据在同步段内缓冲一次性追加一个打包 chunk读取时再拆包——持久性不变收尾时没有任何内容跨越 await 被悬置而行写入比逐 token 追加减少约一个数量级。十一、当前限制与演进方向实时扇出live fanout在隔离in-isolate内这对当前模型是充分的——一个 Durable Object 同一时刻只在一个隔离中执行所有并发读者共享生产者的隔离比隔离活得更久的读者在重连时从自己的游标重放保留是显式的delete()或 cutover 的discard未来工作能力内部按年龄自动清扫、open()上的生产者代数隔离producer-generation fencing、以及从 chat 的恢复协议中抽离传输辅助函数完整设计记录见 design/rfc-streams.md记录着 shipped 设计与实现演化过程chunk 日志机制如何从 chat 词汇中剥离、为何流即恢复证据等。十二、小结何时使用 Streams如果你符合以下任一场景agents/streams就是为 Lifecycle Object 提供持久化增量输出的现成答案流式输出需要断线续传先重放再实时跟随EventSource天然续传?from兜底生产者可能中途崩溃游标就是恢复证据close/error幂等配合 Tasks 一步做到“续传而非重做”内容最终要“变成别的东西”用 cutover 把流收尾、业务写入、行删除放进一个事务杜绝残留清扫聊天 / 长 turn 输出chat 栈已经跑在这套能力之上打包模式 事务化 cutover 是经过生产级 e2e 验证的成熟范式。在写作本文时上述结论均以仓库当前状态为准docs/agents/streams.md、packages/agents/src/streams/*、packages/agents/src/chat/resumable-stream.ts、packages/agents/src/tests/capabilities/streams.ts与packages/agents/src/e2e-tests/stream-cutover-crash.test.ts接口仍属实验性升级时请对照对应版本的 CHANGELOG。【免费下载链接】agentsBuild and deploy AI Agents on Cloudflare项目地址: https://gitcode.com/GitHub_Trending/agents1/agents创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表