Forever Chat:基于 Cloudflare Durable Object 的永不中断 AI 流式对话(多 Provider 恢复实战)
2026/9/18 22:17:29 网站建设 项目流程

Forever Chat:基于 Cloudflare Durable Object 的永不中断 AI 流式对话(多 Provider 恢复实战)

【免费下载链接】agentsBuild and deploy AI Agents on Cloudflare项目地址: https://gitcode.com/GitHub_Trending/agents1/agents

本文围绕开源仓库agentsexperimental/forever-chat示例,系统讲解如何在 Cloudflare Workers 上构建"永不中断"的 AI 流式聊天应用:利用AIChatAgentrunFiber持久化纤维机制为每一轮对话自动续命(keepAlive),在 Durable Object(DO)被驱逐(eviction)后通过onChatRecovery按 Provider 差异化恢复,并用continueLastTurn()无缝衔接被中断的助手消息。读完本文,你将掌握三种主流 LLM Provider(Workers AI / OpenAI / Anthropic)的恢复策略选择、可运行的完整示例搭建方法,以及底层keepAliverunFiber、SQLite 持久化的实现原理。

背景:为什么 AI 聊天需要"持久化流式"

Durable Objects 有三种被驱逐的典型原因(详见 experimental/forever.md):

  1. 空闲超时:约 70~140 秒没有传入请求或打开的 WebSocket;
  2. 代码更新 / 运行时重启:非确定性,每天约 1~2 次;
  3. Alarm 处理器超时:15 分钟。

对 AI Agent 而言,驱逐发生在活跃工作期间是灾难性的:

  • 到 LLM Provider 的上游 HTTP/SSE 连接被永久切断,无法在中途恢复 OpenAI 或 Anthropic 的流;
  • 内存态(流缓冲、部分响应、循环计数器)全部丢失;
  • 已连接的客户端看到流无故停止;
  • 多轮 Agent 循环(工具调用、推理链、编排)完全丢失位置。

最常见的场景是 Agent 连续执行多轮 LLM 调用——每一轮只有几秒到几分钟,但整个会话可能持续 15~30 分钟以上。驱逐可能发生在两轮之间(丢失循环位置)或流式中途(丢失正在进行的生成),两者都必须被处理。forever-chat正是为此设计的端到端演示。

示例概览:forever-chat展示了什么

根据 experimental/forever-chat/README.md,该示例的核心演示内容包括:

  • AIChatAgent与 Think 中的常开恢复纤维(always-on recovery fibers):每一轮对话都被包进runFiber
  • 流式期间的 keepAlive:长时间 LLM 响应期间 DO 保持存活,不会空闲被驱逐;
  • onChatRecovery:驱逐后按 Provider 差异化恢复;
  • continueLastTurn():无缝继续被中断的助手消息内联输出;
  • 多 Provider 支持 + 下拉选择器

整体入口实现位于 experimental/forever-chat/src/server.ts,前端界面在 experimental/forever-chat/src/client.tsx。

三层架构:keepAlive → runFiber → Chat Recovery

仓库的设计文档 experimental/forever.md 把整个机制分成两层内建于Agent基类的能力,加上AIChatAgent在其上的第三层封装:

