Think 默认将聊天轮次包装在可恢复的 fiber 中(chatRecovery = true)。若 Durable Object 在流式传输中途被驱逐,Think 会重建已缓冲的分块、持久化部分输出,并调度 assistant 轮次的续传或未回答用户轮次的重试。
当 chatRecovery 为 true 时,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/observability 在 chat 通道上获取;对话记录修复在 transcript 通道上发出。请参阅 可观测性。
需要提供商特定恢复时重写 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
}
}| 字段 | 类型 | 描述 |
|---|---|---|
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 毫秒 |
| 字段 | 类型 | 描述 |
|---|---|---|
persist |
boolean? |
是否持久化部分 assistant 消息 |
continue |
boolean? |
Agent 达到稳定状态后是否通过 continueLastTurn() 自动继续 |
persist: true 时,部分消息会被保存。continue: true 时,Think 在 Agent 达到稳定状态后调用 continueLastTurn()。
对于流开始前的中断(ctx.streamId === "" 且 ctx.partialText === "",但最新持久化消息仍是未回答的用户消息),除非 continue 为 false,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-available、output-error 或 output-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)限定;若压缩无法缩短历史或预算用尽,溢出通过 onChatError 以 classification: "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 是否处于稳定状态的方法 — 无待处理工具结果、无待处理审批、无活跃轮次。
若任何 assistant 消息有待处理工具调用(无结果的工具或待处理审批),返回 true。
protected hasPendingInteraction(): boolean返回 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." }],
},
]);
}