使用 AIChatAgent 与 useAgentChat 构建 AI 驱动的聊天界面。消息自动持久化到 SQLite,断开时流可恢复,tool 调用在服务端与客户端均可工作。
@cloudflare/ai-chat 包提供两个主要 API:
| 导出 | 导入 | 用途 |
|---|---|---|
AIChatAgent |
@cloudflare/ai-chat |
带消息持久化与流式传输的服务端 Agent 类 |
useAgentChat |
@cloudflare/ai-chat/react |
构建聊天 UI 的 React hook |
高级辅助函数还可从 @cloudflare/ai-chat/react、@cloudflare/ai-chat/types 与 agents/chat 获取;完整导出列表见 导出。
基于 AI SDK ↗ 与 Cloudflare Durable Objects,你将获得:
- 自动消息持久化 — 对话存储在 SQLite,重启后保留
- 可恢复流式传输 — 断开客户端从中断处恢复,无数据丢失
- 实时同步 — 消息通过 WebSocket 广播到所有已连接客户端
- 工具支持 — 服务端、客户端与人机协同(human-in-the-loop)工具模式
- 数据部分(Data part) — 在文本旁向消息附加带类型的 JSON(引用、进度、用量)
- 行大小保护 — 消息接近 SQLite 限制时自动压缩(compaction)
npm install @cloudflare/ai-chat agents ai workers-ai-providerimport { AIChatAgent } from "@cloudflare/ai-chat";
import { createWorkersAI } from "workers-ai-provider";
import { streamText, convertToModelMessages } from "ai";
export class ChatAgent extends AIChatAgent {
async onChatMessage() {
// Use any provider such as workers-ai-provider, openai, anthropic, google, etc.
const workersai = createWorkersAI({ binding: this.env.AI });
const result = streamText({
model: workersai("@cf/zai-org/glm-4.7-flash"),
messages: await convertToModelMessages(this.messages),
});
return result.toUIMessageStreamResponse();
}
}import { AIChatAgent } from "@cloudflare/ai-chat";
import { createWorkersAI } from "workers-ai-provider";
import { streamText, convertToModelMessages } from "ai";
export class ChatAgent extends AIChatAgent {
async onChatMessage() {
// Use any provider such as workers-ai-provider, openai, anthropic, google, etc.
const workersai = createWorkersAI({ binding: this.env.AI });
const result = streamText({
model: workersai("@cf/zai-org/glm-4.7-flash"),
messages: await convertToModelMessages(this.messages),
});
return result.toUIMessageStreamResponse();
}
}import { useAgent } from "agents/react";
import { useAgentChat } from "@cloudflare/ai-chat/react";
function Chat() {
const agent = useAgent({ agent: "ChatAgent" });
const { messages, sendMessage, status } = useAgentChat({ agent });
return (
<div>
{messages.map((msg) => (
<div key={msg.id}>
<strong>{msg.role}:</strong>
{msg.parts.map((part, i) =>
part.type === "text" ? <span key={i}>{part.text}</span> : null,
)}
</div>
))}
<form
onSubmit={(e) => {
e.preventDefault();
const input = e.currentTarget.elements.namedItem("input");
sendMessage({ text: input.value });
input.value = "";
}}
>
<input name="input" placeholder="Type a message..." />
<button type="submit" disabled={status !== "ready"}>
Send
</button>
</form>
</div>
);
}import { useAgent } from "agents/react";
import { useAgentChat } from "@cloudflare/ai-chat/react";
function Chat() {
const agent = useAgent({ agent: "ChatAgent" });
const { messages, sendMessage, status } = useAgentChat({ agent });
return (
<div>
{messages.map((msg) => (
<div key={msg.id}>
<strong>{msg.role}:</strong>
{msg.parts.map((part, i) =>
part.type === "text" ? <span key={i}>{part.text}</span> : null,
)}
</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={status !== "ready"}>
Send
</button>
</form>
</div>
);
}// wrangler.jsonc
{
"ai": { "binding": "AI" },
"durable_objects": {
"bindings": [{ "name": "ChatAgent", "class_name": "ChatAgent" }],
},
"migrations": [{ "tag": "v1", "new_sqlite_classes": ["ChatAgent"] }],
}需要 new_sqlite_classes 迁移 — AIChatAgent 使用 SQLite 做消息持久化与流分块缓冲。
sequenceDiagram
participant Client as Client (useAgentChat)
participant Agent as AIChatAgent
participant DB as SQLite
Client->>Agent: CF_AGENT_USE_CHAT_REQUEST (WebSocket)
Agent->>DB: Persist messages
Agent->>Agent: onChatMessage()
loop Streaming response
Agent-->>Client: CF_AGENT_USE_CHAT_RESPONSE (chunks)
Agent->>DB: Buffer chunks
end
Agent->>DB: Persist final message
Agent-->>Client: CF_AGENT_CHAT_MESSAGES (broadcast to all clients)
- 客户端通过 WebSocket 发送消息
AIChatAgent将消息持久化到 SQLite 并调用你的onChatMessage方法- 你的方法返回流式
Response(通常来自streamText) - Chunks 通过 WebSocket 实时回传
- 流完成后,最终消息会被持久化并广播到所有连接
扩展 agents 包中的 Agent。管理对话状态、持久化与流式传输。
import { AIChatAgent } from "@cloudflare/ai-chat";
export class ChatAgent extends AIChatAgent {
// Access current messages
// this.messages: UIMessage[]
// Limit stored messages (optional)
maxPersistedMessages = 200;
async onChatMessage(onFinish, options) {
// onFinish: callback for streamText (cleanup is automatic)
// options.abortSignal: cancel signal
// options.body: custom data from client
// options.continuation: true for continuation turns
// Return a Response (streaming or plain text)
}
}import { AIChatAgent } from "@cloudflare/ai-chat";
export class ChatAgent extends AIChatAgent {
// Access current messages
// this.messages: UIMessage[]
// Limit stored messages (optional)
maxPersistedMessages = 200;
async onChatMessage(onFinish, options?) {
// onFinish: callback for streamText (cleanup is automatic)
// options.abortSignal: cancel signal
// options.body: custom data from client
// options.continuation: true for continuation turns
// Return a Response (streaming or plain text)
}
}这是你需要重写的主要方法。它接收对话上下文并应返回 Response。
流式响应(最常见):
export class ChatAgent extends AIChatAgent {
async onChatMessage() {
const workersai = createWorkersAI({ binding: this.env.AI });
const result = streamText({
model: workersai("@cf/zai-org/glm-4.7-flash"),
system: "You are a helpful assistant.",
messages: await convertToModelMessages(this.messages),
});
return result.toUIMessageStreamResponse();
}
}export class ChatAgent extends AIChatAgent {
async onChatMessage() {
const workersai = createWorkersAI({ binding: this.env.AI });
const result = streamText({
model: workersai("@cf/zai-org/glm-4.7-flash"),
system: "You are a helpful assistant.",
messages: await convertToModelMessages(this.messages),
});
return result.toUIMessageStreamResponse();
}
}纯文本响应:
export class ChatAgent extends AIChatAgent {
async onChatMessage() {
return new Response("Hello! I am a simple agent.", {
headers: { "Content-Type": "text/plain" },
});
}
}访问自定义 body 数据与 request ID:
export class ChatAgent extends AIChatAgent {
async onChatMessage(_onFinish, options) {
const { timezone, userId } = options?.body ?? {};
// Use these values in your LLM call or business logic
// options.requestId — unique identifier for this chat request,
// useful for logging and correlating events
console.log("Request ID:", options?.requestId);
if (options?.continuation) {
// This turn continues a previous assistant message after a tool result,
// continueLastTurn(), or recovery.
}
}
}options.continuation 在工具结果或审批后的自动续传、调用 continueLastTurn() 或恢复轮次时为 true。可用它选择不同模型、调整系统提示词,或在续传轮次跳过昂贵的上下文组装。
从 SQLite 加载的当前对话历史。这是 AI SDK 的 UIMessage 对象数组。每次交互后消息会自动持久化。
限制 SQLite 中存储的消息数量。超出限制时,最旧的消息会被删除。这仅控制存储,不影响发送给 LLM 的内容。
export class ChatAgent extends AIChatAgent {
maxPersistedMessages = 200;
}export class ChatAgent extends AIChatAgent {
maxPersistedMessages = 200;
}要控制发送给模型的内容,请使用 AI SDK 的 pruneMessages():
import { streamText, convertToModelMessages, pruneMessages } from "ai";
export class ChatAgent extends AIChatAgent {
async onChatMessage() {
const workersai = createWorkersAI({ binding: this.env.AI });
const result = streamText({
model: workersai("@cf/zai-org/glm-4.7-flash"),
messages: pruneMessages({
messages: await convertToModelMessages(this.messages),
reasoning: "before-last-message",
toolCalls: "before-last-2-messages",
}),
});
return result.toUIMessageStreamResponse();
}
}import { streamText, convertToModelMessages, pruneMessages } from "ai";
export class ChatAgent extends AIChatAgent {
async onChatMessage() {
const workersai = createWorkersAI({ binding: this.env.AI });
const result = streamText({
model: workersai("@cf/zai-org/glm-4.7-flash"),
messages: pruneMessages({
messages: await convertToModelMessages(this.messages),
reasoning: "before-last-message",
toolCalls: "before-last-2-messages",
}),
});
return result.toUIMessageStreamResponse();
}
}控制 AIChatAgent 是否在调用 onChatMessage 前等待 MCP 服务器连接就绪。这确保 this.mcp.getAITools() 返回完整工具集,尤其在 Durable Object 休眠后连接在后台恢复时。
| 值 | 行为 |
|---|---|
{ timeout: 10_000 } |
最多等待 10 秒(默认) |
{ timeout: N } |
最多等待 N 毫秒 |
true |
无限等待直至所有连接就绪 |
false |
不等待(0.2.0 之前的旧行为) |
export class ChatAgent extends AIChatAgent {
// Default — waits up to 10 seconds
// waitForMcpConnections = { timeout: 10_000 };
// Wait forever
waitForMcpConnections = true;
// Disable waiting
waitForMcpConnections = false;
}export class ChatAgent extends AIChatAgent {
// Default — waits up to 10 seconds
// waitForMcpConnections = { timeout: 10_000 };
// Wait forever
waitForMcpConnections = true;
// Disable waiting
waitForMcpConnections = false;
}更低级控制时,可在 onChatMessage 内直接调用 this.mcp.waitForConnections()。
控制聊天轮次已活跃或排队时,重叠用户提交的行为。
export class ChatAgent extends AIChatAgent {
messageConcurrency = "queue";
}export class ChatAgent extends AIChatAgent {
messageConcurrency = "queue";
}| 策略 | 行为 |
|---|---|
"queue"(默认) |
排队每个提交并按顺序处理 |
"latest" |
仅保留最新重叠提交;被取代的提交仍会持久化用户消息,但不启动模型轮次 |
"merge" |
排队重叠提交,然后在最新排队轮次运行前,将其末尾的用户消息合并为一个组合轮次 |
"drop" |
完全忽略重叠提交。消息不会持久化。 |
{ strategy: "debounce", debounceMs?: number } |
带静默窗口的 trailing-edge latest(默认 750ms) |
此设置仅适用于 sendMessage() 提交。重新生成、工具续传、审批、清空与编程式 saveMessages() 调用仍保持现有串行行为。
persistMessages 将消息存储在 SQLite 并广播给所有已连接客户端,但不会触发模型轮次。用于向对话注入消息而不启动新响应。
saveMessages 持久化消息 并触发 onChatMessage() 产生新响应。它会等待任何活跃聊天轮次完成后再启动,因此定时或编程式消息不会与进行中的流重叠。
// Store messages without triggering a response
await this.persistMessages(messages);
// Store messages AND trigger onChatMessage
const { requestId, status } = await this.saveMessages(messages);// Store messages without triggering a response
await this.persistMessages(messages);
// Store messages AND trigger onChatMessage
const { requestId, status } = await this.saveMessages(messages);saveMessages 接受消息数组,或从最新持久化的 this.messages 推导下一消息列表的函数。多次调用排队时使用函数形式,以免基线过期:
await this.saveMessages((messages) => [
...messages,
{
id: crypto.randomUUID(),
role: "user",
parts: [{ type: "text", text: "Summarize the latest data" }],
createdAt: new Date(),
},
]);await this.saveMessages((messages) => [
...messages,
{
id: crypto.randomUUID(),
role: "user",
parts: [{ type: "text", text: "Summarize the latest data" }],
createdAt: new Date(),
},
]);saveMessages 返回 { requestId, status, error? },其中 status 为 "completed"(轮次已运行)、"error"(流报错)、"skipped"(开始前聊天被清空)或 "aborted"(外部 AbortSignal 在完成前取消)。status 为 "error" 时,error 包含流错误消息(如有)。
传入 options.signal 可从聊天 agent 外部取消编程式轮次。父工具调用需取消子 agent 轮次且不知内部生成的 request ID 时很有用:
const controller = new AbortController();
const result = await this.saveMessages(
(messages) => [...messages, syntheticUserMessage],
{ signal: controller.signal },
);
if (result.status === "aborted") {
// Partial chunks already streamed are persisted.
}const controller = new AbortController();
const result = await this.saveMessages(
(messages) => [...messages, syntheticUserMessage],
{ signal: controller.signal },
);
if (result.status === "aborted") {
// Partial chunks already streamed are persisted.
}continueLastTurn() 接受相同的 options.signal 参数。AbortSignal 无法跨 Durable Object RPC 边界,因此在调用 saveMessages() 或 continueLastTurn() 的 Durable Object 内构造 controller。信号仅在内存中;若 Durable Object 在轮次中途休眠且启用了 chatRecovery,恢复的轮次会在没有原始信号的情况下运行。
聊天轮次产生并持久化 assistant 消息后调用。此 hook 运行前轮次锁已释放,因此在内部调用 saveMessages 是安全的。对会持久化 assistant 消息的轮次路径触发:WebSocket 聊天请求、saveMessages 与自动续传。若轮次在产生任何 assistant 部分前失败,错误会通过原始请求抛出。
export class ChatAgent extends AIChatAgent {
async onChatResponse(result) {
if (result.status === "completed") {
console.log("Turn completed:", result.requestId);
}
if (result.status === "error") {
console.error("Turn failed:", result.error);
}
}
}import type { ChatResponseResult } from "@cloudflare/ai-chat";
export class ChatAgent extends AIChatAgent {
protected async onChatResponse(result: ChatResponseResult) {
if (result.status === "completed") {
console.log("Turn completed:", result.requestId);
}
if (result.status === "error") {
console.error("Turn failed:", result.error);
}
}
}ChatResponseResult 包含:
| 字段 | 类型 | 描述 |
|---|---|---|
message |
UIMessage |
本次轮次的最终 assistant 消息 |
requestId |
string |
与本次轮次关联的 request ID |
continuation |
boolean |
本次轮次是否为之前 assistant 轮次的续传 |
status |
"completed" | "error" | "aborted" |
轮次如何结束 |
error |
string | undefined |
status 为 "error" 时的错误消息 |
重写此方法可在消息持久化到存储前应用自定义转换。此 hook 在内置清理(剥离 OpenAI metadata、截断 Anthropic 提供商执行的工具载荷、过滤空 reasoning 部分)之后运行。
export class ChatAgent extends AIChatAgent {
sanitizeMessageForPersistence(message) {
return {
...message,
parts: message.parts.map((part) => {
if (
"output" in part &&
typeof part.output === "string" &&
part.output.length > 1000
) {
return { ...part, output: "[redacted]" };
}
return part;
}),
};
}
}export class ChatAgent extends AIChatAgent {
protected sanitizeMessageForPersistence(message: UIMessage): UIMessage {
return {
...message,
parts: message.parts.map((part) => {
if (
"output" in part &&
typeof part.output === "string" &&
part.output.length > 1000
) {
return { ...part, output: "[redacted]" };
}
return part;
}),
};
}
}这些方法帮助协调编程式轮次并等待待处理交互。
assistant 消息等待客户端工具结果或审批时返回 true。
if (this.hasPendingInteraction()) {
console.log("Waiting for user to approve or provide tool output");
}if (this.hasPendingInteraction()) {
console.log("Waiting for user to approve or provide tool output");
}等待对话完全稳定——无活跃流、无待处理的客户端工具交互、无排队的续传轮次。稳定时返回 true;待处理交互在超时前仍未完成则返回 false。
const stable = await this.waitUntilStable({ timeout: 30_000 });
if (stable) {
console.log("All turns complete, safe to proceed");
}const stable = await this.waitUntilStable({ timeout: 30_000 });
if (stable) {
console.log("All turns complete, safe to proceed");
}对服务端驱动流程与 saveMessages 配合尤其有用:
await this.saveMessages((messages) => [...messages, syntheticUserMessage]);
await this.waitUntilStable({ timeout: 60_000 });
// The assistant has finished respondingawait this.saveMessages((messages) => [...messages, syntheticUserMessage]);
await this.waitUntilStable({ timeout: 60_000 });
// The assistant has finished responding中止活跃轮次并使排队的续传失效。内置 CF_AGENT_CHAT_CLEAR 处理程序会自动调用,必要时也可手动调用。
重写 onConnect 与 onClose 以添加自定义逻辑。流恢复与消息同步由框架处理:
export class ChatAgent extends AIChatAgent {
async onConnect(connection, ctx) {
// Your custom logic (e.g., logging, auth checks)
console.log("Client connected:", connection.id);
// Stream resumption and message sync are handled automatically
}
async onClose(connection, code, reason, wasClean) {
console.log("Client disconnected:", connection.id);
// Connection cleanup is handled automatically
}
}export class ChatAgent extends AIChatAgent {
async onConnect(connection, ctx) {
// Your custom logic (e.g., logging, auth checks)
console.log("Client connected:", connection.id);
// Stream resumption and message sync are handled automatically
}
async onClose(connection, code, reason, wasClean) {
console.log("Client disconnected:", connection.id);
// Connection cleanup is handled automatically
}
}destroy() 方法会取消所有待处理的聊天请求并清理流状态。Durable Object 被驱逐时会自动调用,必要时也可手动调用。
当用户在聊天 UI 中点击「stop」时,客户端发送 CF_AGENT_CHAT_REQUEST_CANCEL 消息。服务端将其传播到 options 中的 abortSignal:
export class ChatAgent extends AIChatAgent {
async onChatMessage(_onFinish, options) {
const result = streamText({
model: workersai("@cf/zai-org/glm-4.7-flash"),
messages: await convertToModelMessages(this.messages),
abortSignal: options?.abortSignal, // Pass through for cancellation
});
return result.toUIMessageStreamResponse();
}
}export class ChatAgent extends AIChatAgent {
async onChatMessage(_onFinish, options) {
const result = streamText({
model: workersai("@cf/zai-org/glm-4.7-flash"),
messages: await convertToModelMessages(this.messages),
abortSignal: options?.abortSignal, // Pass through for cancellation
});
return result.toUIMessageStreamResponse();
}
}子类也可在 Durable Object 内取消轮次:
protected abortRequest(requestId: string, reason?: unknown): void
protected abortAllRequests(): void已知 request ID 时使用 abortRequest()。要取消当前任意轮次,使用单用途辅助方法 abortAllRequests()。编程式轮次若可在调用点传入 signal,优先使用 SaveMessagesOptions.signal。
自动流恢复(useAgentChat 上的 resume 选项)是客户端重连恢复——客户端断开并重连时恢复活动流。它不涵盖 Durable Object 驱逐:若模型调用进行中 Worker 进程或 Durable Object 被驱逐,流本身会丢失。chatRecovery 处理该情况。
Durable Object 在流中途被驱逐(代码更新、不活动超时、资源限制)时,LLM 连接永久断开,内存中的流式状态丢失。chatRecovery 将每个聊天轮次包裹在 runFiber() 中,在流式传输期间提供自动 keepAlive,并在重启时提供恢复 hook。
export class ChatAgent extends AIChatAgent {
chatRecovery = true;
}export class ChatAgent extends AIChatAgent {
override chatRecovery = true;
}AIChatAgent 默认 chatRecovery 为 false,因此现有聊天 agent 除非主动开启,否则仅获得客户端重连与可恢复流行为。Think 默认为 true。
启用后,每次 onChatMessage 调用在 fiber 内运行。若 agent 在流中途被驱逐,fiber 行会保留在 SQLite 中。下次激活时,框架检测被中断的 fiber,从缓冲的流分块重建部分响应,并调用 onChatRecovery。
也可将 chatRecovery 设为配置对象,以限制恢复,并在恢复无法成功时自定义终态体验:
export class ChatAgent extends AIChatAgent {
chatRecovery = {
maxAttempts: 10,
stableTimeoutMs: 10_000,
terminalMessage: "The assistant was interrupted and could not recover.",
// Primary stuck-turn bound. Resets on every progress-bearing attempt, so a
// turn that keeps producing content survives unbounded interruption.
noProgressTimeoutMs: 5 * 60 * 1000,
// Runaway-loop guard. Defaults to Infinity (no cap). Set a finite value to
// seal a turn that keeps emitting content but never converges.
maxRecoveryWork: 200,
// Caller policy consulted from the second recovery attempt onward. Return
// false to stop recovery. This is where you enforce a token/cost budget.
// Note: this is called as `config.shouldKeepRecovering(ctx)`, so it is not
// bound to the agent instance — track spend in your own store keyed by the
// incident.
async shouldKeepRecovering(ctx) {
return (await getSpendForTurn(ctx.recoveryRootRequestId)) < MAX_SPEND;
},
async onExhausted(ctx) {
console.warn("Chat recovery exhausted", ctx.incidentId, ctx.reason);
},
};
}export class ChatAgent extends AIChatAgent {
override chatRecovery = {
maxAttempts: 10,
stableTimeoutMs: 10_000,
terminalMessage: "The assistant was interrupted and could not recover.",
// Primary stuck-turn bound. Resets on every progress-bearing attempt, so a
// turn that keeps producing content survives unbounded interruption.
noProgressTimeoutMs: 5 * 60 * 1000,
// Runaway-loop guard. Defaults to Infinity (no cap). Set a finite value to
// seal a turn that keeps emitting content but never converges.
maxRecoveryWork: 200,
// Caller policy consulted from the second recovery attempt onward. Return
// false to stop recovery. This is where you enforce a token/cost budget.
// Note: this is called as `config.shouldKeepRecovering(ctx)`, so it is not
// bound to the agent instance — track spend in your own store keyed by the
// incident.
async shouldKeepRecovering(ctx) {
return (await getSpendForTurn(ctx.recoveryRootRequestId)) < MAX_SPEND;
},
async onExhausted(ctx) {
console.warn("Chat recovery exhausted", ctx.incidentId, ctx.reason);
},
};
}chatRecovery 对象接受以下配置选项:
| 字段 | 默认值 | 描述 |
|---|---|---|
maxAttempts |
10 |
进入终态耗尽前的尝试上限。有前进进度时重置,因此捕获紧密的无进度告警循环,而不是健康的长轮次。 |
stableTimeoutMs |
10_000 |
恢复尝试等待 isolate 达到稳定状态的时长,超时后重新调度。 |
terminalMessage |
通用消息 | 放弃恢复时向用户显示的消息。 |
noProgressTimeoutMs |
300_000(5 分钟) |
主要的卡住轮次边界:事件可无前进进度的最长时间,超时后封存(no_progress_timeout)。每次有进度的尝试都会重置,因此持续产生内容的轮次可在中断后无限继续。 |
maxRecoveryWork |
Infinity |
失控循环防护。事件开始后产生的内容/工具单元上限,仍在推进的轮次也会被封存。默认无上限。 |
shouldKeepRecovering |
— | 从第二次恢复尝试起咨询的调用方策略。返回 false 停止恢复。用于强制执行 token 或成本预算。ctx.work 是粗粒度分段计数,不是 token,因此需自行跟踪真实消耗。 |
onExhausted |
— | 放弃恢复前调用一次,在终态消息送达之前。检查 ctx.reason 了解原因。 |
ChatRecoveryProgressContext(传给 shouldKeepRecovering 的 ctx)包含以下字段:
| 字段 | 类型 | 描述 |
|---|---|---|
incidentId |
string |
此恢复事件的稳定 ID。 |
requestId |
string |
当前续传的 request ID(每条链式续传会变化)。 |
recoveryRootRequestId |
string |
整个续传链的稳定 ID — 按事件跟踪预算的正确键。 |
attempt |
number |
此事件的尝试次数(此 hook 运行时 ≥ 2)。 |
maxAttempts |
number |
配置的尝试上限。 |
recoveryKind |
"retry" | "continue" |
恢复是重试未回答的用户轮次,还是继续部分 assistant 轮次。 |
work |
number |
事件打开后产生的内容/工具分段的粗粒度单调计数(不是 token)。 |
ageMs |
number |
自事件首次中断起的墙上时钟毫秒数。 |
进行中的轮次不会被框架自行终止——只要持续有前进进度,可在无限中断后继续(例如密集部署窗口)。恢复仅因以下 ctx.reason 之一而封存:
no_progress_timeout— 无进度窗口内没有前进进度(卡住的轮次)。max_attempts_exceeded— 尝试上限消耗在紧密的无进度告警循环上。work_budget_exceeded— 轮次持续产生内容但超过maxRecoveryWork(失控循环)。recovery_aborted— 你的shouldKeepRecoveringhook 返回false。stable_timeout— 恢复尝试等待稳定状态持续超时直至预算耗尽(极端抖动)。
轮次可能暂停在无法自行完成的客户端交互上:客户端工具调用(无服务端 execute、由客户端回放结果的工具),或 approval-requested 部分。此类轮次在等待人类,而不是卡住。
交互待处理期间,轮次豁免所有恢复预算。无进度窗口、尝试上限、maxRecoveryWork 与 shouldKeepRecovering 均暂停。用户因部署中断后需数分钟回答提示也不会触发封存。恢复会停放该轮次而不是失败,用户最终审批或工具结果通过正常续传路径恢复。
此豁免仅适用于客户端。execute() 中途被杀的服务端工具是真正的孤立项,不豁免,通过对话记录修复恢复。
通过可观测性监控终态耗尽:
import { subscribe } from "agents/observability";
const unsubscribe = subscribe("chat", (event) => {
if (event.type === "chat:recovery:exhausted") {
console.error("Chat recovery exhausted", event.payload);
}
});import { subscribe } from "agents/observability";
const unsubscribe = subscribe("chat", (event) => {
if (event.type === "chat:recovery:exhausted") {
console.error("Chat recovery exhausted", event.payload);
}
});重写以实现提供商特定的恢复。默认行为持久化部分响应并通过 continueLastTurn() 调度续传。
export class ChatAgent extends AIChatAgent {
chatRecovery = true;
async onChatRecovery(ctx) {
console.log(`Recovered ${ctx.partialText.length} chars of partial text`);
// Default: persist partial + schedule continuation
return {};
}
}import type {
ChatRecoveryContext,
ChatRecoveryOptions,
} from "@cloudflare/ai-chat";
export class ChatAgent extends AIChatAgent {
override chatRecovery = true;
override async onChatRecovery(
ctx: ChatRecoveryContext,
): Promise<ChatRecoveryOptions> {
console.log(`Recovered ${ctx.partialText.length} chars of partial text`);
// Default: persist partial + schedule continuation
return {};
}
}ChatRecoveryContext:
| 字段 | 类型 | 描述 |
|---|---|---|
incidentId |
string |
此恢复事件的稳定 ID |
attempt |
number |
此事件的当前尝试次数,从 1 开始 |
maxAttempts |
number |
进入终态耗尽前的配置尝试上限 |
recoveryKind |
"retry" | "continue" |
恢复是重试未回答的用户轮次,还是继续部分 assistant 轮次 |
streamId |
string |
被中断流的 ID |
requestId |
string |
原始聊天请求的 ID |
partialText |
string |
驱逐前生成的文本 |
partialParts |
MessagePart[] |
驱逐前生成的消息部分(text、reasoning、tool call) |
recoveryData |
unknown | null |
来自 this.stash() 的数据 — 完全由用户控制 |
messages |
ChatMessage[] |
完整对话历史 |
lastBody |
Record<string, unknown> | undefined |
原始请求正文 |
lastClientTools |
ClientToolSchema[] | undefined |
原始请求的客户端工具 schema |
createdAt |
number |
被中断轮次开始时的 epoch 毫秒 |
ChatRecoveryOptions:
| 字段 | 默认值 | 描述 |
|---|---|---|
persist |
true |
将部分响应保存为 assistant 消息 |
continue |
true |
通过 continueLastTurn() 调度续传 |
常见返回值:
{}— 持久化部分结果并自动续传(默认,适用于支持 assistant prefill 的提供商){ continue: false }— 持久化部分结果但不自动续传(自行处理续传){ persist: false, continue: false }— 不持久化未结算的剩余部分,自行处理一切(例如从提供商检索已完成响应)
已结算的工作永不丢弃:persist: false 仅抑制对没有可丢失已结算内容的部分结果的持久化。已携带已结算工具结果(已完成、常为非幂等工作)的部分结果无论如何都会持久化,应用不会意外丢弃已完成的工具调用——也无需为安全而使用 { persist: true }。
若在写入任何流分块前发生恢复,则没有可续传的部分 assistant 消息。若最新持久化消息仍是中断轮次的未回答用户消息,框架会自动重试该轮次,除非 continue 为 false。
使用 ctx.createdAt 跳过过期恢复:
override async onChatRecovery(
ctx: ChatRecoveryContext,
): Promise<ChatRecoveryOptions> {
if (Date.now() - ctx.createdAt > 2 * 60 * 1000) {
return { continue: false };
}
return {};
}通过保存的请求正文重新调用 onChatMessage,追加到最后一条 assistant 消息。响应作为续传流——追加到现有 assistant 消息,而不是新消息。不创建合成的用户消息。
protected continueLastTurn(
body?: Record<string, unknown>,
options?: SaveMessagesOptions,
): Promise<SaveMessagesResult>;默认恢复路径会自动调用。也可从调度回调或其他入口点手动调用。可选 body 参数会覆盖本次续传已保存的请求正文。传入 options.signal 可在续传运行时取消。
在 onChatMessage 内使用 this.stash() 持久化提供商特定的恢复数据。暂存数据存在 fiber 的 SQLite 行中,与 agent 状态分离,在 onChatRecovery 中作为 ctx.recoveryData 可用。
export class ChatAgent extends AIChatAgent {
chatRecovery = true;
async onChatMessage(_onFinish, options) {
const result = streamText({
model: openai("gpt-5.4"),
messages: await convertToModelMessages(this.messages),
providerOptions: { openai: { store: true } },
includeRawChunks: true,
onChunk: ({ chunk }) => {
if (chunk.type === "raw") {
const raw = chunk.rawValue;
if (raw?.type === "response.created" && raw.response?.id) {
this.stash({ responseId: raw.response.id });
}
}
},
});
return result.toUIMessageStreamResponse();
}
}export class ChatAgent extends AIChatAgent {
override chatRecovery = true;
async onChatMessage(_onFinish, options) {
const result = streamText({
model: openai("gpt-5.4"),
messages: await convertToModelMessages(this.messages),
providerOptions: { openai: { store: true } },
includeRawChunks: true,
onChunk: ({ chunk }) => {
if (chunk.type === "raw") {
const raw = chunk.rawValue as {
type?: string;
response?: { id?: string };
};
if (raw?.type === "response.created" && raw.response?.id) {
this.stash({ responseId: raw.response.id });
}
}
},
});
return result.toUIMessageStreamResponse();
}
}正确策略取决于提供商是否支持 assistant prefill,以及断开后响应是否在服务端继续:
| 提供商 | 策略 | Token 成本 |
|---|---|---|
| Workers AI | continueLastTurn() — 模型通过 assistant prefill 继续 |
低 |
| OpenAI (Responses API) | 按 ID 检索已完成响应 — 不浪费 token | 零 |
| Anthropic | 持久化部分结果,发送合成的 user message 以继续 | 中 |
轮次恢复期间,agent 会广播 cf_agent_chat_recovering 状态帧,客户端可显示「recovering…」指示,而不是看起来像冻结。该标志在调度恢复续传时设置,并在每个终态结果时清除,因此指示不会一直转圈。通过 useAgentChat 的 isRecovering 标志消费(见返回值)。该信号仅作提示,且向后兼容——不理解它的客户端会忽略。
对话记录修复(transcript repair)——修复孤立的工具调用(保留为出错结果而不是删除,以便记录在重启后仍在,且模型不会静默重跑该工具),以及在调用提供商前规范化格式错误或缺失的工具输入——会在 transcript 可观测性通道上发出。
聊天恢复如何融入更广的长时间运行 Agent 场景,请参阅长时间运行 Agent:恢复被中断的 LLM 流。底层 fiber API 请参阅 Durable Execution。
通过 WebSocket 连接 AIChatAgent 的 React hook。使用原生 WebSocket 传输包装 AI SDK 的 useChat。
import { useAgent } from "agents/react";
import { useAgentChat } from "@cloudflare/ai-chat/react";
function Chat() {
const agent = useAgent({ agent: "ChatAgent" });
const {
messages,
sendMessage,
clearHistory,
addToolOutput,
addToolApprovalResponse,
setMessages,
status,
isStreaming,
isServerStreaming,
isToolContinuation,
isRecovering,
} = useAgentChat({ agent });
// ...
}import { useAgent } from "agents/react";
import { useAgentChat } from "@cloudflare/ai-chat/react";
function Chat() {
const agent = useAgent({ agent: "ChatAgent" });
const {
messages,
sendMessage,
clearHistory,
addToolOutput,
addToolApprovalResponse,
setMessages,
status,
isStreaming,
isServerStreaming,
isToolContinuation,
isRecovering,
} = useAgentChat({ agent });
// ...
}| 选项 | 类型 | 默认值 | 描述 |
|---|---|---|---|
agent |
ReturnType<typeof useAgent> |
必需 | 来自 useAgent 的 Agent 连接 |
onToolCall |
({ toolCall, addToolOutput }) => void |
— | 处理客户端工具执行 |
tools |
Record<string, AITool> |
— | 高级:从浏览器动态注册由客户端执行的工具 |
autoContinueAfterToolResult |
boolean |
true |
客户端工具结果与审批后自动续传对话 |
resume |
boolean |
true |
重连时启用自动流恢复 |
cancelOnClientAbort |
boolean |
false |
通用客户端流中止或清理发生时取消服务端轮次 |
body |
object | () => object |
— | 随每次请求发送的自定义数据 |
prepareSendMessagesRequest |
(options) => { body?, headers? } |
— | 高级的按请求自定义 |
getInitialMessages |
(options) => Promise<UIMessage[]> or null |
— | 自定义初始消息加载器。设为 null 可完全跳过 HTTP fetch(直接提供 messages 时有用) |
syncMessagesToServer |
boolean |
true |
为 true 时,setMessages 将对话记录推送到服务端。以服务端为权威存储的主机应设为 false,使 setMessages 仅更新本地视图 |
| 属性 | 类型 | 描述 |
|---|---|---|
messages |
UIMessage[] |
当前对话消息 |
sendMessage |
(message) => void |
发送消息 |
clearHistory |
() => void |
清空对话(客户端与服务端) |
addToolOutput |
({ toolCallId, output }) => void |
为客户端工具提供输出 |
addToolApprovalResponse |
({ id, approved }) => void |
批准或拒绝需要审批的工具 |
setMessages |
(messages | updater) => void |
直接设置消息(同步到服务端) |
status |
string |
"ready"、"submitted"、"streaming" 或 "error" |
isStreaming |
boolean |
agent 正在流式传输或等待活跃客户端工具时为 true |
isServerStreaming |
boolean |
服务端发起的流或活跃客户端工具阶段进行中时为 true |
isToolContinuation |
boolean |
工具结果或审批后自动续传运行中为 true |
isRecovering |
boolean |
持久轮次恢复中(被中断并正在恢复)时为 true。与 isStreaming 不同 — 恢复中的轮次尚未产生 token。渲染「recovering…」提示;多数 UI 将 isStreaming || isRecovering 视为「忙碌」 |
UI 需区分全新用户提交与工具结果后的续传时,使用 isToolContinuation。例如,仅在 status === "submitted" && !isToolContinuation 时显示输入中指示,而 isStreaming 为 true 时始终禁用加载控件。
AIChatAgent 支持三种工具模式,均使用 AI SDK 的 tool() 函数:
| 模式 | 运行位置 | 何时使用 |
|---|---|---|
| 服务端 | Server(自动) | API 调用、数据库查询、计算 |
| 客户端 | Browser(通过 onToolCall) |
地理定位、剪贴板、相机、本地存储 |
| 审批 | Server(用户审批后) | 支付、删除、外部操作 |
带 execute 函数的工具在服务端自动运行:
import { streamText, convertToModelMessages, tool, stepCountIs } from "ai";
import { z } from "zod";
export class ChatAgent extends AIChatAgent {
async onChatMessage() {
const workersai = createWorkersAI({ binding: this.env.AI });
const result = streamText({
model: workersai("@cf/zai-org/glm-4.7-flash"),
messages: await convertToModelMessages(this.messages),
tools: {
getWeather: tool({
description: "Get weather for a city",
inputSchema: z.object({ city: z.string() }),
execute: async ({ city }) => {
const data = await fetchWeather(city);
return { temperature: data.temp, condition: data.condition };
},
}),
},
stopWhen: stepCountIs(5),
});
return result.toUIMessageStreamResponse();
}
}import { streamText, convertToModelMessages, tool, stepCountIs } from "ai";
import { z } from "zod";
export class ChatAgent extends AIChatAgent {
async onChatMessage() {
const workersai = createWorkersAI({ binding: this.env.AI });
const result = streamText({
model: workersai("@cf/zai-org/glm-4.7-flash"),
messages: await convertToModelMessages(this.messages),
tools: {
getWeather: tool({
description: "Get weather for a city",
inputSchema: z.object({ city: z.string() }),
execute: async ({ city }) => {
const data = await fetchWeather(city);
return { temperature: data.temp, condition: data.condition };
},
}),
},
stopWhen: stepCountIs(5),
});
return result.toUIMessageStreamResponse();
}
}在服务端定义不带 execute 的工具,然后在客户端用 onToolCall 处理。适用于需要浏览器 API 的工具。
服务端:
tools: {
getLocation: tool({
description: "Get the user's location from the browser",
inputSchema: z.object({}),
// No execute — the client handles it
});
}tools: {
getLocation: tool({
description: "Get the user's location from the browser",
inputSchema: z.object({}),
// No execute — the client handles it
});
}客户端:
const { messages, sendMessage } = useAgentChat({
agent,
onToolCall: async ({ toolCall, addToolOutput }) => {
if (toolCall.toolName === "getLocation") {
const pos = await new Promise((resolve, reject) =>
navigator.geolocation.getCurrentPosition(resolve, reject),
);
addToolOutput({
toolCallId: toolCall.toolCallId,
output: { lat: pos.coords.latitude, lng: pos.coords.longitude },
});
}
},
});const { messages, sendMessage } = useAgentChat({
agent,
onToolCall: async ({ toolCall, addToolOutput }) => {
if (toolCall.toolName === "getLocation") {
const pos = await new Promise((resolve, reject) =>
navigator.geolocation.getCurrentPosition(resolve, reject),
);
addToolOutput({
toolCallId: toolCall.toolCallId,
output: { lat: pos.coords.latitude, lng: pos.coords.longitude },
});
}
},
});LLM 调用 getLocation 时流暂停。onToolCall 回调触发,你的代码提供输出,对话继续。
浏览器在运行时决定可用工具的 SDK 或平台,向 useAgentChat 传入 tools 对象。带客户端 execute 函数的工具会自动序列化并发送到服务端。服务端上 options.clientTools 与 createToolsFromClientSchemas() 仍支持此动态工具模式。
对执行前需要用户确认的工具使用 needsApproval。
服务端:
tools: {
processPayment: tool({
description: "Process a payment",
inputSchema: z.object({
amount: z.coerce.number(),
recipient: z.string(),
}),
needsApproval: async ({ amount }) => amount > 100,
execute: async ({ amount, recipient }) => charge(amount, recipient),
});
}客户端:
import { getToolName, isToolUIPart } from "ai";
import {
getToolApproval,
getToolCallId,
getToolPartState,
} from "@cloudflare/ai-chat/react";
const { messages, addToolApprovalResponse } = useAgentChat({ agent });
// Render pending approvals from message parts
{
messages.map((msg) =>
msg.parts
.filter(
(part) =>
isToolUIPart(part) && getToolPartState(part) === "waiting-approval",
)
.map((part) => (
<div key={getToolCallId(part)}>
<p>Approve {getToolName(part)}?</p>
<button
onClick={() => {
const approval = getToolApproval(part);
if (!approval) return;
addToolApprovalResponse({
id: approval.id,
approved: true,
});
}}
>
Approve
</button>
<button
onClick={() => {
const approval = getToolApproval(part);
if (!approval) return;
addToolApprovalResponse({
id: approval.id,
approved: false,
});
}}
>
Reject
</button>
</div>
)),
);
}用户拒绝 tool 时,addToolApprovalResponse({ id, approved: false }) 将 tool state 设为带通用消息的 output-denied。要给 LLM 更具体的拒绝原因,请改用 addToolOutput 并设置 state: "output-error":
const { addToolOutput } = useAgentChat({ agent });
// Reject with a custom error message
addToolOutput({
toolCallId: part.toolCallId,
state: "output-error",
errorText: "User declined: insufficient budget for this quarter",
});const { addToolOutput } = useAgentChat({ agent });
// Reject with a custom error message
addToolOutput({
toolCallId: part.toolCallId,
state: "output-error",
errorText: "User declined: insufficient budget for this quarter",
});这向 LLM 发送带自定义错误文本的 tool_result,使其能适当响应(例如建议替代方案或提出澄清问题)。
启用 autoContinueAfterToolResult(默认)时,addToolApprovalResponse(approved: false)会自动 continue 对话。state: "output-error" 的 addToolOutput 不会自动 continue — 若希望 LLM 响应错误,之后调用 sendMessage()。
更多模式请参阅 Human-in-the-loop。
使用 body 选项为每次 chat 请求包含自定义数据:
const { messages, sendMessage } = useAgentChat({
agent,
body: {
timezone: Intl.DateTimeFormat().resolvedOptions().timeZone,
userId: currentUser.id,
},
});const { messages, sendMessage } = useAgentChat({
agent,
body: {
timezone: Intl.DateTimeFormat().resolvedOptions().timeZone,
userId: currentUser.id,
},
});动态值请使用函数:
body: () => ({
token: getAuthToken(),
timestamp: Date.now(),
});body: () => ({
token: getAuthToken(),
timestamp: Date.now(),
});在 server 上访问这些字段:
export class ChatAgent extends AIChatAgent {
async onChatMessage(_onFinish, options) {
const { timezone, userId } = options?.body ?? {};
// ...
}
}export class ChatAgent extends AIChatAgent {
async onChatMessage(_onFinish, options) {
const { timezone, userId } = options?.body ?? {};
// ...
}
}高级 per-request 自定义(自定义 header、每次请求不同 body)请使用 prepareSendMessagesRequest:
const { messages, sendMessage } = useAgentChat({
agent,
prepareSendMessagesRequest: async ({ messages, trigger }) => ({
headers: { Authorization: `Bearer ${await getToken()}` },
body: { requestedAt: Date.now() },
}),
});const { messages, sendMessage } = useAgentChat({
agent,
prepareSendMessagesRequest: async ({ messages, trigger }) => ({
headers: { Authorization: `Bearer ${await getToken()}` },
body: { requestedAt: Date.now() },
}),
});数据部分允许在文本旁向消息附加带类型的 JSON——进度指示、来源引用、token 用量或 UI 所需的任何结构化数据。
使用 createUIMessageStream 与 writer.write() 从 server 发送 data part:
import {
streamText,
convertToModelMessages,
createUIMessageStream,
createUIMessageStreamResponse,
} from "ai";
export class ChatAgent extends AIChatAgent {
async onChatMessage() {
const workersai = createWorkersAI({ binding: this.env.AI });
const stream = createUIMessageStream({
execute: async ({ writer }) => {
const result = streamText({
model: workersai("@cf/zai-org/glm-4.7-flash"),
messages: await convertToModelMessages(this.messages),
});
// Merge the LLM stream
writer.merge(result.toUIMessageStream());
// Write a data part — persisted to message.parts
writer.write({
type: "data-sources",
id: "src-1",
data: { query: "agents", status: "searching", results: [] },
});
// Later: update the same part in-place (same type + id)
writer.write({
type: "data-sources",
id: "src-1",
data: {
query: "agents",
status: "found",
results: ["Agents SDK docs", "Durable Objects guide"],
},
});
},
});
return createUIMessageStreamResponse({ stream });
}
}import {
streamText,
convertToModelMessages,
createUIMessageStream,
createUIMessageStreamResponse,
} from "ai";
export class ChatAgent extends AIChatAgent {
async onChatMessage() {
const workersai = createWorkersAI({ binding: this.env.AI });
const stream = createUIMessageStream({
execute: async ({ writer }) => {
const result = streamText({
model: workersai("@cf/zai-org/glm-4.7-flash"),
messages: await convertToModelMessages(this.messages),
});
// Merge the LLM stream
writer.merge(result.toUIMessageStream());
// Write a data part — persisted to message.parts
writer.write({
type: "data-sources",
id: "src-1",
data: { query: "agents", status: "searching", results: [] },
});
// Later: update the same part in-place (same type + id)
writer.write({
type: "data-sources",
id: "src-1",
data: {
query: "agents",
status: "found",
results: ["Agents SDK docs", "Durable Objects guide"],
},
});
},
});
return createUIMessageStreamResponse({ stream });
}
}| 模式 | 方式 | 持久化? | 用例 |
|---|---|---|---|
| 协调(Reconciliation) | 相同 type + id → 原地更新 |
是 | 渐进状态(searching → found) |
| 追加(Append) | 无 id 或不同 id → 追加 |
是 | 日志条目、多条引用 |
| 瞬时(Transient) | transient: true → 不加入 message.parts |
否 | 临时状态(思考指示) |
瞬时部分实时广播给已连接客户端,但排除在 SQLite 持久化与 message.parts 之外。使用 onData 回调消费。
非瞬时数据部分出现在 message.parts 中。使用 UIMessage 泛型为其标注类型:
import { useAgentChat } from "@cloudflare/ai-chat/react";
const { messages } = useAgentChat({ agent });
// Typed access — no casts needed
for (const msg of messages) {
for (const part of msg.parts) {
if (part.type === "data-sources") {
console.log(part.data.results); // string[]
}
}
}import { useAgentChat } from "@cloudflare/ai-chat/react";
import type { UIMessage } from "ai";
type ChatMessage = UIMessage<
unknown,
{
sources: { query: string; status: string; results: string[] };
usage: { model: string; inputTokens: number; outputTokens: number };
}
>;
const { messages } = useAgentChat<unknown, ChatMessage>({ agent });
// Typed access — no casts needed
for (const msg of messages) {
for (const part of msg.parts) {
if (part.type === "data-sources") {
console.log(part.data.results); // string[]
}
}
}瞬时数据部分不在 message.parts 中。请改用 onData 回调:
const [thinking, setThinking] = useState(false);
const { messages } = useAgentChat({
agent,
onData(part) {
if (part.type === "data-thinking") {
setThinking(true);
}
},
});const [thinking, setThinking] = useState(false);
const { messages } = useAgentChat<unknown, ChatMessage>({
agent,
onData(part) {
if (part.type === "data-thinking") {
setThinking(true);
}
},
});在服务端,用 transient: true 写入瞬时部分:
writer.write({
transient: true,
type: "data-thinking",
data: { model: "glm-4.7-flash", startedAt: new Date().toISOString() },
});writer.write({
transient: true,
type: "data-thinking",
data: { model: "glm-4.7-flash", startedAt: new Date().toISOString() },
});onData 在所有代码路径触发 — 新消息、流恢复与跨标签广播。
客户端断开并重连时,stream 会自动恢复。无需配置——开箱即用。
stream 活跃时:
- 所有 chunk 生成时缓冲在 SQLite 中
- 若 client 断开,server 继续 streaming 与缓冲
- client 重连时接收所有缓冲 chunk 并恢复 live streaming
默认情况下,通用 client stream abort 或 cleanup 保留在 browser 本地,因此 server turn 继续运行且可稍后恢复。显式调用 stop() 仍会取消 server turn:
const { messages, stop } = useAgentChat({ agent });
return <button onClick={stop}>Stop</button>;const { messages, stop } = useAgentChat({ agent });
return <button onClick={stop}>Stop</button>;应用有意让 browser 生命周期拥有 server 生命周期时(例如 request-lifetime 或省 token 流程)设置 cancelOnClientAbort: true。无论此选项如何,显式 stop() 始终取消 server 工作。
用 resume: false 禁用:
const { messages } = useAgentChat({ agent, resume: false });const { messages } = useAgentChat({ agent, resume: false });Workers SQLite 行有 2 MB 的硬性上限。为保持在该限制以下,AIChatAgent 会在序列化消息约 1.8 MB 时开始 compaction,例如 tool 返回非常大的 output 时:
- Tool output 压缩 — 大型 tool output 替换为 LLM 友好摘要,指示 model 建议重跑 tool
- 文本截断 — tool compaction 后 message 仍过大时,text part 截断并附说明
Compacted message 包含 metadata.compactedToolOutputs,client 可检测并优雅显示。
存储(maxPersistedMessages)与 LLM 上下文相互独立:
| 关注点 | 控制 | 范围 |
|---|---|---|
| SQLite 存储多少 message | maxPersistedMessages |
持久化 |
| model 看到什么 | pruneMessages() |
LLM 上下文 |
| 行大小限制 | 自动 compaction | 单条 message |
export class ChatAgent extends AIChatAgent {
async onChatMessage() {
const result = streamText({
model: workersai("@cf/zai-org/glm-4.7-flash"),
messages: pruneMessages({
// LLM context limit
messages: await convertToModelMessages(this.messages),
reasoning: "before-last-message",
toolCalls: "before-last-2-messages",
}),
});
return result.toUIMessageStreamResponse();
}
}export class ChatAgent extends AIChatAgent {
async onChatMessage() {
const result = streamText({
model: workersai("@cf/zai-org/glm-4.7-flash"),
messages: pruneMessages({
// LLM context limit
messages: await convertToModelMessages(this.messages),
reasoning: "before-last-message",
toolCalls: "before-last-2-messages",
}),
});
return result.toUIMessageStreamResponse();
}
}AIChatAgent 可与任意 AI SDK 兼容 provider 配合。由 server 代码决定使用哪个 model — client 无需手动更改。
import { createWorkersAI } from "workers-ai-provider";
const workersai = createWorkersAI({ binding: this.env.AI });
const result = streamText({
model: workersai("@cf/zai-org/glm-4.7-flash"),
messages: await convertToModelMessages(this.messages),
});import { createWorkersAI } from "workers-ai-provider";
const workersai = createWorkersAI({ binding: this.env.AI });
const result = streamText({
model: workersai("@cf/zai-org/glm-4.7-flash"),
messages: await convertToModelMessages(this.messages),
});import { createOpenAI } from "@ai-sdk/openai";
const openai = createOpenAI({ apiKey: this.env.OPENAI_API_KEY });
const result = streamText({
model: openai.chat("gpt-4o"),
messages: await convertToModelMessages(this.messages),
});import { createOpenAI } from "@ai-sdk/openai";
const openai = createOpenAI({ apiKey: this.env.OPENAI_API_KEY });
const result = streamText({
model: openai.chat("gpt-4o"),
messages: await convertToModelMessages(this.messages),
});import { createAnthropic } from "@ai-sdk/anthropic";
const anthropic = createAnthropic({ apiKey: this.env.ANTHROPIC_API_KEY });
const result = streamText({
model: anthropic("claude-sonnet-4-20250514"),
messages: await convertToModelMessages(this.messages),
});import { createAnthropic } from "@ai-sdk/anthropic";
const anthropic = createAnthropic({ apiKey: this.env.ANTHROPIC_API_KEY });
const result = streamText({
model: anthropic("claude-sonnet-4-20250514"),
messages: await convertToModelMessages(this.messages),
});由于 onChatMessage 让你完全控制 streamText 调用,可直接使用任意 AI SDK 功能。以下模式均开箱即用——无需特殊 AIChatAgent 配置。
使用 prepareStep ↗ 在多步 agent 循环的各步骤之间更改模型、可用工具或系统提示词:
import { streamText, convertToModelMessages, tool, stepCountIs } from "ai";
import { z } from "zod";
export class ChatAgent extends AIChatAgent {
async onChatMessage() {
const result = streamText({
model: cheapModel, // Default model for simple steps
messages: await convertToModelMessages(this.messages),
tools: {
search: searchTool,
analyze: analyzeTool,
summarize: summarizeTool,
},
stopWhen: stepCountIs(10),
prepareStep: async ({ stepNumber, messages }) => {
// Phase 1: Search (steps 0-2)
if (stepNumber <= 2) {
return {
activeTools: ["search"],
toolChoice: "required", // Force tool use
};
}
// Phase 2: Analyze with a stronger model (steps 3-5)
if (stepNumber <= 5) {
return {
model: expensiveModel,
activeTools: ["analyze"],
};
}
// Phase 3: Summarize
return { activeTools: ["summarize"] };
},
});
return result.toUIMessageStreamResponse();
}
}import { streamText, convertToModelMessages, tool, stepCountIs } from "ai";
import { z } from "zod";
export class ChatAgent extends AIChatAgent {
async onChatMessage() {
const result = streamText({
model: cheapModel, // Default model for simple steps
messages: await convertToModelMessages(this.messages),
tools: {
search: searchTool,
analyze: analyzeTool,
summarize: summarizeTool,
},
stopWhen: stepCountIs(10),
prepareStep: async ({ stepNumber, messages }) => {
// Phase 1: Search (steps 0-2)
if (stepNumber <= 2) {
return {
activeTools: ["search"],
toolChoice: "required", // Force tool use
};
}
// Phase 2: Analyze with a stronger model (steps 3-5)
if (stepNumber <= 5) {
return {
model: expensiveModel,
activeTools: ["analyze"],
};
}
// Phase 3: Summarize
return { activeTools: ["summarize"] };
},
});
return result.toUIMessageStreamResponse();
}
}prepareStep 在每个 step 前运行,可返回 model、activeTools、toolChoice、system 与 messages 的 override。用于:
- 切换 model — 简单 step 用廉价 model,推理 step 升级
- 分阶段 tool — 限制每个 step 可用的 tool
- 管理上下文 — prune 或 transform message 以保持在 token 限制内
- 强制 tool 调用 — 使用
toolChoice: { type: "tool", toolName: "search" }要求特定 tool
使用 wrapLanguageModel ↗ 在不修改 chat 逻辑的情况下添加 guardrail、RAG、缓存或 logging:
import { streamText, convertToModelMessages, wrapLanguageModel } from "ai";
const guardrailMiddleware = {
wrapGenerate: async ({ doGenerate }) => {
const { text, ...rest } = await doGenerate();
// Filter PII or sensitive content from the response
const cleaned = text?.replace(/\b\d{3}-\d{2}-\d{4}\b/g, "[REDACTED]");
return { text: cleaned, ...rest };
},
};
export class ChatAgent extends AIChatAgent {
async onChatMessage() {
const model = wrapLanguageModel({
model: baseModel,
middleware: [guardrailMiddleware],
});
const result = streamText({
model,
messages: await convertToModelMessages(this.messages),
});
return result.toUIMessageStreamResponse();
}
}import { streamText, convertToModelMessages, wrapLanguageModel } from "ai";
import type { LanguageModelV3Middleware } from "@ai-sdk/provider";
const guardrailMiddleware: LanguageModelV3Middleware = {
wrapGenerate: async ({ doGenerate }) => {
const { text, ...rest } = await doGenerate();
// Filter PII or sensitive content from the response
const cleaned = text?.replace(/\b\d{3}-\d{2}-\d{4}\b/g, "[REDACTED]");
return { text: cleaned, ...rest };
},
};
export class ChatAgent extends AIChatAgent {
async onChatMessage() {
const model = wrapLanguageModel({
model: baseModel,
middleware: [guardrailMiddleware],
});
const result = streamText({
model,
messages: await convertToModelMessages(this.messages),
});
return result.toUIMessageStreamResponse();
}
}AI SDK 包含内置 middleware:
extractReasoningMiddleware— 呈现 DeepSeek R1 等 model 的 chain-of-thoughtdefaultSettingsMiddleware— 应用默认 temperature、max tokens 等simulateStreamingMiddleware— 为非 streaming model 添加 streaming
多个 middleware 按顺序组合:middleware: [first, second] 应用为 first(second(model))。
在 tool 内使用 generateObject ↗ 进行结构化数据提取:
import {
streamText,
generateObject,
convertToModelMessages,
tool,
stepCountIs,
} from "ai";
import { z } from "zod";
export class ChatAgent extends AIChatAgent {
async onChatMessage() {
const result = streamText({
model: myModel,
messages: await convertToModelMessages(this.messages),
tools: {
extractContactInfo: tool({
description:
"Extract structured contact information from the conversation",
inputSchema: z.object({
text: z.string().describe("The text to extract contact info from"),
}),
execute: async ({ text }) => {
const { object } = await generateObject({
model: myModel,
schema: z.object({
name: z.string(),
email: z.string().email(),
phone: z.string().optional(),
}),
prompt: `Extract contact information from: ${text}`,
});
return object;
},
}),
},
stopWhen: stepCountIs(5),
});
return result.toUIMessageStreamResponse();
}
}import {
streamText,
generateObject,
convertToModelMessages,
tool,
stepCountIs,
} from "ai";
import { z } from "zod";
export class ChatAgent extends AIChatAgent {
async onChatMessage() {
const result = streamText({
model: myModel,
messages: await convertToModelMessages(this.messages),
tools: {
extractContactInfo: tool({
description:
"Extract structured contact information from the conversation",
inputSchema: z.object({
text: z.string().describe("The text to extract contact info from"),
}),
execute: async ({ text }) => {
const { object } = await generateObject({
model: myModel,
schema: z.object({
name: z.string(),
email: z.string().email(),
phone: z.string().optional(),
}),
prompt: `Extract contact information from: ${text}`,
});
return object;
},
}),
},
stopWhen: stepCountIs(5),
});
return result.toUIMessageStreamResponse();
}
}工具可将工作委派给具有独立上下文的聚焦子调用。使用 ToolLoopAgent ↗ 定义可复用 agent,然后从工具的 execute 调用:
import {
ToolLoopAgent,
streamText,
convertToModelMessages,
tool,
stepCountIs,
} from "ai";
import { z } from "zod";
// Define a reusable research agent with its own tools and instructions
const researchAgent = new ToolLoopAgent({
model: researchModel,
instructions: "You are a research assistant. Be thorough and cite sources.",
tools: { webSearch: webSearchTool },
stopWhen: stepCountIs(10),
});
export class ChatAgent extends AIChatAgent {
async onChatMessage() {
const result = streamText({
model: orchestratorModel,
messages: await convertToModelMessages(this.messages),
tools: {
deepResearch: tool({
description: "Research a topic in depth",
inputSchema: z.object({
topic: z.string().describe("The topic to research"),
}),
execute: async ({ topic }) => {
const { text } = await researchAgent.generate({
prompt: topic,
});
return { summary: text };
},
}),
},
stopWhen: stepCountIs(5),
});
return result.toUIMessageStreamResponse();
}
}import {
ToolLoopAgent,
streamText,
convertToModelMessages,
tool,
stepCountIs,
} from "ai";
import { z } from "zod";
// Define a reusable research agent with its own tools and instructions
const researchAgent = new ToolLoopAgent({
model: researchModel,
instructions: "You are a research assistant. Be thorough and cite sources.",
tools: { webSearch: webSearchTool },
stopWhen: stepCountIs(10),
});
export class ChatAgent extends AIChatAgent {
async onChatMessage() {
const result = streamText({
model: orchestratorModel,
messages: await convertToModelMessages(this.messages),
tools: {
deepResearch: tool({
description: "Research a topic in depth",
inputSchema: z.object({
topic: z.string().describe("The topic to research"),
}),
execute: async ({ topic }) => {
const { text } = await researchAgent.generate({
prompt: topic,
});
return { summary: text };
},
}),
},
stopWhen: stepCountIs(5),
});
return result.toUIMessageStreamResponse();
}
}research agent 在独立上下文中运行——其 token budget 与 orchestrator 分离。仅 summary 返回父 model。
默认情况下,tool part 在 execute 返回前显示为 loading。使用 async generator(async function*)在 tool 仍在工作时向 client stream 进度更新:
deepResearch: tool({
description: "Research a topic in depth",
inputSchema: z.object({
topic: z.string().describe("The topic to research"),
}),
async *execute({ topic }) {
// Preliminary result — the client sees "searching" immediately
yield { status: "searching", topic, summary: undefined };
const { text } = await researchAgent.generate({ prompt: topic });
// Final result — sent to the model for its next step
yield { status: "done", topic, summary: text };
},
});deepResearch: tool({
description: "Research a topic in depth",
inputSchema: z.object({
topic: z.string().describe("The topic to research"),
}),
async *execute({ topic }) {
// Preliminary result — the client sees "searching" immediately
yield { status: "searching", topic, summary: undefined };
const { text } = await researchAgent.generate({ prompt: topic });
// Final result — sent to the model for its next step
yield { status: "done", topic, summary: text };
},
});每个 yield 实时更新客户端上的工具部分(preliminary: true)。最后 yield 的值成为模型看到的最终输出。
此模式适用于:
- 任务需探索大量会膨胀主上下文的信息
- 希望为长时间运行的工具显示实时进度
- 希望并行化独立研究(多个工具调用并发运行)
- 不同子任务需要不同模型或系统提示词
更多请参阅 AI SDK Agents 文档 ↗、Subagents ↗ 与 Preliminary Tool Results ↗。
多个客户端连接到同一 Agent 实例时,消息会自动广播到所有连接。若一个客户端发送消息,所有其他已连接客户端都会收到更新后的消息列表。
Client A ──── sendMessage("Hello") ────▶ AIChatAgent
│
persist + stream
│
Client A ◀── CF_AGENT_USE_CHAT_RESPONSE ──────┤
Client B ◀── CF_AGENT_CHAT_MESSAGES ──────────┘发起客户端接收流式响应。所有其他客户端通过 CF_AGENT_CHAT_MESSAGES 广播接收最终消息。
| 导入路径 | 导出 |
|---|---|
@cloudflare/ai-chat |
AIChatAgent, createToolsFromClientSchemas, ClientToolSchema, ChatRecoveryContext, ChatRecoveryOptions, ChatRecoveryConfig, ChatRecoveryExhaustedContext, ResolvedChatRecoveryConfig, lifecycle types |
@cloudflare/ai-chat/react |
useAgentChat, extractClientToolSchemas, getToolPartState, getToolCallId, getToolInput, getToolOutput, getToolApproval |
@cloudflare/ai-chat/types |
MessageType, OutgoingMessage, IncomingMessage |
agents/chat |
共享高级 chat 原语,如 SaveMessagesResult、SaveMessagesOptions、CHAT_MESSAGE_TYPES、ROW_MAX_BYTES 与 isReplayChunk() |
聊天协议通过 WebSocket 使用带类型的 JSON 消息:
| 消息 | 方向 | 用途 |
|---|---|---|
CF_AGENT_USE_CHAT_REQUEST |
客户端 → 服务端 | 发送聊天消息 |
CF_AGENT_USE_CHAT_RESPONSE |
服务端 → 客户端 | 流式响应分块 |
CF_AGENT_CHAT_MESSAGES |
服务端 → 客户端 | 广播更新后的消息 |
CF_AGENT_CHAT_CLEAR |
双向 | 清空对话 |
CF_AGENT_CHAT_REQUEST_CANCEL |
客户端 → 服务端 | 取消活跃流 |
CF_AGENT_TOOL_RESULT |
客户端 → 服务端 | 提供工具输出 |
CF_AGENT_TOOL_APPROVAL |
客户端 → 服务端 | 批准或拒绝工具 |
CF_AGENT_MESSAGE_UPDATED |
服务端 → 客户端 | 通知消息更新 |
CF_AGENT_STREAM_RESUMING |
服务端 → 客户端 | 通知流恢复 |
CF_AGENT_STREAM_RESUME_REQUEST |
客户端 → 服务端 | 请求流恢复检查 |
CF_AGENT_STREAM_RESUME_ACK |
服务端 → 客户端 | 从游标恢复流 |
CF_AGENT_STREAM_RESUME_NONE |
服务端 → 客户端 | 无可恢复的流 |
以下 API 已弃用,使用时会输出 console 警告,并将在未来版本中移除。
| 已弃用 | 替代 | 说明 |
|---|---|---|
addToolResult({ toolCallId, result }) |
addToolOutput({ toolCallId, output }) |
为与 AI SDK 术语一致而重命名 |
detectToolsRequiringConfirmation() |
在 tool 定义上使用 needsApproval |
审批现为 per-tool,非全局 filter |
toolsRequiringConfirmation option |
在单个 tool 上使用 needsApproval |
Per-tool 审批替代全局列表 |
若从早期版本升级,请将已弃用调用替换为替代 API。已弃用 API 仍可用,但将在未来 major 版本移除。
createToolsFromClientSchemas()、extractClientToolSchemas() 与 useAgentChat 上的 tools 选项仍支持动态客户端工具。它们是高级 API,不是已弃用 API。