pg-workflows

API reference

WorkflowEngine, WorkflowClient, workflow(), refs, context, types, and errors.

WorkflowEngine

Registers workflows, runs workers, and executes steps. Use it in a worker service or in a single service that both starts and runs workflows.

import { WorkflowEngine } from 'pg-workflows'

Constructor

const engine = new WorkflowEngine({
  connectionString: 'postgres://postgres:postgres@localhost:5432/postgres',
  workflows: [sendInvoice],
})
OptionTypeDescription
connectionStringstringThe engine creates a pg.Pool and closes it on stop().
poolpg.PoolUse an existing pool. You own its lifecycle.
workflowsWorkflowDefinition[]Registered when start() runs.
loggerWorkflowLoggerDefaults to console.warn and console.error.
bossPgBossYour own pg-boss instance. When omitted, the engine creates one in the pgboss_v12_pgworkflow schema.

Pass exactly one of connectionString and pool.

Lifecycle

MethodDescription
start(asEngine = true, { batchSize = 1, heartbeatSeconds = 30 }?)Starts pg-boss, runs migrations, registers workflows, starts WORKFLOW_RUN_WORKERS workers, and registers schedules. With asEngine: false, it does everything except start workers and schedules. Calling it again is a no-op.
stop()Unschedules recurring workflows, stops pg-boss, and closes the pool if the engine created it.
registerWorkflow(definition)Parses the handler's steps and registers it. Throws if the ID is already registered or the schedule is invalid.
unregisterWorkflow(workflowId)Removes a definition and its schedule.
unregisterAllWorkflows()Removes every definition and schedule.

If you call startWorkflow before start(), the engine calls start(false): the run is created and enqueued, but this engine runs no workers.

Runs

MethodReturnsDescription
startWorkflow(ref, input, options?)WorkflowRunStarts a run from a typed WorkflowRef.
startWorkflow({ workflowId, input, resourceId?, idempotencyKey?, options? })WorkflowRunStarts a run by ID. Throws WorkflowRunInProgressError for a singleton workflow that already has a run in progress.
pauseWorkflow({ runId, resourceId? })WorkflowRunPauses a pending or running run.
resumeWorkflow({ runId, resourceId?, options? })WorkflowRunResumes a paused run. Throws if the run is not paused. Does nothing while the run waits on a child workflow.
cancelWorkflow({ runId, resourceId? })WorkflowRunCancels a pending, running, or paused run.
triggerEvent({ runId, eventName, data?, resourceId?, options? })WorkflowRunDelivers an event to a run waiting in step.waitFor.
fastForwardWorkflow({ runId, data?, resourceId? })WorkflowRunCompletes the step the run is paused on. See Fast-forward.

options on startWorkflow is a StartWorkflowOptions. options on resumeWorkflow and triggerEvent accepts { expireInSeconds }.

Queries

MethodReturnsDescription
getRun({ runId, resourceId? })WorkflowRunThrows WorkflowRunNotFoundError if the run doesn't exist or belongs to another resource.
getRuns({ resourceId?, workflowId?, statuses?, limit?, startingAfter?, endingBefore? }){ items, nextCursor, prevCursor, hasMore, hasPrev }Lists runs, newest first. limit defaults to 20. Pass nextCursor as startingAfter for the next page, and prevCursor as endingBefore for the previous one.
getWorkflowLastRun({ workflowId, resourceId? })WorkflowRun | nullThe most recently created run of any status.
getStats({ resourceId?, workflowId? })WorkflowRunStatsRun counts per status.
checkProgress({ runId, resourceId? })WorkflowRunProgressThe run plus completedSteps, totalSteps, and completionPercentage. Throws if the workflow is not registered on this engine.

Every method that takes resourceId also matches on it. See Resource ID.

WorkflowClient

Starts and manages runs without loading workflow handlers. Use it in API services that should not import worker code.

import { WorkflowClient } from 'pg-workflows/client'

const client = new WorkflowClient({
  connectionString: 'postgres://postgres:postgres@localhost:5432/postgres',
})

The constructor takes connectionString or pool, plus optional logger and boss, with the same meaning as on WorkflowEngine.

The client has the same run and query methods as the engine, except getWorkflowLastRun. It never executes steps. It also has start(), which connects and runs migrations (called automatically on first use), and stop().

Differences from the engine:

  • checkProgress can't count a workflow's steps without the definition. It returns totalSteps: 0 and completionPercentage: 0 until the run completes, then 100. completedSteps is accurate.
  • singleton must be set on the ref or in startWorkflow options, because the client can't read it from the definition. See Singleton workflows.

workflow()

import { workflow } from 'pg-workflows'

const definition = workflow(id, handler, options)
ParameterTypeDescription
idstringUnique workflow ID, up to 256 characters.
handler(context: WorkflowContext) => Promise<unknown>The workflow body. Its return value becomes run.output.
optionsWorkflowOptionsinputSchema, retries, timeout, priority, singleton, schedule, timezone. See the options table.

workflow.use(plugin)

Returns a new workflow factory whose handlers get the plugin's step methods and wrap middleware. Plugins compose in the order you call .use(). @pg-workflows/otel is built on this.

import { otelPlugin } from '@pg-workflows/otel'
import { workflow } from 'pg-workflows'

const tracedWorkflow = workflow.use(otelPlugin())

workflow.ref(id, options?)

Same as createWorkflowRef, with the generics in <TInput, TOutput> order.

WorkflowRef

A workflow ID and input schema with no handler. Import refs in API services to start runs with typed input. Call a ref with a handler to get a full definition in the worker.

