# One job queue instead of five ways to pick up work

> Replacing ad-hoc background work in a support AI platform, and an n8n enrichment flow, with Postgres job queues that claim work atomically, retry in one place and recover from crashed workers.

Source: https://simasrazinskas.com/work/durable-job-queue  
Published: 2026-06-30  
Reviewed: 2026-10-07

## At a glance

- **Problem**: A support AI platform had grown five different ways to pick up background work, and drafting had no durable record at all, so a deploy could drop a draft and a slow action could run twice.
- **Constraint**: Some of the work posts public replies and cancels subscriptions, and the help-desk API does not deduplicate writes, so a retry had to be safe without help from the receiving side.
- **What I built**: One Postgres job table with an atomic claim, a dedupe index, heartbeats, per-attempt deadlines and a single retry policy, plus idempotency keys and a per-subscription ledger on the actions themselves, and a leased Go queue built to replace an n8n enrichment flow.
- **Result**: Drafts survive restarts, every job has a visible state and error, and the reclaim bug that produced duplicate public comments was removed at the source.
- **Stack**: TypeScript · Go · Postgres · LISTEN/NOTIFY · BigQuery · Zendesk API · n8n

## Overview

A consumer subscription business supervises an AI support agent that drafts replies to Zendesk tickets and, once approved, executes actions such as posting the reply and canceling subscriptions.
Around it runs background work: drafting with a model, executing approved actions, evaluation runs, scheduled scans that start new drafts, and enrichment of tickets with commerce data from BigQuery.

Each piece of that work had grown its own answer to the same question: how does a worker take a job without another worker taking it too.
The answers differed, and some were wrong in ways that only showed up under load or during a deploy.

## Five mechanisms, one table

Before the change there were five patterns:

- Evaluation runs used `FOR UPDATE SKIP LOCKED` with a heartbeat, and were correct.
- Automatic drafting used an advisory lock plus a cooldown ledger.
- A poller that superseded stale drafts used an advisory lock plus a high-water-mark row.
- Drafting itself had a TTL lease that its own documentation said gated nothing.
- Execution used a two-phase reclaim whose lease was shorter than the work it protected.

Drafting also had no durable representation.
It ran as a fire-and-forget promise inside the HTTP process, progress lived in an in-memory registry, and the browser guessed the outcome with a 150-second timeout.
A deploy dropped the draft and left a spinner.

All of it moved into one `job` table with a kind, a key, a state, an attempt budget, a per-attempt deadline, a heartbeat and a progress column.

- **Claim** is one `UPDATE` over `FOR UPDATE SKIP LOCKED`; the state change, attempt count and first heartbeat land together.
- **Dedupe** is a partial unique index on `(kind, key)` over pending and running rows. A colliding enqueue is a no-op that returns nothing, so a double click cannot create two jobs.
- **Wakeup** is `NOTIFY` on enqueue, with a 5-second poll as the backstop for lost notifications.
- **Liveness** is `heartbeat_at`, refreshed every 15 seconds on a timer separate from the work. A reaper returns jobs with no heartbeat for 60 seconds to pending, or fails them when the budget is spent.
- **Deadline** is armed per attempt when the job is claimed, so time spent waiting in the queue does not eat the work's budget. A drafting job gets 300 seconds and 2 attempts.
- **Fencing** applies to every write after the claim: it matches on id, claimant and `running`, and reports whether it won. A worker that lost its claim gets a boolean, never an exception.

Retry lives in exactly one place, the settle step:

```sql
UPDATE job
SET state  = CASE WHEN $outcome = 'failed' AND attempts < max_attempts
                  THEN 'pending' ELSE $outcome END,
    run_at = CASE WHEN $outcome = 'failed' AND attempts < max_attempts
                  THEN now() + make_interval(secs =>
                       least(300, 5 * 2 ^ attempts) * (0.5 + random() * 0.5))
                  ELSE run_at END,
    claimed_by = NULL, heartbeat_at = NULL, error = $error
WHERE id = $id AND claimed_by = $worker AND state = 'running'
RETURNING state;
```

*Exponential backoff from the row's own attempt count, capped at 5 minutes, with half-to-full jitter*

Handlers are not allowed to wrap their own retry loop around a model call.
Progress lives on the row, and terminal jobs are pruned after 30 days.

