Most of what Connect does happens after the HTTP request returns. Syncing a scheduling system, indexing a drawing set well over 100 MB, transcribing a site walk, drafting a weekly client update, scoring yesterday's agent answers: all of it runs in a separate worker process that pulls work out of Postgres. There is no Redis, no message broker and no hosted queue service. The queue is a set of ordinary tables, and the thing that keeps two workers from doing the same job is a lease.
This article explains how that worker is built, what a lease is and why we renew it at a third of its length, how crash recovery works without a reaper process, and how a long indexing job resumes from a checkpoint on a different machine. If you run Postgres already and need reliable background jobs, this pattern will probably take you further than you expect.
Why a separate process, and why Postgres
Connect runs as three processes from one codebase: the API server, the worker, and a third process for MCP and the builders (the deployment article covers the layout). The worker is separate for one blunt reason: a long sync blocks its loop. That is the throughput model, by design. If we need more throughput we start more worker processes. We do not add concurrency inside one.
We learned the cost of mixing these the hard way. Schedule Builder turns, which take minutes of model calls and heavy CPM work, originally ran in the process that also served MCP and HTTP. They starved the event loop until status and list calls timed out mid-build. Moving them to a database-backed job that the worker claims fixed it.
Postgres as the queue gives us things a broker would not:
- The job and its effects commit together. Enqueueing a run happens in the same database as the connection it belongs to, so there is no dual-write problem.
- Constraints do the concurrency control. A partial unique index guarantees at most one active run per connection. Two workers and three browser tabs racing to enqueue still produce one run.
- Job state is queryable. The UI reads run status, attempt history and per-table results with plain SQL.
- One less system to operate. We already back up, monitor and secure Postgres.
The tradeoff is throughput. A Postgres queue polled every two seconds is not built for thousands of jobs per second. Our jobs are minutes long and arrive at human pace, so that ceiling is far away.
The tables
The main queue is connection_runs. The columns that matter for leasing:
| Column | Purpose |
|---|---|
status | waiting, queued, running, succeeded, partial, failed, cancelled |
attempts | Incremented on every claim |
lease_owner | The worker id holding the run |
lease_expires_at | When another worker may take it |
queued_at | Also used as "not before" for retries with backoff |
cancel_requested_at | Set by the UI; the worker notices and stops |
results | JSON with per-table counts and the attempt history |
Two partial unique indexes enforce the invariants: one active (queued or running) run per connection, and one pending purge per connection. A partial index on (queued_at, id, lease_expires_at) over the non-terminal statuses keeps the claim query fast as finished runs accumulate.
Next to it sits a workers table where each process upserts its id and a heartbeat every 15 seconds, which shows which workers are alive. The schema also carries a single-row maintenance_state table from early in the project, enforced by a CHECK (id) on a boolean primary key. It is a neat trick for singleton rows, though the current worker code does not read it.
Claiming work with SKIP LOCKED
Here is the core pattern, simplified. This is illustrative SQL, not our production query:
WITH candidate AS (
SELECT r.id
FROM job_runs r
WHERE (r.status = 'queued' AND r.queued_at <= now())
OR (r.status = 'running' AND r.lease_expires_at < now())
ORDER BY r.queued_at, r.id
FOR UPDATE SKIP LOCKED
LIMIT 1
)
UPDATE job_runs r
SET status = 'running',
attempts = attempts + 1,
lease_owner = $1,
lease_expires_at = now() + make_interval(secs => $2)
FROM candidate
WHERE r.id = candidate.id
RETURNING r.*;
Three things make this work.
FOR UPDATE SKIP LOCKEDlets several workers run the query at the same moment. Each locks a different row and skips rows another transaction holds, so nobody blocks and nobody double-claims.- The claim and the state change are one statement. There is no SELECT-then-UPDATE window for a second worker to slip into.
- An expired lease is claimable. That second WHERE branch is the entire crash-recovery story. If a worker dies mid-run, its lease lapses within two minutes and the next poll from any worker picks the run up. There is no reaper process to build, deploy or forget about.
Our real claim query adds a rate limiter. A run is not claimable if another run in the same workspace is running with a live lease and either one is a purge, or both are for the same provider. Different providers sync in parallel freely; two connections to the same provider in one workspace never sync at once. Evaluating that with NOT EXISTS has a race: two workers could each check before either writes "running." So the query uses FOR UPDATE OF r, t SKIP LOCKED, locking the run row and the workspace row. Locking the workspace serializes claim decisions for that tenant across workers.
A third claim path handles deletion. Purging a connection inserts a run in a waiting state if a sync is still running. The claim query promotes a waiting run once nothing else is active for that connection, so a delete drains cleanly behind an in-flight sync instead of fighting it.
Holding the lease: renew at a third
The lease is 120 seconds. The worker renews it every 40 seconds. The renewal is scoped to the owner:
UPDATE job_runs
SET lease_expires_at = now() + interval '120 seconds'
WHERE id = $1 AND lease_owner = $2 AND status = 'running';
Why a third? It gives three chances to renew before the lease expires, so one slow query or a brief database hiccup does not lose the run. Renew too rarely and a single missed renewal hands your job to someone else. Renew too often and you are spending connections on bookkeeping.
The important part is what happens when a renewal updates zero rows. It means this worker no longer owns the run, most likely because it was slow and another worker reclaimed it. The worker must stop immediately. Two workers writing the same connection's records is exactly the corruption the lease exists to prevent. In code, the run timeout, a lost lease and a cancel request all abort the same AbortSignal, so the sync code has one thing to listen to.
This is also why the worker runs one sync at a time, awaited inline. Parallel runs inside one process fail in an ugly chain: parallel runs exhaust the connection pool, which starves lease renewal, which loses the lease, which lets a second worker write the same connection. The server's pool is 20 connections; the worker's is 4. Scale by process count.
Cancellation and control
Cancelling a queued run is immediate: one UPDATE moves it to cancelled. Cancelling a running run only sets cancel_requested_at. The worker polls a control query every two seconds:
- No row matching
id AND lease_owner AND status = 'running'means the lease was lost. - A row with
cancel_requested_atset means the user cancelled. - Otherwise, carry on.
One query detects both conditions. The worker aborts the signal and records the run as cancelled, keeping whatever results it gathered.
Retries with jittered backoff
A run gets three attempts. When a sync fails with a retryable error, or returns a partial result whose failure is retryable (one endpoint briefly unavailable), the worker requeues the same row instead of inserting a new one. That keeps attempt history and per-table results accumulating on one run, which is what an operator wants to see.
The backoff is 5 seconds, 30 seconds, then 2 minutes, each multiplied by a random factor between 0.8 and 1.2. Jitter keeps a provider outage from producing a synchronized retry stampede when it recovers. A provider's own Retry-After always wins. Requeueing sets queued_at into the future, so the claim query's queued_at <= now() naturally enforces the delay.
The requeue also re-checks the world. If the connection has since moved to deleting or disconnected, or a cancel was requested, the requeue updates nothing and the run ends as cancelled. A non-retryable partial result (an endpoint that is permanently broken) is sealed as partial rather than retried forever.
Scheduled syncs use the same machinery. Each poll enqueues due connections with FOR UPDATE SKIP LOCKED on the connections table, and advances next_run_at from now, not from the old value. A connection that was down for six hours on a 15-minute schedule gets one run, not twenty-four backlogged ones.
One poll, many job kinds
Every two seconds, the worker's poll walks through a fixed list of stages. Each one claims a bounded number of items and is wrapped in its own try/catch, so one failing stage never skips the ones after it or fails a sync.
| Stage | Per poll | Notes |
|---|---|---|
| Due-sync scheduling | up to 250 connections | Enqueue only |
| File deletions | 10 | A failed delete leaves the asset in a deleting state for retry |
| File catalogs and ZIP expansion | 1 each | A stuck expansion is reclaimed after a stale window |
| File indexing | 3 | Leased, resumable (below) |
| Recording transcription, summaries, indexing | 1, 2, 5 | Transcription is the expensive one |
| Job walk transcription, vision, reports, expiry | 1, 2, 2, 10 | Reports fan out to indexing and a draft update |
| Agent evals and quality guardrails | 5 | Rides the same cadence |
| Alerts, escalations, webhooks, notifications | 20 to 50 | Delivery work |
| Regression runs | 1 | Capped at 20 golden items per run so it cannot hold the poll |
| Update generation and cadences | 3 | Drafts only, never auto-published |
| Schedule Builder turns | 2 | Reclaims stale running turns first |
| Connection syncs | up to 20 | Claimed and run one at a time |
Several handlers self-throttle internally (hourly or daily, on a durable marker), so calling them every poll costs a timestamp read. The small per-poll limits are the real scheduling policy: they keep any single kind of work from monopolizing the loop.
Most of these job kinds use the same claim shape: UPDATE ... WHERE id = (SELECT ... FOR UPDATE SKIP LOCKED LIMIT 1) RETURNING .... Not all use a full lease. ZIP expansion uses a simpler "processing since" timestamp and reclaims anything stuck for 15 minutes. That is enough for a job that either finishes quickly or is safe to redo.
Resumable indexing from a checkpoint
File indexing used to be fire-and-forget: claim, extract, embed the whole document in one call. A large drawing set never made it into the knowledge base, and a job that ran for minutes had no progress, no cancel and no way to survive a worker restart.
We fixed it by giving the file ingestion table the same lease and cancel columns as connection_runs, plus two more: processed_units and total_units. The indexer now works like this:
- Claim the next ingestion that is queued, or running with an expired lease. A reclaimed row keeps its
processed_units. - Extract the document page by page and build one ordered list of chunks. Chunking is deterministic, so a reclaimed job rebuilds the identical list.
- Record
total_unitsand start atprocessed_units, skipping everything a previous worker already embedded. - For each batch of 200 chunks: check control (lost or cancelled), embed and upsert the batch, write the new
processed_units, renew the lease. - At the end, delete any chunks past the new tail that a previous, longer attempt left behind, then re-check that the file's project scope did not change mid-run before marking it ready.
The checkpoint is written after the batch is upserted, so the worst case on a crash is re-embedding one batch. Upserts are keyed by chunk index, so that redo is harmless. Progress comes free: processed_units / total_units is the progress bar.
There are still ceilings. Indexing caps raw bytes, pages and chunks so a multi-gigabyte blob is never loaded whole into worker memory, and image-only PDF pages yield no text because this path has no OCR.
Testing it against real Postgres
The worker's integration test runs against a real database in CI (a Postgres service container un-skips it). It builds the real modules, stubs only the connector's extraction, and calls worker.poll() directly instead of starting the loop. It checks three behaviours: a queued run is claimed, executed and published, and rows that disappear upstream are tombstoned on the next sync; a failed extraction is retried in place on the same row with its status back to queued and attempts at one; and a connection delete drains through a purge run until the connection row is gone.
Calling poll() directly is the trick that makes this testable. The loop is just while running: poll(); sleep, so testing one poll tests the logic without timers.
Why this matters if you're building something similar
- Postgres is a good job queue for minute-long jobs.
FOR UPDATE SKIP LOCKEDplus a single UPDATE-RETURNING claim is the whole core. - Make expired leases claimable and skip the reaper. Crash recovery becomes a WHERE clause.
- Renew at a third of the lease and stop on a failed renewal. A worker that keeps writing after losing its lease is worse than one that crashes.
- Scope every write by lease owner. Finish, requeue and renew all check
lease_owner, so a stale worker's writes match nothing. - Let unique indexes do concurrency control. A partial unique index on active jobs beats any application lock.
- Checkpoint long jobs with deterministic units. If the work can be rebuilt identically, store a counter, not the work.
- Bound every stage per poll. Small limits are the scheduler.
Where to go next
- How we built Connect: architecture of an enterprise AI platform
- Postgres schema design for enterprise SaaS
- Resumable large-file uploads to Azure Blob
- A RAG pipeline for construction documents
- The connector framework behind Procore OAuth integrations
Need background processing that survives crashes and restarts? Talk to us about your build.
Frequently asked questions
Can Postgres really replace a message queue for background jobs?
For jobs that take seconds to minutes and arrive at human pace, yes. FOR UPDATE SKIP LOCKED gives safe concurrent claims, and you get transactions, constraints and queryable job state for free. It is not built for thousands of jobs per second.
What happens if a worker crashes mid-job?
Its lease stops being renewed and expires within two minutes. The claim query treats a running job with an expired lease as claimable, so the next worker to poll picks it up and, for indexing, resumes from the saved checkpoint.
Why renew the lease at one third of its duration?
It gives three chances to renew before expiry, so one slow query or a brief database hiccup does not hand the job to another worker, without spending many connections on bookkeeping.
Why run only one sync at a time per worker process?
Parallel runs in one process can exhaust its small connection pool, which starves lease renewal, loses the lease and lets a second worker write the same data. Throughput scales by adding worker processes.
How are failed jobs retried?
The same row is requeued with its queued_at set into the future: 5 seconds, 30 seconds, then 2 minutes, each jittered by plus or minus 20 percent, for up to three attempts. A provider's Retry-After takes precedence.
Next step
Want this workflow automated?
Tell us the handoff, approval, or reminder your team repeats. We'll scope the automation, the approvals it needs, and who owns it.
Prefer email? charley@buildflows.ai
Get the next guide in your inbox
Field Notes: practical guides and new walkthroughs, about once a month.
Field Notes
Practical guides and new walkthroughs on construction data and automation, roughly monthly.
Keep learning
AI agents & MCP · October 9, 2026
How We Built Connect: Architecture of an Enterprise AI Platform
The pillar of our Connect architecture series: a layer-by-layer map of an enterprise agentic AI platform, the three-process shape it runs as, the design principles that kept recurring, and links to every deep-dive article.
Custom applications · October 9, 2026
Postgres Schema Design for an Enterprise AI Platform: 120+ Tables, No ORM, Verified Migrations
A tour of the Postgres schema behind Connect, why it uses direct SQL instead of an ORM, how it splits work with DocumentDB, and the migration discipline that stops the server booting on a mismatched database.
Custom applications · October 9, 2026
Resumable Large File Uploads to Azure Blob for Construction Files
An engineering walkthrough of Connect's Files layer: signed direct-to-blob uploads, server-side sessions that let multi-gigabyte uploads pause and resume, idempotent ZIP expansion, provenance and the handoff to indexing.