Inference Buffer:用 Durable Object 为长时推理流构建跨驱逐的持久响应缓冲
Inference Buffer用 Durable Object 为长时推理流构建跨驱逐的持久响应缓冲【免费下载链接】agentsBuild and deploy AI Agents on Cloudflare项目地址: https://gitcode.com/GitHub_Trending/agents1/agents本文介绍当前仓库中experimental/inference-buffer实验项目——一个验证「AI Gateway 作为长时推理调用持久响应缓冲」构想RFC #1257的独立 Cloudflare Worker Durable Object 原型。它解决的是 Cloudflare Agents 场景下 Agent DO 被驱逐时输出令牌被白白浪费的问题把 Provider 连接与 Agent 生命周期解耦让已生成的分块持久化在 SQLite 中Agent 重启后无缝续传。读完本文你将掌握这套缓冲原型的完整架构、五种状态的状态机、双 Provider 模式、无轮询通知信号机制、两条恢复数据流、五个 HTTP API以及如何在本地运行与测试它。问题背景DO 被驱逐时已计费的令牌被浪费在 Cloudflare Workers 上运行 AI Agent 时每个会话对应一个 Durable ObjectDO。当 Agent 正在流式推理过程中DO 遭遇驱逐例如代码部署触发版本更新时该 DO 到推理提供商的在途 HTTP 请求会被直接掐断。Agents SDK 的 fiber 恢复系统runFiber、ResumableStream、onChatRecovery虽然能恢复客户端侧状态但恢复的前提是发起一次全新的推理调用——这意味着 Provider 已经生成并且已经计费的令牌全部作废客户等于为输出令牌付了两次钱。这个问题在 agentic loop 场景下会成倍放大一次回合内如果包含多个串行工具调用每一次中断都会浪费该回合到目前为止已经生成的全部输出令牌。随着模型能力的提升浪费的代价也在陡增例如 README 中给出的对比gpt-4.1的一次重试成本约为gpt-4.1-mini的 4 倍。解决方案把 Provider 连接从 Agent 生命周期中剥离核心思路是在 Agent 与推理 Provider 之间插入一个独立的缓冲层┌──────────────┐ ┌──────────────────┐ ┌──────────────┐ │ Agent DO │──proxy───▶│ Inference Buffer │──fetch───▶│ Provider │ │ │◀─stream───│ (separate DO) │◀─stream──│ (OpenAI, │ │ │ │ │ │ Anthropic, │ │ [evicted] │ │ [keeps reading] │ │ Workers AI) │ │ │ │ [stores chunks │ │ │ │ [restarts] │ │ in SQLite] │ │ │ │ │──resume──▶│ │ │ │ │ │◀─replay───│ │ │ │ └──────────────┘ └──────────────────┘ └──────────────┘关键点在于Buffer Worker 是一个独立部署。当用户部署新的 Agent 代码时Agent 的 DO 会被驱逐但 Buffer 的 DO 不会。Provider 连接存活在 Buffer 的执行上下文中通过ctx.waitUntil挂起因此它能比调用方活得更久。从源码experimental/inference-buffer/src/index.ts可以确认这一设计动机The buffer Worker is a SEPARATE deployment from the agent Worker. When the user deploys new agent code, their agent DOs are evicted but the buffer DOs are not. The provider connection lives in the buffers execution context (viactx.waitUntil), so it survives the agents eviction.这正是整个原型的成立基础Provider 连接的生命周期归属于 Buffer 这个独立 Worker而不是归属于随时可能被驱逐的 Agent DO。与 AI Gateway 的映射原型在验证什么本原型是「AI Gateway 原生提供持久缓冲能力」的 PoC。README 给出了一份逐项对应表本原型AI Gateway 对应实现独立 Worker DOAI Gateway 基础设施本身就在请求路径上POST /proxyX-Provider-URL头现有代理流程新增X-AI-Gateway-Durable-Id头以显式开启SQLite 分块存储专门构建的流存储内存 持久化 spilloverGET /resume?fromN新端点GET /gateway/buffer/{id}?offset{n}GET /drain快照读取同上附加?snapshottrue标志POST /ack显式确认或通过 TTL 过期隐式清理5 分钟 TTL alarm可配置 TTL默认 5 分钟可按请求覆盖Agent → Buffer 的 Service BindingAI Gateway 已是 Cloudflare 服务无需 binding只需一个 URLX-Provider-Type: workers-ai用于 AI bindingAI Gateway 已路由 Workers AI——这是 no-op这个映射阐明了为何 AI Gateway 是承载该能力的最佳位置它是托管基础设施不会因为用户代码变更而被重新部署。原型验证了两件事——缓冲概念本身可行且 fiber 系统的onChatRecovery钩子提供了干净的集成点。AI Gateway 需要向 SDK 暴露的六项能力从 Agents SDK 的视角README 提出了六条具体的 API 需求通过请求头显式开启缓冲在任意被代理的请求上携带X-AI-Gateway-Durable-Id: id。存在该头时AI Gateway 将响应持久化缓冲不存在时保持原有的直通行为。Resume/replay 端点GET /gateway/buffer/{id}?from{chunk-index}——当缓冲仍在流式写入时尾部跟随实时流阻塞直到完成或回放已完成的缓冲。返回与 Provider 原始返回一致的原始 SSE 字节。状态端点GET /gateway/buffer/{id}/status→{ status: streaming|completed|interrupted|error, chunkCount: N }。格式无关的存储只存原始字节。格式转换原始 Provider SSE → AI SDK UIMessage 格式由 SDK 侧的流式转换器完成。这让 AI Gateway 不依赖任何库或框架。TTL默认 5 分钟可通过X-AI-Gateway-Buffer-TTL头按请求覆盖。清理自动进行不要求显式 ack不过提供 ack 端点以便急切清理也是很好的补充。Workers AI 直通对 Workers AI 模型AI Gateway 已处理路由。缓冲应对所有上游 Provider 一视同仁——只要携带 Durable-Id 头即开启。其中「格式无关存储 原始字节」这一条是理解整个设计的分水岭缓冲层永远不做 SSE 解析解析工作全部交给 SDK 侧的 Provider 解析器完成从而避免了格式转换逻辑在缓冲层腐烂。架构与实现原理核心技巧waitUntil后台消费 SQLite 尾部跟随当 Agent 调用/proxy时Buffer DO 执行三步动作源码见 experimental/inference-buffer/src/index.ts 的_startBuffering自己打开 Provider 连接HTTP Provider 走fetch()Workers AI 走this.env.AI.run()启动后台任务ctx.waitUntil(this._consumeProvider(reader))逐块读取 SSE 分块并写入 SQLite向 Agent 返回一个ReadableStream该流尾部跟随SQLite——按分块到达顺序持续吐出数据。如果 Agent 断开DO 被驱逐后台任务照常消费 Provider 流。当 Agent 重启并调用/resume?fromN时它会拿到一个从第 N 个分块开始的新尾部流。_consumeProvider的实现展示了关键细节src/index.ts用TextDecoder以{stream: true}模式解码每一块写入buffer_chunks表更新_chunkCount然后_notify()流结束后落尾 flush decoder 内部残留字节将状态改为completed并写入元数据表随后通过setAlarm(Date.now() 5 * 60_000)安排 5 分钟 TTL 清理。任何异常路径都会把状态置为error。两种 Provider 模式HTTP ProviderOpenAI、Anthropic 等Buffer 收到完整请求URL、头、体通过fetch()转发并缓冲流式响应。Authorization、Content-Type等头原样透传缓冲专属头X-Provider-URL、X-Buffer-*被剥离X-Provider-Type、X-AI-Model、Host、cf-connecting-ip同样不转发见 src/index.ts。Provider 返回非 2xx 时Buffer 将错误响应原样透传。Workers AIBuffer 持有自己的AIbinding配置见 wrangler.jsonc 中的ai: { binding: AI, remote: true }。调用方发送X-Provider-Type: workers-aiX-AI-Model: cf/model-nameBuffer 直接调用this.env.AI.run(model, body)src/index.ts。这样既不需要 API token又保留了 binding 的既有优势自动路由、无出口流量。如果模型返回的是非流式对象Buffer 会将其包装成 SSE 格式的data: {...}\n\ndata: [DONE]\n\n流从而统一所有 Provider 的消费路径。无轮询的通知信号机制Resume 的尾部流不轮询SQLite。_consumeProvider后台任务在每次插入分块后解析一个信号 promise_notify()尾部流的pull()先尝试从 SQLite 读出新分块读不到就await this._signal.promise阻塞被唤醒后再把可用的分块全部吐出src/index.ts。这个机制的安全性是由 Durable Object 的单线程执行模型保证的SQL 插入与信号解析发生在同一个同步块内因此被唤醒的尾部流一定能在 SQLite 中看到新分块不存在竞态窗口。Buffer 生命周期状态机idle ──[/proxy]──▶ streaming ──[provider done]──▶ completed ──[/ack or TTL]──▶ idle │ │ │ [DO evicted] │ [DO evicted restart] ▼ ▼ interrupted completed (restored) │ └──[TTL alarm]──▶ idle五种状态的定义类型见 src/index.tsidle无活动缓冲。/resume返回 404。streamingProvider 连接存活分块正在写入。此时再调/proxy返回 409已激活。completedProvider 已结束全部分块在 SQLite 中。Resume 可完整回放。interruptedProvider 流活跃期间 Buffer DO 被驱逐Provider 连接丢失。已存储的分块仍可读取调用方会拿到部分数据。error流式过程中 Provider 返回错误。值得注意的实现细节DO 构造函数中的_restore()src/index.ts会从 SQLite 元数据表恢复状态。如果持久化的状态是streaming说明上一个进程的 Provider 连接已随旧进程死亡于是标记为interrupted并立即安排 5 分钟 TTL alarm——否则这个半截缓冲会永久泄漏因为负责设置 alarm 的_consumeProvider已经随旧进程一起消亡了。恢复数据流两条路径Buffer 存储的是原始 Provider 字节。恢复时 SDK 根据 Buffer 状态走两条路径路径一streaming / completed —— Replay Model 组合回放Buffer (/resume) ──▶ Replay Model (real provider with replayFetch) ──▶ streamText() ──▶ clientReplay Model 的做法是用真实 Provider 模型ai-sdk/openai、ai-sdk/anthropic或workers-ai-provider创建一个自定义fetch让这个fetch返回 Buffer 的/resume响应而不是去调真实 API。Provider 自己维护的 SSE 解析器负责把字节流转换为LanguageModelV3StreamPartstreamText()原生处理工具执行、推理与toUIMessageStreamResponse()。零自定义 SSE 解析——如果 Provider 更新了格式Replay 会自动跟随。README 给出的三种 Provider 同一模式// All three providers use the same pattern: createModel: (fetch) createOpenAI({ apiKey: replay, fetch })(gpt-5.4); createModel: (fetch) createAnthropic({ apiKey: replay, fetch })(claude-sonnet-4-6); createModel: (fetch) createWorkersAI({ accountId: replay, apiKey: replay, fetch })( cf/moonshotai/kimi-k2.7-code );源码层面的印证在 experimental/forever-chat/src/replay-model.ts 的createReplayModel它是 Provider 无关的接受调用方传入的createModel工厂replayFetch通过 Service Binding 访问https://buffer/resume?id${bufferId}from0。同时它被设计成单次使用——第一次doStream回放 BufferstreamText工具调用步骤循环中的后续调用返回一个空 finish 流让循环干净地终止。路径二interrupted / error —— 累积解析 持久化 续写Buffer (/drain) ──▶ SSE Parsers ──▶ { text, reasoning, toolCalls } ──▶ execute tools ──▶ persist ──▶ continueLastTurn当 Buffer 的 Provider 连接已死Buffer DO 被驱逐时没有实时流可供回放。此时累积式解析器experimental/forever-chat/src/sse-parsers.ts从已存储的分块中提取文本、推理内容和工具调用。恢复过程中会执行服务端工具带审批检查部分响应被持久化然后continueLastTurn生成剩余部分。parseProviderStream按 Provider 分发实现了多种 SSE 方言的解析源码注释与实现均可在 sse-parsers.ts 中核对OpenAI同时兼容 Chat Completionschoices[0].deltadelta.tool_calls[]按索引增量累积delta.reasoning_content/delta.reasoning与 Responses APIresponse.output_text.delta、response.output_item.added、response.function_call_arguments.delta两种格式Anthropiccontent_block_delta的text_delta、thinking_delta思维块需先由content_block_start标记 type 为thinking、工具调用由content_block_starttool_useinput_json_delta增量拼接Workers AI先尝试 OpenAI 兼容格式无结果时回退到原生data: {response:text}格式。API 参考所有端点都需要通过?idbuffer-id查询参数或X-Buffer-Id请求头提供 buffer IDWorker 入口的路由逻辑见 src/index.ts先解析 bufferId再用env.INFERENCE_BUFFER.idFromName(bufferId)按名字取 DO 实例。缺失 ID 时返回 400并附带端点清单提示。POST /proxy向推理 Provider 转发请求并缓冲响应。HTTP ProviderPOST /proxy?idfiber-abc123 X-Provider-URL: https://api.openai.com/v1/chat/completions Authorization: Bearer sk-... Content-Type: application/json {model: gpt-4.1, messages: [...], stream: true}Workers AIPOST /proxy?idfiber-abc123 X-Provider-Type: workers-ai X-AI-Model: cf/moonshotai/kimi-k2.7-code Content-Type: application/json {messages: [...], stream: true}返回SSE 流响应头带X-Buffer-Status: streaming。源码中的校验细节HTTP 模式缺X-Provider-URL返回 400Workers AI 模式缺X-AI-Model返回 400body 非合法 JSON 返回 400AI.run抛错返回 502缓冲已处于streaming时返回 409。GET /resume从指定偏移回放已缓冲的分块。若 Provider 仍活跃则阻塞直到流完成。GET /resume?idfiber-abc123from0响应头X-Buffer-Status取值为streaming、completed、interrupted、error之一。from为回放起点分块索引缺省为 0。GET /drain快照读取——返回当前已存储的所有分块后立即关闭。不阻塞。用于 Buffer 处于 interrupted 或 error 状态时的部分恢复。GET /drain?idfiber-abc123from0源码还额外返回X-Buffer-Chunk-Count头src/index.ts方便调用方判断拿到了多少比例的数据。GET /status{ status: completed, chunkCount: 142 }POST /ack确认收到并触发提前清理。若仍在 streaming 则返回 409。源码实现src/index.ts会清空buffer_chunks与buffer_meta两张表、重置为idle并deleteAlarm()。打包路径三种演进方向当下作为 RFC 交付物原型experimental/inference-buffer/experimental/forever-chat/演示了完整的端到端流程。forever-chat示例入口见 experimental/forever-chat/src/server.ts展示了 Buffer 与 OpenAI、Anthropic、Workers AI 三种 Provider 的集成包括恢复时的流式回放以及对照的其他恢复策略Workers AI 用continueLastTurn()合并文本与推理、OpenAI 用 Responses API 的store: true取回完整响应、Anthropic 用合成用户消息续写。若 AI Gateway 落地该能力Buffer Worker 变得多余。Replay Model 模式用自定义fetch读取 Buffer 的真实 Provider 组合作为内部工具迁入cloudflare/ai-chat。用户用一个开关即可开启export class MyAgent extends AIChatAgentEnv { override durableBuffer true; // routes inference through AI Gateway buffer }SDK 负责模型包装所有 Provider 的自定义 fetch——Workers AI 也通过 workers-ai-provider 新增的fetch选项支持、通过 Replay Model 组合做恢复编排、通过streamText的原生步骤循环处理工具调用。AI Gateway 负责缓冲、存储、TTL、清理。全程不需要任何自定义 SSE 解析——真正的 Provider 解析器完成所有格式转换。若 AI Gateway 暂不落地把 Buffer Worker 发布为部署模板。用户将其与自己的 Agent Worker 一起部署并添加 Service Binding。SDK 的集成点指向该 bindingoverride durableBuffer { binding: this.env.INFERENCE_BUFFER };等 Gateway 就绪后用户只需把 binding 换成 Gateway 配置——Agent 代码零改动。本地运行与测试进入实验目录并启动cd experimental/inference-buffer pnpm start # starts at http://localhost:8787package.json中start脚本为wrangler devdeploy为wrangler deploytypes为wrangler types env.d.ts --include-runtime false。测试脚本有两个均需 Buffer Worker 运行在localhost:8686pnpm start -- --port 8686test-e2e.shexperimental/inference-buffer/test-e2e.sh完整快乐路径——用内嵌的 Node HTTP mock 服务作为 SSE Provider先完整代理一条 10 分块的流校验/status再从第 5 块/resume回放后半段最后/ack清理。test-eviction.shexperimental/inference-buffer/test-eviction.sh驱逐模拟——这是关键场景。启动一个每 300ms 发一块的慢 Providercurl 代理 1 秒后killcurl 进程模拟 DO 驱逐随后分两次查/status验证 Buffer 在后台继续消费Provider 仍被调用且只有一次Provider 结束后从第 3 块/resume拿到全部错过的数据——脚本最后输出的信息正是DONE — buffer survived caller disconnect。生产化考量README 明确列出原型为验证目的已够用、但 AI Gateway 落地时需要解决的工程问题存储后端原型使用 DO SQLite。AI Gateway 应构建专用流存储内存 持久化 spillover以获得更低延迟。缓冲大小上限当前无上限。长 agentic 回合尤其涉及代码生成可能产生很大响应。多租户隔离原型中 Buffer ID 是全局 UUID。AI Gateway 需要按账号/Gateway 做作用域隔离。可观测性追踪缓冲命中率、未命中率、部分恢复率以及每次恢复省下的令牌数。成本模型缓冲能力计入现有 AI Gateway 定价还是独立计费层级Provider 特定优化OpenAI Responses API 的previous_response_id与 Anthropic 未来可能的 resume 支持都可以与原始字节缓冲叠加利用。总结experimental/inference-buffer是一个小而完整的技术原型它用「独立部署的 Worker 持有waitUntil后台任务 SQLite 分块存储 信号通知尾部流」四件事证明了推理流缓冲可以完全剥离 Provider 连接与 Agent DO 的生命周期让 Agent 在被驱逐后能以「零重复调用、零浪费令牌」的方式续传输出。它与forever-chat示例一起Replay Model 组合 累积式 SSE 解析器为 AI Gateway 未来原生承载该能力提供了可直接迁移的实现路径和明确的 API 契约。对于研究 Cloudflare Agents 长时推理可靠性、fiber 恢复机制或 AI Gateway 能力的开发者这是一个值得通读源码的参考实现。【免费下载链接】agentsBuild and deploy AI Agents on Cloudflare项目地址: https://gitcode.com/GitHub_Trending/agents1/agents创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
上一篇/下一篇内容由系统自动关联
返回资讯列表 →