层级原语用途
1keepAlive()通过 alarm 心跳防止空闲驱逐
2runFiber()持久化执行——注册进 SQLite、可检查点、可恢复
3聊天恢复(AIChatAgent每一轮聊天都跑在纤维里,中断被检测、恢复有界,onChatRecovery支持 Provider 级续写策略

forever-chat就是把第三层完整落地的最小可运行项目。

运行 forever-chat

README 给出了完整的启动流程(需在仓库根目录执行):

npm install cd experimental/forever-chat cp .env.example .env # add your API keys npm start

几点说明:

  • Workers AI 开箱即用:它直接使用AIbinding(见 wrangler.jsonc 中的"ai": { "binding": "AI", "remote": true }),无需任何 API Key;

  • OpenAI 与 Anthropic 需要在.env中配置密钥,模板见 .env.example:

    OPENAI_API_KEY=your-openai-api-key ANTHROPIC_API_KEY=your-anthropic-api-key
  • npm start实际执行vite dev(见 package.json 的 scripts),由 Vite + Cloudflare 插件(vite.config.ts)驱动本地 workerd 运行时;

  • npm run deploy则执行vite build && wrangler deploy发布到线上;npm test运行 vitest 单元测试(覆盖 replay-model)。

wrangler 配置中还声明了两个值得注意的点:

  • Durable ObjectForeverChatAgent注册为 SQLite 类("new_sqlite_classes": ["ForeverChatAgent"]),这是runFibercf_agents_runs落盘的基础;
  • Service BindingINFERENCE_BUFFER绑定到inference-buffer服务,用于可选的"推理缓冲"恢复路径;
  • secretsOPENAI_API_KEYANTHROPIC_API_KEY为线上部署必需。

多 Provider 恢复策略

README 的核心是一张三 Provider 对比表,它是理解整个示例的钥匙:

ProviderModelRecovery strategy
Workers AIkimi-k2.7-codePersist partial + continue viacontinueLastTurn()(text + reasoning 合并进既有 block)
OpenAIgpt-5.4Retrieve completed response via Responses API(store: true)——零浪费 token
Anthropicclaude-sonnet-4.6Persist partial + continue via synthetic user message(恢复时禁用 reasoning)

对应到 server.ts 中的_providerRecovery(),三种策略的代码分支逻辑清晰:

if (provider === "anthropic") { await this.schedule(0, "_continueWithUserMessage", undefined, { idempotent: true }); return { continue: false }; } if (provider === "workersai" || !provider) { await this.schedule(0, "_continueWorkersAI", undefined, { idempotent: true }); return { continue: false }; } if (provider !== "openai") return {}; return this._openAIResponsesRecovery(ctx);

下面逐一展开三种策略的原理与实现。

Workers AI:continueLastTurn()无缝续写

Workers AI(以及未指定 Provider 时的默认路径)走_continueWorkersAI()

async _continueWorkersAI() { const ready = await this.waitUntilStable({ timeout: 10_000 }); if (!ready) return; this._lastBody = { ...this._lastBody, recovering: true }; await this.continueLastTurn(); }

其背后的机制(来自 forever.md):

  1. continueLastTurn()找到this.messages中最后一条助手消息;
  2. 用保存的_lastBody_lastClientTools重新调用onChatMessage
  3. 把新流式输出作为**续写(continuation)**追加到已有助手消息,而不是新建一条;
  4. 不产生合成用户消息——用户看到的体验是"被中断的消息直接从停下的地方继续变长"。

同时,由于 Workers AI 走 prefill 续写,代码把recovering: true写进_lastBodyonChatMessage会在instructions里追加RECOVERY_SUFFIX(" Do not think or reason — continue the text output directly."),提示模型跳过推理直接续写文本。即使模型仍输出了推理内容,框架也会自动把 reasoning 合并进已有的 reasoning block。

OpenAI:Responses API 按 ID 取回完整响应

OpenAI 恢复策略是最特别的:由于 OpenAI Responses API 支持服务端继续生成,连接断开后生成仍在服务端继续。因此恢复时不需要重新请求,只需用流式期间 stash 下来的responseId去取回完整响应,实现零浪费 token

关键实现分三处:

  1. 流式期间 stash responseId_getChunkHandlers):
