# Architectures (/docs/architectures) pg-workflows runs in one of two architectures: * **Single service.** One process starts runs and executes them. * **Microservices (web and worker).** Web or API services start and manage runs with a lightweight client. Worker services load the handlers and execute steps. Both architectures use the same database, and you can move from one to the other without migrating data. ## Single service [#single-service] One process registers workflows, runs workers, and starts runs: ```typescript import { WorkflowEngine, workflow } from 'pg-workflows' import { z } from 'zod' const onboardUser = workflow( 'onboard-user', async ({ step, input }) => { const user = await step.run('create-account', async () => { return { id: `usr_${input.email}`, email: input.email } }) await step.run('send-welcome', async () => { console.log(`Sending welcome email to ${user.email}`) return { sent: true } }) return { userId: user.id } }, { inputSchema: z.object({ email: z.email() }) }, ) const engine = new WorkflowEngine({ connectionString: process.env.DATABASE_URL ?? 'postgres://postgres:postgres@localhost:5432/postgres', workflows: [onboardUser], }) await engine.start() // Call this from a request handler const run = await engine.startWorkflow({ workflowId: 'onboard-user', input: { email: 'alice@example.com' }, }) ``` To scale, run more copies of the process. Every engine that has called `start()` pulls runs from the same queue. Each engine runs [`WORKFLOW_RUN_WORKERS`](/docs/reference/configuration#environment-variables) workers (default 3). ## Microservices (web and worker) [#microservices-web-and-worker] Use this architecture when the web or API service shouldn't load worker dependencies such as LLM SDKs or heavy processing libraries. Both services share a file of **workflow refs**: an ID and input schema, with no handler code. ``` shared/workflows.ts (refs: ID + input schema) │ │ ▼ ▼ ┌───────────────┐ ┌───────────────┐ │ API service │ │ Worker service│ │ WorkflowClient│ │ WorkflowEngine│ │ start, pause, │ │ + handlers │ │ resume, query │ │ execute steps │ └───────┬───────┘ └───────┬───────┘ │ │ └────────┬────────┘ ▼ PostgreSQL ``` ### 1. Define refs [#1-define-refs] ```typescript // shared/workflows.ts import { createWorkflowRef } from 'pg-workflows/client' import { z } from 'zod' export const onboardUser = createWorkflowRef('onboard-user', { inputSchema: z.object({ email: z.email() }), }) ``` ### 2. Start runs from the API service [#2-start-runs-from-the-api-service] ```typescript // api-service.ts import { WorkflowClient } from 'pg-workflows/client' import { onboardUser } from './shared/workflows' const client = new WorkflowClient({ connectionString: process.env.DATABASE_URL ?? 'postgres://postgres:postgres@localhost:5432/postgres', }) // `input` is typed and validated against the ref's schema const run = await client.startWorkflow(onboardUser, { email: 'alice@example.com' }) const current = await client.getRun({ runId: run.id }) console.log(current.status) ``` ### 3. Execute runs in the worker service [#3-execute-runs-in-the-worker-service] ```typescript // worker-service.ts import { WorkflowEngine } from 'pg-workflows' import { onboardUser } from './shared/workflows' const onboardUserDefinition = onboardUser(async ({ step, input }) => { const user = await step.run('create-account', async () => { return { id: `usr_${input.email}`, email: input.email } }) await step.run('send-welcome', async () => { console.log(`Sending welcome email to ${user.email}`) return { sent: true } }) return { userId: user.id } }) const engine = new WorkflowEngine({ connectionString: process.env.DATABASE_URL ?? 'postgres://postgres:postgres@localhost:5432/postgres', workflows: [onboardUserDefinition], }) await engine.start() ``` `pg-workflows/client` contains the client, refs, types, and errors. It does not include the engine or the handler parser. The client doesn't load definitions, which has two consequences: * **Singleton.** Set `singleton: true` on the ref, not only on the worker's definition. See [Singleton workflows](/docs/concepts/singleton-workflows). * **Progress.** `client.checkProgress` can't report `totalSteps` or a percentage for an in-progress run. Use `getRun` and read `status` and `currentStepId`, or call `checkProgress` on an engine that has the workflow registered. A runnable version of this architecture is in [`examples/node/microservices`](https://github.com/SokratisVidros/pg-workflows/tree/main/examples/node/microservices). ## Dashboard [#dashboard] The engine has no built-in UI. [`@pg-workflows/ui`](/docs/ui) is a separate package of React components and hooks for browsing runs and step timelines. ``` ┌───────────┐ HTTP ┌──────────────┐ pg-workflows ┌────────────┐ │ Browser │ ────────► │ Your API │ ───────────────► │ PostgreSQL │ │ React + │ │ getRuns and │ WorkflowEngine │ │ │ ui pkg │ │ getRun │ or Client │ │ └───────────┘ └──────────────┘ └────────────┘ ``` The browser never connects to the database. Expose list and get endpoints on top of `getRuns` and `getRun`, then point the UI's `WorkflowRunsProvider` at them. Setup is in the [UI package README](/docs/ui). Keep `@pg-workflows/ui` out of worker services. Its peer dependencies are React, Tailwind, and TanStack Query. The engine's only peer dependency is `pg`. # Child workflows (/docs/concepts/child-workflows) `step.invokeChildWorkflow` starts another workflow, pauses the parent, and returns the child's output when the child completes. The parent holds no worker while it waits. ```typescript import { createWorkflowRef, workflow } from 'pg-workflows' import { z } from 'zod' type ReceiptOutput = { receiptId: string } const receiptInput = z.object({ orderId: z.string() }) // Explicit generics turn off inference, so pass the schema type as the second one const sendReceiptRef = createWorkflowRef('send-receipt', { inputSchema: receiptInput, }) export const sendReceipt = sendReceiptRef(async ({ step, input }) => { return step.run('email-receipt', async () => ({ receiptId: `rcpt_${input.orderId}` })) }) export const checkout = workflow( 'checkout', async ({ step, input }) => { const receipt = await step.invokeChildWorkflow('send-receipt', sendReceiptRef, { orderId: input.orderId, }) return { receiptId: receipt.receiptId } // receipt: ReceiptOutput }, { inputSchema: z.object({ orderId: z.string() }) }, ) ``` Register both `checkout` and `sendReceipt` with the engine. You can also invoke by ID and type the output with a generic: ```typescript const receipt = await step.invokeChildWorkflow('send-receipt', { workflowId: 'send-receipt', input: { orderId: input.orderId }, }) ``` Behavior: * The child starts once per parent step. Its output is saved on the parent like any step result. * If the child fails or is cancelled, the parent step throws, and the parent's own `retries` apply. * The child inherits the parent's priority unless the call or the child's definition sets one. * **Cancelling the parent does not cancel the child.** The child runs to its own end state. The same applies when the parent fails, completes, or times out while the child is running. * `resumeWorkflow()` and `fastForwardWorkflow()` do nothing while the parent waits on a child. # Events (/docs/concepts/events) `step.waitFor` pauses the run until an event with the matching name arrives. The paused run holds no worker. ```typescript const payment = await step.waitFor('wait-for-payment', { eventName: 'payment-completed', schema: z.object({ amount: z.number() }), }) // payment: { amount: number } ``` Send the event from anywhere that has an engine or client: ```typescript await engine.triggerEvent({ runId: run.id, eventName: 'payment-completed', data: { amount: 99 }, }) ``` With `timeout` (in milliseconds), the step resolves to `undefined` if no event arrives in time: ```typescript const payment = await step.waitFor('wait-for-payment', { eventName: 'payment-completed', timeout: 24 * 60 * 60 * 1000, }) if (!payment) { return { status: 'expired' } } ``` `schema` sets the TypeScript type of the result. It is not checked at runtime, so validate `data` before calling `triggerEvent` if it comes from an untrusted source. Without `schema`, the result is `unknown`. # Input validation (/docs/concepts/input-validation) `inputSchema` accepts any [Standard Schema](https://github.com/standard-schema/standard-schema) validator, including Zod, Valibot, and ArkType. `startWorkflow` validates the input before creating the run, and the handler's `input` is typed from the schema. **Zod:** ```typescript import { workflow } from 'pg-workflows' import { z } from 'zod' const onboarding = workflow( 'user-onboarding', async ({ step, input }) => { // input: { email: string; name: string } await step.run('greet', async () => `Welcome, ${input.name}`) }, { inputSchema: z.object({ email: z.email(), name: z.string() }) }, ) ``` **Valibot:** ```typescript import { workflow } from 'pg-workflows' import * as v from 'valibot' const onboarding = workflow( 'user-onboarding', async ({ step, input }) => { // input: { email: string; name: string } await step.run('greet', async () => `Welcome, ${input.name}`) }, { inputSchema: v.object({ email: v.pipe(v.string(), v.email()), name: v.string() }) }, ) ``` **No schema.** `input` is `unknown` and isn't validated, so narrow it yourself: ```typescript const refund = workflow('refund', async ({ step, input }) => { const { orderId } = input as { orderId: string } await step.run('refund-order', async () => ({ refunded: orderId })) }) ``` # Pause and resume (/docs/concepts/pause-and-resume) `step.pause` pauses the run at that point until something calls `resumeWorkflow`: ```typescript const publishPost = workflow('publish-post', async ({ step }) => { await step.run('render-preview', async () => ({ rendered: true })) await step.pause('editor-approval') await step.run('publish', async () => ({ published: true })) }) ``` ```typescript await engine.resumeWorkflow({ runId: run.id }) ``` `engine.pauseWorkflow({ runId })` pauses a pending or running run from outside. `engine.cancelWorkflow({ runId })` cancels a pending, running, or paused run. ## Fast-forward [#fast-forward] `fastForwardWorkflow` completes whatever step the run is paused on. It's intended for tests, debugging, and support tooling. It does nothing if the run isn't paused. ```typescript // Paused on waitFor: `data` becomes the event payload (default `{}`) await engine.fastForwardWorkflow({ runId: run.id, data: { approved: true } }) // Paused on delay or waitUntil: skips the wait await engine.fastForwardWorkflow({ runId: run.id }) ``` | Paused on | Effect | | ----------------------------------- | --------------------------------------------------------------- | | `step.waitFor()` | Sends the event with `data` (default `{}`). | | `step.delay()` / `step.waitUntil()` | Ends the wait. | | `step.poll()` | Resolves the poll with `data` as its result. | | `step.pause()` | Same as `resumeWorkflow()`. | | `step.invokeChildWorkflow()` | Nothing. The child's outcome decides when the parent continues. | # Polling (/docs/concepts/polling) `step.poll` calls a function on an interval until it returns a truthy value or the timeout passes. Return `false` to keep polling. ```typescript const exportJob = workflow( 'wait-for-export', async ({ step, input }) => { const result = await step.poll( 'wait-for-file', async () => { const response = await fetch(`https://api.example.com/exports/${input.exportId}`) const body = (await response.json()) as { status: string; url?: string } return body.status === 'ready' ? { url: body.url } : false }, { interval: '1 minute', timeout: '24 hours' }, ) if (result.timedOut) { return { status: 'expired' } } return { status: 'ready', url: result.data.url } }, { inputSchema: z.object({ exportId: z.string() }) }, ) ``` | Option | Default | Description | | ---------- | ------- | ----------------------------------------------------------------------------- | | `interval` | `'30s'` | Time between checks. The minimum is 30 seconds. | | `timeout` | none | Stops polling and returns `{ timedOut: true }`. Omit it to poll indefinitely. | The run is paused between checks and holds no worker. # Priorities (/docs/concepts/priorities) `priority` orders runs in the queue. Higher values run first. | Value | Integer | | ----------- | ------------- | | `'high'` | `100` | | `'normal'` | `0` (default) | | `'low'` | `-100` | | any integer | as given | Set a default on the workflow, and override it for a single run: ```typescript const billing = workflow('billing', handler, { priority: 'high' }) await engine.startWorkflow({ workflowId: 'billing', input: {}, options: { priority: 'low' }, }) ``` The effective priority is `startWorkflow` option, then the workflow's `priority`, then `'normal'`. It's resolved once when the run is created and saved as `run.priority`. Resumes, retries, and poll checks reuse the saved value. Child workflows use the call's `options.priority`, then the child definition's `priority`, then the parent run's priority. # Recurring schedules (/docs/concepts/recurring-schedules) Set `schedule` to start a run on a recurring basis. It accepts a cron expression, a duration string, or a duration object. ```typescript workflow('weekday-report', handler, { schedule: '0 9 * * 1-5', timezone: 'America/New_York' }) workflow('every-5-minutes', handler, { schedule: '5m' }) workflow('hourly', handler, { schedule: '1 hour' }) workflow('daily', handler, { schedule: { days: 1 } }) ``` * **Cron vs. duration.** A string of 5 or 6 space-separated fields that uses only cron characters (`0-9 * / , - ? L W #`) is parsed as cron. Anything else is parsed as a duration. * **Durations must divide evenly.** A duration is converted to cron, so it must be whole minutes that divide 60, whole hours that divide 24, or exactly one day. `'23m'` and `'7h'` throw at registration. Use a cron expression for those. * **`timezone`** applies to cron only. The default is UTC. * **A worker must be running.** Scheduled runs are started by an engine that has the workflow registered and has called `engine.start()`. A scheduled run has `ctx.schedule.timestamp`, the time the schedule fired. A run started with `startWorkflow` has `ctx.schedule === undefined`. For incremental syncs, use the last completed run as a cursor. Read it inside a step so the value stays fixed when the run resumes or retries: ```typescript import { WorkflowStatus, workflow } from 'pg-workflows' const syncOrders = workflow( 'sync-orders', async ({ step, schedule, workflowId }) => { const since = await step.run('read-cursor', async () => { const { items } = await engine.getRuns({ workflowId, statuses: [WorkflowStatus.COMPLETED], limit: 1, }) return (items[0]?.completedAt ?? new Date(0)).toISOString() }) const orders = await step.run('fetch-orders', async () => { const response = await fetch(`https://shop.example.com/orders?updated_since=${since}`) return (await response.json()) as { id: string }[] }) return { firedAt: schedule?.timestamp.toISOString(), since, synced: orders.length } }, { schedule: '5m', singleton: true }, ) ``` Don't use `getWorkflowLastRun` for this. It returns the most recently created run of any status, which inside a scheduled run is the current run. **Overlap.** With `singleton: true`, a scheduled fire is skipped while a previous run is pending or running. Without it, scheduled runs can overlap. # Resource IDs and idempotency (/docs/concepts/resource-id-and-idempotency) ## Resource ID [#resource-id] `resourceId` ties a run to an entity in your app, such as a user, tenant, or order. * **Query.** `getRuns({ resourceId })` lists the runs for that entity. * **Scope.** When you pass `resourceId` to `getRun`, `pauseWorkflow`, `resumeWorkflow`, `cancelWorkflow`, `triggerEvent`, and the other run methods, the query also matches on `resource_id`. A run that belongs to a different resource returns `WorkflowRunNotFoundError`. Use this for tenant isolation. ```typescript const run = await engine.startWorkflow({ workflowId: 'send-invoice', resourceId: 'tenant_42', input: { orderId: 'ord_1', total: 99 }, }) const { items } = await engine.getRuns({ resourceId: 'tenant_42' }) ``` `resourceId` is optional everywhere. Omit it to address runs by `runId` alone. ## Idempotency key [#idempotency-key] Pass `idempotencyKey` when the same start can be requested twice, for example on a double click, a client retry, or an at-least-once webhook. A second `startWorkflow` with the same key returns the existing run and enqueues nothing. ```typescript const first = await engine.startWorkflow({ workflowId: 'send-invoice', input: { orderId: 'ord_1', total: 99 }, idempotencyKey: 'send-invoice:ord_1', }) const second = await engine.startWorkflow({ workflowId: 'send-invoice', input: { orderId: 'ord_1', total: 99 }, idempotencyKey: 'send-invoice:ord_1', }) second.id === first.id // true ``` Keys are unique across the whole table, not per workflow or resource, and can be up to 256 characters. Prefix them with the workflow ID. The input of a duplicate call is ignored. # Retries and timeouts (/docs/concepts/retries-and-timeouts) When a handler throws, the run is retried up to `retries` times (default `0`). Each retry runs the handler from the top, and completed steps return their saved result, so only the failed step and the steps after it run again. The number of the current attempt is `attempt` on the handler context and `retryCount` on the run. Retries are scheduled by [pg-boss](https://pgboss.io/) with exponential backoff: roughly 1s, 2s, 4s, 8s, and so on, with up to ±50% jitter. After the last attempt fails, the run's status is `failed` and `error` holds the message. ```typescript const syncCustomer = workflow( 'sync-customer', async ({ step, input, attempt }) => { return step.run('push-to-crm', async () => { const response = await fetch('https://crm.example.com/customers', { method: 'POST', body: JSON.stringify({ id: input.customerId, attempt }), }) if (!response.ok) throw new Error(`CRM returned ${response.status}`) return { synced: true } }) }, { inputSchema: z.object({ customerId: z.string() }), retries: 5 }, ) ``` Override `retries` for a single run with `startWorkflow({ options: { retries } })`. `timeout` (milliseconds, on the workflow or in `startWorkflow` options) is saved on the run as `timeoutAt`. The engine does not currently fail a run that passes `timeoutAt`. To bound a wait, use the `timeout` option of `step.waitFor` or `step.poll`. A single execution of the handler is also bounded by the job expiry, `WORKFLOW_RUN_EXPIRE_IN_SECONDS` (default 300). Split work that takes longer into several steps. See [Configuration](/docs/reference/configuration#environment-variables). # Singleton workflows (/docs/concepts/singleton-workflows) `singleton: true` allows at most one pending or running run of a workflow ID. Starting a second run throws `WorkflowRunInProgressError`. ```typescript import { WorkflowRunInProgressError, workflow } from 'pg-workflows' const nightlySync = workflow( 'nightly-sync', async ({ step }) => { await step.run('pull', async () => ({ pulled: true })) }, { singleton: true }, ) await engine.startWorkflow({ workflowId: 'nightly-sync', input: {} }) try { await engine.startWorkflow({ workflowId: 'nightly-sync', input: {} }) } catch (error) { if (error instanceof WorkflowRunInProgressError) { // the first run is still pending or running } } ``` A run releases the slot when it pauses (`waitFor`, `pause`, `delay`, `waitUntil`, `poll`), completes, fails, or is cancelled. Resuming a paused run while another run holds the slot also throws `WorkflowRunInProgressError`. `WorkflowClient` doesn't load workflow definitions, so it can't see `singleton` on the definition. Set it on the ref instead: ```typescript const nightlySyncRef = workflow.ref('nightly-sync', { singleton: true }) await client.startWorkflow(nightlySyncRef, {}) ``` # Timers (/docs/concepts/timers) `step.waitUntil` pauses until a date. `step.delay` pauses for a duration, and `step.sleep` is an alias for it. A date in the past resumes immediately. ```typescript await step.waitUntil('send-at-launch', new Date('2026-12-01T09:00:00Z')) await step.waitUntil('send-at-launch', '2026-12-01T09:00:00Z') await step.waitUntil('send-at-launch', { date: new Date('2026-12-01T09:00:00Z') }) await step.delay('cool-off', '3 days') await step.delay('cool-off', { days: 3 }) await step.delay('ramp-up', '2 days 12 hours') await step.sleep('backoff', '1 hour') ``` A duration is a string (`'90s'`, `'2h'`, `'3 days'`) or an object with any of `weeks`, `days`, `hours`, `minutes`, and `seconds`. # Workflows and steps (/docs/concepts/workflows) The snippets on this page assume an `engine` like the one in the [quickstart](/docs/quickstart), with each workflow passed to `workflows: [...]`. A workflow is an async function with an ID. Every durable operation inside it goes through `step`. ```typescript 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](/docs/concepts/input-validation). | | `retries` | `number` | `0` | Retry attempts after a failure. See [Retries and timeouts](/docs/concepts/retries-and-timeouts). | | `timeout` | `number` (ms) | none | Recorded as `run.timeoutAt`. See [Retries and timeouts](/docs/concepts/retries-and-timeouts). | | `priority` | `'high' \| 'normal' \| 'low' \| number` | `'normal'` | Queue priority. See [Priorities](/docs/concepts/priorities). | | `singleton` | `boolean` | `false` | At most one pending or running run. See [Singleton workflows](/docs/concepts/singleton-workflows). | | `schedule` | cron string, duration string, or duration object | none | Starts runs on a recurring schedule. See [Recurring schedules](/docs/concepts/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](/docs/reference/api#workflowcontext). 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 [#steps] `step.run` executes a function once and saves its return value on the run. ```typescript 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 [#how-a-run-executes] PostgreSQL is both the job queue (through [pg-boss](https://pgboss.io/)) 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`. # AI and agent workflows (/docs/guides/ai-agents) LLM calls are slow, cost money, and fail with 429s and 500s. Running them as workflow steps gives you: * **Saved results.** Each `step.run` result is saved. If the process crashes or the run retries, completed LLM calls are not repeated. * **Retries.** A thrown error retries the run with exponential backoff, resuming at the failed step. * **Human review.** `step.waitFor` pauses the run until a reviewer responds. A paused run holds no worker or connection, whether it waits minutes or days. * **Inspectable state.** Every step's output is in the run's `timeline`, so you can see what the agent produced up to the point it failed. ## Setup [#setup] The snippets call an `llm` helper that wraps your model provider's SDK and returns the reply as a string. Declare it once: ```typescript // llm.ts export declare const llm: { chat(params: { model: string; messages: { role: 'system' | 'user'; content: string }[] }): Promise embed(text: string): Promise } export declare const vectorStore: { search(embedding: number[], options: { topK: number }): Promise<{ id: string; text: string }[]> } ``` Return plain strings and objects from LLM steps. Results are stored as `jsonb`, so an SDK response object with class instances or methods doesn't round-trip. ## Multi-step agent [#multi-step-agent] A planning call produces a list of tasks. Each task is its own step, so a crash after task 3 of 5 resumes at task 4. ```typescript import { workflow } from 'pg-workflows' import { z } from 'zod' import { llm } from './llm' const researchAgent = workflow( 'research-agent', async ({ step, input }) => { const tasks = await step.run('create-plan', async () => { const reply = await llm.chat({ model: 'gpt-4o', messages: [ { role: 'user', content: `Return a JSON array of {"id": string, "description": string} research tasks for: ${input.topic}`, }, ], }) return JSON.parse(reply) as { id: string; description: string }[] }) const findings: string[] = [] for (const task of tasks) { const finding = await step.run(`research-${task.id}`, async () => { return llm.chat({ model: 'gpt-4o', messages: [{ role: 'user', content: `Research: ${task.description}` }], }) }) findings.push(finding) } const report = await step.run('synthesize', async () => { return llm.chat({ model: 'gpt-4o', messages: [{ role: 'user', content: `Synthesize these findings:\n\n${findings.join('\n\n')}` }], }) }) return { tasks, report } }, { inputSchema: z.object({ topic: z.string() }), retries: 3 }, ) ``` ## Human review [#human-review] Generate a draft, wait for a reviewer, then publish or revise. ```typescript import { workflow } from 'pg-workflows' import { z } from 'zod' import { llm } from './llm' const contentPipeline = workflow( 'ai-content-pipeline', async ({ step, input }) => { const draft = await step.run('generate-draft', async () => { return llm.chat({ model: 'gpt-4o', messages: [{ role: 'user', content: `Write a blog post about: ${input.topic}` }], }) }) const review = await step.waitFor('human-review', { eventName: 'content-reviewed', timeout: 7 * 24 * 60 * 60 * 1000, schema: z.object({ approved: z.boolean(), feedback: z.string().optional() }), }) if (!review) { return { status: 'expired', content: draft } } if (review.approved) { return { status: 'published', content: draft } } const revision = await step.run('revise-draft', async () => { return llm.chat({ model: 'gpt-4o', messages: [ { role: 'user', content: `Revise this draft based on the feedback.\n\nDraft:\n${draft}\n\nFeedback:\n${review.feedback}`, }, ], }) }) return { status: 'revised', content: revision } }, { inputSchema: z.object({ topic: z.string() }), retries: 3 }, ) ``` Send the reviewer's decision from your API: ```typescript await engine.triggerEvent({ runId, eventName: 'content-reviewed', data: { approved: false, feedback: 'Make the intro more engaging' }, }) ``` With `timeout`, `review` is `undefined` if no event arrives within 7 days. `schema` types the event data but doesn't validate it, so check untrusted input before calling `triggerEvent`. ## Retrieval-augmented generation [#retrieval-augmented-generation] Embed the query, retrieve documents, answer, then check the answer against the sources. ```typescript import { workflow } from 'pg-workflows' import { z } from 'zod' import { llm, vectorStore } from './llm' const ragAgent = workflow( 'rag-agent', async ({ step, input }) => { const embedding = await step.run('embed-query', async () => { return llm.embed(input.query) }) const documents = await step.run('search-docs', async () => { return vectorStore.search(embedding, { topK: 10 }) }) const context = documents.map((doc) => doc.text).join('\n') const answer = await step.run('generate-answer', async () => { return llm.chat({ model: 'gpt-4o', messages: [ { role: 'system', content: `Answer using only these documents:\n${context}` }, { role: 'user', content: input.query }, ], }) }) const factCheck = await step.run('fact-check', async () => { return llm.chat({ model: 'gpt-4o', messages: [ { role: 'user', content: `Documents:\n${context}\n\nAnswer:\n${answer}\n\nList any claims in the answer that the documents don't support.`, }, ], }) }) return { answer, factCheck, sources: documents.map((doc) => doc.id) } }, { inputSchema: z.object({ query: z.string() }), retries: 3 }, ) ``` ## Limits to plan for [#limits-to-plan-for] * **A step can repeat after a crash.** If the process dies after an LLM call returns but before its result is saved, the call runs again on retry. Pass an idempotency key to APIs that charge or send. * **Execution time is capped.** A single handler execution is limited by [`WORKFLOW_RUN_EXPIRE_IN_SECONDS`](/docs/reference/configuration#environment-variables) (default 300). The limit resets at every pause, so an agent that waits for review is unaffected, but a long chain of slow calls without a pause can hit it. Raise the limit or split the chain into [child workflows](/docs/concepts/child-workflows). * **The workflow `timeout` isn't enforced.** It's saved as `run.timeoutAt`. To bound a wait, use `timeout` on `step.waitFor` or `step.poll`. * **Code between steps must be deterministic.** The handler runs from the top on every resume. Keep LLM calls, randomness, and `Date.now()` inside steps. # Examples (/docs/guides/examples) Each snippet defines a workflow. To run one, pass it to `workflows: [...]` on an engine like the one in the [quickstart](/docs/quickstart), then call `engine.startWorkflow`. Full runnable scripts are in [`examples/node`](https://github.com/SokratisVidros/pg-workflows/tree/main/examples/node): ```bash git clone https://github.com/SokratisVidros/pg-workflows.git cd pg-workflows && bun install cd examples/node DATABASE_URL=postgres://postgres:postgres@localhost:5432/postgres bun run example:basic ``` | Script | Shows | | ------------------------------------------------------------ | --------------------------------------------------------------------------------------------------- | | `example:basic` | Sequential steps and `checkProgress` | | `example:approval-flow` | `waitFor` with `triggerEvent` | | `example:timeout` | `waitFor` timing out, and bringing your own [pg-boss](https://pgboss.io/) | | `example:polling` | `step.poll` against a simulated payment API | | `example:cron` | A recurring schedule | | `example:microservices:worker` / `example:microservices:api` | The [microservices (web and worker)](/docs/architectures#microservices-web-and-worker) architecture | ## Conditional steps [#conditional-steps] Branch on a saved result outside the step. The engine detects steps inside `if` blocks and reports them as conditional when you register the workflow. ```typescript import { workflow } from 'pg-workflows' import { z } from 'zod' const upgradeAccount = workflow( 'upgrade-account', async ({ step, input }) => { const account = await step.run('load-account', async () => { return { id: input.accountId, plan: 'premium' as 'free' | 'premium' } }) if (account.plan === 'premium') { await step.run('grant-premium-features', async () => { return { granted: true } }) } return { plan: account.plan } }, { inputSchema: z.object({ accountId: z.string() }) }, ) ``` ## Loops [#loops] Give each iteration its own step ID. A retry skips the items that already finished. ```typescript import { workflow } from 'pg-workflows' import { z } from 'zod' const resizeImages = workflow( 'resize-images', async ({ step, input }) => { const images = await step.run('list-images', async () => { return input.imageIds.map((id) => ({ id, url: `https://cdn.example.com/${id}.png` })) }) const resized: string[] = [] for (const image of images) { const result = await step.run(`resize-${image.id}`, async () => { return { url: image.url.replace('.png', '@2x.png') } }) resized.push(result.url) } return { resized } }, { inputSchema: z.object({ imageIds: z.array(z.string()) }) }, ) ``` Each handler execution is limited by [`WORKFLOW_RUN_EXPIRE_IN_SECONDS`](/docs/reference/configuration#environment-variables) (default 300). For a long list, split the work across [child workflows](/docs/concepts/child-workflows) instead. ## Follow-up after a delay [#follow-up-after-a-delay] `step.delay` pauses the run for a duration. The run holds no worker while it waits and survives restarts. ```typescript import { workflow } from 'pg-workflows' import { z } from 'zod' const trialReminder = workflow( 'trial-reminder', async ({ step, input }) => { await step.run('send-welcome', async () => { console.log(`Welcome email to ${input.email}`) return { sent: true } }) await step.delay('wait-for-trial-midpoint', '7 days') await step.run('send-reminder', async () => { console.log(`Reminder email to ${input.email}`) return { sent: true } }) }, { inputSchema: z.object({ email: z.email() }) }, ) ``` ## Poll an external API [#poll-an-external-api] Return `false` to keep polling, or a value to finish. `result.timedOut` tells you which way it ended. ```typescript import { workflow } from 'pg-workflows' import { z } from 'zod' const awaitPayment = workflow( 'await-payment', async ({ step, input }) => { const result = await step.poll( 'wait-for-payment', async () => { const response = await fetch(`https://payments.example.com/payments/${input.paymentId}`) const payment = (await response.json()) as { status: string; amount: number } return payment.status === 'succeeded' ? payment : false }, { interval: '1 minute', timeout: '24 hours' }, ) if (result.timedOut) { await step.run('cancel-order', async () => ({ cancelled: input.orderId })) return { status: 'cancelled' } } return { status: 'paid', amount: result.data.amount } }, { inputSchema: z.object({ orderId: z.string(), paymentId: z.string() }) }, ) ``` ## Retry a flaky call [#retry-a-flaky-call] A thrown error fails the attempt. With `retries: 3`, the run is retried up to three times with exponential backoff, and steps that already succeeded are not repeated. ```typescript import { workflow } from 'pg-workflows' import { z } from 'zod' const syncInventory = workflow( 'sync-inventory', async ({ step, input, attempt }) => { const stock = await step.run('fetch-stock', async () => { const response = await fetch(`https://warehouse.example.com/sku/${input.sku}`) if (!response.ok) throw new Error(`Warehouse returned ${response.status} on attempt ${attempt}`) return (await response.json()) as { available: number } }) return { sku: input.sku, available: stock.available } }, { inputSchema: z.object({ sku: z.string() }), retries: 3 }, ) ``` ## Check progress [#check-progress] `checkProgress` returns the run plus step counts. Call it on an engine that has the workflow registered. ```typescript const run = await engine.startWorkflow({ workflowId: 'resize-images', input: { imageIds: ['a', 'b', 'c'] }, }) const progress = await engine.checkProgress({ runId: run.id }) console.log(progress.status, `${progress.completedSteps}/${progress.totalSteps}`, `${progress.completionPercentage}%`) ``` To show runs in a React app, use [`@pg-workflows/ui`](/docs/ui) behind your own API. Don't call the engine from the browser. # Introduction (/docs) pg-workflows runs durable workflows on the PostgreSQL you already have. Each step's result is saved, a retried run skips the steps that already finished, and a run can pause for an event, a timer, or a polled condition. There's no Redis, broker, or scheduler to run. ## Start here [#start-here] ## Features [#features] | Feature | API | | ---------------------------------------------------------------------------------- | --------------------------------------------------------------------------------------------------------- | | [Durable steps](/docs/concepts/workflows#steps) | `step.run(id, fn)` | | [Wait for external events](/docs/concepts/events) | `step.waitFor(id, { eventName, timeout?, schema? })` + `engine.triggerEvent()` | | [Timers](/docs/concepts/timers) | `step.delay(id, '3 days')`, `step.waitUntil(id, date)` | | [Polling](/docs/concepts/polling) | `step.poll(id, fn, { interval, timeout })` | | [Manual pause and resume](/docs/concepts/pause-and-resume) | `step.pause(id)`, `engine.resumeWorkflow()` | | [Child workflows](/docs/concepts/child-workflows) | `step.invokeChildWorkflow(id, ref, input)` | | [Recurring schedules](/docs/concepts/recurring-schedules) | `workflow(id, fn, { schedule: '0 9 * * 1-5' })` | | [Retries](/docs/concepts/retries-and-timeouts) | `workflow(id, fn, { retries: 3 })` | | [Priorities](/docs/concepts/priorities) | `workflow(id, fn, { priority: 'high' })` | | [One run at a time](/docs/concepts/singleton-workflows) | `workflow(id, fn, { singleton: true })` | | [Deduplicated starts](/docs/concepts/resource-id-and-idempotency#idempotency-key) | `startWorkflow({ idempotencyKey })` | | [Tenant scoping](/docs/concepts/resource-id-and-idempotency#resource-id) | `resourceId` on every run and every API call | | [Typed input](/docs/concepts/input-validation) | Any [Standard Schema](https://github.com/standard-schema/standard-schema) library (Zod, Valibot, ArkType) | | [Microservices (web and worker)](/docs/architectures#microservices-web-and-worker) | `WorkflowClient` from `pg-workflows/client` | ## Packages [#packages] | Package | Purpose | | ------------------------------------------------------------ | -------------------------------------------------------------------------- | | [`pg-workflows`](https://www.npmjs.com/package/pg-workflows) | The engine and client | | [`@pg-workflows/ui`](/docs/ui) | React dashboard, components, and hooks. Try it with `npx @pg-workflows/ui` | | [`@pg-workflows/otel`](/docs/observability/tracing) | OpenTelemetry spans for workflow runs and steps | ## Requirements [#requirements] * Node.js >= 18 * PostgreSQL >= 10 * `pg` >= 8 (peer dependency). [`pg-boss`](https://pgboss.io/) ships with the engine and needs no setup. ## Acknowledgments [#acknowledgments] [Temporal](https://temporal.io/), [Inngest](https://www.inngest.com/), [Trigger.dev](https://trigger.dev/), and [DBOS](https://www.dbos.dev/) pioneered the durable execution patterns this project builds on. # Install with an agent (/docs/install-with-agent) pg-workflows publishes an install skill: instructions written for coding agents. The agent reads your codebase, picks an architecture, installs and verifies the engine, mounts the dashboard for your stack, and offers to add tracing. It asks before touching your database. ## Copy the prompt [#copy-the-prompt] Paste this into your agent, from the root of your project: The agent fetches [`/skill.md`](/skill.md) and follows it. Nothing is installed in the agent itself. ## Install the skill [#install-the-skill] To keep the skill available in every session, install it into your agent instead: ```bash npx skills add SokratisVidros/pg-workflows --skill pg-workflows-install ``` Then ask the agent to "add pg-workflows to this project". The skill's source is in [`skills/pg-workflows-install`](https://github.com/SokratisVidros/pg-workflows/tree/main/skills/pg-workflows-install). ## What the agent does [#what-the-agent-does] 1. **Inspects the project**: package manager, services and their frameworks, hosting, the Postgres connection, React and Tailwind, and any existing OpenTelemetry setup. 2. **Picks an architecture** and tells you why: | Architecture | When | | ------------------- | -------------------------------------------------------------------------------------------------------------------- | | Monolith | One long-running service starts and executes workflows. | | Web app plus worker | One codebase whose web tier is serverless or shouldn't carry heavy workflow dependencies. A worker process joins it. | | Microservices | Several services start runs with `WorkflowClient`. One worker service owns every workflow. | 3. **Installs the engine** and runs a workflow end to end against your database, showing you the output. 4. **Adds the dashboard** with the adapter for your stack: Next.js (App Router or Pages Router), Express, Fastify, Hono, TanStack Start, Bun, or Deno. The dashboard's API goes behind your existing auth. 5. **Offers tracing** with `@pg-workflows/otel`, reusing your OpenTelemetry SDK if you have one. 6. **Reports** every file it changed and what each process needs in production. The same steps are documented for people in [Architectures](/docs/architectures), [UI components](/docs/ui), and [Tracing](/docs/observability/tracing). ## For agents [#for-agents] * [`/skill.md`](/skill.md): the install skill * [`/llms.txt`](/llms.txt): an index of these docs * [`/llms-full.txt`](/llms-full.txt): every page in one file * Every page has a **Copy Markdown** button, and its Markdown at `/llms.mdx/docs//content.md` # Monitoring runs (/docs/observability/monitoring) ## Run state [#run-state] Every run's status, step outputs, and error are in `workflow_runs`. Query them from code: | Method | Use | | ----------------------------------------------- | ------------------------------------------------------------------------- | | `getRun({ runId })` | Status, `output`, `error`, and the `timeline` of step results for one run | | `checkProgress({ runId })` | Completed and total steps, and a percentage | | `getRuns({ statuses, workflowId, resourceId })` | Lists runs, for example every failed run of one workflow | | `getStats({ workflowId })` | Run counts per status | See the [API reference](/docs/reference/api#queries). ## Dashboard [#dashboard] [`@pg-workflows/ui`](/docs/ui) is a React dashboard for browsing runs and step timelines. Try it against your database with: ```bash npx @pg-workflows/ui --database-url=postgres://postgres:postgres@localhost:5432/postgres ``` # Tracing (/docs/observability/tracing) OpenTelemetry tracing for [pg-workflows](https://github.com/SokratisVidros/pg-workflows). `otelPlugin` emits one span per workflow execution and one per step. It's a separate package so the engine has no OpenTelemetry dependency. ## Quickstart [#quickstart] This prints spans to the console. It assumes a Postgres at `DATABASE_URL`. The [engine quickstart](/docs/quickstart) has a Docker command for one. ```bash npm install pg-workflows pg @pg-workflows/otel @opentelemetry/api @opentelemetry/sdk-node npm install -D tsx ``` Save as `traced.ts`: ```typescript import { NodeSDK, tracing } from '@opentelemetry/sdk-node' import { otelPlugin } from '@pg-workflows/otel' import { WorkflowEngine, WorkflowStatus, workflow } from 'pg-workflows' const sdk = new NodeSDK({ traceExporter: new tracing.ConsoleSpanExporter() }) sdk.start() const tracedWorkflow = workflow.use(otelPlugin()) const checkout = tracedWorkflow('checkout', async ({ step }) => { const charge = await step.run('charge', async () => ({ chargeId: 'ch_123' })) await step.run('send-receipt', async () => ({ sentFor: charge.chargeId })) return charge }) async function main() { const engine = new WorkflowEngine({ connectionString: process.env.DATABASE_URL ?? 'postgres://postgres:postgres@localhost:5432/postgres', workflows: [checkout], }) await engine.start() const run = await engine.startWorkflow({ workflowId: 'checkout', input: {} }) let result = await engine.getRun({ runId: run.id }) while (result.status === WorkflowStatus.PENDING || result.status === WorkflowStatus.RUNNING) { await new Promise((resolve) => setTimeout(resolve, 200)) result = await engine.getRun({ runId: run.id }) } await engine.stop() await sdk.shutdown() } main() ``` ```bash npx tsx traced.ts ``` The output includes a `pg_workflows.workflow.run` span and two `pg_workflows.step.run` spans. In production, replace `ConsoleSpanExporter` with your exporter, for example `OTLPTraceExporter` from `@opentelemetry/exporter-trace-otlp-http`. ## Spans [#spans] ``` pg_workflows.workflow.run ├── pg_workflows.step.run ├── pg_workflows.step.waitFor ├── pg_workflows.step.delay ├── pg_workflows.step.waitUntil ├── pg_workflows.step.pause ├── pg_workflows.step.poll └── pg_workflows.step.invokeChildWorkflow ``` `step.sleep` is an alias for `step.delay` and emits `pg_workflows.step.delay`. **One trace per execution.** A run that pauses (`waitFor`, `delay`, `pause`, and so on) and later resumes produces a new trace for each execution. Correlate them with the `workflow.id` and `workflow.run_id` attributes. **Step spans are recorded when the step finishes.** Each step span gets its start time from when the step began, but the span object is created only after the step returns or throws. Spans your code creates inside a `step.run` callback, including auto-instrumented HTTP or database calls, are therefore children of `pg_workflows.workflow.run`, not of the step span. ## Attributes [#attributes] | Span | Attributes | | --------------------------- | -------------------------------------------------------------------------------------------------------------------------------------------------------------------- | | `pg_workflows.workflow.run` | `workflow.id`, `workflow.run_id`, `workflow.attempt` (same as `run.retryCount`), `workflow.resource_id` (when set), and anything returned by the `attributes` option | | `pg_workflows.step.` | `step.id`, `step.type` (a `StepType` value) | On success, a span's status is `OK`. On error, the plugin calls `recordException(error)` and sets the status to `ERROR` with the error message. ## Replayed steps [#replayed-steps] When a run resumes, the handler runs from the top, and completed steps return their saved output. The plugin emits no span for these replays. It uses `isStepCached(context.timeline, stepId)` from `pg-workflows`. A step counts as cached when: * its output is in the timeline, or * it's a `step.invokeChildWorkflow` whose child run has already been created. This covers a parent that re-enters the step while the child is still running. Two exceptions: * **`step.poll`** emits a span on every check, because each check is a real attempt. * **`step.run` returning `undefined`** emits no span. Return a value (for example `{ sent: true }`) from steps you want traced. ## Options [#options] ```typescript import { trace } from '@opentelemetry/api' import { otelPlugin } from '@pg-workflows/otel' otelPlugin({ tracer: trace.getTracer('billing-worker'), spanNamePrefix: 'billing', attributes: (ctx) => ({ 'tenant.id': ctx.resourceId ?? 'none' }), }) ``` | Option | Type | Default | Description | | ---------------- | ---------------------------------------------------------- | --------------------------------- | --------------------------------------------------------------------------------------------------------------- | | `tracer` | `Tracer` | `trace.getTracer('pg-workflows')` | The tracer that creates spans. | | `spanNamePrefix` | `string` | `'pg_workflows'` | Replaces `pg_workflows` in every span name. | | `attributes` | `(ctx: WorkflowContext) => Record` | none | Extra attributes for the `workflow.run` span. Receives the handler context, including `input` and `resourceId`. | ## Errors [#errors] When a step or handler throws, the plugin records the exception, sets the span status to `ERROR`, and rethrows the original error. Retries and failure handling in the engine are unchanged. A thrown non-`Error` value (`throw 'msg'`) is wrapped in an `Error` for the span only. The original value is rethrown. ## Composing plugins [#composing-plugins] `otelPlugin` uses the `wrap(context, next)` middleware hook that any plugin can implement. With several plugins, the first one passed to `.use()` is the outermost wrap: ```typescript import { otelPlugin } from '@pg-workflows/otel' import { type WorkflowPlugin, workflow } from 'pg-workflows' const timingPlugin: WorkflowPlugin = { name: 'timing', methods: () => ({}), wrap: async (ctx, next) => { const startedAt = Date.now() try { return await next() } finally { ctx.logger.log(`${ctx.workflowId} execution took ${Date.now() - startedAt}ms`) } }, } // timingPlugin wraps the whole execution, including the workflow.run span const instrumentedWorkflow = workflow.use(timingPlugin).use(otelPlugin()) ``` ## Context propagation [#context-propagation] Nested spans need a context manager that follows `await`. `NodeSDK` from `@opentelemetry/sdk-node` registers one. If you set up OpenTelemetry by hand, install `@opentelemetry/context-async-hooks` and register it: ```typescript import { context } from '@opentelemetry/api' import { AsyncHooksContextManager } from '@opentelemetry/context-async-hooks' context.setGlobalContextManager(new AsyncHooksContextManager().enable()) ``` ## Not supported yet [#not-supported-yet] Most of these need the trace context stored with the run, and are expected to ship together. * **Metrics.** Only traces are emitted. * **Linking executions.** A resumed run starts a new root trace, not a continuation of the previous one. * **Child workflow traces.** A child run starts its own root trace. * **Caller context.** The trace of the request that called `startWorkflow` isn't propagated into the run. * **Dead-letter failures.** When retries run out, the run is marked failed outside the plugin chain, so that final transition has no span. The error is already on the last execution's span. * **Sampling.** The plugin uses your `TracerProvider`'s sampler. ## Migrating from pg-workflows 0.15 and earlier [#migrating-from-pg-workflows-015-and-earlier] `otelPlugin` used to be exported from `pg-workflows`. Install `@pg-workflows/otel` and change the import. Options and spans are unchanged. ```diff -import { workflow, otelPlugin } from 'pg-workflows' +import { otelPlugin } from '@pg-workflows/otel' +import { workflow } from 'pg-workflows' ``` ## Requirements [#requirements] * `pg-workflows` >= 0.16.0 (peer dependency) * `@opentelemetry/api` ^1.9.0 (peer dependency) * An OpenTelemetry SDK with an async context manager. See [Context propagation](/docs/observability/tracing#context-propagation). # Quickstart (/docs/quickstart) **1. Start Postgres** (skip if you already have one): ```bash docker run -d --name pg-workflows-db -e POSTGRES_PASSWORD=postgres -p 5432:5432 postgres:17 ``` **2. Create a project:** ```bash mkdir workflows-demo && cd workflows-demo npm init -y npm install pg-workflows pg zod npm install -D tsx ``` **3. Save this as `index.ts`:** ```typescript import { WorkflowEngine, WorkflowStatus, workflow } from 'pg-workflows' import { z } from 'zod' const greetUser = workflow( 'greet-user', async ({ step, input }) => { const user = await step.run('load-user', async () => { return { name: input.name, signedUpAt: new Date().toISOString() } }) const message = await step.run('build-message', async () => { return `Welcome, ${user.name}!` }) return { message } }, { inputSchema: z.object({ name: z.string() }) }, ) async function main() { const engine = new WorkflowEngine({ connectionString: process.env.DATABASE_URL ?? 'postgres://postgres:postgres@localhost:5432/postgres', workflows: [greetUser], }) await engine.start() const run = await engine.startWorkflow({ workflowId: 'greet-user', input: { name: 'Ada' }, }) let result = await engine.getRun({ runId: run.id }) while (result.status === WorkflowStatus.PENDING || result.status === WorkflowStatus.RUNNING) { await new Promise((resolve) => setTimeout(resolve, 200)) result = await engine.getRun({ runId: run.id }) } console.log(result.status, result.output) await engine.stop() } main() ``` **4. Run it:** ```bash npx tsx index.ts ``` After the engine's startup logs, you should see: ``` completed { message: 'Welcome, Ada!' } ``` `engine.start()` creates its tables on first run. Each `step.run` result is saved on the run's row in `workflow_runs`. When a run is retried (set `retries` on the workflow), completed steps return their saved result instead of running again. ## Wait for an event [#wait-for-an-event] `step.waitFor` pauses the run until your code calls `triggerEvent`. A paused run holds no worker and no connection. ```typescript import { WorkflowEngine, WorkflowStatus, workflow } from 'pg-workflows' import { z } from 'zod' const approveExpense = workflow( 'approve-expense', async ({ step, input }) => { await step.run('notify-manager', async () => { return { notified: `manager of ${input.employee}` } }) const review = await step.waitFor('wait-for-review', { eventName: 'expense-reviewed', schema: z.object({ approved: z.boolean() }), }) return { amount: input.amount, approved: review.approved } }, { inputSchema: z.object({ employee: z.string(), amount: z.number() }) }, ) const sleep = (ms: number) => new Promise((resolve) => setTimeout(resolve, ms)) async function main() { const engine = new WorkflowEngine({ connectionString: process.env.DATABASE_URL ?? 'postgres://postgres:postgres@localhost:5432/postgres', workflows: [approveExpense], }) await engine.start() const run = await engine.startWorkflow({ workflowId: 'approve-expense', input: { employee: 'ada', amount: 120 }, }) // In a real app, this is a separate request: an approval button, a webhook, and so on. while ((await engine.getRun({ runId: run.id })).status !== WorkflowStatus.PAUSED) await sleep(200) await engine.triggerEvent({ runId: run.id, eventName: 'expense-reviewed', data: { approved: true }, }) let result = await engine.getRun({ runId: run.id }) while (result.status !== WorkflowStatus.COMPLETED && result.status !== WorkflowStatus.FAILED) { await sleep(200) result = await engine.getRun({ runId: run.id }) } console.log(result.status, result.output) // completed { amount: 120, approved: true } await engine.stop() } main() ``` # API reference (/docs/reference/api) ## WorkflowEngine [#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. ```typescript import { WorkflowEngine } from 'pg-workflows' ``` ### Constructor [#constructor] ```typescript 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`](/docs/reference/api#workflowlogger) | Defaults to `console.warn` and `console.error`. | | `boss` | `PgBoss` | Your own [pg-boss](https://pgboss.io/) instance. When omitted, the engine creates one in the `pgboss_v12_pgworkflow` schema. | Pass exactly one of `connectionString` and `pool`. ### Lifecycle [#lifecycle] | Method | Description | | ------------------------------------------------------------------- | ---------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | | `start(asEngine = true, { batchSize = 1, heartbeatSeconds = 30 }?)` | Starts [pg-boss](https://pgboss.io/), runs migrations, registers `workflows`, starts [`WORKFLOW_RUN_WORKERS`](/docs/reference/configuration#environment-variables) 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](https://pgboss.io/), 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 [#runs] | Method | Returns | Description | | ------------------------------------------------------------------------------ | ------------- | ---------------------------------------------------------------------------------------------------------------------------------------------------------- | | `startWorkflow(ref, input, options?)` | `WorkflowRun` | Starts a run from a typed [`WorkflowRef`](/docs/reference/api#workflowref). | | `startWorkflow({ workflowId, input, resourceId?, idempotencyKey?, options? })` | `WorkflowRun` | Starts a run by ID. Throws `WorkflowRunInProgressError` for a [singleton](/docs/concepts/singleton-workflows) 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](/docs/concepts/pause-and-resume#fast-forward). | `options` on `startWorkflow` is a [`StartWorkflowOptions`](/docs/reference/api#startworkflowoptions). `options` on `resumeWorkflow` and `triggerEvent` accepts `{ expireInSeconds }`. ### Queries [#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](/docs/concepts/resource-id-and-idempotency#resource-id). ## WorkflowClient [#workflowclient] Starts and manages runs without loading workflow handlers. Use it in API services that should not import worker code. ```typescript 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](/docs/concepts/singleton-workflows). ## workflow() [#workflow] ```typescript 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` | The workflow body. Its return value becomes `run.output`. | | `options` | `WorkflowOptions` | `inputSchema`, `retries`, `timeout`, `priority`, `singleton`, `schedule`, `timezone`. See the [options table](/docs/concepts/workflows). | ### workflow\.use(plugin) [#workflowuseplugin] 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`](/docs/observability/tracing) is built on this. ```typescript import { otelPlugin } from '@pg-workflows/otel' import { workflow } from 'pg-workflows' const tracedWorkflow = workflow.use(otelPlugin()) ``` ### workflow\.ref(id, options?) [#workflowrefid-options] Same as [`createWorkflowRef`](/docs/reference/api#workflowref), with the generics in `` order. ## WorkflowRef [#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. ```typescript 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: ```typescript await client.startWorkflow(sendInvoiceRef, { orderId: 'ord_1', total: 99 }) ``` Worker service: ```typescript 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: ```typescript const receiptSchema = z.object({ orderId: z.string() }) const sendReceiptRef = createWorkflowRef<{ receiptId: string }, typeof receiptSchema>('send-receipt', { inputSchema: receiptSchema, }) ``` ## WorkflowContext [#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` | 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](/docs/concepts/recurring-schedules). | ### Step methods [#step-methods] | Method | Returns | Guide | | ---------------------------------------------------------------------------------------------------------- | ------------------------------------------------------------- | --------------------------------------------------- | | `step.run(stepId, fn)` | `Promise` | [Steps](/docs/concepts/workflows#steps) | | `step.waitFor(stepId, { eventName, schema? })` | `Promise` | [Events](/docs/concepts/events) | | `step.waitFor(stepId, { eventName, timeout, schema? })` | `Promise` | [Events](/docs/concepts/events) | | `step.waitUntil(stepId, date \| isoString \| { date })` | `Promise` | [Timers](/docs/concepts/timers) | | `step.delay(stepId, duration)` | `Promise` | [Timers](/docs/concepts/timers) | | `step.sleep(stepId, duration)` | `Promise` | Alias for `delay` | | `step.pause(stepId)` | `Promise` | [Pause and resume](/docs/concepts/pause-and-resume) | | `step.poll(stepId, fn, { interval?, timeout? })` | `Promise<{ timedOut: false, data: T } \| { timedOut: true }>` | [Polling](/docs/concepts/polling) | | `step.invokeChildWorkflow(stepId, ref, input, options?)` | `Promise` | [Child workflows](/docs/concepts/child-workflows) | | `step.invokeChildWorkflow(stepId, { workflowId, input, resourceId?, idempotencyKey?, options? })` | `Promise` | [Child workflows](/docs/concepts/child-workflows) | A `duration` is a string (`'90s'`, `'2h'`, `'3 days'`) or `{ weeks?, days?, hours?, minutes?, seconds? }`. ## Types [#types] ### WorkflowRun [#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` | Saved step results. | | `retryCount` / `maxRetries` | `number` | | | `priority` | `number` | Resolved priority. See [Priorities](/docs/concepts/priorities). | | `singleton` | `boolean` | | | `idempotencyKey` | `string \| null` | | | `parentRunId`, `parentStepId`, `parentResourceId` | `string \| null` | Set on child runs. | | `jobId` | `string \| null` | The [pg-boss](https://pgboss.io/) job ID. | | `createdAt`, `updatedAt` | `Date` | | | `pausedAt`, `resumedAt`, `completedAt`, `timeoutAt`, `scheduledAt` | `Date \| null` | | `WorkflowRunProgress` is `WorkflowRun` plus `completedSteps`, `totalSteps`, and `completionPercentage`. ### StartWorkflowOptions [#startworkflowoptions] | Option | Type | Description | | ----------------- | ------------------ | -------------------------------------------------------------------------------------------------------- | | `retries` | `number` | Overrides the definition's `retries`. | | `timeout` | `number` (ms) | Saved as `run.timeoutAt`. Not enforced. See [Retries and timeouts](/docs/concepts/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 [#workflowstatus] ```typescript enum WorkflowStatus { PENDING = 'pending', RUNNING = 'running', PAUSED = 'paused', COMPLETED = 'completed', FAILED = 'failed', CANCELLED = 'cancelled', } ``` `WorkflowRunStats` is `Record`. ### WorkflowPriority [#workflowpriority] `'high' | 'normal' | 'low' | number`. Named levels map to `100`, `0`, and `-100`. ### WorkflowLogger [#workflowlogger] ```typescript interface WorkflowLogger { log(message: string): void error(message: string, ...args: unknown[]): void } ``` ## Errors [#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 [#ui] React components and hooks for browsing runs are in the separate [`@pg-workflows/ui`](/docs/ui) package. They talk to your API over HTTP, which calls `getRuns` and `getRun`. Nothing from `@pg-workflows/ui` is exported by `pg-workflows`. # Configuration (/docs/reference/configuration) ## Environment variables [#environment-variables] The engine reads these at startup. All are optional. | Variable | Default | Description | | -------------------------------- | ------- | ----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | | `WORKFLOW_RUN_WORKERS` | `3` | Number of concurrent workers each `WorkflowEngine` runs in-process. Each worker handles one run execution at a time. | | `WORKFLOW_RUN_EXPIRE_IN_SECONDS` | `300` | Maximum time for a single handler execution. An execution that runs longer is failed and retried. Override per call with `options.expireInSeconds` on `startWorkflow`, `resumeWorkflow`, and `triggerEvent`. | | `WORKFLOW_RUN_HEARTBEAT_SECONDS` | `30` | How often workers report that an execution is alive. If a worker process dies, its run is detected and retried after roughly this interval plus 60 seconds, instead of waiting for the full expiry. Minimum `10`. | The engine does not read `DATABASE_URL`. Pass the connection string or a `pg.Pool` to the constructor: ```typescript import { WorkflowEngine } from 'pg-workflows' const engine = new WorkflowEngine({ connectionString: process.env.DATABASE_URL ?? 'postgres://postgres:postgres@localhost:5432/postgres', }) ``` ## Database objects [#database-objects] `engine.start()` runs migrations and creates: | Object | Purpose | | ---------------------------------------------------------------------------- | --------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | | `public.workflow_runs` | One row per run: status, input, output, error, and the `timeline` of step results. | | `workflow_runs.resource_id` index | Lookups by [resource ID](/docs/concepts/resource-id-and-idempotency#resource-id). | | `workflow_runs.idempotency_key` unique partial index | [Idempotent starts](/docs/concepts/resource-id-and-idempotency#idempotency-key). | | Unique partial index on `workflow_id` for pending and running singleton runs | [Singleton workflows](/docs/concepts/singleton-workflows). | | `pgboss_v12_pgworkflow` schema | [pg-boss](https://pgboss.io/) job queue tables. The schema is isolated so it does not collide with another [pg-boss](https://pgboss.io/) installation in the same database. | ## Retries [#retries] Retries are scheduled by [pg-boss](https://pgboss.io/) with exponential backoff: `2^retryCount` seconds (about 1s, 2s, 4s, 8s), with up to ±50% jitter. A failed attempt includes a thrown error, an execution that passes `WORKFLOW_RUN_EXPIRE_IN_SECONDS`, and a worker that stops sending heartbeats. When the last attempt fails, the run is marked `failed`. See [Retries and timeouts](/docs/concepts/retries-and-timeouts). ## Dependencies [#dependencies] * `pg` is a peer dependency. Install it alongside `pg-workflows`. * [`pg-boss`](https://pgboss.io/) is a regular dependency. It is installed with the engine and needs no setup. To use your own [pg-boss](https://pgboss.io/) configuration, pass a `boss` instance to the `WorkflowEngine` or `WorkflowClient` constructor. ## Requirements [#requirements] * Node.js >= 18 * PostgreSQL >= 10 * `pg` >= 8 * A [Standard Schema](https://github.com/standard-schema/standard-schema) library (Zod, Valibot, ArkType) if you use `inputSchema` # Components (/docs/ui/components) ## `` [#workflowrunsdashboard] The whole dashboard in one component: live toggle, status counts, filters, runs table, pagination, and run detail. It creates its own React Query client and provider. WorkflowRunsDashboard ```tsx import { WorkflowRunsDashboard } from '@pg-workflows/ui' export default function Page() { return } ``` The components below are the pieces of the dashboard. They need a [`WorkflowRunsProvider`](/docs/ui/reference#workflowrunsprovider) above them. Each snippet is an excerpt from the full page in [Compose the components](/docs/ui/components#compose-the-components). ## `` [#runstable] The list of runs, with step progress, IDs, status, timestamps, and duration. RunsTable ```tsx ``` ## `` [#rundetail] One run: header, details, lifecycle actions, step timeline, and input and output. It loads the run and calls the actions itself. RunDetail ```tsx setRunId(null)} /> ``` ## `` [#statussummary] A count per status. Click a count to filter by that status. StatusSummary ```tsx filter({ statuses: [status] })} /> ``` ## `` [#filterbar] Search, plus status, workflow, date, and duration filters. FilterBar ```tsx ``` ## `` [#pagination] Previous and next controls for the cursor-paginated list. Pagination ```tsx setFilters({ endingBefore: runs.data?.prevCursor ?? undefined, startingAfter: undefined }) } onNext={() => setFilters({ startingAfter: runs.data?.nextCursor ?? undefined, endingBefore: undefined }) } /> ``` ## `` [#livetoggle] A Live / Paused switch for polling. You own the state and pass the matching interval to the provider. LiveToggle ```tsx ``` ## `` [#statusbadge] A status label. StatusBadge ```tsx ``` ## Compose the components [#compose-the-components] The same wiring `` uses, as one client component. Filter changes reset the page cursor, and the client-side filters (search, date, duration) and sort are applied to the current page before rendering. ```tsx 'use client' import { applyClientFilters, createFetchClient, FilterBar, LiveToggle, Pagination, RunDetail, type RunFilters, RunsTable, StatusSummary, sortRuns, useRunFilters, useWorkflowRunStats, useWorkflowRuns, WorkflowRunsProvider, } from '@pg-workflows/ui' import { QueryClient, QueryClientProvider } from '@tanstack/react-query' import { useState } from 'react' const client = createFetchClient({ baseUrl: '/workflow-runs' }) export function RunsPage() { const [queryClient] = useState(() => new QueryClient()) const [live, setLive] = useState(true) return (
setLive((v) => !v)} />
) } function Runs({ live, onToggleLive }: { live: boolean; onToggleLive: () => void }) { const { filters, setFilters, clearFilters, hasActiveFilters, serverParams } = useRunFilters() const runs = useWorkflowRuns(serverParams) const stats = useWorkflowRunStats({ workflowId: filters.workflowId }) const [runId, setRunId] = useState(null) const filter = (partial: Partial) => setFilters({ ...partial, startingAfter: undefined, endingBefore: undefined }) const rows = sortRuns( applyClientFilters(runs.data?.items ?? [], filters), filters.sort, filters.dir, ) if (runId) return setRunId(null)} /> return ( <> filter({ statuses: [status] })} /> setFilters({ endingBefore: runs.data?.prevCursor ?? undefined, startingAfter: undefined }) } onNext={() => setFilters({ startingAfter: runs.data?.nextCursor ?? undefined, endingBefore: undefined }) } /> ) } ``` *** # Custom UI with hooks (/docs/ui/hooks) Wrap your tree in `WorkflowRunsProvider` once. Every hook reads from it and returns a plain [TanStack Query](https://tanstack.com/query) result, so loading, error, and refetch states work as usual. ```tsx 'use client' import { createFetchClient, WorkflowRunsProvider } from '@pg-workflows/ui' import { QueryClient, QueryClientProvider } from '@tanstack/react-query' import type { ReactNode } from 'react' const queryClient = new QueryClient() const client = createFetchClient({ baseUrl: '/workflow-runs' }) export function Providers({ children }: { children: ReactNode }) { return ( {children} ) } ``` ## `useWorkflowRuns()`: list runs [#useworkflowruns-list-runs] ```tsx import { useWorkflowRuns } from '@pg-workflows/ui' export function FailedRuns() { const { data, isLoading } = useWorkflowRuns({ limit: 20, statuses: ['failed'] }) if (isLoading) return

Loading…

return (
    {data?.items.map((run) => (
  • {run.workflowId}: {run.status}
  • ))}
) } ``` ## `useWorkflowRun()`: one run [#useworkflowrun-one-run] Polls until the run reaches a terminal status. ```tsx import { useWorkflowRun } from '@pg-workflows/ui' export function RunStatus({ runId }: { runId: string }) { const { data: run } = useWorkflowRun(runId) return (

{run?.workflowId} is {run?.status}

) } ``` ## `useWorkflowRunStats()`: counts by status [#useworkflowrunstats-counts-by-status] ```tsx import { useWorkflowRunStats } from '@pg-workflows/ui' export function RunCounts() { const { data: stats } = useWorkflowRunStats() return (

{stats?.failed ?? 0} failed, {stats?.running ?? 0} running

) } ``` ## `useRunActions()`: control a run [#userunactions-control-a-run] ```tsx import { useRunActions } from '@pg-workflows/ui' export function RunControls({ runId }: { runId: string }) { const { cancel, resume, trigger } = useRunActions() return ( <> ) } ``` After each successful action, the run and every runs list refetch. ## `useRunFilters()`: filter state [#userunfilters-filter-state] Holds the filter, sort, and cursor state. Pass `serverParams` to `useWorkflowRuns`. ```tsx import { useRunFilters, useWorkflowRuns } from '@pg-workflows/ui' export function RunsByWorkflow() { const { filters, setFilters, serverParams } = useRunFilters() const runs = useWorkflowRuns(serverParams) return ( <>

{runs.data?.items.length ?? 0} runs on this page

) } ``` *** # Dashboard setup (/docs/ui) React components, hooks, and HTTP adapters for [pg-workflows](https://github.com/SokratisVidros/pg-workflows). Render the full dashboard, compose the individual components, or build your own UI on the hooks. All components are built with [Base UI](https://base-ui.com/), so they are accessible and accept a `render` prop to swap the underlying element. See [Styling & customization](/docs/ui/reference#styling--customization). Workflow runs dashboard ## Quickstart [#quickstart] ```bash npx @pg-workflows/ui ``` Open [http://127.0.0.1:3777](http://127.0.0.1:3777). The CLI connects to `postgres://localhost:5432/postgres`. To use another database, pass `--database-url` or set `DATABASE_URL`: ```bash npx @pg-workflows/ui --database-url=postgres://user:pass@localhost:5432/mydb ``` The dashboard lists runs and sends lifecycle actions (cancel, pause, resume, fast-forward, trigger). It registers no workflows of its own. On start it runs the engine migrations, so it creates the pg-workflows tables and [pg-boss](https://pgboss.io/) schema in the target database if they are missing. See [CLI](/docs/ui/reference#cli) for every flag and caveat. *** ## Add it to your app [#add-it-to-your-app] This walkthrough adds the dashboard to a Next.js App Router app with Tailwind CSS v4. For other servers, see [Server adapters](/docs/ui/reference#server-adapters). [`examples/dashboard`](https://github.com/SokratisVidros/pg-workflows/tree/main/examples/dashboard) is a complete working app. **1. Install** ```bash npm install @pg-workflows/ui @tanstack/react-query pg-workflows pg ``` If you don't have an app yet, `npx create-next-app@latest --ts --tailwind --app` creates one with Tailwind v4 and the `@/` import alias that the snippets below use. **2. Create one engine per process** ```ts // lib/engine.ts import { WorkflowEngine } from 'pg-workflows' import { workflows } from './workflows' // your workflow definitions, as an array // Cache the engine on globalThis. Next re-evaluates modules on hot reload, // and each new engine opens another pool and another set of workers. const globalForEngine = globalThis as unknown as { engine?: WorkflowEngine } export function getEngine() { if (!globalForEngine.engine) { const connectionString = process.env.DATABASE_URL if (!connectionString) throw new Error('DATABASE_URL is not set') globalForEngine.engine = new WorkflowEngine({ connectionString, workflows }) } return globalForEngine.engine } ``` `engine.start()` also starts queue workers, so this process executes workflow runs. Register every workflow definition here: a worker that picks up a run for an unregistered workflow fails that job with `Workflow not found`. Registered definitions also let the runs table show progress against each workflow's total step count. **3. Mount the API** with one optional catch-all route: ```ts // app/workflow-runs/[[...path]]/route.ts import { createAppRouterHandler } from '@pg-workflows/ui/next' import { getEngine } from '@/lib/engine' export const { GET, POST } = createAppRouterHandler({ engine: getEngine }) ``` Pass `getEngine` itself, not `getEngine()`. The handler calls it when a request arrives, so `next build` can import the route without a database connection. The first request awaits `engine.start()`. **4. Add the styles** to `app/globals.css`: ```css @import 'tailwindcss'; @import '@pg-workflows/ui/styles.css'; @source '../node_modules/@pg-workflows/ui/dist'; ``` The `@source` path is relative to the CSS file. It lets Tailwind generate the utility classes that the components use. **5. Render the dashboard** ```tsx // app/page.tsx import { WorkflowRunsDashboard } from '@pg-workflows/ui' export default function Page() { return } ``` Set `DATABASE_URL` in `.env.local`, run `npm run dev`, and open [http://localhost:3000](http://localhost:3000). *** # API reference (/docs/ui/reference) * [Entry points](/docs/ui/reference#entry-points) * [CLI](/docs/ui/reference#cli) * Components: [`WorkflowRunsDashboard`](/docs/ui/reference#workflowrunsdashboard) · [`RunsTable`](/docs/ui/reference#runstable) · [`RunDetail`](/docs/ui/reference#rundetail) · [`StatusSummary`](/docs/ui/reference#statussummary) · [`FilterBar`](/docs/ui/reference#filterbar) · [`Pagination`](/docs/ui/reference#pagination) · [`LiveToggle`](/docs/ui/reference#livetoggle) · [`StatusBadge`](/docs/ui/reference#statusbadge) * [`WorkflowRunsProvider`](/docs/ui/reference#workflowrunsprovider) * Hooks: [`useWorkflowRuns`](/docs/ui/reference#useworkflowrunsparams) · [`useWorkflowRun`](/docs/ui/reference#useworkflowrunid) · [`useWorkflowRunStats`](/docs/ui/reference#useworkflowrunstatsparams) · [`useRunActions`](/docs/ui/reference#userunactions) · [`useRunFilters`](/docs/ui/reference#userunfiltersinitial) · [`useWorkflowRunsClient`](/docs/ui/reference#useworkflowrunsclient) * [Helpers](/docs/ui/reference#helpers) * [`createFetchClient`](/docs/ui/reference#createfetchclientoptions) * [Server adapters](/docs/ui/reference#server-adapters) * [HTTP API](/docs/ui/reference#http-api) * [Security](/docs/ui/reference#security) * [Styling & customization](/docs/ui/reference#styling--customization) * [Architecture](/docs/ui/reference#architecture) ## Entry points [#entry-points] Client code never imports the server entries, so a browser bundle does not pull in the engine. | Import | Contents | Runs | | ----------------------------- | ------------------------------------------------------------------------------------------ | ---------------- | | `@pg-workflows/ui` | Components, hooks, provider, helpers, and a re-export of `createFetchClient` and its types | client | | `@pg-workflows/ui/client` | `createFetchClient` and types, no React | client or server | | `@pg-workflows/ui/server` | `createWorkflowRunsApi`, `toFetchHandler`, `toNodeHandler`, `HttpError`, `toErrorResponse` | server only | | `@pg-workflows/ui/next` | `createAppRouterHandler`, `createPagesApiHandler`, `createRouteHandlers` | server only | | `@pg-workflows/ui/tailwind` | Tailwind preset for a subset of the `pgw-*` color tokens | build | | `@pg-workflows/ui/styles.css` | CSS variables (light and dark), Tailwind `@theme` tokens, and component styles | client | | `pg-workflows-ui` (bin) | Standalone localhost dashboard | CLI | Peer dependencies: `react >= 18`, `react-dom >= 18`, `@tanstack/react-query >= 5`, `tailwindcss ^4`, `pg-workflows >= 0.13.0`. `pg-workflows` in turn needs `pg`. ## CLI [#cli] ```bash npx @pg-workflows/ui [--database-url=] [--port=3777] ``` | Flag | Default | Description | | ---------------- | --------------------------------------------------------- | -------------------------- | | `--database-url` | `DATABASE_URL`, else `postgres://localhost:5432/postgres` | Postgres connection string | | `--port` | `3777` | Port to listen on | | `-h`, `--help` | | Print usage | * The server binds to `127.0.0.1` only. It has no authentication and no `resolveContext`, so anyone who can reach the port can read and change every run. Do not expose it. * It starts an engine with no registered workflows. `engine.start()` runs migrations if needed. Because no workflows are registered, the runs table shows step progress only from each run's timeline. * `engine.start()` also starts queue workers on the shared run queue. A job those workers pick up fails with `Workflow not found`, because the CLI has no definitions. Prefer it for local databases, and embed the components in your app when it runs against a live queue. * It serves the API under `/workflow-runs` and the prebuilt dashboard on every other path. ## Components [#components] Every component takes `className`, `style`, and `render` (see [Styling & customization](/docs/ui/reference#styling--customization)) and forwards a `ref` to its root element. `StatusSummary` is the exception when it renders nothing. ### `WorkflowRunsDashboard` [#workflowrunsdashboard] Self-contained dashboard. It creates its own `QueryClient` (with query retries off) and `WorkflowRunsProvider`. Pass exactly one of `baseUrl` or `client`. | Prop | Type | Default | Description | | ---------------- | ------------------------------ | ----------------------------------- | ---------------------------------------------------------------------------------------------------------------------------------- | | `baseUrl` | `string` | | Prefix of the mounted API routes. Mutually exclusive with `client`. Read once on mount. | | `client` | `WorkflowRunsClient` | | A client, usually from `createFetchClient`. Mutually exclusive with `baseUrl`. Read once on mount. | | `pollIntervalMs` | `number` | `5000` while Live, `0` while paused | Refresh interval. When set, it overrides the Live toggle: the button still switches, but the interval stays fixed. | | `selectedRunId` | `string \| null` | | Controlled selection. Pair with `onSelectRun` and your router for deep links. When omitted, the dashboard tracks selection itself. | | `onSelectRun` | `(id: string \| null) => void` | | Called when a row is opened (`id`) or the detail view is closed (`null`). | The workflow filter lists the workflow IDs on the current page. Style state: `{ selected: boolean }`. ### `RunsTable` [#runstable] | Prop | Type | Default | Description | | --------------- | ---------------------- | ------- | ------------------------------------------------------------- | | `runs` | `WorkflowRun[]` | | Rows to render, in order. | | `onSelectRun` | `(id: string) => void` | | Called when a row is clicked, or on Enter or Space. | | `selectedRunId` | `string \| null` | | Highlights the matching row. | | `isLoading` | `boolean` | `false` | While `runs` is empty, shows "Loading…" instead of "No runs". | Columns: Workflow, Run ID (copyable), Resource ID, Status, Started, Completed, Duration. The Workflow cell shows step progress for running, paused, and failed runs, using the larger of the timeline's step count and `totalSteps`. Style state: `{ empty: boolean, loading: boolean }`. ### `RunDetail` [#rundetail] Must render under `WorkflowRunsProvider`. It calls `useWorkflowRun(runId)` and `useRunActions()`. | Prop | Type | Default | Description | | -------- | ------------ | ------- | ------------------------------------------------------------------------------- | | `runId` | `string` | | Run to load. | | `onBack` | `() => void` | | Renders a back control that calls this. When omitted, there is no back control. | It renders: * A header with the workflow ID, run ID, and step progress. * A details grid: workflow, resource ID, status, timestamps, duration, retries, priority, job, and error. Priority `100`, `0`, and `-100` display as `high`, `normal`, and `low`. * The actions: Cancel, Pause, Resume, Fast-forward, and Trigger. All are disabled once the run is terminal. Pause is enabled only while the run is `running`, and Resume only while it is `paused`. Each action shows a success or error message. * The step timeline, and input and output JSON. A failed run's error is shown above the steps. Fast-forward sends no `data`. The engine completes the current wait step only when the run is paused on one, and otherwise returns the run unchanged. Trigger always sends an event named `resume` with no data. To send another event, call `useRunActions().trigger` yourself. Style state: `{ phase: 'loading' | 'error' | 'ready', status?: string }`. ### `StatusSummary` [#statussummary] | Prop | Type | Default | Description | | ---------------- | -------------------------------------------- | ------- | -------------------------------------------------------------------- | | `counts` | `Partial>` | | Counts by status, usually `useWorkflowRunStats().data`. | | `onSelectStatus` | `(status: WorkflowRunStatus) => void` | | Called when a count is clicked. | | `trailing` | `ReactNode` | | Rendered after the counts. | | `stat` | `{ className?, style?, render? }` | | Style hooks for each count button. They receive that button's state. | Renders one button per status with a count above zero. When every count is `0`, it renders nothing, so there is no element for `ref`. Style state: `{ empty: boolean }` (always `false` while rendered). ### `FilterBar` [#filterbar] | Prop | Type | Description | | ------------------ | ---------------------------------------- | ---------------------------------------------------------------------------------------------------------------------------------------------- | | `filters` | `RunFilters` | Current filters, from `useRunFilters`. | | `hasActiveFilters` | `boolean` | Enables the Clear control. | | `workflowIds` | `string[]` | Options for the workflow filter. | | `onFiltersChange` | `(partial: Partial) => void` | Called with the changed fields. Reset `startingAfter` and `endingBefore` here, or the next request stays on a cursor from the previous filter. | | `onClear` | `() => void` | Called by the Clear control. | Style state: `{ active: boolean }`. ### `Pagination` [#pagination] | Prop | Type | Default | Description | | ------------ | ------------ | ------- | ----------------------------------------- | | `hasPrev` | `boolean` | | Enables Prev. | | `hasNext` | `boolean` | | Enables Next. | | `onPrev` | `() => void` | | Called by Prev. | | `onNext` | `() => void` | | Called by Next. | | `isFetching` | `boolean` | `false` | Disables both buttons while a page loads. | Style state: `{ hasPrev, hasNext, fetching }`. ### `LiveToggle` [#livetoggle] A [Base UI](https://base-ui.com/) `Toggle`. | Prop | Type | Default | Description | | -------------- | ------------ | ------- | --------------------------------------------------------------------------- | | `isLive` | `boolean` | | Pressed state. | | `isFetching` | `boolean` | | Sets `data-fetching` while a request is in flight. | | `onToggle` | `() => void` | | Called on press. Pass `pollIntervalMs={isLive ? 5000 : 0}` to the provider. | | `nativeButton` | `boolean` | `true` | Set to `false` when `render` is not a `