octaflow
v0.25.0
Published
Durable DAG workflow engine: Zod-typed steps over Postgres + pg-boss, with an optional AI add-on (token/cost/quota). Layered single package.
Maintainers
Readme
octaflow
Durable workflows for TypeScript that run on the Postgres you already have.
Declare a DAG of Zod-typed steps. The engine runs each step as soon as its dependencies complete, persists every transition, retries failures, and picks up where it left off after a restart. There is no workflow server to operate, no control plane, and no vendor — it's a library you import, not a platform you adopt.
📖 Read the docs →
Quick start · Concepts · Postgres & pg-boss · API reference
const wf = buildWorkflow({
type: 'publish-article',
inputSchema: z.object({ draftId: z.string() }),
steps: { fetchDraft, summarize, translate, publish },
});summarize and translate both depend on fetchDraft, and publish waits for both — so the
engine runs the middle pair concurrently and fans back in, without you scheduling anything:
┌── summarize ──┐
fetchDraft ──┤ ├── publish
└── translate ──┘That shape isn't inferred from a trace — it is the workflow. The DAG is a plain value you can walk before anything runs:
for (const s of wf.definition.steps) {
console.log(`${s.key.padEnd(12)} type=${s.type.padEnd(10)} deps: ${s.dependencies?.join(', ') || '—'}`);
}fetchDraft type=fetch deps: —
summarize type=summarize deps: fetchDraft
translate type=translate deps: fetchDraft
publish type=publish deps: summarize, translateThis is the core design choice. Flow is declarative, where Temporal, Inngest, Trigger.dev and DBOS are imperative: there you write a function, and the graph exists only as the trace of what it did. Both models are durable — they trade off differently, and the tradeoff is spelled out below.
Watch it run
A fuller pipeline than the snippet above — dynamic fan-out over two map steps at once, a
concurrency cap so only 2 of 6 images encode at a time, a flaky step retried, a sub-workflow,
a step suspended on an external event under a deadline, a durable sleep, and a branch where
the arm the review didn't choose is skipped but the join still fires. Then a second workflow
fails and its completed steps roll back in reverse.
Every line is an engine transition emitted through the FlowObserver seam —
the same one you'd point at OpenTelemetry or an events table — not a console.log in a handler.
npx tsx scripts/demo.ts # reproduce itContents
How it compares
All of these give you durable execution. They differ in what you operate and how the workflow is expressed.
| | Flow | Temporal | Inngest | Trigger.dev | BullMQ | |---|---|---|---|---|---| | Infra you run | Postgres | a cluster (frontend, history, matching, worker) + a DB | none — engine is hosted | Postgres + Redis (+ ClickHouse recommended at volume) | Redis | | Self-host | it's a library | yes, MIT | no — engine and dashboard are Inngest-hosted | yes, Apache-2.0 | it's a library | | Model | declarative — a static DAG value | imperative workflow code | imperative step functions | imperative tasks | job queue; flows are trees | | Fan-in / diamond deps | yes | yes | yes | yes | no — a job can't be shared by two branches | | Inspect before running | yes — the DAG is a value | no — the graph is the execution trace | no | no | yes — the flow tree is data | | Web dashboard | none — build one | yes | yes | yes | via third-party UIs | | Languages | TypeScript | polyglot SDKs | TS, Python, Go, Kotlin | TS | Node (+ ports) | | Maturity | pre-1.0 | mature | mature | mature | mature, widely deployed |
The honest summary: Flow is the smallest thing that is still a real workflow engine. If you already run Postgres, it adds no infrastructure — the queue (pg-boss) is Postgres too. You give up the dashboards, the polyglot SDKs, and the operational maturity that the others have earned.
When not to use it
Reach for something else if:
- Your control flow is genuinely dynamic. A declarative DAG is fixed at definition time.
Flow softens this a lot —
whenguards and joins (if/else over a static graph),defineMapStep(runtime-sized fan-out), sub-workflows, andwaitForEvent— but the set of steps is still fixed up front. If your process is "loop until a human approves, branching on whatever they typed," an imperative durable function will express it more naturally: Flow can pick a branch, not invent a step. - You want a UI out of the box. Flow ships a wire-safe projection
(
toPublicWorkflow) and lifecycle events, not a dashboard. You build it. - You need non-TypeScript workers. The DAG and its schemas are TypeScript values.
- The caller is waiting. Every step is a persisted transition, and in production a queue hop:
a step starts when a worker fetches it — on the next poll, unless
burstWhenBatchFullis keeping it busy — and adds engine time on top of your handler. Anything inside a request/response — validate, transform, answer — is a function call. Put a workflow behind a request only to start it, and let the page poll or subscribe for the outcome. - One job with retries. A single step that must run later and try again is a queue job; pg-boss on its own does that with less. A DAG earns its keep when steps depend on each other, run in parallel, or must be resumed and read back.
- You can't run Postgres, or you need throughput past what a Postgres-backed queue gives you.
- You need a support contract, or an API frozen by a 1.0 promise. This is pre-1.0 and 0.x minors can break.
Flow fits best when the work is a known pipeline — ingest → enrich → summarize → publish — that must survive crashes, retry sanely, and stay legible to the next person who reads it.
Performance
Reproduce with npx tsx scripts/bench.ts (Docker required — it starts Postgres 17 via
Testcontainers). Workload: 200 workflows × 6 steps in a root → 4 parallel → join diamond.
Handlers are no-ops, so this measures what the engine costs per step — claiming it,
reading dependency outputs, persisting the transition, recomputing readiness — not your work.
Engine + Postgres store (in-process dispatcher), per-step latency:
| concurrency | steps/sec | p50 | p95 | p99 | |---|---|---|---|---| | 1 | 1,031 | 1.0 ms | 2.1 ms | 2.9 ms | | 4 | 1,932 | 2.1 ms | 3.9 ms | 4.8 ms | | 16 | 2,108 | 7.1 ms | 12.7 ms | 15.9 ms | | 64 | 2,270 | 26.8 ms | 46.8 ms | 64.0 ms |
End-to-end through pg-boss workers — the full production path, batch 25:
| workers | burstWhenBatchFull | concurrency | steps/sec |
|---|---|---|---|
| 1 | off | 1 | 50 |
| 1 | on | 1 | 274 |
| 1 | on | 8 | 646 |
| 4 | on | 8 | 902 |
That first row is not a ceiling, it's a polling artifact — and the fix is configuration, not
architecture. A worker drains a batch in milliseconds, then waits out the 0.5 s interval, so
burstWhenBatchFull is the setting that matters: it keeps fetching while batches come back
full. concurrency (steps run at once from one batch) then compounds on top — but on its own,
without burst, it changes nothing at all, because the wait, not the work, is the bottleneck.
Budget connections before raising concurrency: each in-flight step holds one, so
workers × concurrency must fit your pool and Postgres max_connections.
How to read this. Measured on an M-series Mac with Postgres in Docker, which has markedly slower disk I/O than a Linux host — expect better on a real server. These are an order of magnitude and a scaling shape, not a score. A Redis-backed job queue will beat these numbers outright, because it isn't writing a durable transition per step to a relational database; that write is the feature. And in any real workflow, handler time dwarfs the 1–3 ms of engine overhead, so the practical question is usually whether ~1 ms per transition is acceptable next to what your steps actually do.
Features
| | Capability |
|---|---|
| 🧩 | Typed DAG — Zod-validated input/output per step; dependency outputs are typed |
| ⚡ | Auto-parallelism — dependency-free steps run concurrently; a step starts when all its deps complete |
| 🔁 | Retry & timeout — per-step maxAttempts, fixed/exponential backoff, wall-clock timeout |
| 💤 | Durable sleep — hold a step in the queue for N ms (survives restarts) |
| 🔀 | Conditional branching — when guards skip a step and its branch; join: 'any' converges |
| ⏳ | Deadlines — a budget on a suspended step (fail, or continue with a stand-in answer) and on a whole run |
| ♻️ | Retry a failed run — retryWorkflow resumes from the failure point; completed steps keep their output |
| 🚦 | Concurrency & rate limiting — per-step-type caps and token buckets via a pluggable gate |
| ⏰ | Cron / scheduled starts — fire workflows on a schedule (pg-boss) |
| 🔑 | Start idempotency — a dedup key collapses double-clicks / overlapping ticks |
| 🗺️ | Dynamic fan-out / map — spawn one child step per item of a runtime-sized list |
| ⏸️ | Signals / waitForEvent — suspend a step until an external event (resumeStep), with an optional deadline |
| 🪆 | Sub-workflows — a step starts a child workflow and awaits its result |
| ↩️ | Saga compensation — run rollback handlers in reverse order on failure |
| 🚑 | Crash recovery — a step whose worker died is re-queued while its attempt budget lasts |
| 💓 | Heartbeats — a long step proves it is alive, so a dead one is caught in seconds, not minutes |
| 🔭 | Observability — lifecycle events (run history), per-step spans, and createMeterObserver turning the events into OpenTelemetry-style metrics — all pluggable, all structural |
| 🤖 | AI add-on — instrumented models, token/cost capture, quota, daily rollups |
| 🧱 | Pluggable everything — WorkflowStore, Dispatcher, StepGate, FlowObserver, hooks |
Installation
pnpm add octaflow zodzod is a required peer. The heavy dependencies are optional peers — install only
what the layers you import need:
# Postgres store / gate / event sink
pnpm add pg
# the same tables as Drizzle column sets, for a host that lets drizzle-kit own its migrations
pnpm add drizzle-orm
# pg-boss dispatcher, workers, cron scheduler
pnpm add pg-boss
# the AI add-on
pnpm add ai @ai-sdk/providerPure in-memory usage (great for tests and single-process apps) needs nothing beyond
zod— the engine,defineStep/buildWorkflow, and the in-memory store are all in the core.
The Postgres tables
octaflow/postgres owns five tables and ships their DDL: flowStoreDdl() (flow_workflow,
flow_workflow_step), flowEventDdl() (flow_step_event) and flowGateDdl()
(flow_rate_bucket, flow_step_lease). Each is idempotent CREATE TABLE IF NOT EXISTS, takes
a schema, and is meant to be pasted into your own migration system; applySchema(pool, ddl)
is the dev and test shortcut.
The Drizzle columns
The rule for every table-owning capability on the platform: ./postgres ships the store over
the structural executor plus its xDdl(), and ./drizzle ships the same table as a spreadable
column set, so a host on Drizzle declares it in its own schema and drizzle-kit owns the
migration. Where the store's queries are one SQL statement by nature — every one of octaflow's
is — ./drizzle ships the columns only, and the store stays the SQL one. So
octaflow/drizzle exports flowWorkflowColumns, flowWorkflowStepColumns,
flowStepEventColumns, flowRateBucketColumns and flowStepLeaseColumns (drizzle-orm an
optional peer) and nothing that runs a query. The foreign keys, the composite keys, the unique
constraint and the indexes are yours to declare on the pgTable; each column set's doc lists
the ones its DDL creates, and two of them the store relies on:
import { sql } from 'drizzle-orm';
import { foreignKey, index, pgTable, primaryKey, unique, uniqueIndex } from 'drizzle-orm/pg-core';
import { flowWorkflowColumns, flowWorkflowStepColumns, flowStepLeaseColumns } from 'octaflow/drizzle';
export const flowWorkflow = pgTable('flow_workflow', flowWorkflowColumns, (t) => [
index('flow_workflow_partition_status_idx').on(t.partitionKey, t.status),
index('flow_workflow_partition_type_idx').on(t.partitionKey, t.type),
index('flow_workflow_partition_entity_idx').on(t.partitionKey, t.entityRef),
index('flow_workflow_parent_idx').on(t.parentWorkflowId),
// The deadline sweep reads only live runs that have one.
index('flow_workflow_deadline_idx').on(t.deadlineAt).where(sql`deadline_at IS NOT NULL`),
// `start` with an idempotency key upserts against this index: without it a repeated start is a second run.
uniqueIndex('flow_workflow_idempotency_idx').on(t.partitionKey, t.idempotencyKey).where(sql`idempotency_key IS NOT NULL`),
]);
export const flowWorkflowStep = pgTable('flow_workflow_step', flowWorkflowStepColumns, (t) => [
foreignKey({ columns: [t.workflowId], foreignColumns: [flowWorkflow.id] }).onDelete('cascade'),
foreignKey({ columns: [t.parentStepId], foreignColumns: [t.id] }).onDelete('cascade'),
unique('flow_workflow_step_workflow_id_key_unique').on(t.workflowId, t.key),
index('flow_workflow_step_workflow_idx').on(t.workflowId),
index('flow_workflow_step_status_idx').on(t.workflowId, t.status),
index('flow_workflow_step_parent_idx').on(t.parentStepId),
]);
export const flowStepLease = pgTable('flow_step_lease', flowStepLeaseColumns, (t) => [
primaryKey({ columns: [t.partitionKey, t.stepType, t.stepId] }),
index('flow_step_lease_active_idx').on(t.partitionKey, t.stepType, t.expiresAt),
]);The store stays plain SQL, bound to the schema you hand createPgWorkflowStore,
createPgStepGate and createPgEventSink — and, for the gate, to the tables names — so
declare the tables under the same schema and names. A test in the package holds every column
set and its DDL equal, column for column: a column added to one side alone fails the build.
Documentation
Full docs — concepts, every feature, production wiring and the API reference — live at octaflow.octabits.io.
| | |
|---|---|
| Quick start | a runnable workflow in one file, no database |
| Concepts | step, workflow, registry, store, dispatcher, engine, partition |
| Defining steps | defineStep and its variants |
| Retry & timeout | attempt budgets, backoff, and how a failure is classified |
| Fan-out & map | one child step per item of a runtime list |
| Branching | when guards and join rules — if/else over a static DAG |
| Signals · Sub-workflows · Saga | suspend, nest, and roll back |
| Deadlines | budgets for a suspended step and for a whole run |
| Heartbeats | liveness for long steps, and interrupting one that was cancelled |
| Postgres & pg-boss | production wiring, workers, DLQ, cron |
| Cancellation & recovery | cancelling a run, sweeping steps a crash left behind, and retrying a failed run |
| Observability · Live progress | lifecycle events, spans, and streaming them to a browser |
| The AI add-on | token/cost capture, quota, usage rollups |
| Extending | custom stores, dispatchers and gates |
| API reference | every export, by entry point |
Examples
Runnable, focused examples live in examples/ — see examples/README.md.
| # | File | Shows |
|---|---|---|
| 01 | 01-in-memory-quickstart.ts | minimal setup + run loop |
| 02 | 02-dag-parallel-fan-in.ts | parallel branches + fan-in (diamond DAG) |
| 03 | 03-retry-timeout.ts | per-step retry + timeout |
| 04 | 04-durable-sleep.ts | durable delay between steps |
| 05 | 05-concurrency-rate-limit.ts | in-memory StepGate |
| 06 | 06-start-idempotency.ts | dedup key collapses duplicate starts |
| 07 | 07-dynamic-map.ts | runtime fan-out / map |
| 08 | 08-wait-for-event.ts | suspend + resumeStep |
| 09 | 09-sub-workflows.ts | child workflow compose + await |
| 10 | 10-saga-compensation.ts | reverse-order rollback on failure |
| 11 | 11-observability.ts | observer events + tracer spans |
| 12 | 12-postgres-pgboss-production.ts | full pg store + gate + event sink + pg-boss + cron |
| 13 | 13-ai-workflow.ts | AI add-on (instrumented model + cost) |
| 14 | 14-live-progress.ts | FlowObserver → SSE fan-out (build your own dashboard) |
| 15 | 15-conditional-branching.ts | when guards + a join: 'any' convergence |
| 16 | 16-deadlines-and-retry.ts | wait deadlines, run deadlines, retryWorkflow, heartbeats |
The in-memory examples (01–11, 14–16) share a small driver, examples/runtime.ts,
that builds an engine over the in-memory store and an in-process queue you drain.
Contributing
Bug reports with a runnable reproduction are the most useful thing you can send; PRs are
welcome. See CONTRIBUTING.md for the setup, the layer rules the lint
enforces, and the correctness requirements a custom WorkflowStore has to meet.
Status
Pre-1.0 — developed in the octabits platform monorepo
under packages/octaflow
(a standalone repository from 2026-07-14 to 2026-09-06; that history is merged in).
Independently versioned and published as the unscoped octaflow. The API is stable but may
still see breaking changes in 0.x minors.
