Use WorkflowInstance.subscribe() to receive events without polling status(). A subscription first delivers events recorded before you subscribed. After delivering these retained events, the subscription waits for new events as the instance runs.
You can subscribe immediately after creating an instance. To subscribe later, retrieve the instance with get(). Subscriptions remain available during the instance retention period.
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>;The subscription ends when the instance emits workflow_completed, workflow_errored, or workflow_terminated. After a terminal event, each later next() call returns done: true.
Set filter to limit next() results to specific event types.
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;
}A subscription ends even when its filter excludes a terminal event. In that case, next() returns done: true without the event.
Each event includes an eventId. To resume after a remote procedure call (RPC) fails, store the last processed event ID. Then pass that ID as 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);
}The cursor identifies the last processed event. The subscription starts with the first event whose eventId is greater than the cursor.
For steps marked as sensitive, the step_completed event sets output to "[REDACTED]".
A subscription holds a Workers RPC resource. Disposing the subscription stops event delivery, clears its state, and releases its resources.
Declare the subscription with using for automatic disposal when the scope exits, or call subscription[Symbol.dispose]() in a finally block. For more information, refer to RPC lifecycle.
The public type definition shows the fields available on each event:
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" }
);The following sections describe when each event is emitted.
| Event type | Emitted when |
|---|---|
workflow_queued |
The instance enters the execution queue |
workflow_started |
The instance starts |
workflow_running |
The instance starts or resumes execution |
workflow_paused |
The instance pauses |
workflow_waiting_for_pause |
The instance waits for current work before pause |
workflow_waiting |
The instance enters waiting state |
workflow_completed |
The instance completes successfully |
workflow_errored |
The instance ends with an error |
workflow_terminated |
The instance is terminated |
| Event type | Emitted when |
|---|---|
step_started |
A step.do() call starts |
step_completed |
A step.do() call completes |
step_errored |
A step.do() call errors |
attempt_started |
A step attempt starts |
attempt_completed |
A step attempt completes |
attempt_errored |
A step attempt errors |
| Event type | Emitted when |
|---|---|
sleep_started |
A step.sleep() or step.sleepUntil() call starts |
sleep_completed |
A sleep finishes |
wait_started |
A step.waitForEvent() call starts |
wait_completed |
A matching event reaches step.waitForEvent() |
wait_timed_out |
A step.waitForEvent() call times out |
| Event type | Emitted when |
|---|---|
rollback_started |
The Workflow starts a rollback |
rollback_step_started |
A rollback handler starts |
rollback_step_completed |
A rollback handler completes |
rollback_step_errored |
A rollback handler errors |
rollback_attempt_started |
A rollback attempt starts |
rollback_attempt_completed |
A rollback attempt completes |
rollback_attempt_errored |
A rollback attempt errors |
rollback_completed |
All required rollback handlers complete |
rollback_errored |
The rollback operation errors |
For method signatures and option types, refer to WorkflowInstance.subscribe().