Cloudflare Agents 服务端驱动消息:用 saveMessages 与 onChatResponse 构建自主 Agent 工作流
【免费下载链接】agentsBuild and deploy AI Agents on Cloudflare项目地址: https://gitcode.com/GitHub_Trending/agents1/agents
导读
在典型的聊天流程中,用户发送消息、Agent 做出响应;但真实业务中 Agent 往往需要自己主动行动——定时提醒触发、Webhook 到达、邮件到达、队列处理完毕,或 Agent 检查自己的回答后决定继续追问。Cloudflare Agents 项目(agents、ai-chat 等包)通过saveMessages、submitMessages、persistMessages、onChatResponse与waitUntilStable等原语,提供了完整的服务端驱动消息(server-driven messages)能力。本文将围绕 docs/agents/server-driven-messages.md 展开,先讲清几个核心原语的分工与取舍,再逐一带你实现 Cron 定时摘要、队列处理、邮件/Webhook 触发、静默注入上下文、以及"自我续答"式链式推理,并深入源码说明waitUntilStable、messageConcurrency、流式取消(AbortSignal)的底层行为。读完本文,你将能够在 Cloudflare Workers/Durable Objects 上构建无需人工参与的自主 Agent 工作流,并能在前端准确区分"用户发起"与"服务端发起"的流式状态。
核心原语总览
服务端驱动消息的四个核心原语来自 AIChatAgent 与 Think 两个包,职责各不相同:
| 原语 | 作用 |
|---|---|
saveMessages | 注入消息并触发 LLM 响应——服务端的sendMessage等价物,可await,返回时模型已响应、消息已持久化 |
submitMessages | 持久化地接纳一个 Think turn 用于异步执行,之后可随时查询状态(Think 包提供) |
persistMessages | 只存储消息、不触发响应——用于静默注入上下文 |
onChatResponse | 在任何响应完成时触发回调,包括非你发起的响应 |
isServerStreaming | 客户端侧标志:服务端发起的流处于活跃状态时为true |
其中saveMessages、persistMessages、onChatResponse、waitUntilStable均定义在 packages/ai-chat/src/index.ts 的AIChatAgent类上;submitMessages是 Think 包的持久化接纳原语,实现在 packages/think/src/think.ts,用于"先记账、后执行"的异步编排。
saveMessagesvspersistMessages:要不要触发模型
saveMessages会把消息持久化到 SQLite,并触发onChatMessage开启一轮新的 LLM 响应。它是可等待的——await返回之后,模型已经响应完毕、消息也已落库。从源码看,saveMessages的完整链路是:先解析消息(支持函数式入参),调用persistMessages(resolvedMessages)落库,再调用_runProgrammaticChatTurn执行一轮编程式 turn,最后返回{ requestId, status, error? }(见 packages/ai-chat/src/index.ts)。
// 可等待:await 返回后,LLM 已响应、消息已持久化 const result = await this.saveMessages([ { id: crypto.randomUUID(), role: "user", parts: [{ type: "text", text: "Summarize today's activity." }], createdAt: new Date() } ]); // result.status === "completed" | "error" | "skipped" | "aborted"persistMessages则只负责存储消息并广播给已连接的客户端,不会启动模型 turn(见 packages/ai-chat/src/index.ts)。它的典型场景是把系统消息、后台数据作为上下文注入会话,但不希望此刻产生回答:
async addBackgroundContext(data: string) { const stable = await this.waitUntilStable({ timeout: 30_000 }); if (!stable) return; await this.persistMessages([ ...this.messages, { id: crypto.randomUUID(), role: "user", parts: [{ type: "text", text: `[Background context]: ${data}` }], createdAt: new Date() } ]); // 消息已存储并广播给客户端,但没有发生 LLM 调用。 }一句话总结:要回答用saveMessages,只要上下文用persistMessages。
saveMessagesvssubmitMessages:同步等待还是持久接纳
submitMessages()在存储待处理工作后立即返回,消息只在 submission 真正开始执行时才追加到会话Session。它接受可序列化的UIMessage[],不接受saveMessages((messages) => ...)那种函数形式。这使其非常适合 Webhook 处理器、RPC 调用方、或对超时限制严格的父 Worker——调用方需要一个快速的持久化收据、幂等重试和后续状态查询:
const submission = await this.submitMessages( [ { id: crypto.randomUUID(), role: "user", parts: [ { type: "text", text: `Webhook event: ${JSON.stringify(payload)}` } ] } ], { idempotencyKey: payload.id } ); return Response.json({ submissionId: submission.submissionId, status: submission.status, accepted: submission.accepted });两者分工的决策标准:
- 用
saveMessages():调用方可以等待模型 turn 完整跑完。比如内部队列循环、Cron 回调。 - 用
submitMessages():调用方需要快速的持久收据 + 幂等重试 + 稍后状态检查,且不能长时间阻塞(如严格超时的 Webhook)。 - 在 Think 之外,若持久的"工作单元"是外层应用任务(一次性地接受 Webhook、恢复 provider 状态、发布可见回复、记录恢复策略),则应使用
startFiber()。submitMessages()拥有 Think 的对话接纳权,而 managed fiber 负责该 turn 周围的外部副作用。
关于 Think 的 turn API 选择(原始
chat()调用 vs agent tools),可参考 docs/think/index.md 中"Choosing a turn API"一节的专门对比。
saveMessagesvsonChatResponse:谁来触发
- 用
saveMessages的场景是你控制触发时机——调度回调、Webhook、邮件处理器,即由你决定何时注入消息。 - 用
onChatResponse的场景是你需要对非你触发的响应做出反应——用户发起的消息、工具审批后的自动续答、或框架替你运行的任何 turn。
onChatResponse是AIChatAgent上的一个可覆写钩子(默认空实现),定义见 packages/ai-chat/src/index.ts,其回调参数类型ChatResponseResult与SaveMessagesOptions定义在 packages/ai-chat/src/chat/lifecycle.ts。
waitUntilStable:服务端注入前的安全阀
从调度回调、Webhook、邮件处理器等非聊天入口读取this.messages或调用saveMessages之前,必须先调用waitUntilStable()。它会等待对话完全稳定,即同时满足:
- 没有进行中的 LLM 流
- 没有待处理的客户端工具交互(用户尚未提供的工具结果或审批)
- 没有排队中的续答 turn
waitUntilStable的实现位于 packages/ai-chat/src/index.ts,接受可选的{ timeout?, pendingInteraction? }配置。它返回true表示稳定;若在超时前仍有未决交互,则返回false;若当前没有未决项,则立即返回。
const stable = await this.waitUntilStable({ timeout: 30_000 }); if (!stable) { // 对话被用户交互阻塞,或有一个 30 秒内未完成的进行中流。 console.warn("Conversation not stable, skipping server-driven message"); return; } // 现在可以安全地读取 this.messages 并调用 saveMessages。不加这个守卫,你可能会读到过期消息,或与进行中的流重叠。注意,waitUntilStable在内部也被聊天恢复流程使用(如_chatRecoveryContinue通过recoveryConfig.stableTimeoutMs调用它,并用hasPendingClientInteraction收紧未决判定),这说明"稳定"语义是框架自身的核心依赖。
从服务端触发响应:四种典型场景
场景一:Cron 定时任务
一个每天早晨总结活动的 Digest Agent。Cron 调度默认是幂等的,所以在onStart中调用schedule()是安全的——Durable Object 多次重启不会产生重复调度。
import { AIChatAgent } from "@cloudflare/ai-chat"; export class DigestAgent extends AIChatAgent { async onChatMessage() { // ... your LLM call } async onStart() { await this.schedule("0 9 * * *", "dailyDigest"); } async dailyDigest() { const stable = await this.waitUntilStable({ timeout: 30_000 }); if (!stable) { console.warn("Conversation not stable, skipping daily digest"); return; } await this.saveMessages((messages) => [ ...messages, { id: crypto.randomUUID(), role: "user", parts: [ { type: "text", text: "Summarize what happened since your last digest." } ], createdAt: new Date() } ]); // 此时 LLM 已响应、消息已持久化。 } }这里使用了saveMessages的函数形式——saveMessages((messages) => [...])——它在执行时读取最新已持久化消息,避免多个调用排队时基于过期基线追加(比如 Webhook 快速连续到达)。Cron 语法与schedule()的更多细节见 scheduling.md。
场景二:处理队列
当你控制触发时机时,一个简单的循环就是最清晰的模式:
async processQueue() { for (const task of this.taskQueue) { const stable = await this.waitUntilStable({ timeout: 30_000 }); if (!stable) { console.warn("Conversation not stable, stopping queue processing"); break; } await this.saveMessages((messages) => [ ...messages, { id: crypto.randomUUID(), role: "user", parts: [{ type: "text", text: task }], createdAt: new Date() } ]); // LLM 已响应。this.messages 已更新。进入下一轮。 } this.taskQueue = []; }不需要任何特殊钩子——saveMessages在整个 turn 完成之后才返回。
场景三:邮件触发
把收到的邮件内容转成一条用户消息注入会话:
async onEmail(email: AgentEmail) { const stable = await this.waitUntilStable({ timeout: 30_000 }); if (!stable) { console.warn("Conversation not stable, cannot process email"); return; } const subject = email.headers.get("subject") ?? "(no subject)"; const body = await new Response(email.raw).text(); await this.saveMessages((messages) => [ ...messages, { id: crypto.randomUUID(), role: "user", parts: [ { type: "text", text: `Email from ${email.from}: ${subject}\n\n${body}` } ], createdAt: new Date() } ]); }场景四:Webhook 触发
Webhook 处理器需要先确认对话稳定,稳定则注入消息并返回ok;否则返回503 Agent is busy:
async onRequest(request: Request): Promise<Response> { const url = new URL(request.url); if (url.pathname.endsWith("/webhook") && request.method === "POST") { const stable = await this.waitUntilStable({ timeout: 30_000 }); if (!stable) { return new Response("Agent is busy", { status: 503 }); } const payload = await request.json(); try { await this.saveMessages((messages) => [ ...messages, { id: crypto.randomUUID(), role: "user", parts: [ { type: "text", text: `Webhook event: ${JSON.stringify(payload)}` } ], createdAt: new Date() } ]); return new Response("ok"); } catch (error) { console.error("Failed to process webhook:", error); return new Response("Internal error", { status: 500 }); } } return super.onRequest(request); }对非自己发起的响应做出反应:onChatResponse
onChatResponse在每一个完成的 turn 之后触发——无论它是用户发起的消息、saveMessages调用,还是自动续答。当你需要观察或反应任意来源的响应时使用它。项目在 packages/ai-chat/src/index.ts 附近展示了它的调用方式:turn 完成、锁释放后排队触发;钩子抛错会被捕获并打日志,不会影响主流程。
用法一:广播状态
import { AIChatAgent, type ChatResponseResult } from "@cloudflare/ai-chat"; export class ChatAgent extends AIChatAgent { async onChatMessage() { // ... your LLM call } protected async onChatResponse(result: ChatResponseResult) { if (result.status === "completed") { this.broadcast(JSON.stringify({ streaming: false })); } } }用法二:分析上报
protected async onChatResponse(result: ChatResponseResult) { try { await fetch("https://analytics.example.com/event", { method: "POST", body: JSON.stringify({ requestId: result.requestId, status: result.status, continuation: result.continuation }) }); } catch (error) { console.error("Analytics reporting failed:", error); } }用法三:链式推理(自我续答)
Agent 检查自己的回答并决定是否继续。这对用户发起的消息同样有效——你无法预测用户问什么,但可以对 Agent 说了什么做出反应:
protected async onChatResponse(result: ChatResponseResult) { if (result.status !== "completed") return; const lastText = result.message.parts .filter((p) => p.type === "text") .map((p) => p.text) .join(""); if (lastText.includes("[NEEDS_MORE_RESEARCH]")) { await this.saveMessages((messages) => [ ...messages, { id: crypto.randomUUID(), role: "user", parts: [{ type: "text", text: "Continue your research." }], createdAt: new Date() } ]); } }从onChatResponse内部调用saveMessages时,内层 turn 会跑完、saveMessages返回;当前onChatResponse返回后,框架会为内层响应再次触发onChatResponse,直到没有更多排队工作。框架绝不会嵌套调用onChatResponse——结果被顺序排空。源码中对应的机制是_pendingChatResponseResults数组与_insideResponseHook重入守卫(见 packages/ai-chat/src/index.ts),确保钩子在 turn 锁释放后才触发、且不会递归。
用法四:响应式队列处理
当队列项随时可能被外部事件(用户消息、Webhook)追加时,onChatResponse让你在每一次响应之后(不管谁触发的)排空队列:
protected async onChatResponse(result: ChatResponseResult) { if (result.status === "completed" && this.taskQueue.length > 0) { const next = this.taskQueue.shift()!; await this.saveMessages((messages) => [ ...messages, { id: crypto.randomUUID(), role: "user", parts: [{ type: "text", text: next }], createdAt: new Date() } ]); } }ChatResponseResult字段
onChatResponse收到的ChatResponseResult类型定义在 packages/ai-chat/src/chat/lifecycle.ts:
| 字段 | 类型 | 说明 |
|---|---|---|
message | UIMessage | 最终确定的 assistant 消息 |
requestId | string | 本 turn 的唯一 ID |
continuation | boolean | 若为自动续答则为true |
status | "completed" \| "error" \| "aborted" | turn 如何结束 |
error | string \| undefined | 当status为"error"时的错误详情 |
客户端:检测服务端发起的流
当服务端通过saveMessages触发流时,AI SDK 的status会停留在"ready",因为请求不是客户端发起的。useAgentChat钩子额外提供了两个标志位:
| 标志 | 追踪内容 |
|---|---|
status | AI SDK 生命周期:"submitted"、"streaming"、"ready"、"error"——仅针对客户端发起的请求 |
isServerStreaming | 服务端发起的流活跃时为true |
isStreaming | 客户端或服务端任一流式活跃即为true——用作通用指示器 |
UI 通用场景用isStreaming(禁用发送按钮、显示加载指示器);只有需要区分用户发起还是服务端发起时才用isServerStreaming(例如显示不同的提示文案"Agent is working in the background..."):
import { useAgent } from "agents/react"; import { useAgentChat } from "@cloudflare/ai-chat/react"; function Chat() { const agent = useAgent({ agent: "ChatAgent" }); const { messages, sendMessage, isStreaming, isServerStreaming } = useAgentChat({ agent }); return ( <div> {messages.map((m) => ( <div key={m.id}>{/* render message */}</div> ))} {isServerStreaming && <div>Agent is working in the background...</div>} {!isServerStreaming && isStreaming && <div>Agent is responding...</div>} <form onSubmit={(e) => { e.preventDefault(); const input = e.currentTarget.elements.namedItem( "input" ) as HTMLInputElement; sendMessage({ text: input.value }); input.value = ""; }} > <input name="input" placeholder="Type a message..." /> <button type="submit" disabled={isStreaming}> Send </button> </form> </div> ); }当用户在空闲状态下,服务端驱动的响应到达时,已连接的客户端会实时看到新消息出现。isStreaming标志会随流运行经历false → true → false的转换,因此发送按钮等 UI 元素会自动禁用再恢复。仓库中 server-initiated-stream.test.ts 即为服务端发起流的端到端测试,验证了该行为。
与messageConcurrency的交互
AIChatAgent上的messageConcurrency设置控制重叠的用户提交如何表现("queue"、"latest"、"merge"、"drop"、"debounce"),它只作用于sendMessage()——即客户端发起的用户消息。saveMessages()无论messageConcurrency如何设置,始终使用串行化(排队)行为:服务端驱动的消息永远不会被丢弃、合并或防抖,它们总是排队并按顺序执行。MessageConcurrency类型定义在 packages/ai-chat/src/chat/lifecycle.ts,注释明确写明:只有submit-message请求受此设置影响,重新生成、工具续答、审批、清空、编程式saveMessages与continueLastTurn均保持既有串行行为。
与其他 Agent 原语的组合
| 原语 | 组合方式 |
|---|---|
schedule() | 调度一个回调,回调中调用saveMessages——见上文 Cron 示例 |
submitMessages() | 调用方等不起saveMessages()完成时,持久接纳一个 Think turn |
queue() | 排队一个调用saveMessages的方法以延迟处理 |
runWorkflow() | 启动 Workflow;用AgentWorkflow.agentRPC 调用触发saveMessages或submitMessages的方法 |
onEmail() | 把邮件内容转换为聊天消息并调用saveMessages |
onRequest() | 处理 Webhook 并调用saveMessages |
this.broadcast() | 从onChatResponse广播自定义状态 |
取消一个服务端驱动的 turn
通过options.signal可以从外部取消编程式 turn,而无需知道内部生成的 request id:
async runLongTask(query: string, abortSignal: AbortSignal) { const result = await this.saveMessages( [{ id: crypto.randomUUID(), role: "user", parts: [{ type: "text", text: query }] }], { signal: abortSignal } ); if (result.status === "aborted") { // 信号在流中途中止。部分已流出的 chunk 仍然被持久化。 } }信号中止时的行为:
- 推理循环的信号被中止(与
chat-request-cancel走同一路径); - 中止前已流出的部分 chunk 被持久化;
saveMessages以{ status: "aborted" }解析;onChatResponse以status: "aborted"触发。
预中止的信号会在任何模型工作开始前短路。源码层面,saveMessages的SaveMessagesOptions.signal会把外部信号链接到 turn 的 abort registry 控制器(见 packages/ai-chat/src/chat/lifecycle.ts 与_abortRegistry.linkExternal调用)。
已知限制
- 信号不能跨越 Durable Object 边界。
AbortSignal不是 RPC 可序列化类型。控制器必须在调用saveMessages的 DO 内部构造。对于 Think 子 Agent 编排,应使用 Agent Tools:runAgentTool()会把父级 abort 桥接到子运行中。对于更底层的自定义 RPC,让子端返回ReadableStream,由父端取消它——workerd 会把取消传播回源的cancel回调。 - Hibernation 会丢失监听器。信号只存在于内存中。DO 重启后,持久恢复通常不带原始信号调用
continueLastTurn(),因此重启后触发的 abort 不会生效。对于流开始前的中断,恢复逻辑可以自动重试最近一条未答复的用户消息(顶层 Agent 与子 Agent 均如此)。如果取消意图必须在重启后存活,请把取消意图持久化到 agent state 或 SQL,在onChatRecovery()中检查它并返回{ continue: false }。持久恢复无法被禁用。
这正是 agent-tool 编排的集成点:父 Agent 的 AI SDK abort 信号需要传播进子 DO 的saveMessages调用。
重要注意事项速查
saveMessages是可等待的。返回后 LLM 已响应、消息已持久化。当你控制触发时机时使用它。submitMessages是持久化接纳。它在 turn 被接受后返回,而非 LLM 响应后。当超时歧义会使重试不安全时使用它。- 使用
saveMessages的函数形式。saveMessages((messages) => [...messages, newMsg])在执行时读取最新已持久化消息,避免多个调用排队时基于过期基线。 submitMessages只接受可序列化消息。它接收UIMessage[],这样被接纳的工作可以在执行前持久化存储。persistMessages不触发响应。用它静默注入上下文或系统消息。onChatResponse用于反应非你发起的 turn。适用于用户发起的消息、自动续答、或任何不是你亲自调用saveMessages的 turn。onChatResponse不嵌套。在onChatResponse内调用saveMessages时,内层 turn 完成后onChatResponse会顺序再次触发——而不是递归。- 消息在
onChatResponse触发前就已持久化。如果钩子执行期间 Durable Object 被驱逐,对话在 SQLite 中仍然安全——只有钩子回调本身丢失。 - 注入前先
waitUntilStable()。从调度回调、Webhook 或其他非聊天入口进入时,务必先调用,避免与进行中的流或未决工具交互重叠。 - 客户端会在
onChatResponse运行前看到done: true。服务端钩子不会延迟客户端。 saveMessages接受options.signal以支持外部取消。在把上游AbortSignal(例如父 Agent 的 AI SDK toolexecute)转发进子 DO 的聊天 turn 时非常有用。messageConcurrency不影响saveMessages。服务端驱动的消息总是排队并按顺序执行。
总结
服务端驱动消息把 Cloudflare Agents 从"一问一答"的被动系统升级为可自主行动的工作流引擎:saveMessages提供了可等待的"注入并回答",submitMessages提供了面向严格超时调用方的持久接纳,persistMessages支持静默上下文注入,onChatResponse让你对任意来源的响应(包括自触发续答)做出反应,而waitUntilStable与串行化的messageConcurrency交互保证了并发安全。客户端通过isStreaming/isServerStreaming可以精确呈现服务端发起的后台工作。这套原语的实现细节——从saveMessages的持久化加 turn 调度链路,到onChatResponse的非嵌套顺序排空,再到 AbortSignal 的边界限制——都可以在 packages/ai-chat/src/index.ts、packages/ai-chat/src/chat/lifecycle.ts 与 packages/think/src/think.ts 中逐一验证。
【免费下载链接】agentsBuild and deploy AI Agents on Cloudflare项目地址: https://gitcode.com/GitHub_Trending/agents1/agents
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考