Postgres as the Timer Wheel: Scheduling Workflows Without a Message Broker

Orch8 keeps every scheduled workflow as one row with a timestamp and lets engine nodes claim due rows with FOR UPDATE SKIP LOCKED. This is the claim query as it exists in the code, the Postgres rule that broke the first version of it, how a dead node is told apart from a slow step, and where the design stops being a good idea.

By Oleksii Vasylenko, Technical Lead · Published · Updated · 14 min read

Where this comes from. Orch8 is a durable workflow engine I built solo in Rust: every crate, every SDK, every integration. Ten-crate workspace, PostgreSQL and SQLite backends, server and on-device execution. The details below are the actual implementation, not a reference architecture.

A workflow engine has to wake work up at the right moment, and it has to do that across several engine nodes without two of them running the same step. When I started Orch8, the two designs I kept seeing were a message broker with delayed delivery, or an in-memory timer wheel that gets rebuilt from the database after a restart. I did not want to operate a broker, and I did not want a second copy of scheduling state living in process memory. So the Postgres job queue in this article is what I built instead, using FOR UPDATE SKIP LOCKED as the only coordination primitive.

The model is small. A scheduled workflow is a row in task_instances with state = 'scheduled' and a next_fire_at timestamp. "What is due?" is a range scan on that column. A million waiting workflows use no engine memory, because they are a million rows. Restart the process and nothing is lost, because nothing was in memory.

The hard part is not finding due rows. It is letting several nodes claim from the same table at the same instant without any two of them taking the same row, without a lock table of my own, and without one large tenant filling every batch. The rest of this article is about those three things and the bugs I hit on the way.

Some vocabulary, because the query below uses it. In Orch8 a workflow definition is a sequence: a versioned JSON document that lists blocks. A block is usually a step, and a step names a handler such as http_request, llm_call, or a handler you registered yourself. A run of a sequence is an instance. The engine is one Rust binary that talks to Postgres (or SQLite for tests and on-device use) through a StorageBackend trait, so the claim query has a Postgres version and a SQLite version and both have to behave the same way.

CREATE TABLE IF NOT EXISTS task_instances (
    id              UUID PRIMARY KEY,
    sequence_id     UUID NOT NULL REFERENCES sequences(id),
    tenant_id       TEXT NOT NULL,
    state           TEXT NOT NULL DEFAULT 'scheduled',
    next_fire_at    TIMESTAMPTZ,
    priority        SMALLINT NOT NULL DEFAULT 1,
    context         JSONB NOT NULL DEFAULT '{}',
    updated_at      TIMESTAMPTZ NOT NULL DEFAULT NOW(),
    ...
);
migrations/002_create_task_instances.sql, the columns that matter to the scheduler. Everything else on the row is workflow data.

An instance moves through a handful of states: scheduled (waiting for its fire time), running (a node has claimed it), waiting (parked on a human approval or an external event), and the terminal states completed, failed, and cancelled. The context column holds the snapshot of the run so far. Because resume is "load the row and continue at the next block", the scheduler never has to replay anything. That decision is covered in the snapshot-versus-replay article linked at the end.

The engine runs a tick loop in orch8-engine/src/scheduler.rs. With the defaults from orch8-types/src/config.rs, a tick fires every 100 milliseconds, claims up to 256 due instances, and executes them under a semaphore that allows 128 in flight at once. One tick, end to end, looks like this:

  1. run_tick_loop wakes on a tokio interval. If the previous tick is still busy the missed tick is skipped, not queued.
  2. process_tick checks how many semaphore permits are free and asks storage for at most that many rows. There is no point claiming rows it cannot run.
  3. claim_due_instances runs the SKIP LOCKED query below inside one transaction, flips the claimed rows to running, and commits. From this moment other nodes cannot see those rows as due.
  4. Two batch queries fetch every pending signal and every already-completed block id for the whole batch. This is the prefetch step, and it exists because the first version did it per instance.
  5. For each instance, a permit is taken, a task is spawned, and a heartbeat task starts touching the row's updated_at while the step is in flight.
  6. process_instance loads the sequence from an in-process cache, then execute_step_loop runs every block that is ready. It only stops when a block asks for a delay, parks on input, fails, or the sequence ends.
  7. On the way out, the instance is written back as scheduled with a new next_fire_at, waiting, completed, or failed. The heartbeat task is cancelled and the permit is released.
