Workflows and steps
Define a workflow, configure it, and make each unit of work durable with step.run.
The snippets on this page assume an engine like the one in the quickstart, with each workflow passed to workflows: [...].
A workflow is an async function with an ID. Every durable operation inside it goes through step.
import { workflow } from 'pg-workflows'
import { z } from 'zod'
export const sendInvoice = workflow(
'send-invoice',
async ({ step, input }) => {
const invoice = await step.run('create-invoice', async () => {
return { id: `inv_${input.orderId}`, total: input.total }
})
return invoice
},
{
inputSchema: z.object({ orderId: z.string(), total: z.number() }),
retries: 3,
},
)| Option | Type | Default | Description |
|---|---|---|---|
inputSchema | Standard Schema | none | Validates input at startWorkflow and types it in the handler. See Input validation. |
retries | number | 0 | Retry attempts after a failure. See Retries and timeouts. |
timeout | number (ms) | none | Recorded as run.timeoutAt. See Retries and timeouts. |
priority | 'high' | 'normal' | 'low' | number | 'normal' | Queue priority. See Priorities. |
singleton | boolean | false | At most one pending or running run. See Singleton workflows. |
schedule | cron string, duration string, or duration object | none | Starts runs on a recurring schedule. See Recurring schedules. |
timezone | IANA zone | 'UTC' | Time zone for cron schedules. |
The handler receives input, step, runId, workflowId, resourceId, attempt (zero-based retry count), timeline, logger, and schedule. The full type is in the API reference.
The handler runs from the top each time the run resumes after a pause or a retry. Completed steps return their saved result without running again, so the code between steps must be deterministic. Put anything with side effects, randomness, or the current time inside a step.
Steps
step.run executes a function once and saves its return value on the run.
const user = await step.run('create-user', async () => {
return { id: 'usr_123', email: input.email }
})- IDs must be unique within a workflow. In loops, build the ID from the item:
step.run(`charge-$\{order.id\}`, ...). - Return JSON-serializable values. Results are stored as
jsonb. ADatecomes back as an ISO string when the step is replayed, and class instances lose their prototype. Return plain objects, and usetoISOString()for dates. - A step can run more than once if the process crashes mid-step, before its result is saved. Make external side effects idempotent, for example by passing an idempotency key to your payment provider.
How a run executes
PostgreSQL is both the job queue (through pg-boss) and the state store (workflow_runs).
- Start.
startWorkflowinserts aworkflow_runsrow with statusrunningand enqueues a job. - Execute. A worker picks up the job and calls the handler.
- Save each step. Each
step.runresult is written to the run'stimelinebefore the handler moves on. - Pause.
waitFor,pause,delay,waitUntil,poll, andinvokeChildWorkflowsave the run aspausedand end the execution. No worker or connection is held while it waits. - Resume. An event,
resumeWorkflow, a timer, or a finished child enqueues a new job. The handler runs from the top, and completed steps return their saved results. - Finish. The handler's return value is saved as
outputand the run iscompleted. If the handler throws, the run is retried up toretriestimes, then markedfailed.