Skip to content

非公式本サイトは非公式の日本語ドキュメントであり、Cloudflare 公式サイトではありません。最新情報はdevelopers.cloudflare.comをご確認ください。

自律応答

最終更新 Markdown で表示Agent セットアップ

人間の操作なしに、サーバー側からメッセージを送り、LLM の応答を起動します。定期フォローアップ、キュー処理、メール起点の応答、自律エージェントワークフローに使います。

概要

通常のチャットフローでは、ユーザーがメッセージを送り、エージェントが応答します。ただしエージェントは、自分で動く必要もよくあります。定期リマインダーの発火、webhook の到着、ワークフローの完了、自身の応答を確認したあとの継続などです。

主なプリミティブは次のとおりです。

プリミティブ 役割
saveMessages メッセージを注入し、LLM を起動します。sendMessage のサーバー側相当です
submitMessages Think ターンを非同期実行向けに耐久的に受け付け、あとから確認します
startFiber ターン周辺のアプリケーション所有の副作用を耐久的に受け付けます
persistMessages 応答を起動せずにメッセージを保存します。コンテキストを静かに注入する場合に使います
onChatResponse 自分で起動していないものも含め、応答が完了したときに反応します
isServerStreaming クライアント側フラグです。サーバー起点のストリームが稼働中なら true です

saveMessagespersistMessages

saveMessages はメッセージを SQLite に永続化し、かつ 新しい LLM 応答のために onChatMessage を起動します。await できます。返ってきた時点で、LLM は応答済みで、メッセージは永続化されています。

persistMessages はメッセージを保存し、接続中のクライアントへ配信しますが、モデルターンは起動しません。会話にコンテキスト(システムメッセージやバックグラウンドデータなど)を注入し、応答は始めない場合に使います。

saveMessagessubmitMessages

呼び出し側がモデルターンの完了を待てる場合は saveMessages() を使います。

呼び出し側が、素早い耐久的な受付、べき等な再試行、あとからの状態確認を必要とする場合は、Think と一緒に submitMessages() を使います。タイムアウト制限が厳しい webhook ハンドラー、RPC 呼び出し元、親 Workers に向いています。

const submission = await this.submitMessages(
	[
		{
			id: crypto.randomUUID(),
			role: "user",
			parts: [
				{ type: "text", text: `Webhook event: ${JSON.stringify(payload)}` },
			],
		},
	],
	{ idempotencyKey: payload.id },
);

return Response.json({
	submissionId: submission.submissionId,
	status: submission.status,
	accepted: submission.accepted,
});
const submission = await this.submitMessages(
	[
		{
			id: crypto.randomUUID(),
			role: "user",
			parts: [
				{ type: "text", text: `Webhook event: ${JSON.stringify(payload)}` },
			],
		},
	],
	{ idempotencyKey: payload.id },
);

return Response.json({
	submissionId: submission.submissionId,
	status: submission.status,
	accepted: submission.accepted,
});

submitMessages() は、保留中の作業を先に保存し、実行が始まった時点で会話 Session にメッセージを追加します。受け付けるのはシリアライズ可能な UIMessage[] です。saveMessages((messages) => ...) が対応している関数形式は使えません。

Think の外で、webhook を一度だけ受け付ける、プロバイダー状態を復元する、ユーザーに見える返信を投稿する、復旧ポリシーを記録する、といった周囲のアプリケーションジョブが耐久単位である場合は startFiber() を使います。submitMessages() は Think の会話受付を担当し、managed fiber はそのターン周辺の外部副作用を担当します。

Think API の全体は submitMessages() を参照してください。

saveMessagesonChatResponse の使い分け

トリガーを自分で制御できる場合は saveMessages を使います。 スケジュールコールバック、webhook、メールハンドラー、メッセージ注入のタイミングを決める任意のメソッドです。

自分で起動していない応答に反応する必要がある場合は onChatResponse を使います。 ユーザー起点のメッセージ、ツール承認後の自動継続、フレームワークが代わりに実行した任意のターンです。

waitUntilStable

スケジュールコールバック、webhook、メールハンドラー、その他チャット以外のエントリポイントから this.messages を読む、または saveMessages を呼ぶ前に、必ず waitUntilStable() を呼び出します。

