十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

深度拆解 Cloudflare agents 框架的 Agent 类:从 Durable Object 到全功能 AI Agent 基类

深度拆解 Cloudflare agents 框架的 Agent 类:从 Durable Object 到全功能 AI Agent 基类 深度拆解 Cloudflare agents 框架的 Agent 类从 Durable Object 到全功能 AI Agent 基类【免费下载链接】agentsBuild and deploy AI Agents on Cloudflare项目地址: https://gitcode.com/GitHub_Trending/agents1/agentsagents库的核心是导出的Agent类它直接继承 Cloudflare 的DurableObject以继承即获得全部能力的方式为开发者提供有状态、可调度、可观测的 Agent 原语。本文以官方文档《Demystifying the Agent class》为骨架结合本仓库源码逐层拆解Agent的三层结构Durable Object → Lifecycle → Agent讲透状态持久化、callableRPC、任务队列、调度、MCP 客户端、邮件收发、Fiber 可恢复执行等每个内置特性的底层原理与实战写法帮助你在开始写自己的 Agent 之前建立完整心智模型避开常见陷阱。什么是 AgentAgent直接继承 Cloudflare 的DurableObject因此每一个 Agent 实例都是一个全局可寻址、单线程执行的分布式计算单元自带持久化的 KV / SQLite 存储。如果你对 Durable Object 还不熟悉建议先阅读 What are Durable Objects 了解这一平台原语。每个 Agent 内部组合compose了一个Lifecycle实例。Lifecycle 负责安装面向平台的请求request、告警alarm和休眠 WebSockethibernating WebSocket入口点而 Agent 本身向开发者提供语义化的回调与高层特性DurableObject └── Agent └── owns Lifecycle从源码看Agent 类的定义确实以extends DurableObjectEnv开始并在字段初始化时通过Lifecycle.installEnv, Props(this, ...)构造 Lifecycle再以new WebSockets({ handlers: {...} })挂载 WebSocket 子系统见 index.ts。理解这个三层结构是理解整个框架的关键。Layer 0Durable Object 平台原语这一层不展开讲 Durable Object 的全部细节但要清楚它暴露了哪些原语以及上两层如何利用它们。constructor签名固定必须经由 Namespace 创建constructor(ctx: DurableObjectState, env: Env) {}Workers 运行时每次初始化 DO 时都会调用构造函数处理内部事务这带来两点约束构造函数每次 DO 初始化都会被调用但签名是固定的开发者不能通过构造函数增加或修改参数开发者不能手动new这个类必须通过绑定 API 经由 DurableObjectNamespace 创建实例。RPC公共方法即远程方法只要 DO 类继承自内置类型DurableObject其公共方法就会被暴露为 RPC 方法开发者可以从 Worker 中通过 DurableObjectStub 调用// 该实例可能是活跃的、休眠的、未初始化的甚至从未被创建过 const stub env.MY_DO.getByName(foo); // 我们可以调用类的任意公共方法运行时**保证** // 如果实例未激活会先为我们调用构造函数。 await stub.bar();fetch()唯一的外部 HTTP 入口Durable Object 可以从 Worker 接收Request并返回Response。这只能通过fetch方法完成必须由开发者实现。WebSockets一流支持与回调式事件Durable Object 对 WebSockets 提供一等支持DO 可以在fetch中接受从Request传来的 WebSocket 然后忘掉它hibernation。基类提供可覆写的方法作为回调有效替代事件监听器webSocketMessage(ws, message)、webSocketClose(ws, code, reason, wasClean)和webSocketError(ws, error)。export class MyDurableObject extends DurableObject { async fetch(request) { // 创建 WebSocket 连接的两端 const webSocketPair new WebSocketPair(); const [client, server] Object.values(webSocketPair); // 调用 acceptWebSocket() 将 WebSocket 接入 Durable Object // 使其可以收发消息 this.ctx.acceptWebSocket(server); return new Response(null, { status: 101, webSocket: client }); } async webSocketMessage(ws, message) { // 原样回显消息 ws.send(msg); } }Alarms归 Lifecycle 统一所有Lifecycle拥有 Agent 物理 Durable Object 的 alarm因为调度schedule、keep-alive、Fiber、子 Agent 等能力共享同一个 alarm 槽位。因此不要在 Agent 特性中覆写alarm()也不要调用this.ctx.storage.setAlarm()——一个调用者可能覆盖另一个特性的唤醒。应使用this.schedule()做具名的 Agent 回调可复用的、自带持久化工作的能力则通过this.lifecycle.jobs推送任务并实现onJob()由 Lifecycle 驱动到期任务并根据队列状态重新武装 alarm。关于 job queue 的机制参见 Durable Object lifecycle 文档与调度文档。Lifecycle 拥有队列与单一物理 alarm队列按时间戳排序任何队列变更都会自动重新武装 alarm无需显式 rearm 调用alarm 触发时 Lifecycle 先以事件循环驱动到期任务、再运行宿主的onAlarm()、最后根据队列状态重新武装见 lifecycle.md。this.ctx状态与环境的统一入口基类DurableObject将 DurableObjectState 放进this.ctx。其中最重要的是this.ctx.storage。this.ctx.storageDurableObjectStorage 是与 DO 持久化机制交互的主接口同时提供 KV 和 SQLite 两种同步APIconst sql this.ctx.storage.sql; const kv this.ctx.storage.kv; // 同步 SQL 查询示例 const rows sql.exec(SELECT * FROM contacts WHERE country ?, US); // 同步 KV 示例 const token kv.get(someToken);this.ctx.envDO 还能通过this.env访问 Worker 的Env绑定bindings 文档。Layer 1Lifecycle 组合层Agent使用Lifecycle.install(this)构造 lifecycle 并安装面向平台的fetch、alarm、webSocketMessage、webSocketClose、webSocketError处理器。Agent 子类应实现语义回调而非第二套基类方法class MyAgent extends Agent { onStart() { // 每次内存生命周期开始、处理工作之前运行一次 } onRequest(request: Request) { return new Response(Hello from ${request.url}); } onConnect(connection: Connection) { connection.send(connected); } }在 index.ts 中可以看到这些回调的默认实现onStart()、onRequest()默认返回 404、onConnect()、onMessage()、onClose()均为空操作或兜底实现开发者按需覆写即可。Lifecycle 的 WebSocket 始终使用 Cloudflare 的Hibernation API空闲客户端保持连接而 Durable Object 可以离开内存当消息唤醒它时构造函数字段和onStart会再次运行。任何需要跨唤醒保留的数据必须持久化到 storage 或connection.state中。可复用的能力通过this.lifecycle.use(capability)安装能力在 Agent 启动之前启动可以在onRequest之前拦截请求也可以在onAlarm之前处理告警。从 lifecycle.md 看能力按注册顺序作为中间件运行第一个返回Response的能力处理该请求返回undefined则继续传递给下一个能力无人认领的请求最终落到宿主的onRequest。claims默认selective而WebSockets能力是catch-all它认领所有 upgrade。请求调用路径routeAgentRequest(request) └─ named Durable Object stub.fetch(request) └─ lifecycle-installed fetch ├─ lifecycle startup capabilities ├─ host onStart ├─ capability middleware, in registration order │ └─ first Response handles the request └─ host onRequest已唤醒的对象会跳过 startup但仍会把每个请求交给其中间件。Identity命名身份自 2026-03-15 起Workers 将idFromName()或getByName()使用的名称暴露为ctx.id.name包括在 alarm 处理器中。Agent 与 Agent facet 使用命名 IDthis.name即映射这一原生身份。迁移方面lifecycle 可以读取旧版本写入的__ps_name值但绝不会写入重复名称。原始 ID、idFromString()以及超过 1024 字节的名称不提供原生身份2026-03-15 之前创建的 alarm 必须从具名的 fetch 或 RPC 处理器重新调度详见 lifecycle.md。Layer 2Agent——有状态、可调度、可观测Agent直接继承DurableObject组合了上述 lifecycle并提供面向有状态、可调度、可观测 Agent 的强约定原语支持通过 RPC、WebSocket、甚至电子邮件进行通信。this.state 与 this.setState()自动状态持久化Agent的核心特性之一是自动状态持久化。开发者通过泛型参数和initialState仅在存储中不存在状态时使用定义状态形状Agent 负责加载、保存以及通过 lifecycle 管理的 WebSocket 连接广播状态变更。this.state是一个 getter惰性地从存储SQL加载状态。当通过this.setState()更新状态时状态会跨 DO 驱逐持久化——setState()自动序列化状态并写回存储。可以覆写this.onStateChanged来响应状态变更。class MyAgent extends AgentEnv, { count: number } { initialState { count: 0 }; increment() { this.setState({ count: this.state.count 1 }); } onStateChanged(state, source) { console.log(State updated:, state); } }从源码看state getter 首先检查内存缓存_state未命中时执行SELECT state FROM cf_agents_state WHERE id ${STATE_ROW_ID}以行是否存在作为状态是否曾写入的信号这正确处理了null、0、false、等 falsy 值若 JSON 解析失败则回退到initialState若连初始状态都没有则删除损坏数据行以防无限重试。状态存储在cf_agents_stateSQL 表中状态消息以type: cf_agent_state发送客户端和服务端均如此。由于agents提供了 JS 和 React 客户端实时状态同步开箱即用。协议消息CF_AGENT_IDENTITY、CF_AGENT_STATE、CF_AGENT_MCP_SERVERS会在连接时自动发送并在变更时广播。可以在每个连接上覆写shouldSendProtocolMessages(connection, ctx)来抑制这些消息详见 Protocol Message Control。this.sql同步 SQL 模板标签Agent 提供了便捷的sql模板标签来对 DO 的 SQLite 存储执行查询。它构造参数化查询并执行底层就是this.ctx.storage.sql的同步SQL API。从 sql 方法实现 看它把模板片段与?占位符拼接后交给ctx.storage.sql.exec异常时包装为SqlError抛出。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.sqlINSERT 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 与可调用方法callableagents把 Durable Object 的 RPC 又推进了一步通过 WebSocket 实现 RPC因此客户端也能直接调用 Agent 的方法。要让一个方法可以通过 WebSocket 调用使用callable装饰器。方法可以返回可序列化的值也可以使用callable({ streaming: 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在 index.ts 中可以找到 RPC 请求的类型守卫isRPCRequest消息必须是{ type: rpc, id: string, method: string, args: unknown[] }形状。装饰器定义于 callable-decorator.tscallablesFromDecorated(this)负责收集被装饰方法供 WebSockets 能力使用见 index.ts。callable方法同样是 Agent 在 Capn Web RPC 端点上的远程接口——一套接口多种传输详见 lifecycle.md。this.queue 与相关方法内置任务队列Agent 内置任务队列用于延迟执行适合卸载工作或重试操作。可用方法this.queue、this.dequeue、this.dequeueAll、this.dequeueAllByCallback、this.getQueue、this.getQueues。class MyAgent extends Agent { async onConnect() { // 排队一个稍后执行的任务 await this.queue(processTask, { userId: 123 }); } async processTask(payload: { userId: string }, queueItem: QueueItem) { console.log(Processing task for user:, payload.userId); } }任务存储在cf_agents_queuesSQL 表中按创建时间升序FIFO自动刷新执行。任务成功后自动出队。源码细节queue()index.ts用nanoid(9)生成 ID校验回调必须是字符串且必须是 Agent 实例上的函数可选retry选项默认maxAttempts: 3, baseDelayMs: 100, maxDelayMs: 3000随后INSERT OR REPLACE INTO cf_agents_queues并触发_flushQueue()异步排空循环_flushQueue()index.ts用tryN做指数退避重试失败后调用this.onError(e)并在finally中dequeue(row.id)。dequeueAllByCallback(callback)按回调名批量删除getQueues(key, value)则对 payload 反序列化后按键值过滤。this.schedule 与相关方法调度执行Agent 通过包装 DO 的alarm()支持方法的定时执行。可用方法this.schedule、this.getScheduleById、this.listSchedules、this.cancelSchedule以及已弃用的同步版本this.getSchedule/this.getSchedules。调度可以是一次性的、延迟的也可以是循环的使用 cron 表达式。由于 DO 同一时间只允许一个 alarmAgent通过在 SQL 中管理多个调度并复用一个 alarm 来解决这一限制。从源码看调度能力由独立的 Scheduler 能力 提供Agent在组合根处安装readonly scheduler: Scheduler见 index.tsthis.schedule()等稳定方法直接委托给this.scheduler详见 scheduling.md。class MyAgent extends Agent { async foo() { // 指定时间调度 await this.schedule(new Date(2025-12-25T00:00:00Z), sendGreeting, { message: Merry Christmas! }); // 延迟调度单位秒 await this.schedule(60, checkStatus, { check: health }); // cron 表达式调度 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); } }调度以 job 形式存储在cf_agents_jobsSQL 表Lifecycle 拥有的 job 队列中。Cron 调度执行后会自动重新调度自身一次性调度执行后会被删除。调度系统支持四种模式详见 scheduling.md模式语法使用场景延迟this.schedule(60, ...)60 秒后运行定时this.schedule(new Date(...), ...)在指定时间运行Cronthis.schedule(0 8 * * *, ...)循环调度间隔this.scheduleEvery(30, ...)每 30 秒运行一次scheduleEvery()在callback, interval, payload组合上幂等重复调用不会创建重复调度因此可以在onStart()中安全调用onStart每次 DO 唤醒都会执行它还内置重叠预防——回调运行时间超过间隔则跳过下一次执行不是排队并记录警告。this.mcp 与相关方法多服务器 MCP 客户端Agent内置多服务器 MCPModel Context Protocol客户端让 Agent 能够与暴露 MCP 接口的外部服务交互。MCP 客户端的完整文档参见 MCP Client API。class MyAgent extends Agent { async onConnect() { // 添加一个 MCP 服务器 await this.addMcpServer(GitHub, https://mcp.example.com/sse); } }从源码看Agent 以readonly mcp: MCPClientManager暴露该能力index.ts其实现位于 mcp/client支持sse、streamable-http、auto等传输类型并可选 OAuth 客户端提供者DurableObjectOAuthClientProvider见 do-oauth-client-provider.ts。MCP 服务器状态通过CF_AGENT_MCP_SERVERS协议消息同步给客户端。邮件处理发送与接收 EmailAgent 可以使用 Cloudflare 的 Email Service 收发邮件。使用this.sendEmail()配合send_email绑定发送外发邮件class MyAgent extends Agent { callable() async sendWelcome(to: string) { return this.sendEmail({ binding: this.env.EMAIL, to, from: supportyourdomain.com, subject: Welcome!, text: Thanks for signing up. }); } async onEmail(email: AgentEmail) { console.log(Received email from:, email.from); console.log(Subject:, email.headers.get(subject)); await this.replyToEmail(email, { fromName: My Agent, body: Thanks for your email! }); } }将邮件路由到 Agent需要在 Worker 的 email 处理器中使用routeAgentEmailimport { routeAgentEmail } from agents; import { createAddressBasedEmailResolver } from agents/email; export default { async email(message, env, ctx) { await routeAgentEmail(message, env, { resolver: createAddressBasedEmailResolver(my-agent) }); } };关于sendEmail()、入站邮件路由、resolver 以及安全回复流程的完整细节参见 Email Service 指南。补充几个源码可确认的要点sendEmail()会为每条消息自动注入路由头X-Agent-Name、X-Agent-ID并可选用secret做 HMAC-SHA256 签名以便回信路由回同一 Agent 实例wrangler.jsonc中的绑定形如{ send_email: [{ name: EMAIL, remote: true }] }remote: true允许本地开发时调用真实 Email Service API。上下文管理getCurrentAgent()agents用AsyncLocalStorage包裹所有方法在整个请求生命周期内维护上下文让你可以从代码任意位置访问当前 agent、connection、request 或 email取决于正在处理的事件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); } }实现上Agent 通过runInInvocation()index.ts在agentContext.run(store, ...)中建立上下文并叠加withInvocationScope划定调用边界防止 invocation 内的 span 超期存活内部存储位于 internal_context.ts。this.onError统一错误处理Agent.onError同时处理 WebSocket 错误和 Agent 内部的其他错误参数可能是Connection或unknown错误class MyAgent extends Agent { onError(connectionOrError: Connection | unknown, error?: unknown) { if (error) { // WebSocket 连接错误 console.error(Connection error:, error); } else { // 服务器错误 console.error(Server error:, connectionOrError); } // 可选抛出以传播错误 throw connectionOrError; } }this.destroy彻底销毁 Agentthis.destroy()会删除所有表、删除 alarm、清空存储并中止上下文。为确保 DO 被完全驱逐会调用this.ctx.abort()抛出不可捕获的错误会出现在日志中详见 abort 文档。class MyAgent extends Agent { async onStart() { console.log(Agent is starting up...); // 初始化你的 agent } async cleanup() { // 这会抹掉一切 await this.destroy(); } }从源码看destroy()还配套了延迟销毁机制_cf_scheduleDestroy先写入cf_agents_destroy_pending标记再推入cf:destroy主机 job由带完整执行预算的 alarm 调用destroy()从而避免 HTTP 处理器中执行多步 I/O 序列时与响应竞态index.ts 及上方常量注释。this.keepAlive 与 this.keepAliveWhile防止空闲驱逐this.keepAlive()通过持有 alarm 支撑的心跳引用heartbeat ref防止 Durable Object 因不活动而被驱逐。返回一个 disposer 函数用于停止心跳。对于有明确作用域的工作使用this.keepAliveWhile(fn)——函数完成或抛错时自动清理。完整文档见 Keeping the Agent Alive。const dispose await this.keepAlive(); try { // 不得被中断的长时工作 const result await longRunningComputation(); await sendResults(result); } finally { dispose(); }实现要点见 scheduling.mdkeepAlive()使用内存引用计数每次调用 1、disposer 调用 -1计数大于 0 时 Agent 每 30 秒向 Lifecycle 贡献一次唤醒时间。不创建调度行、不发送可观测事件因此对listSchedules()与调度诊断通道不可见。AIChatAgent在流式响应期间自动调用keepAlive()无需手动添加。this.runFiber 与 this.startFiber可恢复执行与受管 Fiberthis.runFiber()运行可检查点checkpointable、带崩溃恢复的工作this.startFiber()则持久化地接受后台工作支持幂等、状态检查、取消以及保留终态记录。const receipt await this.startFiber( process-webhook, async (ctx) { ctx.stash({ webhookId }); await processWebhook(webhookId, { signal: ctx.signal }); }, { idempotencyKey: webhook:${webhookId}, waitForCompletion: true } ); const current await this.inspectFiber(receipt.fiberId); await this.cancelFiber(receipt.fiberId, Superseded); if (current?.status interrupted) { await this.resolveFiber(receipt.fiberId, { status: completed }); } await this.deleteFibers({ status: [completed, error, aborted], settledBefore: new Date(Date.now() - 7 * 24 * 60 * 60 * 1000) });选择指南调用方等待结果时用runFiber()调用方需要立即获得持久化回执、可能用相同幂等键重试时用startFiber()调用方应等待已接受任务达到终态、同时仍要消除重试重复时加上waitForCompletion: true用onFiberRecovered()决定被中断的受管 Fiber对你的应用意味着什么返回一个恢复结果即可更新保留的状态记录resolveFiber()仅用于外部解决interrupted行。Fiber 的状态机为pending | running | completed | aborted | interrupted | errorindex.ts受管 Fiber 的账本保存在cf_agents_fibers表含idempotency_key UNIQUE索引。startFiber()支持fiberId、idempotencyKey、metadata、waitForCompletion等选项StartFiberOptions。路由getAgentByName 与 routeAgentRequest使用getAgentByName获取具名 RPC stub使用routeAgentRequest处理/agents/:class/:name的 HTTP 与 WebSocket 路由const stub await getAgentByName(env.MY_DO, foo); // ... const res await routeAgentRequest(request, env); if (res) return res; return Response(Not found, { status: 404 });路由实现位于 agent-routing.ts其中routeAgentRequest支持路由重试RoutingRetryOptions例如对内存限制重置与代码更新导致的重置进行安全重试isDurableObjectMemoryLimitReset等判定函数见 retries.ts。若要在纯 Lifecycle Object 上使用同样的 URL 路由直接让 Worker 的fetch调用routeAgentRequest即可默认 URL 形态为/agents/:binding/:name见 lifecycle.md。附Agent 静态配置选项Agent子类可通过static options覆写框架级配置源码在 DEFAULT_AGENT_STATIC_OPTIONS 中集中定义了默认值与语义选项默认值说明sendIdentityOnConnecttrue连接时是否向客户端发送身份name、agenthungScheduleTimeoutSeconds30运行中的间隔调度被判定挂起并强制重置的超时秒keepAliveIntervalMs30000keepAlive()alarm 心跳间隔毫秒越小恢复越快但 alarm 越频繁retry{ maxAttempts: 3, baseDelayMs: 100, maxDelayMs: 3000 }schedule()、queue()、this.retry()的默认重试fiberRecoveryHookTimeoutMs10000内部 Fiber 恢复钩子超时fiberRecoveryScanDeadlineMs10000单次中断 Fiber 恢复扫描的软截止fiberRecoveryMaxAgeMs24h未受管中断 Fiber 行在放弃恢复前的最长年龄agentToolReattachNoProgressTimeoutMs120000重新挂接 agent-tool 子运行时无进展预算按进展重置agentToolReattachMaxWindowMsInfinity单次重新挂接的硬墙钟上限detachedMaxBudgetMs24h分离式后台agent-tool 运行的绝对预算上限detachedNoProgressBudgetMs1h分离式运行无进展窗口收到信号后重置maxAlarmMemoryLimitStrikes3alarm 边界内内存限制重置的连续次数上限防无限重试回路覆写方式class SecureAgent extends Agent { static options { sendIdentityOnConnect: false, keepAliveIntervalMs: 15_000, retry: { maxAttempts: 5, baseDelayMs: 200, maxDelayMs: 10_000 } }; }从 index.ts 的实现看_resolvedOptions会合并默认值与子类覆写并缓存静态 options 在 DO 实例生命周期内不变。在仓库中继续深入Agent 类完整源码本文所有特性的实现所在约一万行的核心文件Agent 类定义处与状态 getter 实现Durable Object lifecycleLifecycle 组合、能力中间件、job queue、WebSockets 能力与原生 RPC调度指南四种调度模式、scheduleEvery()幂等、keep-alive、AI 辅助调度Email Service 指南绑定配置、出站发送、入站路由与安全回复Agent 路由实现getAgentByName与routeAgentRequest对应测试状态持久化 state.ts、队列 queue.ts、调度 schedule.ts、Fiber run-fiber.ts、keep-alive keep-alive.ts、邮件 email.ts、可调用方法 callable.ts。小结把Agent理解为带电池的 Durable Object并不夸张底层的生命周期管理、alarm 协调、休眠 WebSocket、SQLite 状态与队列、调度、Fiber 可恢复执行等重活全部由框架托管开发者只需继承Agent并实现onStart、onRequest、onConnect、onMessage等语义回调再按需使用callable、this.state、this.queue、this.schedule、this.mcp、sendEmail等高层能力。理解三层结构与一个 alarm 槽位被 Lifecycle 统一协调这一核心约束是写出健壮、可长期运行的 Agent 应用的前提。【免费下载链接】agentsBuild and deploy AI Agents on Cloudflare项目地址: https://gitcode.com/GitHub_Trending/agents1/agents创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表