跳转到内容
搜索文档

Agent 类内部机制

最后更新 查看 MarkdownAgent 设置

agents 库的核心是 Agent 类。你扩展它、重写少量方法,即可免费获得状态管理、WebSocket、调度、RPC 等能力。本页逐层说明 Agent 的构建方式,帮助你理解底层机制。

此处代码片段仅供说明,不一定代表最佳实践。完整 API 请参阅 API 参考源代码

什么是 Agent?

Agent 类扩展 DurableObject——Agent 就是 Durable Object。若不熟悉 Durable Object,请先阅读什么是 Durable Objects。其核心是全局可寻址(每个实例有唯一 ID)、单线程计算实例,并具备长期存储(键值与 SQLite)。

Agent 类并不直接继承 DurableObject,而是继承 partyserver 包中的 Server,而 Server 继承 DurableObject。可将其理解为分层:DurableObject > Server > Agent

第 0 层:Durable Object

简要了解 Durable Objects 暴露的原语,有助于理解外层如何使用它们。Durable Object 类提供:

constructor

constructor(ctx: DurableObjectState, env: Env) {}

Workers 运行时始终调用构造函数以处理内部事务。这意味着两点:

  1. 虽然每次 Durable Object 初始化时都会调用构造函数,但其签名是固定的。开发者无法从构造函数添加或更新参数。
  2. 开发者不能手动实例化类,而必须通过绑定 API,通过 DurableObjectNamespace 完成。

RPC

编写继承内置类型 DurableObject 的 Durable Object 类时,公开方法会暴露为 RPC 方法,开发者可使用 Worker 中的 DurableObjectStub 调用。

// This instance could've been active, hibernated,
// not initialized or maybe had never even been created!
const stub = env.MY_DO.getByName("foo");

// We can call any public method on the class. The runtime
// ensures the constructor is called if the instance was not active.
await stub.bar();

fetch()

Durable Objects 可接收 Worker 的 Request 并返回 Response,但只能通过开发者必须实现的 fetch 方法完成。

WebSockets

Durable Objects 对 WebSockets 提供一等支持。Durable Object 可接受 fetchRequest 携带的 WebSocket 并随后不再管理它。基类提供开发者可实现的回调方法,有效替代事件监听器的需求。

基类提供 webSocketMessage(ws, message)webSocketClose(ws, code, reason, wasClean)webSocketError(ws , error)API)。

export class MyDurableObject extends DurableObject {
	async fetch(request) {
		// Creates two ends of a WebSocket connection.
		const webSocketPair = new WebSocketPair();
		const [client, server] = Object.values(webSocketPair);

		// Calling `acceptWebSocket()` connects the WebSocket to the Durable Object, allowing the WebSocket to send and receive messages.
		this.ctx.acceptWebSocket(server);

		return new Response(null, {
			status: 101,
			webSocket: client,
		});
	}

	async webSocketMessage(ws, message) {
		ws.send(message);
	}
}

alarm()

HTTP 与 RPC 请求并非 Durable Object 的唯一入口。闹钟允许开发者调度稍后触发的事件。当下一个闹钟到期时,运行时会调用由开发者实现的 alarm() 方法。

可使用 this.ctx.storage.setAlarm() 调度闹钟。更多信息请参阅 Alarms

this.ctx

基类 DurableObjectDurableObjectState 设置到 this.ctx。其中有许多有用的方法与属性,此处聚焦 this.ctx.storage

this.ctx.storage

DurableObjectStorage 是与 Durable Object 持久化机制交互的主接口,包括 KV 与 SQLITE 的同步 API。

const sql = this.ctx.storage.sql;

// Synchronous SQL query
const rows = sql.exec("SELECT * FROM contacts WHERE country = ?", "US");

// Key-value storage
const token = this.ctx.storage.get("someToken");

this.ctx.env

最后,Durable Object 在 this.env 中也有 Worker Env。详见 Bindings

第 1 层:Server(partyserver)

了解 Durable Object 开箱即用能力后,来自 partyserverServer 类会更容易理解。它是有明确主张的 DurableObject 包装,用开发者友好的回调替换底层原语。

Server 不会添加任何存储操作——它仅包装 Durable Object 生命周期。

寻址

partyserver 暴露辅助函数,可按名称寻址 Durable Objects,而无需手动通过绑定。其中包含 URL 路由方案(<your-worker>/servers/:durableClass/:durableName),Agent 层在此基础上构建。