## Leases have to outlive the work

The bug that forced the issue was in execution.
An approved action held a 30-second lease while it worked.
The work itself was a Zendesk write with a 30-second timeout followed by sequential cancellation calls, each with its own 30-second timeout.
When the real work outlasted the lease, a retry or a concurrent caller reclaimed the action and ran all of it again: a duplicate public comment and duplicate cancellations.

The fix was a heartbeat that extends the lease while the claimant is still working.
A lease now expires only for a claimant that stopped, which is the same liveness rule the job table uses.
## Idempotency where the API offers none

Zendesk comment creation is not idempotent; sending the same request twice posts two comments.
So the deduplication has to happen before the call.

Each proposal action gets a key of the form `{proposalId}:{actionIndex}`, and the action table has a unique index on account and key.
An action moves from pending to applying to a terminal state.
Only `applied` deduplicates a re-run.
If the Zendesk write landed and a cancellation failed, the action ends in a separate partial-failure state, and a retry re-runs only the cancellation leg.

Cancellations get their own ledger, with one row per action and subscription.
Before calling the subscription service, the executor reads which subscriptions already succeeded and passes only the rest.
After the call returns or throws, it writes one row per outcome before re-raising, so a crash after a partial success still records it.

Reserving an action takes a per-ticket Postgres advisory lock with the same key the ticket-mirroring sync worker uses, so the reservation serializes against mirror writes for that ticket.
## The enrichment queue and its cutover

An n8n flow used to look up each ticket's customer in BigQuery.
I built a Go service to replace it, around its own leased queue.

The queue has explicit states: `queued`, `leased`, `done`, `failed`, `dead_letter`.
A row becomes visible 10 minutes after a new ticket, to ride out BigQuery replication lag, or 2 minutes after activity on an existing one.
A partial unique index allows one open queued row per ticket, so bursts merge.
A claim sets `lease_owner` and a 60-second `lease_expires_at`, a sweeper recovers expired leases every 30 seconds, and a row moves to `dead_letter` after 5 attempts.

The first sweeper flipped every expired row back to queued in one statement.
When two expired rows belonged to the same ticket, the second hit the unique index, and the sweep failed every 30 seconds.
The fix promotes the newest expired row per ticket and dead-letters the rest with a reason, all in one transaction.

The new worker ran in shadow first.
It wrote only to its own tables, a test failed the build if it ever imported a path that writes to the legacy tables, and a parity job compared order coverage between the n8n-written data and the new data.
Its first production traffic surfaced three bugs: an `int4` overflow once a run billed more than 2.1 GB of BigQuery, a BigQuery job ID written into a UUID column that made committed runs retry and duplicate, and a lookup log that was never wired up.
After the fixes, the first-generation Go worker, its budget ledger and its projection tables were deleted.

## Limitations

- An ambiguous network failure, where Zendesk accepted a write but the response was lost, can still double-post if an operator explicitly re-approves. That was chosen over the alternative, which silently marked actions executed with no writes.
- The two queues are still two: the Go enrichment queue and the TypeScript job table share the pattern, not the code.
- The enrichment queue has a heartbeat API, but its live worker relies on a fixed lease TTL, so a run that outlasts the lease would be reclaimed. It repeats the execution bug in a place where a rerun is cheap.
- Another family of n8n workflows was also rewritten as a Go service with a queue and cron jobs. It ran only in dry-run, the event subscription was never moved from n8n, and the project was scrapped. Its queue deleted each job after dispatch whether or not the call succeeded, copied from the flow it replaced. The lessons were to plan the cutover path before the rewrite, and to fix the old system's failure handling instead of porting it.

## Related case studies

- [Help-desk ingestion that proves nothing went missing](https://simasrazinskas.com/work/idempotent-webhook-ingestion): Every support ticket reaches the database; missing ones are found and repaired automatically.
- [One governed platform for a company's internal AI automation](https://simasrazinskas.com/work/internal-ai-platform): New internal AI tools launch on one shared, access-controlled platform, not from scratch.

## Start with one process

A one-hour call costs €80. Afterwards you get a written plan, whether or not we work together.

[Book a call](https://simasrazinskas.com/book-me) · [simas@simasrazinskas.com](mailto:simas@simasrazinskas.com)
