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],
})| Option | Type | Description |
|---|---|---|
connectionString | string | The engine creates a pg.Pool and closes it on stop(). |
pool | pg.Pool | Use an existing pool. You own its lifecycle. |
workflows | WorkflowDefinition[] | Registered when start() runs. |
logger | WorkflowLogger | Defaults to console.warn and console.error. |
boss | PgBoss | Your own pg-boss instance. When omitted, the engine creates one in the pgboss_v12_pgworkflow schema. |
Pass exactly one of connectionString and pool.
Lifecycle
| Method | Description |
|---|---|
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
| Method | Returns | Description |
|---|---|---|
startWorkflow(ref, input, options?) | WorkflowRun | Starts a run from a typed WorkflowRef. |
startWorkflow({ workflowId, input, resourceId?, idempotencyKey?, options? }) | WorkflowRun | Starts a run by ID. Throws WorkflowRunInProgressError for a singleton workflow that already has a run in progress. |
pauseWorkflow({ runId, resourceId? }) | WorkflowRun | Pauses a pending or running run. |
resumeWorkflow({ runId, resourceId?, options? }) | WorkflowRun | Resumes a paused run. Throws if the run is not paused. Does nothing while the run waits on a child workflow. |
cancelWorkflow({ runId, resourceId? }) | WorkflowRun | Cancels a pending, running, or paused run. |
triggerEvent({ runId, eventName, data?, resourceId?, options? }) | WorkflowRun | Delivers an event to a run waiting in step.waitFor. |
fastForwardWorkflow({ runId, data?, resourceId? }) | WorkflowRun | Completes the step the run is paused on. See Fast-forward. |
options on startWorkflow is a StartWorkflowOptions. options on resumeWorkflow and triggerEvent accepts { expireInSeconds }.
Queries
| Method | Returns | Description |
|---|---|---|
getRun({ runId, resourceId? }) | WorkflowRun | Throws 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 | null | The most recently created run of any status. |
getStats({ resourceId?, workflowId? }) | WorkflowRunStats | Run counts per status. |
checkProgress({ runId, resourceId? }) | WorkflowRunProgress | The 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:
checkProgresscan't count a workflow's steps without the definition. It returnstotalSteps: 0andcompletionPercentage: 0until the run completes, then100.completedStepsis accurate.singletonmust be set on the ref or instartWorkflowoptions, 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)| Parameter | Type | Description |
|---|---|---|
id | string | Unique workflow ID, up to 256 characters. |
handler | (context: WorkflowContext) => Promise<unknown> | The workflow body. Its return value becomes run.output. |
options | WorkflowOptions | inputSchema, 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.
| Option | Type | Description |
|---|---|---|
inputSchema | Standard Schema | Types and validates input. |
singleton | boolean | Applies 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.
| Field | Type | Description |
|---|---|---|
input | inferred from inputSchema | The run's input. unknown without a schema. |
step | StepBaseContext | Step methods, listed below. |
runId | string | The run's ID. |
workflowId | string | The workflow's ID. |
resourceId | string | undefined | The run's resource ID, if set. |
attempt | number | Zero-based retry attempt. Same as run.retryCount. |
timeline | Record<string, unknown> | Saved step results, keyed by step ID. |
logger | WorkflowLogger | The engine's logger. |
schedule | { timestamp: Date } | undefined | Set only for runs started by a schedule. |
Step methods
| Method | Returns | Guide |
|---|---|---|
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
| Field | Type | Description |
|---|---|---|
id | string | Run ID (KSUID). |
workflowId | string | |
resourceId | string | null | |
status | 'pending' | 'running' | 'paused' | 'completed' | 'failed' | 'cancelled' | Compare with WorkflowStatus values. |
input | unknown | |
output | unknown | null | The handler's return value, once completed. |
error | string | null | The last error message. |
currentStepId | string | |
timeline | Record<string, unknown> | Saved step results. |
retryCount / maxRetries | number | |
priority | number | Resolved priority. See Priorities. |
singleton | boolean | |
idempotencyKey | string | null | |
parentRunId, parentStepId, parentResourceId | string | null | Set on child runs. |
jobId | string | null | The pg-boss job ID. |
createdAt, updatedAt | Date | |
pausedAt, resumedAt, completedAt, timeoutAt, scheduledAt | Date | null |
WorkflowRunProgress is WorkflowRun plus completedSteps, totalSteps, and completionPercentage.
StartWorkflowOptions
| Option | Type | Description |
|---|---|---|
retries | number | Overrides the definition's retries. |
timeout | number (ms) | Saved as run.timeoutAt. Not enforced. See Retries and timeouts. |
priority | WorkflowPriority | Overrides the definition's priority. |
expireInSeconds | number | Per-execution limit for this run's first job. Defaults to WORKFLOW_RUN_EXPIRE_IN_SECONDS. |
singleton | boolean | For WorkflowClient, which can't read the definition. |
resourceId | string | Alternative to the top-level resourceId. |
idempotencyKey | string | Alternative 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.
| Class | Thrown when |
|---|---|
WorkflowEngineError | Base class. Has workflowId, runId, cause, and issues (input validation failures). |
WorkflowRunNotFoundError | The run doesn't exist, or doesn't match the given resourceId. |
WorkflowRunInProgressError | A 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.