onChunk: ({ chunk }) => { if (chunk.type !== "raw") return; const raw = chunk.rawValue; if (raw?.type === "response.created" && raw.response?.id) { // 将 responseId 写入 buffer stash 或直接 this.stash({ responseId }) } }, includeRawChunks: true
  1. 请求时开启store: true_getProviderOptions):
return { openai: { store: true, reasoningEffort: "low", reasoningSummary: "auto" } };
  1. 恢复时_openAIResponsesRecovery:从ctx.recoveryDataresponseId,调用GET https://api.openai.com/v1/responses/{responseId},若data.status === "completed",就把完整文本作为新的 assistant 消息直接persistMessages([...this.messages]),然后返回{ persist: false, continue: false }——既不需要持久化部分响应,也不需要续写。

源码注释特别解释了为什么不能复用_persistOrphanedStream:如果 DO 在 chunk 缓冲(10 个 chunk)刷入 SQLite 之前就被驱逐,getStreamChunks()会返回[],导致旧路径失效,所以这里直接手写消息落盘。

Anthropic:合成用户消息续写

Anthropic不支持 assistant prefill,所以continueLastTurn()(本质是把对话以"结尾是部分 assistant 消息"的方式重发)不可用。替代方案是_continueWithUserMessage()

async _continueWithUserMessage() { const ready = await this.waitUntilStable({ timeout: 10_000 }); if (!ready) return; this._lastBody = { ...this._lastBody, recovering: true }; await this.saveMessages((messages) => [ ...messages, { id: crypto.randomUUID(), role: "user", parts: [{ type: "text", text: "Your previous response was interrupted. Please continue exactly where you left off." }], metadata: { synthetic: true } } ]); }

即调度saveMessages追加一条metadata.synthetic = true的用户消息引导模型续写(前端在 client.tsx 中通过message.metadata?.synthetic跳过该合成消息的渲染,用户无感知)。同时恢复调用会通过_getProviderOptions对 Anthropic 传入{ thinking: { type: "disabled" } }(正常对话则为{ thinking: { type: "adaptive" } }),即恢复时禁用 reasoning,只续写文本。

恢复上限:chatRecovery.maxAttempts

README 明确指出:如果反复恢复超过chatRecovery.maxAttempts,框架会持久化配置好的终止消息(terminal message),而不是让该轮对话卡死。这意味着恢复不是无界的——框架通过有界恢复(bounded recovery incidents)避免死循环,并在达到上限时给出确定性收尾。

第四种路径:Inference Buffer 恢复

除了 README 表格里的三种 Provider 策略,server.ts 的头部注释还描述了一条可选的通用缓冲恢复路径,可用于任何 Provider:

  • Inference Buffer (any provider): resume from the durable response buffer — zero wasted tokens, zero duplicate provider calls

onChatRecovery中,如果this.state?.useBuffer为 true,会优先尝试_tryBufferRecovery(ctx),失败才回退到 Provider 特定恢复。其状态分支逻辑:

  • streaming/completed:调度流式回放。先持久化一条空的 assistant 消息(否则continueLastTurn因缺少最后一条 assistant 消息而跳过),然后通过createReplayModel从 buffer 读原始 SSE,交给真实 Provider 的模型实现_getReplayModelapiKey: "buffer-replay"+ 自定义fetch重新构造 OpenAI/Anthropic/Workers AI 模型)做实时格式转换,客户端看到 token 像从 Provider 实时流入一样;
  • interrupted/error:用 sse-parsers.ts 的累加解析器(parseProviderStream)从原始字节中提取文本、reasoning 与工具调用,执行可恢复的服务器端工具,持久化部分响应后调度续写;
  • idle/ 空 buffer:返回null,落回 Provider 特定恢复。

该路径通过自定义 fetch / 伪造 AI binding 把 Provider 调用"劫持"进INFERENCE_BUFFER服务绑定(_routeThroughBuffer为每次调用生成crypto.randomUUID()作为 bufferId 并stash),同时把原始 Provider SSE 完整落盘。这样即便 Agent 的 DO 被驱逐,Provider 连接本身存活在 buffer 服务中,从而可以做到零重复调用。

sse-parsers.ts还系统梳理了三种 Provider 的 SSE 格式差异:

  • OpenAI:Chat Completions(delta.content/delta.tool_calls[])与 Responses API(response.output_text.delta/response.output_item.added/response.function_call_arguments.delta)两种格式都会被parseOpenAIStream识别;
  • Anthropiccontent_block_deltatext_delta)、thinking_delta(推理)与content_block_starttool_use)+input_json_delta(工具参数增量);
  • Workers AI:先按 OpenAI 兼容格式解析,失败再回退原生{"response":"text"}格式。

Replay Model 的设计哲学:零漂移

replay-model.ts 实现了一个巧妙设计:恢复时不自己解析 SSE,而是把 buffer 的原始响应通过replayFetch喂给真实 Provider 的模型(createOpenAI/createAnthropic/createWorkersAI),让官方维护的 SSE 解析器完成格式转换。注释中明确指出:

if @ai-sdk/openai updates how it parses Responses API SSE, the replay model picks it up automatically. Zero drift.

同时它是单次使用的:第一次doStream回放 buffer,后续调用(来自 streamText 工具调用循环)返回一个空的finish流(createEmptyFinishStream)让循环干净终止;doGenerate直接抛错("Replay model is stream-only")。这些行为都被 replay-model.test.ts 中的"second doStream returns empty finish (single-use)""doStream throws on buffer fetch failure""doGenerate throws"等用例验证。

恢复中途的工具调用如何处置

