Skip to content

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

インスタンスイベントを購読する

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

WorkflowInstance.subscribe() を使うと、status() をポーリングせずにイベントを受け取れます。購読は、購読前に記録されたイベントをまず配信します。保持していたイベントを配信したあと、購読はインスタンスの実行中に新しいイベントを待ちます。

インスタンス作成直後に購読できます。後から購読するには、get() でインスタンスを取得します。購読は インスタンス保持期間 のあいだ利用できます。

すべてのイベントを購読する

export default {
	async fetch(_request, env) {
		const instance = await env.MY_WORKFLOW.create({
			params: { reportId: "report-123" },
		});

		using subscription = await instance.subscribe();

		while (true) {
			const result = await subscription.next();
			if (result.done) {
				break;
			}

			console.log(result.value.type, result.value);
		}

		return Response.json({ instanceId: instance.id });
	},
};
interface Env {
	MY_WORKFLOW: Workflow;
}

export default {
	async fetch(_request: Request, env: Env) {
		const instance = await env.MY_WORKFLOW.create({
			params: { reportId: "report-123" },
		});

		using subscription = await instance.subscribe();

		while (true) {
			const result = await subscription.next();
			if (result.done) {
				break;
			}

			console.log(result.value.type, result.value);
		}

		return Response.json({ instanceId: instance.id });
	},
} satisfies ExportedHandler<Env>;

インスタンスが workflow_completedworkflow_errored、または workflow_terminated を発行すると、購読は終了します。終端イベントのあと、以降の next() 呼び出しは done: true を返します。

イベントをフィルタする

filter を設定すると、next() の結果を特定のイベント種類に限定できます。

const instance = await env.MY_WORKFLOW.get("report-123");

using subscription = await instance.subscribe({
	filter: ["workflow_completed", "workflow_errored", "workflow_terminated"],
});

const result = await subscription.next();
if (result.done) {
	throw new Error("The instance ended without a matching event.");
}

switch (result.value.type) {
	case "workflow_completed":
		console.log("Workflow output:", result.value.output);
		break;
	case "workflow_errored":
		console.error("Workflow errored:", result.value.error);
		break;
	case "workflow_terminated":
		console.log("Workflow terminated.");
		break;
}
const instance = await env.MY_WORKFLOW.get("report-123");

using subscription = await instance.subscribe({
	filter: ["workflow_completed", "workflow_errored", "workflow_terminated"],
});

const result = await subscription.next();
if (result.done) {
	throw new Error("The instance ended without a matching event.");
}

switch (result.value.type) {
	case "workflow_completed":
		console.log("Workflow output:", result.value.output);
		break;
	case "workflow_errored":
		console.error("Workflow errored:", result.value.error);
		break;
	case "workflow_terminated":
		console.log("Workflow terminated.");
		break;
}

フィルタが終端イベントを除外していても、購読は終了します。その場合、next() はイベントなしで done: true を返します。

カーソルから再開する

各イベントには eventId が含まれます。リモートプロシージャコール (RPC) が失敗したあとに再開するには、最後に処理したイベント ID を保存します。次に、その ID を cursor として渡します。

using subscription = await instance.subscribe({
	cursor: lastProcessedEventId,
	filter: ["step_completed", "workflow_completed", "workflow_errored"],
});

while (true) {
	const result = await subscription.next();
	if (result.done) {
		break;
	}

	await processEvent(result.value);
	await saveLastProcessedEventId(result.value.eventId);
}
using subscription = await instance.subscribe({
	cursor: lastProcessedEventId,
	filter: ["step_completed", "workflow_completed", "workflow_errored"],
});

while (true) {
	const result = await subscription.next();
	if (result.done) {
		break;
	}

	await processEvent(result.value);
	await saveLastProcessedEventId(result.value.eventId);
}

カーソルは最後に処理したイベントを識別します。購読は、eventId がカーソルより大きい最初のイベントから始まります。

機密の出力

機密としてマークしたステップでは、step_completed イベントの output"[REDACTED]" になります。

購読を破棄する

購読は Workers RPC リソースを保持します。購読を破棄すると、イベント配信が止まり、状態がクリアされ、リソースが解放されます。

スコープを抜けるときに自動破棄されるよう、購読を using で宣言するか、finally ブロックで subscription[Symbol.dispose]() を呼び出します。詳細は RPC ライフサイクル を参照してください。

イベントフィールド

公開型定義は、各イベントで利用できるフィールドを示します。