waitUntilStable() は、会話が完全に安定するまで待ちます。

  • 進行中の LLM ストリームがない
  • 保留中のクライアントツール操作(ユーザーがまだ提供していないツール結果や承認)がない
  • キューに入った継続ターンがない

安定していれば true を返します。保留中の操作が解消される前にタイムアウトした場合は false です。保留がなければ、すぐに返ります。

const stable = await this.waitUntilStable({ timeout: 30_000 });
if (!stable) {
	// The conversation is blocked on a user interaction or an in-flight
	// stream that did not complete within 30 seconds.
	console.warn("Conversation not stable, skipping server-driven message");
	return;
}
// Safe to read this.messages and call saveMessages.
const stable = await this.waitUntilStable({ timeout: 30_000 });
if (!stable) {
	// The conversation is blocked on a user interaction or an in-flight
	// stream that did not complete within 30 seconds.
	console.warn("Conversation not stable, skipping server-driven message");
	return;
}
// Safe to read this.messages and call saveMessages.

このガードがないと、古いメッセージを読んだり、進行中のストリームと重なったりするおそれがあります。

トリガーパターン

Cron スケジュール

毎朝アクティビティを要約する日次ダイジェストエージェントです。Cron スケジュールはデフォルトでべき等なので、onStartschedule() を呼んでも安全です。Durable Object の再起動をまたいで重複は作られません。

import { AIChatAgent } from "@cloudflare/ai-chat";

export class DigestAgent extends AIChatAgent {
	async onChatMessage() {
		// ... your LLM call
	}

	async onStart() {
		await this.schedule("0 9 * * *", "dailyDigest");
	}

	async dailyDigest() {
		const stable = await this.waitUntilStable({ timeout: 30_000 });
		if (!stable) {
			console.warn("Conversation not stable, skipping daily digest");
			return;
		}

		await this.saveMessages((messages) => [
			...messages,
			{
				id: crypto.randomUUID(),
				role: "user",
				parts: [
					{
						type: "text",
						text: "Summarize what happened since your last digest.",
					},
				],
				createdAt: new Date(),
			},
		]);
		// At this point the LLM has responded and the message is persisted.
	}
}
import { AIChatAgent } from "@cloudflare/ai-chat";

export class DigestAgent extends AIChatAgent {
	async onChatMessage() {
		// ... your LLM call
	}

	async onStart() {
		await this.schedule("0 9 * * *", "dailyDigest");
	}

	async dailyDigest() {
		const stable = await this.waitUntilStable({ timeout: 30_000 });
		if (!stable) {
			console.warn("Conversation not stable, skipping daily digest");
			return;
		}

		await this.saveMessages((messages) => [
			...messages,
			{
				id: crypto.randomUUID(),
				role: "user",
				parts: [
					{
						type: "text",
						text: "Summarize what happened since your last digest.",
					},
				],
				createdAt: new Date(),
			},
		]);
		// At this point the LLM has responded and the message is persisted.
	}
}

saveMessages の関数形式 — saveMessages((messages) => [...]) — は、実行時点の最新の永続化済みメッセージを読みます。複数の呼び出しがキューに並んだとき(連続した webhook 到着など)に、古いベースラインを避けるためです。schedule() と cron 構文の詳細は タスクのスケジュール を参照してください。

キューの処理

トリガーを自分で制御できる場合、単純なループが一番わかりやすいパターンです。

async processQueue() {
	for (const task of this.taskQueue) {
		const stable = await this.waitUntilStable({ timeout: 30_000 });
		if (!stable) {
			console.warn("Conversation not stable, stopping queue processing");
			break;
		}

		await this.saveMessages((messages) => [
			...messages,
			{
				id: crypto.randomUUID(),
				role: "user",
				parts: [{ type: "text", text: task }],
				createdAt: new Date(),
			},
		]);
		// LLM has responded. this.messages is updated. Next iteration.
	}
	this.taskQueue = [];
}

特別なフックは不要です。saveMessages はターン全体が完了してから返ります。

メール起点

