跳转到内容
搜索文档

聊天 Agent

最后更新 查看 MarkdownAgent 设置

使用 AIChatAgentuseAgentChat 构建 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/typesagents/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-provider

服务端

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 { 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 配置

// 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)
  1. 客户端通过 WebSocket 发送消息
  2. AIChatAgent 将消息持久化到 SQLite 并调用你的 onChatMessage 方法
  3. 你的方法返回流式 Response(通常来自 streamText
  4. Chunks 通过 WebSocket 实时回传
  5. 流完成后,最终消息会被持久化并广播到所有连接

服务端 API

AIChatAgent

扩展 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)
	}
}

onChatMessage

这是你需要重写的主要方法。它接收对话上下文并应返回 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。可用它选择不同模型、调整系统提示词,或在续传轮次跳过昂贵的上下文组装。

this.messages

从 SQLite 加载的当前对话历史。这是 AI SDK 的 UIMessage 对象数组。每次交互后消息会自动持久化。

maxPersistedMessages

限制 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();
	}
}

waitForMcpConnections

控制 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()

messageConcurrency

控制聊天轮次已活跃或排队时,重叠用户提交的行为。

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() 调用仍保持现有串行行为。

persistMessagessaveMessages

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,恢复的轮次会在没有原始信号的情况下运行。

onChatResponse

聊天轮次产生并持久化 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" 时的错误消息

sanitizeMessageForPersistence

重写此方法可在消息持久化到存储前应用自定义转换。此 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;
			}),
		};
	}
}

轮次生命周期辅助方法

这些方法帮助协调编程式轮次并等待待处理交互。

hasPendingInteraction()

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");
}

waitUntilStable()

等待对话完全稳定——无活跃流、无待处理的客户端工具交互、无排队的续传轮次。稳定时返回 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 responding
await this.saveMessages((messages) => [...messages, syntheticUserMessage]);
await this.waitUntilStable({ timeout: 60_000 });
// The assistant has finished responding

resetTurnState()

中止活跃轮次并使排队的续传失效。内置 CF_AGENT_CHAT_CLEAR 处理程序会自动调用,必要时也可手动调用。

生命周期钩子

重写 onConnectonClose 以添加自定义逻辑。流恢复与消息同步由框架处理:

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 默认 chatRecoveryfalse,因此现有聊天 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(传给 shouldKeepRecoveringctx)包含以下字段:

字段 类型 描述
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 — 你的 shouldKeepRecovering hook 返回 false
  • stable_timeout — 恢复尝试等待稳定状态持续超时直至预算耗尽(极端抖动)。

等待人类的轮次不会被封存

轮次可能暂停在无法自行完成的客户端交互上:客户端工具调用(无服务端 execute、由客户端回放结果的工具),或 approval-requested 部分。此类轮次在等待人类,而不是卡住。

交互待处理期间,轮次豁免所有恢复预算。无进度窗口、尝试上限、maxRecoveryWorkshouldKeepRecovering 均暂停。用户因部署中断后需数分钟回答提示也不会触发封存。恢复会停放该轮次而不是失败,用户最终审批或工具结果通过正常续传路径恢复。

此豁免仅适用于客户端。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);
	}
});

onChatRecovery

重写以实现提供商特定的恢复。默认行为持久化部分响应并通过 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 消息。若最新持久化消息仍是中断轮次的未回答用户消息,框架会自动重试该轮次,除非 continuefalse

使用 ctx.createdAt 跳过过期恢复:

override async onChatRecovery(
	ctx: ChatRecoveryContext,
): Promise<ChatRecoveryOptions> {
	if (Date.now() - ctx.createdAt > 2 * 60 * 1000) {
		return { continue: false };
	}
	return {};
}

continueLastTurn

通过保存的请求正文重新调用 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…」指示,而不是看起来像冻结。该标志在调度恢复续传时设置,并在每个终态结果时清除,因此指示不会一直转圈。通过 useAgentChatisRecovering 标志消费(见返回值)。该信号仅作提示,且向后兼容——不理解它的客户端会忽略。

对话记录修复(transcript repair)——修复孤立的工具调用(保留为出错结果而不是删除,以便记录在重启后仍在,且模型不会静默重跑该工具),以及在调用提供商前规范化格式错误或缺失的工具输入——会在 transcript 可观测性通道上发出。

聊天恢复如何融入更广的长时间运行 Agent 场景,请参阅长时间运行 Agent:恢复被中断的 LLM 流。底层 fiber API 请参阅 Durable Execution

客户端 API

useAgentChat

通过 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(用户审批后) 支付、删除、外部操作

服务端 tool

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();
	}
}

客户端 tool

在服务端定义不带 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.clientToolscreateToolsFromClientSchemas() 仍支持此动态工具模式。

工具审批(人机协同)

对执行前需要用户确认的工具使用 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>
			)),
	);
}

使用 addToolOutput 自定义拒绝消息

用户拒绝 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(默认)时,addToolApprovalResponseapproved: 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() },
	}),
});

数据部分(Data part)

数据部分允许在文本旁向消息附加带类型的 JSON——进度指示、来源引用、token 用量或 UI 所需的任何结构化数据。

写入 data part(服务端)

使用 createUIMessageStreamwriter.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[]
		}
	}
}

使用 onData 处理瞬时部分

瞬时数据部分不在 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 活跃时:

  1. 所有 chunk 生成时缓冲在 SQLite 中
  2. 若 client 断开,server 继续 streaming 与缓冲
  3. 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 时:

  1. Tool output 压缩 — 大型 tool output 替换为 LLM 友好摘要,指示 model 建议重跑 tool
  2. 文本截断 — tool compaction 后 message 仍过大时,text part 截断并附说明

Compacted message 包含 metadata.compactedToolOutputs,client 可检测并优雅显示。

控制 LLM 上下文与存储

存储(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();
	}
}

使用不同的 AI 提供商

AIChatAgent 可与任意 AI SDK 兼容 provider 配合。由 server 代码决定使用哪个 model — client 无需手动更改。

Workers AI (Cloudflare)

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),
});

OpenAI

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),
});

Anthropic

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 前运行,可返回 modelactiveToolstoolChoicesystemmessages 的 override。用于:

  • 切换 model — 简单 step 用廉价 model,推理 step 升级
  • 分阶段 tool — 限制每个 step 可用的 tool
  • 管理上下文 — prune 或 transform message 以保持在 token 限制内
  • 强制 tool 调用 — 使用 toolChoice: { type: "tool", toolName: "search" } 要求特定 tool

语言 model 中间件

使用 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-thought
  • defaultSettingsMiddleware — 应用默认 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();
	}
}

进程内子 agent 委派

工具可将工作委派给具有独立上下文的聚焦子调用。使用 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。

使用 preliminary result 流式传输进度

默认情况下,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 文档SubagentsPreliminary 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 广播接收最终消息。

API 参考

导出

导入路径 导出
@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 原语,如 SaveMessagesResultSaveMessagesOptionsCHAT_MESSAGE_TYPESROW_MAX_BYTESisReplayChunk()

WebSocket 协议

聊天协议通过 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

以下 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。

后续步骤

这篇文档对您有帮助吗?