// Note the await here!
const stub = await getServerByName(env.MY_DO, "foo");

// We can still call RPC methods.
await stub.bar();

URL 方案还启用请求路由器。在 Agent 层,这被重新导出为 routeAgentRequest

  async fetch(request: Request, env: Env, ctx: ExecutionContext) {
    const res = await routeAgentRequest(request, env);

    if (res) return res;

    return new Response("Not found", { status: 404 });
  }

onStart

寻址层允许 Server 暴露 onStart 回调,在 Durable Object 每次启动时(驱逐、休眠或首次创建后)以及任何 fetch 或 RPC 调用之前运行。

class MyServer extends Server {
	onStart() {
		// Some initialization logic that you wish
		// to run every time the DO is started up.
		const sql = this.ctx.storage.sql;
		sql.exec(`...`);
	}
}

onRequestonConnect

Server 已为底层 Durable Object 实现 fetch,并暴露两个开发者可用的回调:onRequestonConnect,分别用于 HTTP 请求与传入 WebSocket 连接(WebSocket 连接默认被接受)。

class MyServer extends Server {
	async onRequest(request: Request) {
		const url = new URL(request.url);

		return new Response(`Hello from ${url.origin}!`);
	}

	async onConnect(conn, ctx) {
		const { request } = ctx;
		const url = new URL(request.url);

		// Connections are a WebSocket wrapper
		conn.send(`Hello from ${url.origin}!`);
	}
}

WebSockets

onConnect 是每个新连接的回调类似,Server 也在 DurableObject 类的默认回调之上提供包装:onMessageonCloseonError

还有 this.broadcast,向所有已连接客户端发送 WebSocket 消息(并无魔法,只是遍历 this.getConnections()!)。

this.name

在 Durable Object 内部很难获取其 namepartyserver 尝试在 this.name 中提供,但并非完美方案。详见 此 GitHub issue

第 2 层:Agent

最后是 Agent 类。Agent 扩展 Server,为有状态、可调度、可观测的 Agent 提供有明确主张的原语,可通过 RPC、WebSocket 甚至邮件通信。

this.statethis.setState()

Agent 的核心特性之一是自动状态持久化。开发者通过泛型参数与 initialState(仅当存储中无状态时使用)定义状态形状,Agent 负责加载、保存与广播状态变更(见上方 Serverthis.broadcast())。

this.state 是惰性从存储(SQL)加载状态的 getter。使用 this.setState() 更新时,状态会在 Durable Object 驱逐后仍被持久化,该方法会自动序列化状态并写回存储。

还可重写 this.onStateChanged 以响应状态变更。

class MyAgent extends Agent<Env, { count: number }> {
	initialState = { count: 0 };

	increment() {
		this.setState({ count: this.state.count + 1 });
	}

	onStateChanged(state, source) {
		console.log("State updated:", state);
	}
}

状态存储在 cf_agents_state SQL 表中。状态消息以 type: "cf_agent_state" 发送(来自客户端与服务端)。由于 agents 提供 JS 与 React 客户端,实时状态更新开箱即用。

this.sql

Agent 提供便捷的 sql 模板标签,用于对 Durable Object 的 SQL 存储执行查询。它构建参数化查询并执行,使用 this.ctx.storage.sql同步 SQL API。

class MyAgent extends Agent {
	onStart() {
		this.sql`
      CREATE TABLE IF NOT EXISTS users (
        id TEXT PRIMARY KEY,
        name TEXT
      )
    `;

		const userId = "1";
		const userName = "Alice";
		this.sql`INSERT INTO users (id, name) VALUES (${userId}, ${userName})`;

		const users = this.sql<{ id: string; name: string }>`
      SELECT * FROM users WHERE id = ${userId}
    `;
		console.log(users); // [{ id: "1", name: "Alice" }]
	}
}

RPC 与可调用方法

agents 在 Durable Objects RPC 基础上进一步实现 WebSocket RPC,使客户端可直接调用 Agent 方法。要通过 WebSocket 使方法可调用,使用 @callable() 装饰器。方法可返回可序列化值或流(使用 @callable({ stream: true }) 时)。

class MyAgent extends Agent {
	@callable({ description: "Add two numbers" })
	async add(a: number, b: number) {
		return a + b;
	}
}

客户端可通过发送 WebSocket 消息调用此方法:

{
	"type": "rpc",
	"id": "unique-request-id",
	"method": "add",
	"args": [2, 3]
}

