耐久ジョブにモデル駆動の推論ステップが 1 つ必要なとき、ThinkWorkflow が Think と Cloudflare Workflows をつなぎます。
プロセスの主導権が Workflow 側にある場合に使います。
- 耐久的な複数ステップのオーケストレーション
- 承認ゲートや長時間の待機
- 再試行可能な決定的な副作用
- 型付きの構造化出力を返す Think ターン
繰り返しのプロンプトは スケジュールタスク に、単純な単発のバックグラウンドターンは submitMessages() に任せます。Workflows は、ステップ自体が重要なジョブ向けです。
@cloudflare/think/workflows からインポートします。
import { ThinkWorkflow } from "@cloudflare/think/workflows";ThinkWorkflow を継承し、run() 内で step.prompt() を呼び出します。
import { z } from "zod";
import { ThinkWorkflow } from "@cloudflare/think/workflows";
const draftSchema = z.object({
title: z.string(),
summary: z.string(),
labels: z.array(z.string()),
});
export class TriageWorkflow extends ThinkWorkflow {
async run(event, step) {
const draft = await step.prompt("triage-issue", {
prompt: `Triage issue #${event.payload.issueNumber}`,
output: draftSchema,
timeout: "3 days",
});
await step.do("apply-labels", async () => {
await this.agent.applyLabels(draft.labels);
});
}
}import { z } from "zod";
import { ThinkWorkflow } from "@cloudflare/think/workflows";
import type { ThinkWorkflowStep } from "@cloudflare/think/workflows";
import type { AgentWorkflowEvent } from "agents/workflows";
const draftSchema = z.object({
title: z.string(),
summary: z.string(),
labels: z.array(z.string()),
});
export class TriageWorkflow extends ThinkWorkflow<TriageAgent, Params> {
async run(event: AgentWorkflowEvent<Params>, step: ThinkWorkflowStep) {
const draft = await step.prompt("triage-issue", {
prompt: `Triage issue #${event.payload.issueNumber}`,
output: draftSchema,
timeout: "3 days",
});
await step.do("apply-labels", async () => {
await this.agent.applyLabels(draft.labels);
});
}
}Think Agent 内から runWorkflow() で Workflow を開始します。
export class TriageAgent extends Think {
async triageIssue(issueNumber) {
return this.runWorkflow(
"TRIAGE_WORKFLOW",
{ issueNumber },
{ metadata: { issueNumber } },
);
}
}export class TriageAgent extends Think<Env> {
async triageIssue(issueNumber: number): Promise<string> {
return this.runWorkflow(
"TRIAGE_WORKFLOW",
{ issueNumber },
{ metadata: { issueNumber } },
);
}
}runWorkflow() は Workflow インスタンスを作成し、run() 内で this.agent に再接続するために ThinkWorkflow が必要とする Agent の識別情報を注入します。Workflows バインディングを直接呼ぶより、こちらを使います。
// Avoid this for Agent workflows. It does not include Agent context.
await this.env.TRIAGE_WORKFLOW.create({ params: { issueNumber } });待機中の Workflow に、人の承認などの外部シグナルが必要なときは、Agent から sendWorkflowEvent() を使います。
await this.sendWorkflowEvent("TRIAGE_WORKFLOW", workflowId, {
type: "approval",
payload: { approved: true },
});step.prompt() は、プロンプト文字列と Zod オブジェクトスキーマを受け取ります。Workflow が Agent を呼ぶ前に、スキーマは JSON Schema へ変換されます。その後 Think は完全なエージェントターンを実行します。Agent は複数ステップにわたってツールを使え、スキーマに一致する引数で内部の final_answer ツールを呼ぶことで構造化結果を返します。これはストリーミングの response_format ではなく通常のツール呼び出しなので、Think が対応するすべてのプロバイダーで動きます。ストリーミングリクエストで JSON Schema 応答を拒否する Workers AI も含みます。Workflow が再開すると、型付きの値を返す前に、元の Zod スキーマでペイロードを再検証します。
JSON Schema に表せない未対応の Zod 機能は、プロンプトステップの作成時に失敗します。Think は不正なモデル出力を黙って修復しません。モデルが有効な final_answer 呼び出しを出さない場合、サブミッションは終端のエラー状態になり、step.prompt() は例外を投げます。
- Agent は先にツールを使えます。
step.prompt()のターンは完全なエージェントターンです。Agent は複数ステップにわたって自身のツールを呼び、その後に final-answer ツールを呼べます。回答前にツールを使う想定なら、少なくともmaxSteps: 2を許可します。maxSteps: 1では最初のステップで回答せざるを得ず、ほかのツールは呼べません。 - 構造化ターンではツール使用が強制されます。 プレーンテキストではなく構造化回答で終えるよう、Think はターンに
toolChoiceを設定します。step.prompt()ターンのbeforeTurnからtoolChoiceを上書きしないでください。上書きすると Agent が final-answer ツールを呼べなくなり、プロンプトが失敗します。 think_final_answerは予約済みです。 Think は構造化結果を運ぶ内部ツールthink_final_answerを注入します。この名前(およびthink_final_answer_*の派生名)は予約されています。呼び出しと結果は永続化された会話から取り除かれるため、トランスクリプトや後続ターンに Think の内部処理は見えません。- モデルはストリーミング中のツール呼び出しに対応している必要があります。 Think はすべてのターンをストリーミングするため、
step.prompt()が動くのは、ストリーミング中に強制ツール呼び出しを確実に出せるモデルだけです。ツール呼び出しが強いモデル(例: OpenAIgpt-4o-mini、Anthropicclaude-haiku-4-5、Workers AI@cf/moonshotai/kimi-k2.6)は動作確認済みです。一部のモデルは、非ストリーミングリクエストでのみ強制toolChoiceを守り、ストリーミング中はプレーンテキストで返信して止まります。例: Workers AI@cf/meta/llama-3.3-70b-instruct-fp8-fast。こうしたモデルではターンがthink_final_answer呼び出しなしで終わり、step.prompt()は失敗します(Model ended the turn without calling the think_final_answer tool)。代わりに、ストリーミング中のツール呼び出しが動くモデルを使います。
呼び出しはブロッキングステップのように読めますが、長寿命の Durable Object RPC は開きっぱなしにしません。
step.do("<name>:submit", ...)が、べき等な Think サブミッションを作成または検索します。- Think は、通常のサブミッションキュー経由で提出されたターンを実行します。
- サブミッションが
completed、error、aborted、skippedになると、Think は保留中の workflow 通知を記録します。 - Think は
sendWorkflowEvent()と Durable Object アラームで通知アウトボックスを排出し、配信が成功するまで続けます。 step.waitForEvent("<name>:wait", ...)が Workflow を再開します。step.prompt()は構造化出力を検証するか、型付きエラーを投げます。
機械可読な出力は、保留中の通知と Workflow イベントのペイロードに載ります。Think はサブミッション台帳に別の output_json 列を持たず、配信後に通知ペイロードを消します。配信後の耐久結果の所有権は Workflow 側にあります。
デフォルトでは、step.prompt() は Workflow の識別情報とステップ名からべき等キーを推論します。
think-workflow:<workflowName>:<workflowId>:<stepName>ループでは、同じステップ名の繰り返しを区別するため、文字列の key を渡します。
await step.prompt("summarize-file", {
key: file.path,
prompt: `Summarize ${file.path}`,
output: summarySchema,
});プロンプト本文は推論キーに含まれません。ただし Think は診断用に、workflow メタデータとプロンプト/設定のフィンガープリントを保存します。
端末イベントを Workflow が待つ時間は timeout で制御します。待機がタイムアウトすると、step.prompt() はデフォルトで Think サブミッションをキャンセルし、ThinkPromptTimeoutError を投げます。
Workflow が待機をやめたあとも Think サブミッションを続けたい場合は、cancelOnTimeout: false を設定します。
繰り返しのプロンプトサブミッションや、決定的なスケジュールハンドラーには getScheduledTasks() を使います。
getScheduledTasks() {
return {
dailySummary: {
schedule: "every day at 09:00",
timezone: "UTC",
prompt: "Generate the daily report."
},
dailyWorkflow: {
schedule: "every day at 09:00",
timezone: "UTC",
retry: { maxAttempts: 3 },
handler: async ({ idempotencyKey, scheduledFor, timezone }) => {
await this.env.REPORT_WORKFLOW.create({
id: idempotencyKey,
params: { scheduledFor, timezone }
});
}
}
};
}呼び出し側があとからサブミッション状態を確認できる、耐久的な単発ターンには submitMessages() を使います。
Agent 内での復旧が必要な、アプリ所有のべき等な Agent ジョブには startFiber() を使います。Think の workflow 通知配信はファイバーを使いません。配信成功までイベントを保持する必要があるため、専用のアウトボックスを使います。
プロセスに複数の決定的ステップ、長時間の待機、人の承認がある場合は Workflows を使います。