当 buffer 部分响应里含有工具调用时,_buildRecoveredParts+_executeRecoveredTool决定如何重建消息:

  • 无需审批的服务器端工具:直接执行,输出作为tool-xxxpart 以output-available状态呈现;
  • 需要审批的工具(如calculateMath.abs(a) > 1000needsApproval返回 true):返回合成错误消息("requires user approval which was interrupted — please try again"),标记resolved: false
  • 纯客户端工具(无execute定义,如getUserTimezone):返回合成错误("Tool could not be executed during recovery — please retry"),标记resolved: false
  • 执行抛出异常:返回{ error: \Tool failed during recovery: ${e}` },但仍标记resolved: true`(错误本身就是确定的输出)。

恢复完成后,未解析(hasUnresolvedTools)的调用会由后续续写轮次决定重试或提示用户。这套逻辑保证了即使中断发生在工具执行阶段,聊天状态也不会"悬空"。

前端如何呈现"恢复"

client.tsx 通过useAgent(来自agents/react)连接ForeverChatAgent,通过useAgentChat(来自@cloudflare/ai-chat/react)发送消息,并渲染消息中的textreasoning(Thinking 块)、工具输出、审批请求(Approve / Reject)等 part。界面细节:

  • Provider 下拉框:切换lastProvider并写入 agent state,禁流式期间切换;
  • Buffer 开关:切换useBuffer,开启后所有 Provider 调用经推理缓冲(对应界面上的"零浪费 token"说明);
  • ConnectionIndicator:绿/黄/红三点显示 WebSocket 连接状态;
  • 合成用户消息(metadata.synthetic)被直接跳过渲染,因此用户感知到的就是"上一条消息继续长出内容"。

如何验证恢复确实发生

README 给出了手工验证方法:

Start a long response, then restart the dev server while the model is still streaming. On the next activation, the agent records one recovery incident and either continues the partial assistant turn or retries the unanswered user turn, depending on where interruption happened.

即:发起一个长响应,在模型仍在流式输出时重启 dev server。下一次激活时,Agent 会记录一次恢复事件,并根据中断点位置决定:

  • 若中断发生在部分助手消息处 → 继续该部分回复(continueLastTurn路径);
  • 若中断发生在用户消息未被回答时 → 重试该未回答的用户轮次;
  • 若恢复次数超过chatRecovery.maxAttempts→ 持久化配置好的终止消息,避免回合卡死。

底层恢复的判定逻辑(来自 forever.md):每一条聊天轮次被包成名为__cf_internal_chat_turn:{requestId}的 fiber;DO 重启后onStart()_checkRunFibers()会扫描cf_agents_runs表,把不在内存活跃集合中的行视为"被中断",调用内部恢复钩子解析requestId、从cf_ai_chat_stream_metadata读取流 chunk,再回调onChatRecoveryrequestId编码在 fiber 名而非 snapshot 中,意味着snapshot 完全是开发者自己的领域——在onChatMessagethis.stash({ responseId })不会覆盖框架数据,恢复时ctx.recoveryData即用户 stash 的内容。

本地 workerd 会把 SQLite 与 alarm 状态持久化到磁盘,因此本地开发与线上行为一致:fiber 运行中杀掉进程(SIGKILL / Ctrl-C)→ 重启后首个请求触发恢复。

底层机制速览:keepAlive 与 runFiber

如果你想把forever-chat的模式复用到自己的 Agent 中,需要理解两层原语(详见 experimental/forever.md):

keepAlive():引用计数的 alarm 心跳。首次持引用时ctx.storage.setAlarm(now + keepAliveIntervalMs)(默认 30 秒,可用static options = { keepAliveIntervalMs: 2_000 }覆盖以便测试);alarm 触发时执行日常清理并继续设置下一个。所有 disposer 调用后 refs 归零,alarm 停止,DO 自然空闲。它不产生 schedule 行,对listSchedules()不可见。

runFiber(name, fn):SQLite 中先INSERTcf_agents_runs(仅 4 列:idnamesnapshotcreated_at),再持有 keepAlive、执行fn(ctx)ctx.stash()同步写检查点、完成后DELETE行。被驱逐后由_checkRunFibers()恢复,onFiberRecovered钩子拿到name+snapshot决定如何重跑。相比旧的 12 列cf_agents_fibers表,新表刻意保持极简——恢复逻辑属于开发者的钩子,而不是框架管理的状态字段。

forever-chat正是把这三层能力组合起来的完整参考实现:keepAlive保活流式响应、runFiber持久化每一轮、onChatRecovery按 Provider 差异化恢复、continueLastTurn无缝续写、Inference Buffer 实现零浪费恢复。无论是想学习持久化 Agent 的实现原理,还是需要一个可直接改造的多 Provider 聊天模板,它都是理想起点。

【免费下载链接】agentsBuild and deploy AI Agents on Cloudflare项目地址: https://gitcode.com/GitHub_Trending/agents1/agents

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询