跳转到内容
搜索文档

子 agent RPC 与编程式轮次

最后更新 查看 MarkdownAgent 设置

Think 既可作顶层 Agent 也可作子 Agent。作子 Agent 时,chat() 运行完整轮次并通过回调流式传输事件。

chat

async chat(
	userMessage: string | UIMessage,
	callback: StreamCallback,
	options?: ChatOptions,
): Promise<void>

StreamCallback

方法 触发时机
onStart(event) 工作开始前;暴露用于取消的 request ID
onEvent(json) 每个流式分块(JSON 序列化的 UIMessageChunk
onDone() 轮次完成且 assistant 消息持久化后
onError(message) 轮次期间出错
onInterrupted() 可选。尝试被中断,已调度的续传(后续 isolate)拥有最终结果 — 非完成、非终态错误。默认无操作

onInterruptedchat() 驱动且被中断并恢复的轮次很重要:RPC promise 干净兑现(isolate 仍存活),仅依赖干净兑现的消费者会误读为成功并最终化已流式传输的部分内容。应视为「未完成、未失败 — 续传拥有答案」:保持通道开放、显示恢复中状态或重新挂接,而非最终化部分内容。部署或驱逐中断会在触发前杀死 isolate(调用方看到传输中断);onInterrupted 覆盖 isolate 内停滞转入恢复的路径。

ChatOptions

字段 描述
signal AbortSignal,用于在流中途取消轮次

工具属于子 agent。用子级的 getTools()、扩展、MCP 工具或客户端工具 schema 定义持久能力。向 chat() 传旧版 options.tools 的调用方会收到警告且值被忽略。

示例:父级调用子级

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

export class ParentAgent extends Think {
	getModel() {
		/* ... */
	}

	async delegateToChild(task) {
		const child = await this.subAgent(ChildAgent, "child-1");

		const chunks = [];
		await child.chat(task, {
			onStart: (event) => {
				console.log("Child started:", event.requestId);
			},
			onEvent: (json) => {
				chunks.push(json);
			},
			onDone: () => {
				console.log("Child completed");
			},
			onError: (error) => {
				console.error("Child failed:", error);
			},
		});

		return chunks;
	}
}

export class ChildAgent extends Think {
	getModel() {
		/* ... */
	}

	getSystemPrompt() {
		return "You are a research assistant. Analyze data and report findings.";
	}
}
import { Think } from "@cloudflare/think";
import type { StreamCallback } from "@cloudflare/think";

export class ParentAgent extends Think<Env> {
	getModel() {
		/* ... */
	}

	async delegateToChild(task: string) {
		const child = await this.subAgent(ChildAgent, "child-1");

		const chunks: string[] = [];
		await child.chat(task, {
			onStart: (event) => {
				console.log("Child started:", event.requestId);
			},
			onEvent: (json) => {
				chunks.push(json);
			},
			onDone: () => {
				console.log("Child completed");
			},
			onError: (error) => {
				console.error("Child failed:", error);
			},
		});

		return chunks;
	}
}

export class ChildAgent extends Think<Env> {
	getModel() {
		/* ... */
	}

	getSystemPrompt() {
		return "You are a research assistant. Analyze data and report findings.";
	}
}

取消子 agent 轮次

onStartcancelChat() 跨子 agent 边界以 RPC 安全方式取消:

let requestId;

const callback = {
	onStart(event) {
		requestId = event.requestId;
	},
	onEvent(json) {
		// Forward stream chunks.
	},
	onDone() {},
	onError(error) {
		console.error(error);
	},
};

const turn = child.chat("Long analysis task", callback);

// Later, from another RPC call or failure handler:
if (requestId) {
	await child.cancelChat(requestId, "client disconnected");
}

await turn;
let requestId: string | undefined;

const callback: StreamCallback = {
	onStart(event) {
		requestId = event.requestId;
	},
	onEvent(json) {
		// Forward stream chunks.
	},
	onDone() {},
	onError(error) {
		console.error(error);
	},
};

const turn = child.chat("Long analysis task", callback);

// Later, from another RPC call or failure handler:
if (requestId) {
	await child.cancelChat(requestId, "client disconnected");
}

await turn;

调用方与被调用方未由 Workers RPC 分离时,也可传 AbortSignal 在流中途取消:

const controller = new AbortController();
setTimeout(() => controller.abort(), 30_000);

await child.chat("Long analysis task", callback, {
	signal: controller.signal,
});
const controller = new AbortController();
setTimeout(() => controller.abort(), 30_000);

await child.chat("Long analysis task", callback, {
	signal: controller.signal,
});

轮次已完成或请求 ID 未知时 cancelChat(requestId, reason?) 为无操作。中止时部分 assistant 消息仍会持久化。

saveMessages

注入消息并触发模型轮次,无需 WebSocket。用于定时响应、webhook 触发轮次、主动 agent 或从 onChatResponse 链式调用。

async saveMessages(
	messages:
		| UIMessage[]
		| ((current: UIMessage[]) => UIMessage[] | Promise<UIMessage[]>),
	options?: SaveMessagesOptions,
): Promise<SaveMessagesResult>

返回 { requestId, status, error? },其中 status"completed""error""skipped""aborted"

status 时机
"completed" 轮次运行至完成。
"error" 轮次已开始但流报告错误。error 含流错误消息(如有)。
"skipped" 轮次中途失效,例如 chat-clear;用户消息已持久化,无模型运行。
"aborted" 通过 options.signalchat-request-cancel 在完成前取消。部分 assistant 分块仍会持久化。

从启动轮次的 Durable Object 传 options.signal 取消程序化轮次。AbortSignal 不能跨 Durable Object RPC,信号不会跨休眠持久化。

静态消息

await this.saveMessages([
	{
		id: crypto.randomUUID(),
		role: "user",
		parts: [{ type: "text", text: "Time for your daily summary." }],
	},
]);
await this.saveMessages([
	{
		id: crypto.randomUUID(),
		role: "user",
		parts: [{ type: "text", text: "Time for your daily summary." }],
	},
]);

函数形式

多个 saveMessages 排队时,函数形式在轮次实际开始时用最新消息运行:

await this.saveMessages((current) => [
	...current,
	{
		id: crypto.randomUUID(),
		role: "user",
		parts: [{ type: "text", text: "Continue your analysis." }],
	},
]);
await this.saveMessages((current) => [
	...current,
	{
		id: crypto.randomUUID(),
		role: "user",
		parts: [{ type: "text", text: "Continue your analysis." }],
	},
]);

定时响应

getScheduledTasks() 触发周期性提示词轮次:

export class MyAgent extends Think {
	getModel() {
		/* ... */
	}

	getScheduledTasks() {
		return {
			dailyReport: {
				schedule: "every day at 09:00",
				timezone: "UTC",
				prompt: "Generate the daily report.",
			},
		};
	}
}
export class MyAgent extends Think<Env> {
	getModel() {
		/* ... */
	}

	getScheduledTasks() {
		return {
			dailyReport: {
				schedule: "every day at 09:00",
				timezone: "UTC",
				prompt: "Generate the daily report.",
			},
		};
	}
}

从 onChatResponse 链式

当前轮次完成后启动后续轮次:

async onChatResponse(result: ChatResponseResult) {
	if (result.status === "completed" && this.needsFollowUp(result.message)) {
		await this.saveMessages([{
			id: crypto.randomUUID(),
			role: "user",
			parts: [{ type: "text", text: "Now summarize what you found." }],
		}]);
	}
}

continueLastTurn

在最新 assistant 消息后运行另一次模型调用,不注入新用户消息。Think 将结果持久化为 continuation: true 的新 assistant 消息;不向现有 assistant 消息追加分块。

protected async continueLastTurn(
	body?: Record<string, unknown>,
	options?: SaveMessagesOptions,
): Promise<SaveMessagesResult>

若最后一条消息不是 assistant 消息,返回 { requestId, status: "skipped" }。可选 body 覆盖此续传的已存储 body。传 options.signal 在续传运行时取消。

abortRequest 与 abortAllRequests

从 Durable Object 内部取消进行中的聊天轮次:

protected abortRequest(requestId: string, reason?: unknown): void
protected abortAllRequests(): void

知道请求 ID 时用 abortRequest()。专用辅助方法应取消当前任意轮次时用 abortAllRequests()。能在调用点传 signal 的编程式轮次优先用 SaveMessagesOptions.signal

这篇文档对您有帮助吗?