async onEmail(email: AgentEmail) {
	const stable = await this.waitUntilStable({ timeout: 30_000 });
	if (!stable) {
		console.warn("Conversation not stable, cannot process email");
		return;
	}

	const subject = email.headers.get("subject") ?? "(no subject)";
	const body = await new Response(email.raw).text();

	await this.saveMessages((messages) => [
		...messages,
		{
			id: crypto.randomUUID(),
			role: "user",
			parts: [
				{
					type: "text",
					text: `Email from ${email.from}: ${subject}\n\n${body}`,
				},
			],
			createdAt: new Date(),
		},
	]);
}

webhook 起点

async onRequest(request: Request): Promise<Response> {
	const url = new URL(request.url);

	if (url.pathname.endsWith("/webhook") && request.method === "POST") {
		const stable = await this.waitUntilStable({ timeout: 30_000 });
		if (!stable) {
			return new Response("Agent is busy", { status: 503 });
		}

		const payload = await request.json();
		try {
			await this.saveMessages((messages) => [
				...messages,
				{
					id: crypto.randomUUID(),
					role: "user",
					parts: [
						{
							type: "text",
							text: `Webhook event: ${JSON.stringify(payload)}`,
						},
					],
					createdAt: new Date(),
				},
			]);
			return new Response("ok");
		} catch (error) {
			console.error("Failed to process webhook:", error);
			return new Response("Internal error", { status: 500 });
		}
	}

	return super.onRequest(request);
}

webhook プロバイダーが素早い応答を期待する場合は、代わりに submitMessages() を使います。プロバイダーへ耐久的な受付通知を返し、同じべき等キーで安全に再試行できます。

async onRequest(request: Request): Promise<Response> {
	if (request.method !== "POST") return super.onRequest(request);

	const payload = await request.json<{ id: string }>();
	const submission = await this.submitMessages(
		[
			{
				id: crypto.randomUUID(),
				role: "user",
				parts: [
					{ type: "text", text: `Webhook event: ${JSON.stringify(payload)}` },
				],
			},
		],
		{ idempotencyKey: payload.id },
	);

	return Response.json({
		submissionId: submission.submissionId,
		accepted: submission.accepted,
		status: submission.status,
	});
}

応答を起動せずにコンテキストを注入する

次のターンで LLM に見せたいメッセージを、今はターンを始めずに追加するには persistMessages を使います。

async addBackgroundContext(data: string) {
	const stable = await this.waitUntilStable({ timeout: 30_000 });
	if (!stable) return;

	await this.persistMessages([
		...this.messages,
		{
			id: crypto.randomUUID(),
			role: "user",
			parts: [{ type: "text", text: `[Background context]: ${data}` }],
			createdAt: new Date(),
		},
	]);
	// Message is stored and broadcast to clients, but no LLM call happens.
}

自分で起動していない応答への反応

onChatResponse は、完了したすべてのターンのあとに発火します。ユーザー起点のメッセージ、saveMessages の呼び出し、自動継続です。起動方法に関係なく応答を観察・反応したい場合に使います。

状態のブロードキャスト

import { AIChatAgent } from "@cloudflare/ai-chat";

export class ChatAgent extends AIChatAgent {
	async onChatMessage() {
		// ... your LLM call
	}

	async onChatResponse(result) {
		if (result.status === "completed") {
			this.broadcast(JSON.stringify({ streaming: false }));
		}
	}
}
import { AIChatAgent, type ChatResponseResult } from "@cloudflare/ai-chat";

export class ChatAgent extends AIChatAgent {
	async onChatMessage() {
		// ... your LLM call
	}

	protected async onChatResponse(result: ChatResponseResult) {
		if (result.status === "completed") {
			this.broadcast(JSON.stringify({ streaming: false }));
		}
	}
}

アナリティクス

protected async onChatResponse(result: ChatResponseResult) {
	try {
		await fetch("https://analytics.example.com/event", {
			method: "POST",
			body: JSON.stringify({
				requestId: result.requestId,
				status: result.status,
				continuation: result.continuation,
			}),
		});
	} catch (error) {
		console.error("Analytics reporting failed:", error);
	}
}

連鎖的な推論

エージェントは自身の応答を確認し、続けるかを決められます。ユーザー起点のメッセージでも同じです。ユーザーが何を聞くかは予測できませんが、エージェントが何と言ったかには反応できます。

