免费获取学习方案
ARTICLE DETAIL

资讯详情

深耕编程基础知识与建站技术分享的一线实战洞察。

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

Forever Chat:基于 Cloudflare Durable Object 的永不中断 AI 流式对话(多 Provider 恢复实战) Forever Chat基于 Cloudflare Durable Object 的永不中断 AI 流式对话多 Provider 恢复实战【免费下载链接】agentsBuild and deploy AI Agents on Cloudflare项目地址: https://gitcode.com/GitHub_Trending/agents1/agents本文围绕开源仓库agents中experimental/forever-chat示例系统讲解如何在 Cloudflare Workers 上构建永不中断的 AI 流式聊天应用利用AIChatAgent的runFiber持久化纤维机制为每一轮对话自动续命keepAlive在 Durable ObjectDO被驱逐eviction后通过onChatRecovery按 Provider 差异化恢复并用continueLastTurn()无缝衔接被中断的助手消息。读完本文你将掌握三种主流 LLM ProviderWorkers AI / OpenAI / Anthropic的恢复策略选择、可运行的完整示例搭建方法以及底层keepAlive、runFiber、SQLite 持久化的实现原理。背景为什么 AI 聊天需要持久化流式Durable Objects 有三种被驱逐的典型原因详见 experimental/forever.md空闲超时约 70140 秒没有传入请求或打开的 WebSocket代码更新 / 运行时重启非确定性每天约 12 次Alarm 处理器超时15 分钟。对 AI Agent 而言驱逐发生在活跃工作期间是灾难性的到 LLM Provider 的上游 HTTP/SSE 连接被永久切断无法在中途恢复 OpenAI 或 Anthropic 的流内存态流缓冲、部分响应、循环计数器全部丢失已连接的客户端看到流无故停止多轮 Agent 循环工具调用、推理链、编排完全丢失位置。最常见的场景是 Agent 连续执行多轮 LLM 调用——每一轮只有几秒到几分钟但整个会话可能持续 1530 分钟以上。驱逐可能发生在两轮之间丢失循环位置或流式中途丢失正在进行的生成两者都必须被处理。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-chatREADME 给出了完整的启动流程需在仓库根目录执行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 KeyOpenAI 与 Anthropic 需要在.env中配置密钥模板见 .env.exampleOPENAI_API_KEYyour-openai-api-key ANTHROPIC_API_KEYyour-anthropic-api-keynpm 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]这是runFiber表cf_agents_runs落盘的基础Service BindingINFERENCE_BUFFER绑定到inference-buffer服务用于可选的推理缓冲恢复路径secretsOPENAI_API_KEY与ANTHROPIC_API_KEY为线上部署必需。多 Provider 恢复策略README 的核心是一张三 Provider 对比表它是理解整个示例的钥匙ProviderModelRecovery strategyWorkers AIkimi-k2.7-codePersist partial continue viacontinueLastTurn()text reasoning 合并进既有 blockOpenAIgpt-5.4Retrieve completed response via Responses APIstore: true——零浪费 tokenAnthropicclaude-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 AIcontinueLastTurn()无缝续写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.mdcontinueLastTurn()找到this.messages中最后一条助手消息用保存的_lastBody与_lastClientTools重新调用onChatMessage把新流式输出作为**续写continuation**追加到已有助手消息而不是新建一条不产生合成用户消息——用户看到的体验是被中断的消息直接从停下的地方继续变长。同时由于 Workers AI 走 prefill 续写代码把recovering: true写进_lastBodyonChatMessage会在instructions里追加RECOVERY_SUFFIX Do not think or reason — continue the text output directly.提示模型跳过推理直接续写文本。即使模型仍输出了推理内容框架也会自动把 reasoning 合并进已有的 reasoning block。OpenAIResponses API 按 ID 取回完整响应OpenAI 恢复策略是最特别的由于 OpenAI Responses API 支持服务端继续生成连接断开后生成仍在服务端继续。因此恢复时不需要重新请求只需用流式期间 stash 下来的responseId去取回完整响应实现零浪费 token。关键实现分三处流式期间 stash responseId_getChunkHandlersonChunk: ({ 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请求时开启store: true_getProviderOptionsreturn { openai: { store: true, reasoningEffort: low, reasoningSummary: auto } };恢复时_openAIResponsesRecovery从ctx.recoveryData取responseId调用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.maxAttemptsREADME 明确指出如果反复恢复超过chatRecovery.maxAttempts框架会持久化配置好的终止消息terminal message而不是让该轮对话卡死。这意味着恢复不是无界的——框架通过有界恢复bounded recovery incidents避免死循环并在达到上限时给出确定性收尾。第四种路径Inference Buffer 恢复除了 README 表格里的三种 Provider 策略server.ts 的头部注释还描述了一条可选的通用缓冲恢复路径可用于任何 ProviderInference 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 的模型实现_getReplayModel用apiKey: 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 格式差异OpenAIChat Completionsdelta.content/delta.tool_calls[]与 Responses APIresponse.output_text.delta/response.output_item.added/response.function_call_arguments.delta两种格式都会被parseOpenAIStream识别Anthropiccontent_block_deltatext_delta、thinking_delta推理与content_block_starttool_useinput_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状态呈现需要审批的工具如calculate中Math.abs(a) 1000时needsApproval返回 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发送消息并渲染消息中的text、reasoningThinking 块、工具输出、审批请求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}的 fiberDO 重启后onStart()里_checkRunFibers()会扫描cf_agents_runs表把不在内存活跃集合中的行视为被中断调用内部恢复钩子解析requestId、从cf_ai_chat_stream_metadata读取流 chunk再回调onChatRecovery。requestId编码在 fiber 名而非 snapshot 中意味着snapshot 完全是开发者自己的领域——在onChatMessage里this.stash({ responseId })不会覆盖框架数据恢复时ctx.recoveryData即用户 stash 的内容。本地 workerd 会把 SQLite 与 alarm 状态持久化到磁盘因此本地开发与线上行为一致fiber 运行中杀掉进程SIGKILL / Ctrl-C→ 重启后首个请求触发恢复。底层机制速览keepAlive 与 runFiber如果你想把forever-chat的模式复用到自己的 Agent 中需要理解两层原语详见 experimental/forever.mdkeepAlive()引用计数的 alarm 心跳。首次持引用时ctx.storage.setAlarm(now keepAliveIntervalMs)默认 30 秒可用static options { keepAliveIntervalMs: 2_000 }覆盖以便测试alarm 触发时执行日常清理并继续设置下一个。所有 disposer 调用后 refs 归零alarm 停止DO 自然空闲。它不产生 schedule 行对listSchedules()不可见。runFiber(name, fn)SQLite 中先INSERT进cf_agents_runs仅 4 列id、name、snapshot、created_at再持有 keepAlive、执行fn(ctx)、ctx.stash()同步写检查点、完成后DELETE行。被驱逐后由_checkRunFibers()恢复onFiberRecovered钩子拿到namesnapshot决定如何重跑。相比旧的 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),仅供参考
返回列表