Quickstart
Run your first durable workflow against Postgres in about two minutes.
1. Start Postgres (skip if you already have one):
docker run -d --name pg-workflows-db -e POSTGRES_PASSWORD=postgres -p 5432:5432 postgres:172. Create a project:
mkdir workflows-demo && cd workflows-demo
npm init -y
npm install pg-workflows pg zod
npm install -D tsx3. Save this as index.ts:
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:
npx tsx index.tsAfter 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
step.waitFor pauses the run until your code calls triggerEvent. A paused run holds no worker and no connection.
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()