protected async onChatResponse(result: ChatResponseResult) {
	if (result.status !== "completed") return;

	const lastText = result.message.parts
		.filter((p) => p.type === "text")
		.map((p) => p.text)
		.join("");

	if (lastText.includes("[NEEDS_MORE_RESEARCH]")) {
		await this.saveMessages((messages) => [
			...messages,
			{
				id: crypto.randomUUID(),
				role: "user",
				parts: [{ type: "text", text: "Continue your research." }],
				createdAt: new Date(),
			},
		]);
	}
}

onChatResponse の内側から saveMessages を呼ぶと、内側のターンが完了するまで実行され、saveMessages が返ります。現在の onChatResponse が返ったあと、フレームワークは内側の応答に対して再度 onChatResponse を発火します。キューに作業がなくなるまで続きます。フレームワークが onChatResponse を入れ子にすることはありません。結果は順次処理されます。

リアクティブなキュー処理

キュー項目が外部イベント(ユーザーメッセージ、webhook)でいつでも追加される場合、onChatResponse を使うと、誰が起動したかに関係なく、応答のたびにキューを消化できます。

protected async onChatResponse(result: ChatResponseResult) {
	if (result.status === "completed" && this.taskQueue.length > 0) {
		const next = this.taskQueue.shift()!;
		await this.saveMessages((messages) => [
			...messages,
			{
				id: crypto.randomUUID(),
				role: "user",
				parts: [{ type: "text", text: next }],
				createdAt: new Date(),
			},
		]);
	}
}

ChatResponseResult のフィールド

フィールド 説明
message UIMessage 確定したアシスタントメッセージ
requestId string このターンの一意 ID
continuation boolean 自動継続なら true
status "completed" | "error" | "aborted" ターンの終了方法
error string | undefined status が "error" のときのエラー詳細

クライアント側: サーバー起点ストリームの検出

サーバーが saveMessages でストリームを起動すると、クライアントがリクエストを始めていないため、AI SDK の status"ready" のままです。useAgentChat フックは、これを扱う追加フラグを 2 つ提供します。

フラグ 追跡対象
status AI SDK のライフサイクル: "submitted""streaming""ready""error"。クライアント起点のリクエストのみ
isServerStreaming サーバー起点のストリームが稼働中なら true
isStreaming クライアントまたはサーバーのいずれかのストリーミングが稼働中なら true。共通インジケーターに使います

ほとんどの UI(送信ボタンの無効化、読み込み表示)には isStreaming を使います。isServerStreaming は、ユーザー起点とサーバー起点のストリームを区別する必要がある場合だけ使います(「Agent is working in the background...」のような別表示など)。

import { useAgent } from "agents/react";
import { useAgentChat } from "@cloudflare/ai-chat/react";

function Chat() {
	const agent = useAgent({ agent: "ChatAgent" });
	const { messages, sendMessage, isStreaming, isServerStreaming } =
		useAgentChat({ agent });

	return (
		<div>
			{messages.map((m) => (
				<div key={m.id}>{/* render message */}</div>
			))}

			{isServerStreaming && <div>Agent is working in the background...</div>}
			{!isServerStreaming && isStreaming && <div>Agent is responding...</div>}

			<form
				onSubmit={(e) => {
					e.preventDefault();
					const input = e.currentTarget.elements.namedItem(
						"input",
					) as HTMLInputElement;
					sendMessage({ text: input.value });
					input.value = "";
				}}
			>
				<input name="input" placeholder="Type a message..." />
				<button type="submit" disabled={isStreaming}>
					Send
				</button>
			</form>
		</div>
	);
}

ユーザーが待機中にサーバー駆動の応答が届くと、接続中のクライアントには新しいメッセージがリアルタイムで表示されます。ストリームの進行に合わせて isStreaming フラグは falsetruefalse と変わるため、送信ボタンなどの UI は自動的に無効化・再有効化されます。

messageConcurrency との関係

AIChatAgentmessageConcurrency 設定は、重なり合うユーザー送信の扱いを制御します("queue""latest""merge""drop""debounce")。この設定が適用されるのは sendMessage() だけです。クライアントからのユーザー起点メッセージです。

saveMessages() は、messageConcurrency の設定に関係なく、常に直列(キュー)動作です。そのため、サーバー駆動メッセージがドロップ、マージ、デバウンスされることはありません。常にキューに並び、順に実行されます。