The scheduler tick loopEvery 100 milliseconds the engine claims due instances with SKIP LOCKED, batch-prefetches signals and completed blocks in two queries, then runs every ready step within one claim cycle.yesnotick — every 100msCLAIMFOR UPDATE SKIP LOCKEDpartial index on next_fire_atBATCH PREFETCHsignals + completed blocks2 queries for the whole batchPROCESSbounded by a semaphorerun every ready stepin one claim cyclenext blockhas a delay?set next_fire_atstate = scheduledcontinue in this cycle
One claim cycle: lock a batch, prefetch what the batch needs in two queries, run every ready step, write the new state back.

FOR UPDATE SKIP LOCKED is the whole trick. It locks the rows it returns and silently skips rows that another transaction has already locked. Two nodes running the same statement at the same instant get disjoint result sets, with no coordination between them. This is the Postgres version from orch8-storage/src/postgres/instances.rs, with the per-tenant cap enabled:

WITH locked AS (
    SELECT id, tenant_id, priority, next_fire_at
    FROM task_instances
    WHERE (next_fire_at IS NULL OR next_fire_at <= $1)
      AND state = 'scheduled'
    ORDER BY priority DESC, next_fire_at ASC NULLS FIRST
    LIMIT $4
    FOR UPDATE SKIP LOCKED
), ranked AS (
    SELECT id, tenant_id, priority, next_fire_at,
           ROW_NUMBER() OVER (
               PARTITION BY tenant_id
               ORDER BY priority DESC, next_fire_at ASC NULLS FIRST
           ) AS rn
    FROM locked
), winners AS (
    SELECT id
    FROM ranked
    WHERE rn <= $3
    ORDER BY priority DESC, next_fire_at ASC NULLS FIRST
    LIMIT $2
)
SELECT task_instances.*
FROM task_instances
JOIN winners USING (id)
ORDER BY task_instances.priority DESC,
         task_instances.next_fire_at ASC NULLS FIRST;

-- same transaction, after concurrency-key filtering:
UPDATE task_instances SET state = 'running', updated_at = $1
WHERE id = ANY($2);
orch8-storage/src/postgres/instances.rs, claim_due. $1 is now, $2 is the batch limit, $3 is max rows per tenant, $4 is the over-select limit.

When max_instances_per_tenant is zero, which is the default, the code takes a shorter path: a plain SELECT ... ORDER BY priority DESC, next_fire_at ASC NULLS FIRST LIMIT $2 FOR UPDATE SKIP LOCKED. The capped version is what a multi-tenant deployment turns on, and it is the version with all the interesting mistakes in it.

Two engine nodes claiming from one table with SKIP LOCKEDThe first node locks rows one to fifty. The second node's identical query silently skips those locked rows and receives rows fifty-one to one hundred. The result sets are disjoint with no coordination.Engine node 2PostgreSQLEngine node 1rows 1-50 are locked,so they are skipped silentlydisjoint sets, no coordination,no lock table of our ownSELECT ... FOR UPDATE SKIPLOCKED1rows 1-50 (now locked)2SELECT ... FOR UPDATE SKIPLOCKED3rows 51-1004
Two nodes issue the same query at the same time. The second node skips the rows the first one locked and gets the next batch instead.

The update to running happens inside the same transaction as the lock. Between the select and the commit, the rows are locked, so a third node cannot claim them either. After the commit they are no longer scheduled, so nobody will even try. There is no lease table, no leader, and no advisory lock for the claim itself. (There is one advisory lock in this path, taken per concurrency_key so that two nodes cannot both admit instances that share a key and overshoot its max_concurrency. That is a different problem and I will leave it out here.)

Two indexes support the query. The first is partial, which is what keeps it small: only scheduled rows enter it, so completed and failed instances, which end up being most of the table, are never indexed. The second, added later in migration 028 when the tenant cap arrived, covers the ordering the window function needs.

CREATE INDEX IF NOT EXISTS idx_instances_fire
    ON task_instances (next_fire_at)
    WHERE state = 'scheduled';

CREATE INDEX IF NOT EXISTS idx_instances_claim_priority
    ON task_instances (state, priority DESC, next_fire_at ASC)
    WHERE state = 'scheduled';
migrations/009_create_indexes.sql and migrations/028_throughput_indexes.sql.

The first version of the capped query did the ranking and the locking in one subselect. It was shorter and it was wrong.

