pg-workflows

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,
  },
)
OptionTypeDefaultDescription
inputSchemaStandard SchemanoneValidates input at startWorkflow and types it in the handler. See Input validation.
retriesnumber0Retry attempts after a failure. See Retries and timeouts.
timeoutnumber (ms)noneRecorded as run.timeoutAt. See Retries and timeouts.
priority'high' | 'normal' | 'low' | number'normal'Queue priority. See Priorities.
singletonbooleanfalseAt most one pending or running run. See Singleton workflows.
schedulecron string, duration string, or duration objectnoneStarts runs on a recurring schedule. See Recurring schedules.
timezoneIANA 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. A Date comes back as an ISO string when the step is replayed, and class instances lose their prototype. Return plain objects, and use toISOString() 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).

  1. Start. startWorkflow inserts a workflow_runs row with status running and enqueues a job.
  2. Execute. A worker picks up the job and calls the handler.
  3. Save each step. Each step.run result is written to the run's timeline before the handler moves on.
  4. Pause. waitFor, pause, delay, waitUntil, poll, and invokeChildWorkflow save the run as paused and end the execution. No worker or connection is held while it waits.
  5. 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.
  6. Finish. The handler's return value is saved as output and the run is completed. If the handler throws, the run is retried up to retries times, then marked failed.

On this page