例如,使用提供的 React 客户端非常简单:

const { stub } = useAgent({ name: "my-agent" });
const result = await stub.add(2, 3);
console.log(result); // 5

this.queue 及相关方法

Agent 包含内置任务队列用于延迟执行。适用于卸载工作或重试操作。可用方法包括 this.queuethis.dequeuethis.dequeueAllthis.dequeueAllByCallbackthis.getQueuethis.getQueues

class MyAgent extends Agent {
	async onConnect() {
		// Queue a task to be executed later
		await this.queue("processTask", { userId: "123" });
	}

	async processTask(payload: { userId: string }, queueItem: QueueItem) {
		console.log("Processing task for user:", payload.userId);
	}
}

任务存储在 cf_agents_queues SQL 表中并自动按序刷新。任务成功后会自动出队。

this.schedule 及相关方法

Agent 支持通过包装 Durable Object 的 alarm() 来调度方法执行。可用方法包括 this.schedulethis.getSchedulethis.getSchedulesthis.cancelSchedule。调度可以是一次性、延迟或周期性(使用 cron 表达式)。

由于 Durable Objects 同一时间只允许一个闹钟,Agent 类通过在 SQL 中管理多个调度并使用单个闹钟来绕过此限制。

class MyAgent extends Agent {
	async foo() {
		// Schedule at a specific time
		await this.schedule(new Date("2025-12-25T00:00:00Z"), "sendGreeting", {
			message: "Merry Christmas!",
		});

		// Schedule with a delay (in seconds)
		await this.schedule(60, "checkStatus", { check: "health" });

		// Schedule with a cron expression
		await this.schedule("0 0 * * *", "dailyTask", { type: "cleanup" });
	}

	async sendGreeting(payload: { message: string }) {
		console.log(payload.message);
	}

	async checkStatus(payload: { check: string }) {
		console.log("Running check:", payload.check);
	}

	async dailyTask(payload: { type: string }) {
		console.log("Daily task:", payload.type);
	}
}

调度存储在 cf_agents_schedules SQL 表中。Cron 调度执行后会自动重新调度,一次性调度会被删除。

this.mcp 及相关方法

Agent 包含多服务器 MCP 客户端,使 Agent 可与暴露 MCP 接口的外部服务交互。MCP 客户端详见 MCP 客户端 API

class MyAgent extends Agent {
	async onStart() {
		// Add an HTTP MCP server (callbackHost only needed for OAuth servers)
		await this.addMcpServer("GitHub", "https://mcp.github.com/mcp", {
			callbackHost: "https://my-worker.example.workers.dev",
		});

		// Add an MCP server via RPC (Durable Object binding, no HTTP overhead)
		await this.addMcpServer("internal-tools", this.env.MyMCP);
	}
}

邮件处理

Agent 可使用 Cloudflare 的 Email Routing 接收并回复邮件。

class MyAgent extends Agent {
	async onEmail(email: AgentEmail) {
		console.log("Received email from:", email.from);
		console.log("Subject:", email.headers.get("subject"));

		const raw = await email.getRaw();
		console.log("Raw email size:", raw.length);

		// Reply to the email
		await this.replyToEmail(email, {
			fromName: "My Agent",
			subject: "Re: " + email.headers.get("subject"),
			body: "Thanks for your email!",
			contentType: "text/plain",
		});
	}
}

要将邮件路由到 Agent,在 Worker 的邮件处理程序中使用 routeAgentEmail

export default {
	async email(message, env, ctx) {
		await routeAgentEmail(message, env, {
			resolver: createAddressBasedEmailResolver("my-agent"),
		});
	},
} satisfies ExportedHandler<Env>;

上下文管理

agentsAsyncLocalStorage 包装所有方法,以在请求生命周期内保持上下文。这样可从代码任意位置访问当前 agent、连接、请求或邮件(取决于正在处理的事件):

import { getCurrentAgent } from "agents";

function someUtilityFunction() {
	const { agent, connection, request, email } = getCurrentAgent();

	if (agent) {
		console.log("Current agent:", agent.name);
	}

	if (connection) {
		console.log("WebSocket connection ID:", connection.id);
	}
}

this.onError

Agent 继承 ServeronError,可用于处理不一定是 WebSocket 错误的错误。它以 Connectionunknown 错误为参数调用。

