Files
Pranay Prakash d07c668495 [world-postgres] Honor $retention: 0 when a run finishes
A run started with `experimental_retention: 0` carries the reserved
`$retention: '0'` attribute, and a World that implements retention is
expected to delete that run's user payloads once it reaches a terminal
state. world-postgres was ignoring the attribute entirely, so the option
was silently a no-op here: the SDK documents the Postgres World as
implementing retention, and it did not.

The purge is one transaction issued after the terminal event row commits
and before the terminal NOTIFY.

- After the event row, because the terminal event is itself
  payload-bearing — `run_completed` carries the run's output — so
  anything earlier leaves that one payload behind.
- Before the NOTIFY, so a waiter woken by it re-reads an already-expired
  run rather than catching the output on its way out.

One transaction is also what makes the ordering rule the Vercel World has
to hand-sequence (data must become unreadable no later than it becomes
unrecoverable) a non-issue here: `expired_at` and the last cleared byte
commit together, so no reader can observe a half-purged run. It cannot,
however, be folded into the terminal run UPDATE itself, because that
UPDATE commits before the event insert — so a crash in between leaves a
terminal run holding data with no `expired_at`, exactly the window the
Vercel World has, and nothing retries it.

Both halves of every payload column are cleared. Each CBOR column has a
legacy JSONB twin beside it and the read paths fall back to the twin
(`value.output ||= value.outputJson`), so clearing only `*_cbor` would
leave the payload in place and resurrect it on the next read. Payload
columns are set through raw SQL `NULL` rather than a JS `null`: the CBOR
codec's `toDriver` would otherwise encode `null` into a one-byte CBOR
value and store that.

`expired_at` is stamped because it is what the CLI and web UI gate their
`<data expired>` rendering on. Without it the deletion is silent, and a
purged run reads back as one that never had any input.

Rows are kept — a purged run stays listable and traceable, matching the
Vercel World's contract — and so is plaintext metadata (status, name,
attributes, errorCode, executionContext, hook tokens). `executionContext`
in particular is excluded deliberately: the reference implementation
purges input/output/error and nothing else, and widening that here would
delete something the contract does not ask us to.

Stream chunk rows are blanked rather than deleted. The reader closes on
the `eof` row and `streams.list()` enumerates from these rows, so
deleting them would turn a finished stream into one that never
terminates and a run's stream list into an empty one. An empty `bytea`
carries no user data and keeps both behaviors.

Every value other than the literal `'0'` keeps the data, including a
well-formed non-zero duration. That is the load-bearing half: the unit
`$retention` is measured in is deliberately undecided, so an SDK that
starts sending a unit-bearing value to a World that predates the decision
must get the safe answer. The failure mode worth engineering against is
not a purge that does not fire, it is a purge that fires on a value
nobody meant as "delete my data".

Purge failures are caught and logged rather than thrown: a run must still
be able to finish. The cost of a failure is a run that keeps data it
asked to have deleted, which is loud in the log and safe on disk.

No migration: `expired_at` has existed since 0002_add_expired_at.sql and
nothing wrote it until now.
2026-08-28 12:27:09 -07:00
..
2025-10-23 12:07:52 +03:00
2026-08-21 22:17:38 -07:00
2025-10-23 12:07:52 +03:00
2026-08-21 22:17:38 -07:00
2025-10-23 12:07:52 +03:00

@workflow/world-postgres

An embedded worker and workflow system backed by PostgreSQL for multi-host self-hosted solutions. This is a reference implementation. A production system might run workers in separate processes with a dedicated queuing system.

Installation

npm install @workflow/world-postgres
# or
pnpm add @workflow/world-postgres
# or
yarn add @workflow/world-postgres

Usage

Basic setup

The PostgreSQL World can be configured by setting the WORKFLOW_TARGET_WORLD environment variable to the package name:

export WORKFLOW_TARGET_WORLD="@workflow/world-postgres"

Configuration

Configure the PostgreSQL world using environment variables:

# Required: PostgreSQL connection string
export WORKFLOW_POSTGRES_URL="postgres://username:password@localhost:5432/database"

# Optional: Job prefix for queue operations
export WORKFLOW_POSTGRES_JOB_PREFIX="myapp"

# Optional: Worker concurrency (default: 10)
export WORKFLOW_POSTGRES_WORKER_CONCURRENCY="10"

