跳转到内容
搜索文档

使用 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() 删除已结算的 completederroraborted 行。除非你显式传递该状态,否则不会删除 interrupted 行,因为 interrupted 行通常需要检查或手动解决。

取消是协作式的。cancelFiber() 记录 aborted 终端状态,并在 fiber 于当前 isolate 中运行时 abort ctx.signal。你的回调应在昂贵工作前后及可见副作用之前检查 ctx.signal.aborted。使用 waitForCompletion: true 的调用方在账本达到 aborted 时返回,即使非协作回调仍在当前 isolate 中运行。

如果 Durable Object 在 fiber 中途被驱逐,保留记录被标记为 interruptedonFiberRecovered() 接收最后检查点。原始闭包无法自动重放;使用 ctx.namectx.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 已消失。 恢复时你只有 namesnapshot。lambda 无法序列化——恢复逻辑必须在钩子中。
  • 非托管 runFiber() 行在钩子成功返回后删除。 若要继续非托管工作,在钩子内再次调用 runFiber()——这会创建新行。
  • 托管 startFiber() 行会保留。 返回 FiberRecoveryResult 将 interrupted 托管 fiber 标记为 completederroraborted 或仍为 interrupted
  • 你控制恢复的含义。 从头重试、从检查点恢复、跳过并通知用户,或什么都不做。框架不强制策略。
  • 如果钩子抛出异常,行会保留(有上限)。 后续启动或闹钟扫描会重试恢复,防止瞬时存储或调度失败。当你想将工作标记为终态而非重试时,请自行捕获应用级错误。始终抛出的钩子会在退避调度上重试(恢复闹钟使用上限 5 分钟的指数延迟,因此不是忙循环),直到行超过 fiberRecoveryMaxAgeMs(默认 24 小时),之后以 fiber:recovery:skippedreason: "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 和局部变量)。
  • Returnsfn 返回的值。如果在完成前 DO 被驱逐,返回值丢失;通过钩子恢复。

startFiber(name, fn, options)

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

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

inspectFiber(fiberId) / inspectFiberByKey(idempotencyKey)

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

listFibers(options)

列出保留的托管 fiber。按 statusname 筛选,使用 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 已结算的 completederroraborted 行。传递 statussettledBeforelimit 缩小清理范围。

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.createdAtrunFiber() 启动时的 epoch 毫秒。与 Date.now() 比较以跳过过旧无法安全重放的恢复。
  • ctx.recoveryReason — 恢复运行原因。驱逐或重启恢复目前始终为 "interrupted"

keepAlive()

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

keepAliveWhile(fn)

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

相关

这篇文档对您有帮助吗?