跳转到内容
搜索文档

持久化恢复

最后更新 查看 MarkdownAgent 设置

Think 默认将聊天轮次包装在可恢复的 fiber 中(chatRecovery = true)。若 Durable Object 在流式传输中途被驱逐,Think 会重建已缓冲的分块、持久化部分输出,并调度 assistant 轮次的续传或未回答用户轮次的重试。

chatRecoverytrue 时,WebSocket 轮次、子 Agent chat() 轮次、持久 submitMessages() 执行、自动续写、saveMessages()continueLastTurn() 均包装在 runFiber 中。

有界恢复

流式停滞看门狗中止(chatStreamStallTimeoutMs)被视为另一种中断:开启 chatRecovery 时,停滞进入同一有界路径——已结算的部分结果被保留并调度续传——因此短暂挂起可自动恢复。持续挂起的提供商会耗尽预算,并通过与部署或驱逐中断相同的耗尽处理进入终态:onExhausted 触发、发出 chat:recovery:exhausted 事件,并显示配置的 terminalMessage(而非原始停滞错误)。

通过将 chatRecovery 设为对象来配置有界恢复:

export class MyAgent extends Think {
	chatRecovery = {
		maxAttempts: 6,
		stableTimeoutMs: 10_000,
		terminalMessage: "The assistant was interrupted and could not recover.",
		async onExhausted(ctx) {
			console.warn("Chat recovery exhausted", ctx.incidentId);
		},
	};

	getModel() {
		/* ... */
	}
}
export class MyAgent extends Think<Env> {
	override chatRecovery = {
		maxAttempts: 6,
		stableTimeoutMs: 10_000,
		terminalMessage: "The assistant was interrupted and could not recover.",
		async onExhausted(ctx) {
			console.warn("Chat recovery exhausted", ctx.incidentId);
		},
	};

	getModel() {
		/* ... */
	}
}

相同的恢复事件可通过 agents/observabilitychat 通道上获取;对话记录修复在 transcript 通道上发出。请参阅 可观测性

onChatRecovery

需要提供商特定恢复时重写 onChatRecovery,例如检索已存储的 OpenAI Responses 结果而非发起新模型调用:

export class MyAgent extends Think {
	chatRecovery = {
		maxAttempts: 10,
		terminalMessage: "The assistant was interrupted. Please try again.",
	};

	async onChatRecovery(ctx) {
		console.log("Recovering chat turn", ctx.incidentId, ctx.attempt);
		return {}; // persist partial output and continue/retry when possible
	}
}
import type {
	ChatRecoveryContext,
	ChatRecoveryOptions,
} from "@cloudflare/think";

export class MyAgent extends Think<Env> {
	override chatRecovery = {
		maxAttempts: 10,
		terminalMessage: "The assistant was interrupted. Please try again.",
	};

	override async onChatRecovery(
		ctx: ChatRecoveryContext,
	): Promise<ChatRecoveryOptions> {
		console.log("Recovering chat turn", ctx.incidentId, ctx.attempt);
		return {}; // persist partial output and continue/retry when possible
	}
}

ChatRecoveryContext

字段 类型 描述
incidentId string 此恢复事件的稳定 ID
attempt number 此事件的当前尝试次数,从 1 开始
maxAttempts number 进入终态耗尽前的配置尝试上限
recoveryKind "retry" | "continue" 恢复将重试未回答的用户轮次,还是继续部分 assistant 轮次
streamId string 被中断轮次的流 ID
requestId string 被中断轮次的 request ID
partialText string 中断前生成的文本
partialParts MessagePart[] 中断前累积的部分
recoveryData unknown | null 轮次期间 this.stash() 的数据
messages UIMessage[] 当前对话历史
lastBody Record<string, unknown>? 被中断轮次的正文
lastClientTools ClientToolSchema[]? 被中断轮次的客户端工具
createdAt number 轮次开始时的 epoch 毫秒

ChatRecoveryOptions

字段 类型 描述
persist boolean? 是否持久化部分 assistant 消息
continue boolean? Agent 达到稳定状态后是否通过 continueLastTurn() 自动继续

persist: true 时,部分消息会被保存。continue: true 时,Think 在 Agent 达到稳定状态后调用 continueLastTurn()

对于流开始前的中断(ctx.streamId === ""ctx.partialText === "",但最新持久化消息仍是未回答的用户消息),除非 continuefalse,Think 会自动重试该轮次。

onChatRecovery(ctx: ChatRecoveryContext): ChatRecoveryOptions {
	if (!ctx.streamId && !ctx.partialText) {
		console.log("Recovering a pre-stream interruption");
	}
	return {};
}

使用 ctx.createdAt 跳过过期的恢复。例如,若被中断的轮次已超过数分钟,返回 { continue: false },以保留部分响应而不启动旧的续传。

恢复预算与限制

