人間の操作なしに、サーバー側からメッセージを送り、LLM の応答を起動します。定期フォローアップ、キュー処理、メール起点の応答、自律エージェントワークフローに使います。
通常のチャットフローでは、ユーザーがメッセージを送り、エージェントが応答します。ただしエージェントは、自分で動く必要もよくあります。定期リマインダーの発火、webhook の到着、ワークフローの完了、自身の応答を確認したあとの継続などです。
主なプリミティブは次のとおりです。
| プリミティブ | 役割 |
|---|---|
saveMessages |
メッセージを注入し、LLM を起動します。sendMessage のサーバー側相当です |
submitMessages |
Think ターンを非同期実行向けに耐久的に受け付け、あとから確認します |
startFiber |
ターン周辺のアプリケーション所有の副作用を耐久的に受け付けます |
persistMessages |
応答を起動せずにメッセージを保存します。コンテキストを静かに注入する場合に使います |
onChatResponse |
自分で起動していないものも含め、応答が完了したときに反応します |
isServerStreaming |
クライアント側フラグです。サーバー起点のストリームが稼働中なら true です |
saveMessages はメッセージを SQLite に永続化し、かつ 新しい LLM 応答のために onChatMessage を起動します。await できます。返ってきた時点で、LLM は応答済みで、メッセージは永続化されています。
persistMessages はメッセージを保存し、接続中のクライアントへ配信しますが、モデルターンは起動しません。会話にコンテキスト(システムメッセージやバックグラウンドデータなど)を注入し、応答は始めない場合に使います。
呼び出し側がモデルターンの完了を待てる場合は 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() を参照してください。
トリガーを自分で制御できる場合は saveMessages を使います。 スケジュールコールバック、webhook、メールハンドラー、メッセージ注入のタイミングを決める任意のメソッドです。
自分で起動していない応答に反応する必要がある場合は onChatResponse を使います。 ユーザー起点のメッセージ、ツール承認後の自動継続、フレームワークが代わりに実行した任意のターンです。
スケジュールコールバック、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 スケジュールはデフォルトでべき等なので、onStart で schedule() を呼んでも安全です。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(),
},
]);
}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(),
},
]);
}
}| フィールド | 型 | 説明 |
|---|---|---|
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 フラグは false → true → false と変わるため、送信ボタンなどの UI は自動的に無効化・再有効化されます。
AIChatAgent の messageConcurrency 設定は、重なり合うユーザー送信の扱いを制御します("queue"、"latest"、"merge"、"drop"、"debounce")。この設定が適用されるのは sendMessage() だけです。クライアントからのユーザー起点メッセージです。
saveMessages() は、messageConcurrency の設定に関係なく、常に直列(キュー)動作です。そのため、サーバー駆動メッセージがドロップ、マージ、デバウンスされることはありません。常にキューに並び、順に実行されます。
| プリミティブ | 組み合わせ方 |
|---|---|
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の実行前に、完了した応答を見ます。 サーバー側フックがクライアントを遅らせることはありません。 messageConcurrencyはsaveMessagesに影響しません。 サーバー駆動メッセージは常にキューされ、順に実行されます。