pg-workflows

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
ScriptShows
example:basicSequential steps and checkProgress
example:approval-flowwaitFor with triggerEvent
example:timeoutwaitFor timing out, and bringing your own pg-boss
example:pollingstep.poll against a simulated payment API
example:cronA recurring schedule
example:microservices:worker / example:microservices:apiThe 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.

On this page