import { createWorkflowRef } from 'pg-workflows/client'
import { z } from 'zod'

export const sendInvoiceRef = createWorkflowRef('send-invoice', {
  inputSchema: z.object({ orderId: z.string(), total: z.number() }),
})

createWorkflowRef(id, options?) is exported from both pg-workflows and pg-workflows/client.

OptionTypeDescription
inputSchemaStandard SchemaTypes and validates input.
singletonbooleanApplies the singleton constraint when starting from WorkflowClient.

API service:

await client.startWorkflow(sendInvoiceRef, { orderId: 'ord_1', total: 99 })

Worker service:

const sendInvoice = sendInvoiceRef(
  async ({ step, input }) => {
    return step.run('create-invoice', async () => ({ id: `inv_${input.orderId}` }))
  },
  { retries: 3 },
)

The second argument accepts every workflow() option except inputSchema.

To type the output of step.invokeChildWorkflow, pass it as the first generic of createWorkflowRef. Explicit generics turn off inference for the others, so pass the schema type too:

const receiptSchema = z.object({ orderId: z.string() })

const sendReceiptRef = createWorkflowRef<{ receiptId: string }, typeof receiptSchema>('send-receipt', {
  inputSchema: receiptSchema,
})

WorkflowContext

The object passed to a handler.

FieldTypeDescription
inputinferred from inputSchemaThe run's input. unknown without a schema.
stepStepBaseContextStep methods, listed below.
runIdstringThe run's ID.
workflowIdstringThe workflow's ID.
resourceIdstring | undefinedThe run's resource ID, if set.
attemptnumberZero-based retry attempt. Same as run.retryCount.
timelineRecord<string, unknown>Saved step results, keyed by step ID.
loggerWorkflowLoggerThe engine's logger.
schedule{ timestamp: Date } | undefinedSet only for runs started by a schedule.

Step methods

MethodReturnsGuide
step.run(stepId, fn)Promise<T>Steps
step.waitFor(stepId, { eventName, schema? })Promise<T>Events
step.waitFor(stepId, { eventName, timeout, schema? })Promise<T | undefined>Events
step.waitUntil(stepId, date | isoString | { date })Promise<void>Timers
step.delay(stepId, duration)Promise<void>Timers
step.sleep(stepId, duration)Promise<void>Alias for delay
step.pause(stepId)Promise<void>Pause and resume
step.poll(stepId, fn, { interval?, timeout? })Promise<{ timedOut: false, data: T } | { timedOut: true }>Polling
step.invokeChildWorkflow(stepId, ref, input, options?)Promise<TOutput>Child workflows
step.invokeChildWorkflow<TOutput>(stepId, { workflowId, input, resourceId?, idempotencyKey?, options? })Promise<TOutput>Child workflows

A duration is a string ('90s', '2h', '3 days') or { weeks?, days?, hours?, minutes?, seconds? }.

Types

WorkflowRun

FieldTypeDescription
idstringRun ID (KSUID).
workflowIdstring
resourceIdstring | null
status'pending' | 'running' | 'paused' | 'completed' | 'failed' | 'cancelled'Compare with WorkflowStatus values.
inputunknown
outputunknown | nullThe handler's return value, once completed.
errorstring | nullThe last error message.
currentStepIdstring
timelineRecord<string, unknown>Saved step results.
retryCount / maxRetriesnumber
prioritynumberResolved priority. See Priorities.
singletonboolean
idempotencyKeystring | null
parentRunId, parentStepId, parentResourceIdstring | nullSet on child runs.
jobIdstring | nullThe pg-boss job ID.
createdAt, updatedAtDate
pausedAt, resumedAt, completedAt, timeoutAt, scheduledAtDate | null

WorkflowRunProgress is WorkflowRun plus completedSteps, totalSteps, and completionPercentage.

StartWorkflowOptions

OptionTypeDescription
retriesnumberOverrides the definition's retries.
timeoutnumber (ms)Saved as run.timeoutAt. Not enforced. See Retries and timeouts.
priorityWorkflowPriorityOverrides the definition's priority.
expireInSecondsnumberPer-execution limit for this run's first job. Defaults to WORKFLOW_RUN_EXPIRE_IN_SECONDS.
singletonbooleanFor WorkflowClient, which can't read the definition.
resourceIdstringAlternative to the top-level resourceId.
idempotencyKeystringAlternative to the top-level idempotencyKey.

WorkflowStatus

enum WorkflowStatus {
  PENDING = 'pending',
  RUNNING = 'running',
  PAUSED = 'paused',
  COMPLETED = 'completed',
  FAILED = 'failed',
  CANCELLED = 'cancelled',
}

WorkflowRunStats is Record<WorkflowStatus, number>.

WorkflowPriority

'high' | 'normal' | 'low' | number. Named levels map to 100, 0, and -100.

WorkflowLogger

interface WorkflowLogger {
  log(message: string): void
  error(message: string, ...args: unknown[]): void
}

Errors

All three are exported from pg-workflows and pg-workflows/client.

ClassThrown when
WorkflowEngineErrorBase class. Has workflowId, runId, cause, and issues (input validation failures).
WorkflowRunNotFoundErrorThe run doesn't exist, or doesn't match the given resourceId.
WorkflowRunInProgressErrorA singleton workflow already has a pending or running run.

UI

React components and hooks for browsing runs are in the separate @pg-workflows/ui package. They talk to your API over HTTP, which calls getRuns and getRun. Nothing from @pg-workflows/ui is exported by pg-workflows.

On this page