他の Agent プリミティブとの組み合わせ

プリミティブ 組み合わせ方
schedule() saveMessages を呼ぶコールバックをスケジュールします。上の cron 例を参照してください
queue() 遅延処理のために saveMessages を呼ぶメソッドをキューします
startFiber() メッセージターン周辺のアプリケーション所有作業を耐久的に受け付け、確認します
runWorkflow() Workflow を開始します。AgentWorkflow.agent RPC で saveMessages または submitMessages を起動するメソッドを呼びます
onEmail() メール内容をチャットメッセージに変換し、saveMessages を呼びます
onRequest() webhook を処理し、saveMessages または submitMessages を呼びます
this.broadcast() onChatResponse からカスタム状態をブロードキャストします

サーバー駆動ターンのキャンセル

同じ Durable Object がターンを開始・制御する場合は、AbortSignal を渡します。

const controller = new AbortController();
const result = await this.saveMessages(
	[
		{
			id: crypto.randomUUID(),
			role: "user",
			parts: [{ type: "text", text: "Run the long analysis." }],
		},
	],
	{ signal: controller.signal },
);

if (result.status === "aborted") {
	// Partial chunks already streamed are persisted.
}
const controller = new AbortController();
const result = await this.saveMessages(
	[
		{
			id: crypto.randomUUID(),
			role: "user",
			parts: [{ type: "text", text: "Run the long analysis." }],
		},
	],
	{ signal: controller.signal },
);

if (result.status === "aborted") {
	// Partial chunks already streamed are persisted.
}

continueLastTurn() は同じ options.signal 引数を受け付けます。AbortSignal オブジェクトは Durable Object の RPC 境界を越えられず、シグナルはメモリ上のみです。ターンの途中で Durable Object がハイバネートすると、耐久復旧は通常、元のシグナルなしで継続します。ストリーム前の中断では、復旧は最新の未回答ユーザーメッセージを再試行できます。再起動後に発火した abort は、復旧済みターンには効きません。

キャンセルを再起動後も残す必要がある場合は、キャンセル意図を永続化します。onChatRecovery() でその状態を読み、{ continue: false } を返すと、再度のモデル呼び出しを止められます。

作業を submitMessages() で受け付けた場合、またはキャンセルが Worker と Durable Object の RPC 境界を越える必要がある場合は、耐久キャンセルに cancelSubmission(submissionId) を使います。

耐久単位を startFiber() で受け付け、キャンセルを Think ターンではなく周囲のアプリケーションジョブに適用する場合は cancelFiber(fiberId) を使います。

重要な注意点

  • saveMessages は await できます。 返ってきた時点で、LLM は応答済みで、メッセージは永続化されています。トリガーを自分で制御できる場合に使います。
  • saveMessages は関数形式を使います。 saveMessages((messages) => [...messages, newMsg]) は実行時点の最新の永続化済みメッセージを読むため、複数呼び出しがキューに並んでも古いベースラインを避けられます。
  • persistMessages は応答を起動しません。 コンテキストやシステムメッセージを静かに注入する場合に使います。
  • onChatResponse は、自分で起動していないターンへの反応用です。 ユーザー起点のメッセージ、自動継続、自分で saveMessages を呼んでいない任意のターンに使います。
  • onChatResponse は入れ子になりません。 onChatResponse の内側から saveMessages を呼ぶと、内側のターンが完了し、onChatResponse は再帰ではなく順次、再度発火します。
  • メッセージは onChatResponse の発火前に永続化されます。 フック実行中に Durable Object が退避しても、会話は SQLite 上で安全です。失われるのはフックのコールバックだけです。
  • 注入前に waitUntilStable() を呼びます。 進行中のストリームや保留中のツール操作との重複を避けるため、スケジュールコールバック、webhook、その他チャット以外のエントリポイントからは必ず呼びます。
  • クライアントは onChatResponse の実行前に、完了した応答を見ます。 サーバー側フックがクライアントを遅らせることはありません。
  • messageConcurrencysaveMessages に影響しません。 サーバー駆動メッセージは常にキューされ、順に実行されます。

次のステップ

Webhook

webhook イベントを受け取り、エージェントインスタンスへルーティングします。

役に立ちましたか?