PostgreSQL does not allow a locking clause in a query that also contains a window function. The documentation says FOR UPDATE requires that returned rows be "clearly identifiable with individual table rows", and a window function breaks that correspondence. Depending on how you shape it you get an error, or a plan whose locking behaviour is not what you assumed. I got the second one first, which is worse, because the tests passed on a single node.

Splitting the query created a second problem. If the locked CTE takes only limit rows and one tenant's backlog fills all of them, the per-tenant cap trims most of them away and the caller gets a nearly empty batch. A fairness rule was starving the scheduler it was meant to protect. The fix is the over-select: lock limit * max_per_tenant rows so that after trimming there are still enough rows from other tenants to fill the batch. The multiplication is saturating (saturating_mul in the Rust) so a silly configuration degrades instead of overflowing. The SQLite twin uses the same pre-limit, for a different reason: SQLite has no SKIP LOCKED, so it serialises claims with BEGIN IMMEDIATE, and ranking the entire due backlog inside that transaction would block every other writer.

The next_fire_at IS NULL branch is a smaller mistake with the same shape. A scheduled row with no fire time is due now. The SQLite backend always treated it that way. The Postgres query originally did not, so rows created without a fire time sat in the table forever. Nothing errored. They were just never due. A pair of tests in orch8-storage/tests/bugs_group_b.rs now pins the behaviour for both backends, including under the tenant cap, because the capped and uncapped paths are two different SQL strings and can drift independently. Two implementations behind one trait need conformance tests that assert the same answer, not just the same function signature.

Claiming rows is the easy half. The obvious next move is to loop over the claimed instances and, for each one, fetch its pending signals and its list of completed block ids. With a batch of 256 that is 512 round trips per tick. At a 100 millisecond tick the database spends its time answering the same two questions over and over.

So process_tick collects the instance ids and calls get_pending_signals_batch and get_completed_block_ids_batch once each, joined with tokio::try_join!, then hands out the results from a map. There is a third batched load, preload_externalized_markers, which hydrates any oversized context values that were moved out to a side table. That one is allowed to fail softly, because the per-step resolver can still fetch them one at a time if it has to.

The related win is executing every ready block within one claim. Early on, process_instance ran one step and wrote the row back, so a five-step workflow with no delays took five ticks and about half a second. Now execute_step_loop keeps going until a block asks to wait. Nothing about persistence changed. The engine just stopped handing the row back to the scheduler between two steps that were both ready.

A claimed instance is running. If the node holding it dies, the row stays running forever and nobody finishes the work. The usual fix is a reaper that finds rows that have been running for longer than some threshold and puts them back to scheduled. Orch8 has one. It runs at startup in recover_stale_instances, and again from a background task every half threshold (with a ten second floor) so an instance orphaned by a panic does not have to wait for a restart.

UPDATE task_instances
SET state = 'scheduled', next_fire_at = NOW(), updated_at = NOW()
WHERE state = 'running'
  AND updated_at < NOW() - make_interval(secs => $1::double precision);
orch8-storage/src/postgres/misc.rs, recover_stale_instances. Note the state filter: running only, never waiting.

That query has a well-known failure mode. A step that takes longer than the threshold, say a large export or a slow provider, looks the same as a dead node. The reaper re-dispatches it while the original node is still working, and now two nodes are executing the same step against the same external system. With the default threshold of 300 seconds this is rare, but rare is not the same as impossible, and the steps that take five minutes are usually the ones you least want to run twice.

The question the reaper needs answered is not "how long has this been running" but "is the node that owns it still alive". So while a step is in flight, the owning node runs a small heartbeat task that touches the row:

let heartbeat_interval = instance_heartbeat_interval(stale_instance_threshold_secs);
tokio::spawn(async move {
    let mut ticker = tokio::time::interval(heartbeat_interval);
    ticker.tick().await; // first tick is immediate; claiming already set updated_at
    loop {
        tokio::select! {
            () = heartbeat_stop.cancelled() => break,
            _ = ticker.tick() => {
                if let Err(e) = heartbeat_storage.heartbeat_instance(instance_id).await {
                    warn!(instance_id = %instance_id, error = %e, "instance heartbeat failed");
                }
            }
        }
    }
});

// orch8-storage/src/postgres/misc.rs
UPDATE task_instances SET updated_at = NOW()
WHERE id = $1 AND state IN ('running', 'waiting');
orch8-engine/src/scheduler.rs, inside the spawned task, and the storage call it makes. The interval is a third of the reaper threshold, so 100 seconds at the defaults.

