从服务端发送消息并触发 LLM 响应,无需人工操作。适用于定时跟进、队列处理、邮件触发响应与自主 Agent 工作流。
典型聊天流程中,用户发送消息,Agent 响应。但 Agent 常需自主行动——定时提醒触发、webhook 到达、workflow 完成,或 Agent 在检查自身响应后决定继续。
关键原语:
| 原语 | 作用 |
|---|---|
saveMessages |
注入消息并触发 LLM——服务端等价于 sendMessage |
submitMessages |
持久接受 Think 轮次以异步执行并稍后检查 |
startFiber |
持久接受轮次周围的应用自有副作用 |
persistMessages |
存储消息但不触发响应——静默注入上下文 |
onChatResponse |
任意响应完成时作出反应,包括不是你发起的 |
isServerStreaming |
客户端标志:服务端发起的流活跃时为 true |
saveMessages 将消息持久化到 SQLite 并触发 onChatMessage 以产生新 LLM 响应。可 await——返回后 LLM 已响应且消息已持久化。
persistMessages 存储消息并广播到已连接客户端,但不触发模型轮次。用于向对话注入上下文(例如系统消息或后台数据)而不启动响应。
调用方可等待模型轮次完成时使用 saveMessages()。
调用方需要快速持久回执、幂等重试与后续状态检查时,对 Think 使用 submitMessages()。适用于 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,
});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,
});submitMessages() 先存储待处理工作,仅在提交开始执行时才将消息追加到对话 Session。它接受可序列化的 UIMessage[] 值,不支持 saveMessages((messages) => ...) 的函数形式。
在 Think 之外,当持久单元是周围的应用任务(例如接受一次 webhook、恢复提供商状态、发布可见回复并记录恢复策略)时,使用 startFiber()。submitMessages() 负责 Think 的对话准入;托管 fiber 负责该轮次周围的外部副作用。
完整 Think API 请参阅 submitMessages()。
在可控触发时使用 saveMessages — 调度回调、webhook、邮件处理程序,或任何由你决定何时注入消息的方法。
在需要响应非你触发的回复时使用 onChatResponse — 用户发起的消息、工具审批后的自动续传,或框架代你运行的任意轮次。
从调度回调、webhook、邮件处理程序或其他非聊天入口点读取 this.messages 或调用 saveMessages 前,务必调用 waitUntilStable()。
waitUntilStable() 等待对话完全稳定:
- 无进行中的 LLM stream
- 无 pending 的 client-tool 交互(用户尚未提供的 tool 结果或审批)
- 无排队的 continuation turn
稳定时返回 true;若在 pending 交互 resolve 前超时则返回 false。若无 pending 项则立即返回。
const stable = await this.waitUntilStable({ timeout: 30_000 });
if (!stable) {
// The conversation is blocked on a user interaction or an in-flight
// stream that did not complete within 30 seconds.
console.warn("Conversation not stable, skipping server-driven message");
return;
}
// Safe to read this.messages and call saveMessages.const stable = await this.waitUntilStable({ timeout: 30_000 });
if (!stable) {
// The conversation is blocked on a user interaction or an in-flight
// stream that did not complete within 30 seconds.
console.warn("Conversation not stable, skipping server-driven message");
return;
}
// Safe to read this.messages and call saveMessages.若无此 guard,可能读到 stale message 或与进行中的 stream 重叠。
每日摘要 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(),
},
]);
// At this point the LLM has responded and the message is persisted.
}
}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(),
},
]);
// At this point the LLM has responded and the message is persisted.
}
}saveMessages 的函数形式 — saveMessages((messages) => [...]) — 在执行时读取最新持久化的 message。这可避免多次调用排队时出现 stale baseline(例如 webhook 快速到达)。更多 schedule() 与 cron 语法请参阅调度任务。
当你控制触发时,简单循环是最清晰的模式:
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 has responded. this.messages is updated. Next iteration.
}
this.taskQueue = [];
}无需特殊 hook — 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(),
},
]);
}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);
}若 webhook 提供商期望快速响应,请改用 submitMessages()。这向提供商提供持久确认,并允许使用相同幂等键安全重试:
async onRequest(request: Request): Promise<Response> {
if (request.method !== "POST") return super.onRequest(request);
const payload = await request.json<{ id: string }>();
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,
accepted: submission.accepted,
status: submission.status,
});
}使用 persistMessages 添加 LLM 在下次 turn 会看到的 message,但不立即启动 turn:
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(),
},
]);
// Message is stored and broadcast to clients, but no LLM call happens.
}onChatResponse 在每次完成的 turn 后触发 — 用户发起的 message、saveMessages 调用与自动 continuation。无论触发方式如何,需要观察或响应回复时使用它。
import { AIChatAgent } from "@cloudflare/ai-chat";
export class ChatAgent extends AIChatAgent {
async onChatMessage() {
// ... your LLM call
}
async onChatResponse(result) {
if (result.status === "completed") {
this.broadcast(JSON.stringify({ streaming: false }));
}
}
}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 可检查自身回复并决定是否继续。对用户发起的 message 同样适用 — 无法预测用户会问什么,但可对其回复做出反应。
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 调用 — 结果按顺序 drain。
当队列项可由外部事件(用户 message、webhook)随时添加时,onChatResponse 可在每次响应后 drain 队列,无论由谁触发:
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(),
},
]);
}
}| 字段 | 类型 | 描述 |
|---|---|---|
message |
UIMessage |
本次 turn 的最终 assistant message |
requestId |
string |
本次 turn 的唯一 ID |
continuation |
boolean |
若为自动 continuation 则为 true |
status |
"completed" | "error" | "aborted" |
turn 如何结束 |
error |
string | undefined |
status 为 "error" 时的错误详情 |
服务端通过 saveMessages 触发 stream 时,AI SDK 的 status 保持 "ready",因为 client 未发起请求。useAgentChat hook 提供两个额外标志处理此情况:
| 标志 | 跟踪内容 |
|---|---|
status |
AI SDK 生命周期:"submitted"、"streaming"、"ready"、"error" — 仅针对 client 发起的请求 |
isServerStreaming |
服务端发起的 stream 活跃时为 true |
isStreaming |
client 或 server stream 任一活跃时为 true — 用作通用指示 |
大多数 UI 场景使用 isStreaming(禁用发送按钮、显示加载指示)。仅在需区分用户发起与服务端发起 stream 时使用 isServerStreaming(例如显示「Agent 正在后台工作…」等不同指示)。
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>
);
}用户 idle 时服务端驱动响应到达,已连接 client 会实时看到新 message。isStreaming 标志随 stream 运行从 false → true → false 变化,发送按钮等 UI 元素会自动禁用与重新启用。
AIChatAgent 上的 messageConcurrency 设置控制重叠用户提交的行为("queue"、"latest"、"merge"、"drop"、"debounce")。此设置仅适用于 sendMessage() — 客户端发起的用户消息。
无论 messageConcurrency 如何,saveMessages() 始终使用串行(排队)行为。即服务端驱动的消息不会被丢弃、合并或防抖 — 始终排队并按顺序执行。
| 原语 | 组合方式 |
|---|---|
schedule() |
调度调用 saveMessages 的回调 — 见上文 cron 示例 |
queue() |
将调用 saveMessages 的方法入队以延迟处理 |
startFiber() |
持久接受并检查消息轮次周围的应用自有工作 |
runWorkflow() |
启动 Workflow;用 AgentWorkflow.agent RPC 调用触发 saveMessages 或 submitMessages 的方法 |
onEmail() |
将邮件内容转为聊天消息并调用 saveMessages |
onRequest() |
处理 webhook 并调用 saveMessages 或 submitMessages |
this.broadcast() |
从 onChatResponse 广播自定义状态 |
同一 Durable Object 启动并控制轮次时,传入 AbortSignal:
const controller = new AbortController();
const result = await this.saveMessages(
[
{
id: crypto.randomUUID(),
role: "user",
parts: [{ type: "text", text: "Run the long analysis." }],
},
],
{ signal: controller.signal },
);
if (result.status === "aborted") {
// Partial chunks already streamed are persisted.
}const controller = new AbortController();
const result = await this.saveMessages(
[
{
id: crypto.randomUUID(),
role: "user",
parts: [{ type: "text", text: "Run the long analysis." }],
},
],
{ signal: controller.signal },
);
if (result.status === "aborted") {
// Partial chunks already streamed are persisted.
}continueLastTurn() 接受相同的 options.signal 参数。AbortSignal 对象无法跨 Durable Object RPC 边界,且 signal 仅在内存中。若 Durable Object 在轮次中途休眠且启用了聊天恢复,恢复的轮次通常在没有原始 signal 的情况下继续;对于流式传输前中断,恢复可自动重试最新未回答的用户消息。重启后触发的中止对恢复的轮次无效。
当工作通过 submitMessages() 接受,或取消必须跨 Worker 与 Durable Object RPC 边界时,使用 cancelSubmission(submissionId) 进行持久取消。
当持久单元通过 startFiber() 接受,且取消应作用于周围的应用任务而非 Think 轮次时,使用 cancelFiber(fiberId)。
saveMessages可等待。 返回后 LLM 已响应且消息已持久化。在可控触发时使用。- 使用
saveMessages的函数形式。saveMessages((messages) => [...messages, newMsg])在执行时读取最新持久化消息,避免多次调用排队时的过期基线。 persistMessages不触发响应。 用于静默注入上下文或系统消息。onChatResponse用于响应非你发起的轮次。 用于用户发起的消息、自动续传,或任何未自行调用saveMessages的轮次。onChatResponse不嵌套。 在onChatResponse内调用saveMessages时,内部轮次完成后再次顺序触发onChatResponse— 非递归。- 消息在
onChatResponse触发前已持久化。 若 Durable Object 在钩子期间被驱逐,对话在 SQLite 中安全 — 仅钩子回调丢失。 - 注入前调用
waitUntilStable()。 从调度回调、webhook 或其他非聊天入口点调用,避免与进行中的流或待处理工具交互重叠。 - 客户端在
onChatResponse运行前看到完成的响应。 服务端钩子不会延迟客户端。 messageConcurrency不影响saveMessages。 服务端驱动的消息始终排队并按顺序执行。