chatRecovery = true 外,也可分配对象以调整恢复允许运行的时长及何时放弃。持续有前进进度的轮次不会被框架自行终止 — 时长不是边界。恢复仅由下表中的限制之一封存。

export class MyAgent extends Think {
	chatRecovery = {
		maxAttempts: 10,
		noProgressTimeoutMs: 5 * 60 * 1000,
		maxRecoveryWork: Infinity,
		terminalMessage: "The assistant was interrupted and could not recover.",
		// Consulted from the second recovery attempt onward. Return false to stop.
		// Called as `config.shouldKeepRecovering(ctx)`, so it is NOT bound to the
		// agent instance — track real token/cost spend in your own store keyed by
		// `ctx.recoveryRootRequestId`.
		async shouldKeepRecovering(ctx) {
			return (await getSpendForTurn(ctx.recoveryRootRequestId)) < MAX_SPEND;
		},
		async onExhausted(ctx) {
			console.warn("Recovery exhausted", ctx.incidentId, ctx.reason);
		},
	};
}
export class MyAgent extends Think<Env> {
	override chatRecovery = {
		maxAttempts: 10,
		noProgressTimeoutMs: 5 * 60 * 1000,
		maxRecoveryWork: Infinity,
		terminalMessage: "The assistant was interrupted and could not recover.",
		// Consulted from the second recovery attempt onward. Return false to stop.
		// Called as `config.shouldKeepRecovering(ctx)`, so it is NOT bound to the
		// agent instance — track real token/cost spend in your own store keyed by
		// `ctx.recoveryRootRequestId`.
		async shouldKeepRecovering(ctx) {
			return (await getSpendForTurn(ctx.recoveryRootRequestId)) < MAX_SPEND;
		},
		async onExhausted(ctx) {
			console.warn("Recovery exhausted", ctx.incidentId, ctx.reason);
		},
	};
}
字段 默认值 描述
maxAttempts 10 尝试上限。有前进进度时重置,因此捕获紧密的无进度告警循环,而不是健康的长轮次。
stableTimeoutMs 10_000 某次尝试等待 isolate 达到稳定状态的时长,超时后重新调度。
noProgressTimeoutMs 300_000(5 分钟) 主要的卡住轮次边界:无前进进度的最长时间,超时后封存。每次有进度的尝试都会重置。
maxRecoveryWork Infinity 失控循环防护:事件打开后产生的内容/工具单元上限,仍在推进的轮次也会被封存。默认无上限。
shouldKeepRecovering 从第二次尝试起咨询的调用方策略。返回 false 停止恢复。token/成本预算的 hook 点(ctx.work 是粗粒度分段计数,不是 token)。
terminalMessage 通用消息 放弃恢复时向用户显示的消息。
onExhausted 放弃恢复时调用一次。检查 ctx.reason

耗尽 hook 上的 ctx.reason 为以下之一:no_progress_timeout(卡住)、max_attempts_exceeded(无进度告警循环)、work_budget_exceeded(失控)、recovery_aborted(你的 shouldKeepRecovering 返回 false)或 stable_timeout(极端抖动)。完整共享参考请参阅 流恢复 — Think 与 @cloudflare/ai-chat 使用相同的恢复配置。

修复中断的工具调用

当轮次在进行中被中断时,对话记录可能包含尚无已结算结果的工具调用。在下次提供商调用前,Think 修复每个此类调用,以免模型静默重跑,也避免提供商以 AI_MissingToolResultsError 拒绝对话记录。默认将中断的调用翻转为出错的工具结果,因此记录保留,转换后仍有工具结果。

重写 repairInterruptedToolPart 以自定义修复后的形态。常见情况是由客户端解析的工具 — 例如没有服务端 execute、通常由用户下一条消息回答的 ask_user 问题。将其转换为纯文本部分可让模型将其视为普通对话而非工具错误,并在压缩中逐字保留问题:

export class MyAgent extends Think {
	repairInterruptedToolPart(part) {
		const record = part;
		if (record.type === "tool-ask_user") {
			const input = record.input;
			if (input?.prompt) {
				return { type: "text", text: input.prompt };
			}
		}
		return super.repairInterruptedToolPart(part);
	}
}
import type { UIMessage } from "ai";

export class MyAgent extends Think<Env> {
	protected override repairInterruptedToolPart(
		part: UIMessage["parts"][number],
	): UIMessage["parts"][number] {
		const record = part as Record<string, unknown>;
		if (record.type === "tool-ask_user") {
			const input = record.input as { prompt?: string } | undefined;
			if (input?.prompt) {
				return { type: "text", text: input.prompt };
			}
		}
		return super.repairInterruptedToolPart(part);
	}
}

这在对话记录修复期间运行 — 在修复后的记录持久化并发送给模型之前 — 因此转换塑造当前轮次,而不仅是下一个。input 已规范化为有效对象。返回的工具部分必须携带已结算结果(output-availableoutput-erroroutput-denied);返回文本等非工具部分也可以。

