跳转到内容
搜索文档

使用 fiber 的持久执行

最后更新 查看 MarkdownAgent 设置

运行可在 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 });
		}
	}
}

为何需要 fiber

Durable Objects 因以下三个原因被驱逐:

  1. 不活动超时 — 约 70–140 秒内无传入请求或打开的 WebSocket
  2. 代码更新 / 运行时重启 — 非确定性,每天 1–2 次
  3. Alarm 处理程序超时 — 15 分钟

驱逐发生在工作进行中时,上游 HTTP 连接(到 LLM 提供商、API、数据库)会永久断开。内存状态——流式缓冲区、部分响应、循环计数器——会丢失。多轮 agent 循环会完全失去位置。

keepAlive() 降低驱逐概率。runFiber() 使驱逐后仍可恢复。

对于应独立于 agent 运行、具有每步重试和多步编排的工作,请改用 Workflows。Fiber 适用于 agent 自身执行的一部分。比较请参阅长时间运行 agent:Workflows 与 agent 内部模式。

keepAlive

通过创建 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 与 runFiber

keepAlive 防止驱逐但不处理恢复。如果 agent 仍 被驱逐(代码更新、alarm 超时、资源限制),任何进行中的工作都会丢失。

runFiber 内部调用 keepAlive 并 将工作持久化到 SQLite 以便恢复。当工作重做成本低或不需要检查点时,单独使用 keepAlive。当工作成本高且需要从断点恢复时,使用 runFiber。

场景 使用
等待慢速 API 调用 keepAlive()
流式传输 LLM 响应(通过 AIChatAgent) 自动(内置)
带中间结果的多步计算 runFiber()
耗时 10 分钟以上的后台研究循环 带 stash() 的 runFiber()
必须恰好接受一次的 webhook 任务 startFiber()

runFiber

带检查点和恢复的持久执行。

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。

子 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,你可在此获得快照及出错的完整上下文。

内联与 fire-and-forget

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

当调用方需要持久接受后台工作、快速返回并安全去重重试时,请使用 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。

使用 stash 的检查点

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

this.stash 与 ctx.stash

两者作用相同。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

多个 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 状态在重启之间持久化到磁盘。

  1. 启动 agent 并触发 fiber(runFiber)
  2. 终止 wrangler 进程(Ctrl-C 或 SIGKILL)
  3. 重启 wrangler
  4. 恢复自动触发——如有请求到达则通过 onStart(),无客户端连接则通过持久化 alarm

API 参考

runFiber(name, fn)

执行持久 fiber。fiber 在 fn 运行前注册到 SQLite,完成(或抛出)后删除。期间持有 keepAlive()。

  • name — fiber 标识符,在 onFiberRecovered 中用于区分 fiber 类型。不唯一——多个 fiber 可共享名称。
  • fn — 接收 FiberContext 的异步函数。闭包自然工作(捕获 this 和局部变量)。
  • Returns — fn 返回的值。如果在完成前 DO 被驱逐,返回值丢失;通过钩子恢复。

startFiber(name, fn, options)

持久接受保留的后台 fiber。返回的 StartFiberResult 包含生成的 fiberId、当前 status、可选 metadata 和 accepted;当现有 fiber 匹配相同幂等键时为 false。

  • name — 托管 fiber 标识符,用于检查和恢复。
  • fn — 接收 FiberContext 的异步函数。函数结果不存储。
  • options.idempotencyKey — 用于去重重试的稳定外部键。
  • options.metadata — 与保留行一起存储的可 JSON 序列化数据。
  • options.waitForCompletion — 在返回前等待终端状态。

inspectFiber(fiberId) / inspectFiberByKey(idempotencyKey)

返回托管 fiber 的保留状态行,无行则返回 null。

listFibers(options)

列出保留的托管 fiber。按 status 或 name 筛选,使用 limit 限制结果集。

cancelFiber(fiberId, reason) / cancelFiberByKey(idempotencyKey, reason)

将托管 fiber 标记为 aborted,并在当前 isolate 中运行时 abort 其内存 ctx.signal。fiber 不存在或已终端时返回 false。

resolveFiber(fiberId, result)

应用级恢复成功后解析 interrupted 托管 fiber。对 pending、running 或已终端行返回 false。

deleteFibers(options)

删除保留的托管 fiber 行。默认 eligible 已结算的 completed、error 和 aborted 行。传递 status、settledBefore 或 limit 缩小清理范围。

stash(data) / ctx.stash(data)

检查点当前 fiber 状态。同步写入 SQLite。每次调用完全替换先前快照。data 必须可 JSON 序列化。

onFiberRecovered(ctx)

agent 重启时为每个孤立 fiber 行调用一次。覆盖以实现恢复。非托管 runFiber() 行在此钩子成功返回后删除;若恢复抛出,行保留供后续扫描,避免短暂失败丢失恢复句柄。托管 startFiber() 行保持保留,可通过返回 FiberRecoveryResult 解析。

  • ctx.id — 唯一 fiber ID
  • ctx.name — 传递给 runFiber() 的名称
  • ctx.status — 托管 fiber 的保留状态
  • ctx.idempotencyKey — 托管 fiber 的幂等键(如提供)
  • ctx.metadata — 托管 fiber 的 metadata(如提供)
  • ctx.snapshot — 最后一次 stash() 数据,或 stash() 从未调用时为 null
  • ctx.createdAt — runFiber() 启动时的 epoch 毫秒。与 Date.now() 比较以跳过过旧无法安全重放的恢复。
  • ctx.recoveryReason — 恢复运行原因。驱逐或重启恢复目前始终为 "interrupted"。

keepAlive()

创建 30 秒 alarm 心跳。返回 disposer 函数。幂等——多次调用 disposer 安全。

keepAliveWhile(fn)

在保持 DO 存活的同时运行异步函数。fn 开始前启动心跳,完成或抛出时停止。返回 fn 返回的值。

相关

这篇文档对您有帮助吗?