Examples
Conditional steps, loops, reminders, polling, retries, and progress.
Each snippet defines a workflow. To run one, pass it to workflows: [...] on an engine like the one in the quickstart, then call engine.startWorkflow. Full runnable scripts are in examples/node:
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 |
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) architecture |
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.
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
Give each iteration its own step ID. A retry skips the items that already finished.
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 (default 300). For a long list, split the work across child workflows instead.
Follow-up after a delay
step.delay pauses the run for a duration. The run holds no worker while it waits and survives restarts.
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
Return false to keep polling, or a value to finish. result.timedOut tells you which way it ended.
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
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.
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
checkProgress returns the run plus step counts. Call it on an engine that has the workflow registered.
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 behind your own API. Don't call the engine from the browser.