A slow but healthy step keeps refreshing updated_at. A dead node stops. The reaper then only recovers rows whose owner went quiet, and long-running work is never pulled out from under a live process. The configuration validator refuses a stale threshold shorter than the tick interval, since that would make every claimed row look dead before it had a chance to heartbeat.

Timeout-based recovery versus lease heartbeatsA threshold alone cannot distinguish a slow healthy step from a dead node, so it re-dispatches live work. A heartbeat touched by the owning node makes silence, not slowness, the signal of death.Lease heartbeat — detect death, not slownessnoyesowning node touchesupdated_at while in flightheartbeatgone stale?healthy but slow — leave itowner really died — recoverTimeout alone — the wrong signalyesstep running 5 minover threshold?reaper re-dispatchestwo nodes run the same step
A threshold measures how long the work has taken. A heartbeat measures whether anyone is still doing it. Only the second one tells you the owner is gone.

A failure this design invites: the tick loop dies while the HTTP server keeps serving. The API accepts new workflows, returns 200s, and nothing ever executes. Every dashboard is green. Work accumulates until someone notices a scheduled report never arrived.

In Orch8 the server creates an engine_ready atomic flag, hands it to both the engine task and the API state, and the engine task clears it when run() returns for any reason. The readiness handler checks the flag first and the database second:

pub(crate) async fn readiness(State(state): State<AppState>) -> impl IntoResponse {
    if !state.engine_ready.load(Ordering::Relaxed) {
        return StatusCode::SERVICE_UNAVAILABLE;
    }
    match state.storage.ping().await {
        Ok(()) => StatusCode::OK,
        Err(_) => StatusCode::SERVICE_UNAVAILABLE,
    }
}
orch8-api/src/health.rs. Liveness is unconditional; readiness is not.

When a process has a background loop doing the real work and a request surface answering health checks, the two must not be able to disagree about whether the process is healthy. The probe has to assert the loop, so the orchestrator pulls the pod instead of leaving a zombie API accepting work it will never run.

Postgres as a queue works well up to a point, and I would rather write down where that point is than find it in production.

  • Claim throughput has a ceiling. Every claim is a write transaction. Tens of thousands per second on one primary is reachable with tuning. Hundreds of thousands is a different architecture. Measure your own workload before trusting anyone's headline number, mine included.
  • Dead tuples accumulate. A high-churn queue table generates a lot of update volume. Autovacuum needs per-table tuning or the partial index bloats and the range scan degrades. This is the failure people actually hit, and it tends to arrive a few weeks after launch.
  • Long transactions are poison. Holding the claim transaction open during a slow HTTP call pins the snapshot horizon and blocks vacuum across the whole database. Claim, commit, then execute. The step itself runs outside any transaction.
  • Fan-out is not free. A broker delivering one message to many consumers is doing something Postgres has to emulate with more rows and more writes.
  • Sub-millisecond scheduling is out of scope. A 100 millisecond tick means up to 100 milliseconds of scheduling latency, plus the claim itself. Fine for business workflows. Wrong for anything that would call itself low latency.

What you get in exchange is the reason the trade usually goes this way: the queue is in the same transaction as your data. Creating an instance and writing the row that caused it is atomic. No outbox, no dual-write gap, no reconciliation between the broker's view of the world and the database's. That one property removes a whole category of bug that broker-based designs spend real effort managing.

For a workflow engine I think the trade is almost always right. Business processes are measured in seconds and days, and surviving a crash matters far more than scheduling precision. If a workload ever outgrows it, the migration is a storage-layer change behind the StorageBackend trait rather than a rewrite of the engine, which is itself a reason to start here.

ConcernPostgres + SKIP LOCKEDDedicated broker
Transactional enqueueAtomic with your dataNeeds an outbox pattern
Operational surfaceA database you already runOne more system to operate
Delayed deliveryA timestamp columnNative, often with limits
Querying pending workOrdinary SQLUsually opaque
Throughput ceilingTens of thousands/sec, tunedMuch higher
Fan-outEmulated with rowsNative
Maintenance riskVacuum and index bloatBroker cluster operations
Scheduling precisionTick-bounded, about 100msSub-millisecond possible

The first row decides most cases. Removing the gap between "job enqueued" and "data committed" is worth more than broker throughput to a workflow engine.

The pillar guide places the scheduler in the full engine: execution model, storage backends, crate layout, and operational design.

Read the durable workflow engine architecture guide →