class MyAgent extends Agent {
	onError(connectionOrError: Connection | unknown, error?: unknown) {
		if (error) {
			// WebSocket connection error
			console.error("Connection error:", error);
		} else {
			// Server error
			console.error("Server error:", connectionOrError);
		}

		// Optionally throw to propagate the error
		throw connectionOrError;
	}
}

this.destroy

this.destroy() 会删除所有表、清除闹钟、清空存储并中止上下文。为确保 Durable Object 被完全驱逐,使用 setTimeout() 异步调用 this.ctx.abort(),以便当前正在执行的处理程序(如已调度任务)在上下文中止前完成清理。

这意味着 this.ctx.abort() 会抛出无法在代码中捕获的错误并出现在日志中,但会在让出事件循环后抛出(详见 abort())。

destroy() 方法可在已调度任务内安全调用。从调度回调内调用时,Agent 会设置内部标志以跳过剩余数据库更新,并将 ctx.abort() 让出给事件循环,确保闹钟处理程序在 Agent 驱逐前干净完成。

class MyAgent extends Agent {
	async onStart() {
		console.log("Agent is starting up...");
		// Initialize your agent
	}

	async cleanup() {
		// This wipes everything!
		await this.destroy();
	}

	async selfDestruct() {
		// Safe to call from within a scheduled task
		await this.schedule(60, "destroyAfterDelay", {});
	}

	async destroyAfterDelay() {
		// This will safely destroy the Agent even when
		// called from within the alarm handler
		await this.destroy();
	}
}

static options

通过在类上重写 static options 配置 agent 行为。所有字段均为可选——默认值在运行时应用。

export class MyAgent extends Agent {
	static options = {
		hibernate: true,
		sendIdentityOnConnect: false,
		retry: { maxAttempts: 5, baseDelayMs: 200, maxDelayMs: 5000 },
	};
}
选项 类型 默认值 描述
hibernate boolean true agent 不活跃时是否休眠。Durable Object 睡眠时 WebSocket 连接保持打开
sendIdentityOnConnect boolean true WebSocket 连接时是否向客户端发送身份(agent 名称、实例名称)。设为 false 可隐藏敏感实例名
hungScheduleTimeoutSeconds number 30 运行中的间隔调度被视为挂起并强制重置前的超时。长时间运行的回调应增大此值
keepAliveIntervalMs number 30000 keepAlive() 闹钟心跳间隔(毫秒)。值越低恢复越快,但闹钟越频繁
retry RetryOptions { maxAttempts: 3, baseDelayMs: 100, maxDelayMs: 3000 } schedule()queue()this.retry() 的默认重试选项。每任务选项会覆盖这些默认值

this.keepAlive()this.keepAliveWhile()

Durable Objects 在不活动一段时间后会被驱逐(通常 70–140 秒内无传入请求、WebSocket 消息或闹钟)。在长时间运行操作期间——流式 LLM 响应、等待外部 API、运行多步计算——agent 可能在执行中途被驱逐。

keepAlive() 创建闹钟心跳以防止驱逐。keepAliveWhile() 包装异步函数并保证清理。

class MyAgent extends Agent {
	async handleLongTask() {
		// Option 1: manual dispose
		const dispose = await this.keepAlive();
		try {
			await longRunningComputation();
		} finally {
			dispose();
		}

		// Option 2: automatic cleanup (recommended)
		const result = await this.keepAliveWhile(async () => {
			return await longRunningComputation();
		});
	}
}

AIChatAgent 在流式 LLM 响应期间内部使用 keepAliveWhile。更多细节请参阅调度任务 — 保持 Agent 存活

路由

Agent 类将寻址辅助函数 重新导出为 getAgentByNamerouteAgentRequest

const stub = await getAgentByName(env.MY_DO, "foo");
await stub.someMethod();

const res = await routeAgentRequest(request, env);
if (res) return res;

return new Response("Not found", { status: 404 });

第 3 层:AIChatAgent

AIChatAgent 类来自 @cloudflare/ai-chat,以有明确主张的层继承 Agent 用于 AI 聊天。它增加自动消息持久化到 SQLite、可恢复流式传输、工具支持(服务端、客户端与人机协同)以及用于构建聊天 UI 的 React 钩子(useAgentChat)。

完整层次为:DurableObject > Server > Agent > AIChatAgent

若构建聊天 agent,请从 AIChatAgent 开始。若需要更低层控制或不构建聊天界面,请直接使用 Agent

这篇文档对您有帮助吗?