上下文窗口溢出恢复

压缩轮次之间 检查 — compactAfter() 在每次 appendMessage() 后运行。但单个很长、工具很多的轮次会在一个 streamText 循环内逐步增长 prompt,并可能在 轮次中间、下次轮次前检查之前超出模型上下文窗口。提供商随后拒绝请求("prompt is too long"context_length_exceeded),轮次否则会以终态失败。

Think 通过 contextOverflow 属性的两层可选、与提供商无关的机制从此恢复。两者默认关闭,因此现有行为不变。两者复用你的会话压缩函数,因此需要配置了 onCompaction()configureSession()。两者都需要 classifyChatError 告诉 Think 哪些错误是溢出 — Think 核心不包含提供商特定匹配。

1. 反应式兜底 — contextOverflow.reactive 当轮次因你分类为 "context_overflow" 的错误失败时,Think 丢弃截断的部分结果、运行 session.compact(),并从压缩后的历史重新运行轮次。部分结果不会持久化:轮次从头重启,因此保留被截断的 assistant 消息会使其孤立在恢复后的答案旁。由 contextOverflow.maxRetries(默认 1)限定;若压缩无法缩短历史或预算用尽,溢出通过 onChatErrorclassification: "context_overflow" 终态抛出 — 不会循环或静默结束。

import { Think, defaultContextOverflowClassifier } from "@cloudflare/think";

export class MyAgent extends Think {
	contextOverflow = { reactive: true };

	// The bundled classifier covers the common providers (Anthropic, OpenAI,
	// Google, Bedrock, …). Assign it directly, or write your own.
	classifyChatError = defaultContextOverflowClassifier;
}
import { Think, defaultContextOverflowClassifier } from "@cloudflare/think";

export class MyAgent extends Think<Env> {
	override contextOverflow = { reactive: true };

	// The bundled classifier covers the common providers (Anthropic, OpenAI,
	// Google, Bedrock, …). Assign it directly, or write your own.
	override classifyChatError = defaultContextOverflowClassifier;
}

2. 主动防护 — contextOverflow.proactive 在提供商错误发生前预防。每步之前,Think 读取上一步模型报告的 usage.inputTokens(与提供商无关),若超过 maxInputTokens * (headroom ?? 0.9),则原地压缩并将重新压缩的历史送入即将到来的步骤。若提供商省略 inputTokens,回退到 usage.totalTokens(安全的高估 — 略早压缩而不是错过阈值)。每轮次最多压缩 proactive.maxCompactions 次(默认 1)— 独立于反应式 maxRetries 预算 — 因此无法缩短的历史不会每步都压缩。

import { Think, defaultContextOverflowClassifier } from "@cloudflare/think";

export class MyAgent extends Think {
	contextOverflow = {
		reactive: true,
		// Compact mid-turn once a step approaches 90% of a 200K window.
		proactive: { maxInputTokens: 200_000 },
	};

	classifyChatError = defaultContextOverflowClassifier;
}
import { Think, defaultContextOverflowClassifier } from "@cloudflare/think";

export class MyAgent extends Think<Env> {
	override contextOverflow = {
		reactive: true,
		// Compact mid-turn once a step approaches 90% of a 200K window.
		proactive: { maxInputTokens: 200_000 },
	};

	override classifyChatError = defaultContextOverflowClassifier;
}

可单独使用任一层,或两者一起:主动防护避免大多数溢出,反应式兜底捕获仍漏网的(例如轮次开始时已超预算,或单个工具结果过大、压缩无法帮助 — 此时干净地进入终态)。两者适用于每个轮次入口路径(WebSocket、子 agent chat()、编程式 saveMessages() / submitMessages()),并发出 chat:context:compacted 可观测性事件

针对真实 Workers AI 模型的可运行演示,请参阅 context-overflow-recovery 示例

稳定性检测

Think 提供检查 Agent 是否处于稳定状态的方法 — 无待处理工具结果、无待处理审批、无活跃轮次。

hasPendingInteraction

若任何 assistant 消息有待处理工具调用(无结果的工具或待处理审批),返回 true

protected hasPendingInteraction(): boolean

waitUntilStable

返回 promise,Agent 达到稳定状态时解析为 true,超时则 false

const stable = await this.waitUntilStable({ timeout: 30_000 });
if (stable) {
	await this.saveMessages([
		{
			id: crypto.randomUUID(),
			role: "user",
			parts: [{ type: "text", text: "Now that you are done, summarize." }],
		},
	]);
}
const stable = await this.waitUntilStable({ timeout: 30_000 });
if (stable) {
	await this.saveMessages([
		{
			id: crypto.randomUUID(),
			role: "user",
			parts: [{ type: "text", text: "Now that you are done, summarize." }],
		},
	]);
}

这篇文档对您有帮助吗?