运行可在 Durable Object 驱逐后存活的工作。runFiber() 在 SQLite 中注册任务,在执行期间保持 agent 存活,允许你使用 stash() 检查点中间状态,并在 agent 在任务中途被驱逐时在下次激活时调用 onFiberRecovered()。
当调用方需要持久接受后台工作、快速返回、安全去重重试、稍后检查状态或取消正在运行的任务时,请使用 startFiber()。
import { Agent } from "agents";
import type { FiberRecoveryContext } from "agents";
class MyAgent extends Agent {
async doWork() {
await this.runFiber("my-task", async (ctx) => {
const step1 = await expensiveOperation();
ctx.stash({ step1 });
const step2 = await anotherExpensiveOperation(step1);
this.setState({ ...this.state, result: step2 });
});
}
async onFiberRecovered(ctx: FiberRecoveryContext) {
if (ctx.name !== "my-task") return;
const snapshot = ctx.snapshot as { step1: unknown } | null;
if (snapshot) {
const step2 = await anotherExpensiveOperation(snapshot.step1);
this.setState({ ...this.state, result: step2 });
}
}
}Durable Objects 因以下三个原因被驱逐:
- 不活动超时 — 约 70–140 秒内无传入请求或打开的 WebSocket
- 代码更新 / 运行时重启 — 非确定性,每天 1–2 次
- Alarm 处理程序超时 — 15 分钟
驱逐发生在工作进行中时,上游 HTTP 连接(到 LLM 提供商、API、数据库)会永久断开。内存状态——流式缓冲区、部分响应、循环计数器——会丢失。多轮 agent 循环会完全失去位置。
keepAlive() 降低驱逐概率。runFiber() 使驱逐后仍可恢复。
对于应独立于 agent 运行、具有每步重试和多步编排的工作,请改用 Workflows。Fiber 适用于 agent 自身执行的一部分。比较请参阅长时间运行 agent:Workflows 与 agent 内部模式。
通过创建 30 秒 alarm 心跳重置不活动计时器,防止空闲驱逐。
class Agent {
keepAlive(): Promise<() => void>;
keepAliveWhile<T>(fn: () => Promise<T>): Promise<T>;
}keepAliveWhile() 是推荐方式——它运行异步函数,并在完成或抛出异常时自动清理心跳:
const result = await this.keepAliveWhile(async () => {
return await slowAPICall();
});如需手动控制,keepAlive() 返回 disposer。完成后务必调用——否则心跳会无限继续:
const dispose = await this.keepAlive();
try {
await longWork();
} finally {
dispose();
}只要持有任何 keepAlive 引用,alarm 每 30 秒触发一次以重置不活动计时器。所有 disposer 被调用后,alarm 停止,DO 可自然进入空闲。
心跳对 listSchedules() 不可见——不会创建 schedule 行。它不会与你自己的 schedule 冲突;alarm 系统通过单个 alarm 槽复用所有 schedule 和 keepAlive 心跳。
默认:30 秒。不活动超时约 70–140 秒,因此 30 秒提供充足余量。通过静态选项覆盖:
class MyAgent extends Agent {
static options = { keepAliveIntervalMs: 2_000 };
}keepAlive 防止驱逐但不处理恢复。如果 agent 仍 被驱逐(代码更新、alarm 超时、资源限制),任何进行中的工作都会丢失。
runFiber 内部调用 keepAlive 并 将工作持久化到 SQLite 以便恢复。当工作重做成本低或不需要检查点时,单独使用 keepAlive。当工作成本高且需要从断点恢复时,使用 runFiber。
| 场景 | 使用 |
|---|---|
| 等待慢速 API 调用 | keepAlive() |
流式传输 LLM 响应(通过 AIChatAgent) |
自动(内置) |
| 带中间结果的多步计算 | runFiber() |
| 耗时 10 分钟以上的后台研究循环 | 带 stash() 的 runFiber() |
| 必须恰好接受一次的 webhook 任务 | startFiber() |
带检查点和恢复的持久执行。
class Agent {
runFiber<T>(name: string, fn: (ctx: FiberContext) => Promise<T>): Promise<T>;
startFiber(
name: string,
fn: (ctx: FiberContext) => Promise<void>,
options?: StartFiberOptions,
): Promise<StartFiberResult>;
inspectFiber(fiberId: string): Promise<FiberInspection | null>;
inspectFiberByKey(idempotencyKey: string): Promise<FiberInspection | null>;
listFibers(options?: ListFibersOptions): Promise<FiberInspection[]>;
cancelFiber(fiberId: string, reason?: string): Promise<boolean>;
cancelFiberByKey(idempotencyKey: string, reason?: string): Promise<boolean>;
deleteFibers(options?: DeleteFibersOptions): Promise<number>;
resolveFiber(fiberId: string, result: FiberRecoveryResult): Promise<boolean>;
stash(data: unknown): void;
onFiberRecovered(
ctx: FiberRecoveryContext,
): Promise<void | FiberRecoveryResult>;
}
type FiberContext = {
id: string;
signal: AbortSignal;
stash(data: unknown): void;
snapshot: unknown | null;
};
type FiberStatus =
| "pending"
| "running"
| "completed"
| "aborted"
| "interrupted"
| "error";
type FiberRecoveryContext = {
id: string;
name: string;
status?: FiberStatus;
idempotencyKey?: string;
metadata?: Record<string, unknown> | null;
snapshot: unknown | null;
createdAt: number;
recoveryReason: "interrupted";
};runFiber("work", fn)
├─ Persist recovery metadata
├─ keepAlive() — heartbeat starts
├─ Execute fn(ctx)
│ ├─ ctx.stash(data) → persist snapshot
│ ├─ ctx.stash(data) → persist snapshot
│ └─ return result
├─ Delete recovery metadata
├─ keepAlive dispose — heartbeat stops
└─ Return result to caller[DO evicted — all in-memory state lost]
On next activation:
├─ Request/connection → onStart() → check for orphaned fibers [primary path]
│ OR
├─ Persisted alarm fires → housekeeping check [fallback path]
Recovery:
├─ Load interrupted fibers from storage
├─ For each interrupted fiber:
│ ├─ Parse snapshot from JSON
│ ├─ Call onFiberRecovered(ctx)
│ └─ Delete recovery metadata after successful recovery
└─ If onFiberRecovered calls runFiber() again → new fiber, normal execution两条恢复路径调用同一钩子。alarm 路径对无传入客户端连接的后台 agent 至关重要——持久化 alarm 可独立唤醒 agent。
Fiber 也可在子 agent 内工作。fiber 行和快照存储在子 agent 自己的 SQLite 数据库中,onFiberRecovered() 以子 agent 作为 this 运行。
子 agent 没有独立的 alarm 槽,因此顶层父级拥有物理心跳。当子 agent 启动 fiber 时,父级跟踪足够元数据以将恢复检查路由回所属子 agent,即使子级没有客户端连接或传入 RPC。
这使恢复保持在子级本地,同时保留父级拥有的单个物理 alarm 槽。恢复的延续可在 facet 内使用 schedule();父级拥有物理 alarm 并将回调路由回子级。
fn(ctx) throws Error
├─ DELETE row from cf_agents_runs
├─ keepAlive dispose
└─ Error propagates to caller (or logged if fire-and-forget)无自动重试。恢复逻辑属于 onFiberRecovered,你可在此获得快照及出错的完整上下文。
runFiber() 支持两种模式:
// Inline — await the result
const result = await this.runFiber("work", async (ctx) => {
return computeExpensiveThing();
});
// Fire-and-forget — caller does not wait
void this.runFiber("background", async (ctx) => {
await longRunningProcess();
});如果在内联 await 期间 DO 被驱逐,调用方已消失。恢复时 onFiberRecovered 触发——无法将结果返回给原始调用方。这是跨进程边界持久执行的固有局限。对于可能超过单个 DO 生命周期的长时间运行工作,当调用方需要保留状态记录、幂等接受或取消时,请使用 startFiber()。
当调用方需要持久接受后台工作、快速返回并安全去重重试时,请使用 startFiber()。它在回调运行前存储保留的 fiber 记录,然后使用与 runFiber() 相同的 keep-alive 和恢复机制在后台启动回调。
const receipt = await this.startFiber(
"reply-to-webhook",
async (ctx) => {
ctx.stash({ webhookId, threadId });
await postReply(threadId);
},
{
idempotencyKey: `webhook:${webhookId}`,
metadata: { threadId },
},
);
if (!receipt.accepted) {
// This webhook was already accepted by an earlier delivery.
}默认情况下,startFiber() 在工作被持久接受后返回。当调用方应保持打开直到接受的 fiber 达到终端状态时,传递 waitForCompletion: true。具有相同幂等键的重复调用在可能时加入活跃内存执行,然后返回 accepted: false 的保留状态。
const result = await this.startFiber("reply-to-webhook", reply, {
idempotencyKey: `webhook:${webhookId}`,
waitForCompletion: true,
});
if (result.status === "error") {
console.error(result.error);
}startFiber() 是持久接受 API,而非返回值 API。它返回托管 fiber 状态,但不返回回调结果。稍后使用 inspectFiber() 或 inspectFiberByKey() 检查状态。
const current = await this.inspectFiberByKey(`webhook:${webhookId}`);
if (current) {
await this.cancelFiber(current.fiberId, "No longer needed");
}
await this.deleteFibers({
status: ["completed", "error", "aborted"],
settledBefore: new Date(Date.now() - 7 * 24 * 60 * 60 * 1000),
});默认情况下,deleteFibers() 删除已结算的 completed、error 和 aborted 行。除非你显式传递该状态,否则不会删除 interrupted 行,因为 interrupted 行通常需要检查或手动解决。
取消是协作式的。cancelFiber() 记录 aborted 终端状态,并在 fiber 于当前 isolate 中运行时 abort ctx.signal。你的回调应在昂贵工作前后及可见副作用之前检查 ctx.signal.aborted。使用 waitForCompletion: true 的调用方在账本达到 aborted 时返回,即使非协作回调仍在当前 isolate 中运行。
如果 Durable Object 在 fiber 中途被驱逐,保留记录被标记为 interrupted,onFiberRecovered() 接收最后检查点。原始闭包无法自动重放;使用 ctx.name、ctx.snapshot 和 metadata 决定是恢复、补偿还是保留记录以供检查。
从 onFiberRecovered() 返回 FiberRecoveryResult 以记录策略决策:
async onFiberRecovered(ctx: FiberRecoveryContext) {
if (ctx.name !== "reply-to-webhook") return;
const snapshot = ctx.snapshot as { webhookId: string; threadId: string };
await postRecoveryMessage(snapshot.threadId);
return {
status: "completed",
snapshot: { ...snapshot, recovered: true },
};
}返回 undefined 使托管 fiber 保持 interrupted。抛出异常使其保持 interrupted 并记录恢复错误以供检查。如果存在过期的 run 行,终端托管 fiber(如 aborted)不会再次恢复。
如果恢复由后续重复 webhook 触发而非 onFiberRecovered(),在应用级恢复成功后使用相同结果形状的 resolveFiber()。resolveFiber() 仅更新当前为 interrupted 的托管 fiber;对 pending、running 或已终端行返回 false。
ctx.stash(data) 同步写入 SQLite。在「我决定保存」与「已保存」之间没有异步间隙。如果在 stash() 返回后发生驱逐,数据保证在 SQLite 中。
每次调用完全替换先前的快照——不是合并。写入你需要的完整恢复状态:
await this.runFiber("research", async (ctx) => {
const steps = ["search", "analyze", "synthesize"];
const completed: string[] = [];
const results: Record<string, unknown> = {};
for (const step of steps) {
results[step] = await executeStep(step);
completed.push(step);
ctx.stash({
completed,
results,
pendingSteps: steps.slice(completed.length),
});
}
});两者作用相同。ctx.stash() 对 fiber ID 使用直接闭包。this.stash() 使用 AsyncLocalStorage 查找当前执行的 fiber——即使并发 fiber 也能正确工作,因为每个 fiber 的 ALS 上下文独立。
this.stash() 便于从无 ctx 访问的嵌套函数调用。在 runFiber 回调外调用会抛出异常。
覆盖 onFiberRecovered 以处理 interrupted fiber。默认实现记录警告并删除行。
class ResearchAgent extends Agent {
async onFiberRecovered(ctx: FiberRecoveryContext) {
if (ctx.name !== "research") return;
const snapshot = ctx.snapshot as {
completed: string[];
results: Record<string, unknown>;
pendingSteps: string[];
} | null;
if (snapshot && snapshot.pendingSteps.length > 0) {
void this.runFiber("research", async (fiberCtx) => {
const { completed, results, pendingSteps } = snapshot;
for (const step of pendingSteps) {
results[step] = await this.executeStep(step);
completed.push(step);
fiberCtx.stash({
completed,
results,
pendingSteps: pendingSteps.slice(pendingSteps.indexOf(step) + 1),
});
}
});
}
}
}要点:
- 原始 lambda 已消失。 恢复时你只有
name和snapshot。lambda 无法序列化——恢复逻辑必须在钩子中。 - 非托管
runFiber()行在钩子成功返回后删除。 若要继续非托管工作,在钩子内再次调用runFiber()——这会创建新行。 - 托管
startFiber()行会保留。 返回FiberRecoveryResult将 interrupted 托管 fiber 标记为completed、error、aborted或仍为interrupted。 - 你控制恢复的含义。 从头重试、从检查点恢复、跳过并通知用户,或什么都不做。框架不强制策略。
- 如果钩子抛出异常,行会保留(有上限)。 后续启动或闹钟扫描会重试恢复,防止瞬时存储或调度失败。当你想将工作标记为终态而非重试时,请自行捕获应用级错误。始终抛出的钩子会在退避调度上重试(恢复闹钟使用上限 5 分钟的指数延迟,因此不是忙循环),直到行超过
fiberRecoveryMaxAgeMs(默认 24 小时),之后以fiber:recovery:skipped(reason: "max_age_exceeded")事件丢弃。设置fiberRecoveryMaxAgeMs: 0会无限保留此类行——恢复在有上限的退避上持续重试,且 Durable Object 在存在不可恢复行时永不空闲驱逐,因此除非你打算自行检查或清除这些行,否则优先使用有限期限。对于托管工作,保留行保持interrupted并记录恢复错误以供检查。
AIChatAgent 基于 fiber 实现 LLM 流式恢复。启用 chatRecovery 时,每个聊天轮次自动包装在 fiber 中。框架处理内部恢复路径并暴露 onChatRecovery 用于提供商特定策略。详情请参阅长时间运行 agent:恢复中断的 LLM 流。
多个 fiber 可同时运行。每个在 SQLite 中有自己的行和快照,并独立调用 keepAlive()(引用计数,因此 DO 在所有 fiber 完成前保持存活)。
void this.runFiber("fetch-data", async (ctx) => {
/* ... */
});
void this.runFiber("process-queue", async (ctx) => {
/* ... */
});恢复时,所有孤立行被迭代,并为每个调用 onFiberRecovered。在恢复钩子中使用 ctx.name 区分 fiber 类型。
在 wrangler dev 中,fiber 恢复与生产环境相同。SQLite 和 alarm 状态在重启之间持久化到磁盘。
- 启动 agent 并触发 fiber(
runFiber) - 终止 wrangler 进程(Ctrl-C 或 SIGKILL)
- 重启 wrangler
- 恢复自动触发——如有请求到达则通过
onStart(),无客户端连接则通过持久化 alarm
执行持久 fiber。fiber 在 fn 运行前注册到 SQLite,完成(或抛出)后删除。期间持有 keepAlive()。
name— fiber 标识符,在onFiberRecovered中用于区分 fiber 类型。不唯一——多个 fiber 可共享名称。fn— 接收FiberContext的异步函数。闭包自然工作(捕获this和局部变量)。- Returns —
fn返回的值。如果在完成前 DO 被驱逐,返回值丢失;通过钩子恢复。
持久接受保留的后台 fiber。返回的 StartFiberResult 包含生成的 fiberId、当前 status、可选 metadata 和 accepted;当现有 fiber 匹配相同幂等键时为 false。
name— 托管 fiber 标识符,用于检查和恢复。fn— 接收FiberContext的异步函数。函数结果不存储。options.idempotencyKey— 用于去重重试的稳定外部键。options.metadata— 与保留行一起存储的可 JSON 序列化数据。options.waitForCompletion— 在返回前等待终端状态。
返回托管 fiber 的保留状态行,无行则返回 null。
列出保留的托管 fiber。按 status 或 name 筛选,使用 limit 限制结果集。
将托管 fiber 标记为 aborted,并在当前 isolate 中运行时 abort 其内存 ctx.signal。fiber 不存在或已终端时返回 false。
应用级恢复成功后解析 interrupted 托管 fiber。对 pending、running 或已终端行返回 false。
删除保留的托管 fiber 行。默认 eligible 已结算的 completed、error 和 aborted 行。传递 status、settledBefore 或 limit 缩小清理范围。
检查点当前 fiber 状态。同步写入 SQLite。每次调用完全替换先前快照。data 必须可 JSON 序列化。
agent 重启时为每个孤立 fiber 行调用一次。覆盖以实现恢复。非托管 runFiber() 行在此钩子成功返回后删除;若恢复抛出,行保留供后续扫描,避免短暂失败丢失恢复句柄。托管 startFiber() 行保持保留,可通过返回 FiberRecoveryResult 解析。
ctx.id— 唯一 fiber IDctx.name— 传递给runFiber()的名称ctx.status— 托管 fiber 的保留状态ctx.idempotencyKey— 托管 fiber 的幂等键(如提供)ctx.metadata— 托管 fiber 的 metadata(如提供)ctx.snapshot— 最后一次stash()数据,或stash()从未调用时为nullctx.createdAt—runFiber()启动时的 epoch 毫秒。与Date.now()比较以跳过过旧无法安全重放的恢复。ctx.recoveryReason— 恢复运行原因。驱逐或重启恢复目前始终为"interrupted"。
创建 30 秒 alarm 心跳。返回 disposer 函数。幂等——多次调用 disposer 安全。
在保持 DO 存活的同时运行异步函数。fn 开始前启动心跳,完成或抛出时停止。返回 fn 返回的值。
- 长时间运行 agent — fiber 如何与 schedule、plan 和异步操作组合
- 调度任务 —
keepAlive详情与 alarm 系统 - 子 agent — 子 agent 内的持久执行与调度
- Workflows — agent 外部的持久多步执行
- Chat agent —
chatRecovery与onChatRecovery