# Optional: Internal pg.Pool max size (default: 10)
export WORKFLOW_POSTGRES_MAX_POOL_SIZE="10"

# Optional: Let the application coordinate shutdown (default: false)
export WORKFLOW_POSTGRES_APPLICATION_MANAGED_SHUTDOWN="1"

# Optional: Maximum Hook minimum retention in days (default: 30)
export WORKFLOW_POSTGRES_HOOK_RETENTION_LIMIT_DAYS="30"

Programmatic usage

You can also create a PostgreSQL world directly in your code:

import { createWorld } from "@workflow/world-postgres";

const world = createWorld({
  connectionString: "postgres://username:password@localhost:5432/database",
  jobPrefix: "myapp", // optional
  queueConcurrency: 50, // optional
  maxPoolSize: 10, // optional, overrides WORKFLOW_POSTGRES_MAX_POOL_SIZE when `pool` is omitted
});

// Or pass an existing pg.Pool (shared with your app Drizzle, etc.); `world.close()` will not end it.
import { Pool } from "pg";
const pool = new Pool({ connectionString: process.env.DATABASE_URL });
const worldFromPool = createWorld({ pool });

Application-managed shutdown

By default, Graphile Worker responds automatically when the application is asked to shut down. If your application already coordinates shutdown, set WORKFLOW_POSTGRES_APPLICATION_MANAGED_SHUTDOWN=1 when selecting the package with WORKFLOW_TARGET_WORLD, or set applicationManagedShutdown: true when calling createWorld() directly. Await world.close() from your shutdown path so Graphile Worker cannot terminate the process as soon as its queue stops, before your application finishes closing dependent resources:

import { createWorld } from '@workflow/world-postgres';

const world = createWorld({
  connectionString: process.env.DATABASE_URL!,
  applicationManagedShutdown: true,
});

await world.start();

Use this option only when your application or framework has its own shutdown hook. Handle cleanup errors there and await world.close() first, then close the workflow HTTP server and any caller-owned pg.Pool.

Closing the world stops the queue from accepting new jobs and waits for active jobs. After Graphile Worker's graceful-shutdown timeout (5s by default), it aborts any workflow HTTP request that is still pending. Graphile Worker then unlocks the same row through its normal failure handling. Graphile counts a delivery attempt when it claims the row, so the aborted delivery consumes that attempt and is retried only if its Graphile attempt budget remains. A one-attempt or final-attempt job is unlocked but not retried. The shutdown handler does not create a replacement row.

An aborted HTTP request does not guarantee that its server-side handler stopped, so workflow and step handlers must continue to tolerate at-least-once execution. Keep the workflow HTTP routes and any caller-owned pool available until world.close() resolves.

Configuration options

Option Type Default Description
connectionString string process.env.WORKFLOW_POSTGRES_URL, process.env.DATABASE_URL, or 'postgres://world:world@localhost:5432/world' Used only when pool is omitted, to construct an internal pool
maxPoolSize number process.env.WORKFLOW_POSTGRES_MAX_POOL_SIZE or pg.Pool default (10) Optional. Sets the internal pg.Pool max size when createWorld() creates the pool
pool pg.Pool Not applicable Optional. When set, used for Drizzle, Graphile Worker, and stream writes. world.close() does not end it.
jobPrefix string process.env.WORKFLOW_POSTGRES_JOB_PREFIX Optional prefix for queue job names
queueConcurrency number 50 Number of concurrent active step executions per process. Must be high enough to cover any parent→child workflow polling in flight because each Run#returnValue await holds a worker slot until the child run terminates.
applicationManagedShutdown boolean false; WORKFLOW_POSTGRES_APPLICATION_MANAGED_SHUTDOWN=1 enables it for the default package configuration Whether the application coordinates shutdown and awaits world.close() instead of Graphile Worker responding automatically.

Environment variables

