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 scheduled workflow is a row with a timestamp
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.
How Orch8 stores and runs a workflow
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(),
...
);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:
run_tick_loopwakes on a tokio interval. If the previous tick is still busy the missed tick is skipped, not queued.process_tickchecks how many semaphore permits are free and asks storage for at most that many rows. There is no point claiming rows it cannot run.claim_due_instancesruns the SKIP LOCKED query below inside one transaction, flips the claimed rows torunning, and commits. From this moment other nodes cannot see those rows as due.- 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.
- For each instance, a permit is taken, a task is spawned, and a heartbeat task starts touching the row's
updated_atwhile the step is in flight. process_instanceloads the sequence from an in-process cache, thenexecute_step_loopruns every block that is ready. It only stops when a block asks for a delay, parks on input, fails, or the sequence ends.- On the way out, the instance is written back as
scheduledwith a newnext_fire_at,waiting,completed, orfailed. The heartbeat task is cancelled and the permit is released.
The claim query
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);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.
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';Why SKIP LOCKED and ROW_NUMBER cannot share a query level
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.
Batch prefetch: two queries per tick, not two per instance
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.
Lease heartbeats: telling a dead node from a slow step
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);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');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.
Readiness has to include the tick loop
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,
}
}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.
Where a Postgres job queue stops being the right answer
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.
Postgres Queue vs Dedicated Broker
| Concern | Postgres + SKIP LOCKED | Dedicated broker |
|---|---|---|
| Transactional enqueue | Atomic with your data | Needs an outbox pattern |
| Operational surface | A database you already run | One more system to operate |
| Delayed delivery | A timestamp column | Native, often with limits |
| Querying pending work | Ordinary SQL | Usually opaque |
| Throughput ceiling | Tens of thousands/sec, tuned | Much higher |
| Fan-out | Emulated with rows | Native |
| Maintenance risk | Vacuum and index bloat | Broker cluster operations |
| Scheduling precision | Tick-bounded, about 100ms | Sub-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.
Related Matching Engine Guides
Why an instance is one row that can be resumed without replaying history.
Where the per-tenant claim cap fits among the other isolation boundaries.
What happens after a claim, when the step calls an external system.
Claim-and-lease semantics extended across machines that do not share a database.
Related reading
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 →