type WorkflowInstanceEvent = {
	instanceId: string;
	eventId: number;
	timestamp: number;
} & (
	| { type: "workflow_queued" }
	| { type: "workflow_started"; params?: unknown }
	| { type: "workflow_running" }
	| { type: "workflow_paused" }
	| { type: "workflow_waiting_for_pause" }
	| { type: "workflow_waiting" }
	| { type: "workflow_completed"; output?: unknown }
	| { type: "workflow_errored"; error: { name: string; message: string } }
	| { type: "workflow_terminated" }
	| {
			type: "step_started";
			stepName: string;
			config?: {
				retries: {
					limit: number;
					delay: WorkflowSleepDuration | "[dynamic]";
					backoff?: "constant" | "linear" | "exponential";
				};
				timeout: WorkflowSleepDuration;
				sensitive?: "output";
			};
	  }
	| { type: "step_completed"; stepName: string; output?: unknown }
	| { type: "step_errored"; stepName: string }
	| { type: "attempt_started"; stepName: string; attempt: number }
	| { type: "attempt_completed"; stepName: string; attempt: number }
	| {
			type: "attempt_errored";
			stepName: string;
			attempt: number;
			retryDelayMs?: number;
			error: { name: string; message: string };
	  }
	| { type: "sleep_started"; stepName: string; durationMs: number }
	| { type: "sleep_completed"; stepName: string }
	| { type: "wait_started"; stepName: string; eventType: string }
	| { type: "wait_completed"; stepName: string }
	| { type: "wait_timed_out"; stepName: string }
	| { type: "rollback_started" }
	| {
			type: "rollback_step_started";
			stepName: string;
			config?: {
				retries: {
					limit: number;
					delay: WorkflowSleepDuration | "[dynamic]";
					backoff?: "constant" | "linear" | "exponential";
				};
				timeout: WorkflowSleepDuration;
				sensitive?: "output";
			};
	  }
	| { type: "rollback_step_completed"; stepName: string }
	| {
			type: "rollback_step_errored";
			stepName: string;
			error: { name: string; message: string };
	  }
	| { type: "rollback_attempt_started"; stepName: string; attempt: number }
	| { type: "rollback_attempt_completed"; stepName: string; attempt: number }
	| {
			type: "rollback_attempt_errored";
			stepName: string;
			attempt: number;
			retryDelayMs?: number;
			error: { name: string; message: string };
	  }
	| { type: "rollback_completed" }
	| { type: "rollback_errored" }
);

次のセクションは、各イベントがいつ発行されるかを説明します。

Workflow ライフサイクルイベント

イベント種類 発行タイミング
workflow_queued インスタンスが実行キューに入る
workflow_started インスタンスが開始する
workflow_running インスタンスが実行を開始または再開する
workflow_paused インスタンスが一時停止する
workflow_waiting_for_pause インスタンスが一時停止前に現在の作業を待つ
workflow_waiting インスタンスが待機状態に入る
workflow_completed インスタンスが正常に完了する
workflow_errored インスタンスがエラーで終了する
workflow_terminated インスタンスが終了させられる

ステップと試行のイベント

イベント種類 発行タイミング
step_started step.do() 呼び出しが開始する
step_completed step.do() 呼び出しが完了する
step_errored step.do() 呼び出しがエラーになる
attempt_started ステップの試行が開始する
attempt_completed ステップの試行が完了する
attempt_errored ステップの試行がエラーになる

Sleep と wait のイベント

イベント種類 発行タイミング
sleep_started step.sleep() または step.sleepUntil() 呼び出しが開始する
sleep_completed sleep が終了する
wait_started step.waitForEvent() 呼び出しが開始する
wait_completed 一致するイベントが step.waitForEvent() に到達する
wait_timed_out step.waitForEvent() 呼び出しがタイムアウトする

Rollback イベント

イベント種類 発行タイミング
rollback_started Workflow が rollback を開始する
rollback_step_started rollback ハンドラーが開始する
rollback_step_completed rollback ハンドラーが完了する
rollback_step_errored rollback ハンドラーがエラーになる
rollback_attempt_started rollback の試行が開始する
rollback_attempt_completed rollback の試行が完了する
rollback_attempt_errored rollback の試行がエラーになる
rollback_completed 必要な rollback ハンドラーがすべて完了する
rollback_errored rollback 操作がエラーになる

メソッドシグネチャとオプション型は WorkflowInstance.subscribe() を参照してください。

役に立ちましたか?