Files
Nathan Rajlich 03455a2979 Carry run identity on step-dispatch messages; drop the blocking runs.get from the queued-step prologue (#3457)
* Carry immutable run identity on step-dispatch messages; drop the blocking runs.get from the consumer prologue

Closes #3456. Every queued step execution paid a runs.get round trip
before its step_started claim — one RTT per branch on the TTLS-critical
path, and under a 256-branch fan-out burst the read amplification drove
that read to p90 ~5.1s (durabench parallel sweeps), smearing branch
starts.

The dispatch sites (node dispatch loop, delayed retries, the suspension
handler's resilient publish, and the quickjs engine's queueStepMessage)
now stamp WorkflowInvokePayload.runContext with the fields the consumer
actually needs — deploymentId, specVersion, startedAt, rootRunId — all
immutable for the life of a run and known from the run row the producer
already holds. A consumer that receives it skips the run fetch: the
run-status early exit is enforced by the step_started claim itself
(RunExpired → gone, terminal step → skipped), guardDeployment takes the
carried identity, and only the fan-out's LAST completer fetches the full
run row, lazily, for its inline replay — once per fan-out instead of
once per branch. The deployment-mismatch re-route now also preserves
stepInput/runContext on the re-enqueued payload.

Messages without runContext (older producers) keep the legacy prologue;
messages are deployment-pinned, so mixed handling within one run cannot
occur.

* Address review: terminal-only lazy status gate, terminal-run start fence in local worlds, prologue telemetry, last-completer coverage

- The last completer's lazy runs.get result is now gated on
  isTerminalWorkflowRunStatus (with a debug log): a stale 'pending' read
  — a run with completed steps has necessarily started — no longer
  silently abandons the fan-out's continuation; it falls through to the
  inline replay, whose next entity write is fenced server-side if the
  run truly ended meanwhile.
- world-local / world-postgres now reject step_started on terminal runs
  even when the step row still reads 'running' (a redelivered start a
  previous delivery claimed): starting work on a finished run is never
  valid, and previously the body re-ran with its outcome unconsumable.
  In-flight steps still write their terminal events unchanged. This
  closes the adapter gap behind the fetch-free prologue's reliance on
  the step_started claim as the run-liveness check, and the prologue
  comment now states the contract precisely.
- workflow.step.dispatch_prologue span attribute ('run_context' |
  'runs_get') makes fetch-free adoption and the saved round trip
  observable during version-skew windows.
- Restated why the eager redelivery re-ensure survives on the
  fetch-free path (no run fetch to overlap; still cheaper than the
  in-band recovery's failed-start round trip).
- New two-phase fan-out coverage: a real replay emits the queued step
  message (asserting the stamped runContext), then its redelivery runs
  as the LAST completer — zero reads before the step, exactly one lazy
  runs.get, run completed; plus the stale-'pending' fall-through and
  the genuinely-terminal skip.
2026-09-11 21:31:15 +00:00
..
2025-10-23 12:07:52 +03:00
2026-09-09 12:12:11 -07:00
2025-10-23 12:07:52 +03:00
2026-09-09 12:12:11 -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"