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_completed、workflow_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_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_started |
step.sleep() または step.sleepUntil() 呼び出しが開始する |
sleep_completed |
sleep が終了する |
wait_started |
step.waitForEvent() 呼び出しが開始する |
wait_completed |
一致するイベントが step.waitForEvent() に到達する |
wait_timed_out |
step.waitForEvent() 呼び出しがタイムアウトする |
| イベント種類 | 発行タイミング |
|---|---|
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() を参照してください。