Variable Description Default
WORKFLOW_TARGET_WORLD Set to "@workflow/world-postgres" to use this world -
WORKFLOW_POSTGRES_URL PostgreSQL connection string DATABASE_URL or 'postgres://world:world@localhost:5432/world'
WORKFLOW_POSTGRES_JOB_PREFIX Prefix for queue job names -
WORKFLOW_POSTGRES_WORKER_CONCURRENCY Number of concurrent workers 50
WORKFLOW_POSTGRES_MAX_POOL_SIZE Internal pg.Pool max size 10
WORKFLOW_POSTGRES_APPLICATION_MANAGED_SHUTDOWN Set to 1 when the application coordinates shutdown and awaits world.close() unset (false)
WORKFLOW_POSTGRES_HOOK_RETENTION_LIMIT_DAYS Maximum Hook minimum retention in days 30

When pool is omitted, maxPoolSize precedence is: createWorld({ maxPoolSize }), then WORKFLOW_POSTGRES_MAX_POOL_SIZE, then the pg.Pool default.

For higher worker concurrency, Graphile Worker recommends setting maxPoolSize to 10 or queueConcurrency + 2, whichever is larger.

Database setup

This package uses PostgreSQL with the following components:

  • Graphile Worker: For queue processing and job management
  • Drizzle ORM: For database operations and schema management
  • pg (node-postgres): For PostgreSQL client connections. Drizzle and Graphile Worker share a pg.Pool, while LISTEN uses a dedicated pg.Client created from the same connection options.

Quick setup with CLI

Set up your database with the included CLI tool:

# npm
npx --package=@workflow/world-postgres bootstrap

# pnpm
pnpm dlx --package @workflow/world-postgres bootstrap

# Yarn
yarn dlx --package @workflow/world-postgres bootstrap

# Bun
bunx --package @workflow/world-postgres bootstrap

The CLI and runtime World automatically load the connection string from:

  1. WORKFLOW_POSTGRES_URL environment variable
  2. DATABASE_URL environment variable
  3. Default: postgres://world:world@localhost:5432/world

Database schema

The setup creates the following tables:

  • workflow_runs: Stores workflow execution runs
  • workflow_events: Stores workflow events
  • workflow_steps: Stores individual workflow steps
  • workflow_hooks: Stores webhook hooks
  • workflow_stream_chunks: Stores streaming data chunks

You can also access the schema programmatically:

import { runs, events, steps, hooks, streams } from '@workflow/world-postgres';
// or
import * as schema from '@workflow/world-postgres/schema';

Make sure your PostgreSQL database is accessible and the user has sufficient permissions to create tables and manage jobs.

Data retention

Postgres World does not yet perform general workflow-run cleanup. After a retained Hook's run ends and its deadline passes, reads treat the Hook as absent and its token can be reused. If the token is never reused, the expired workflow_hooks row remains.

Features

  • Durable storage: Stores workflow runs, events, steps, hooks, and webhooks in PostgreSQL
  • Queue processing: Uses Graphile Worker as the durable queue and executes jobs over the workflow HTTP routes
  • Durable delays: Reschedules waits and retries in PostgreSQL
  • Streaming: Real-time event streaming capabilities
  • Health checks: Built-in connection health monitoring
  • Configurable concurrency: Adjustable worker concurrency for queue processing

Queue behavior

  • Graphile jobs are acknowledged only after execution finishes, or after the worker durably schedules a delayed follow-up job
  • Backlog stays in PostgreSQL when all execution slots are busy
  • Retry and sleep-style delays use Graphile runAt scheduling
  • Workflow orchestration and queued step execution are both sent through /.well-known/workflow/v1/flow

Development

For local development, you can use the included Docker Compose configuration:

# Start PostgreSQL database
docker-compose up -d

# Create and run migrations
pnpm drizzle-kit generate
pnpm drizzle-kit migrate

# Set environment variables for local development
export WORKFLOW_POSTGRES_URL="postgres://world:world@localhost:5432/world"
export WORKFLOW_TARGET_WORLD="@workflow/world-postgres"

Testing

Integration tests use Testcontainers to start a PostgreSQL container. Docker must be installed and running before you run tests.

  • Linux/macOS: Start the Docker daemon (e.g. sudo systemctl start docker or Docker Desktop).
  • WSL2: Use Docker Desktop with WSL2 integration, or run the Docker engine inside WSL and ensure the daemon is started. Verify with docker info.

Then from the package directory:

pnpm build
pnpm test

World selection

To use the PostgreSQL world, set the WORKFLOW_TARGET_WORLD environment variable to the package name:

export WORKFLOW_TARGET_WORLD="@workflow/world-postgres"