Priority 2 of the horizontal-scaling work (#1574). #1755 depends on the primitive described here for its delete and provision jobs.
Why
Background work exists in four places, each with its own mechanism: the segment store purge (claims with FOR UPDATE SKIP LOCKED, one bounded batch per call, every replica runs the loop), the collection registry purge in #1734 (same pattern), the semantic ingestion loop (no claim at all, #1745), and the session deletion queue (in process memory, #1749). #1722 proposes a leased lock service with no consumer. The queue library trial (design/queue_library_trial.md on branch queue-library-trial) measured that a leased claim inside the store solves the registry's vacuum-horizon problem without a lock service. One primitive should serve all of them.
At a glance
| Consumer |
Today |
With the primitive |
Issue |
| Session provision and delete |
in-process asyncio.Queue, boot fan-out |
one job per (session, kind), claimed from any replica |
#1749, #1577, #1755 |
| Semantic ingestion |
every replica processes every dirty set |
one claim per set, lease covers the LLM round, replicas take disjoint sets |
#1745, #1699 |
| Segment store purge, registry purge |
already claim with SKIP LOCKED |
adopt the shared helper when it exists, semantics unchanged |
#1734 |
| Not-yet-ingested episodes |
three independent writes per add, no record of partial failure |
not a per-item queue: a watermark or outbox per (session, subsystem) advanced by a replay job on the same table |
#1738, #1634 |
| Leased lock service |
no consumer |
fold into the claim primitive, or close |
#1722 |
The standard shape
A job table: (id, kind, key, state, attempts, run_after, claimed_by, claimed_until, last_error, payload).
- Claim: one statement, one round trip, no lock held across the work:
UPDATE jobs SET claimed_by = :token, claimed_until = now() + :lease WHERE id = (SELECT id FROM jobs WHERE state = 'pending' AND run_after <= now() ORDER BY run_after, id FOR UPDATE SKIP LOCKED LIMIT 1) RETURNING *.
- Completion: delete the row, or
state = 'done' where history matters. Failure: attempts + 1, run_after = now() + backoff(attempts), last_error; after N attempts state = 'failed', visible to operators and the API.
- Crash recovery: a claim expires at
claimed_until and the next claimer takes it, so handlers are idempotent and at-least-once. Where the handler holds one transaction for the whole step (the segment store purge today), the transaction is the lease and claimed_until is unused. Where the step calls a remote system (vector purge, LLM), the lease is a timestamp and every write of a result carries WHERE claimed_by = :token AND claimed_until > now(), which refuses a late result from an expired claim. The token is a per-claim uuid.
- Not FIFO:
run_after and kind order the work. A partial unique index on (kind, key) WHERE state IN ('pending', 'running') gives one job per key, so a repeated trigger resets run_after instead of inserting a duplicate (the replay job per session and subsystem in design/server_redesign.md).
- Database clock throughout; no server clock enters a comparison.
This is the shape of the PostgreSQL-backed queues in general use (pg-boss, graphile-worker, Oban, Solid Queue, River): SKIP LOCKED claims, a visibility timeout or a transaction-held lock, idempotent at-least-once handlers, backoff, a dead-letter state. A message broker buys nothing here: every job is tied to a row in the same database and would need an outbox anyway.
Not part of this
A general lease lock service (#1722) unless a long-running indivisible job appears; leader election; coordination across databases.
Tracked
🤖 Written by Claude Code (Claude Fable 5.1) on behalf of @edwinyyyu.
Priority 2 of the horizontal-scaling work (#1574). #1755 depends on the primitive described here for its delete and provision jobs.
Why
Background work exists in four places, each with its own mechanism: the segment store purge (claims with
FOR UPDATE SKIP LOCKED, one bounded batch per call, every replica runs the loop), the collection registry purge in #1734 (same pattern), the semantic ingestion loop (no claim at all, #1745), and the session deletion queue (in process memory, #1749). #1722 proposes a leased lock service with no consumer. The queue library trial (design/queue_library_trial.mdon branchqueue-library-trial) measured that a leased claim inside the store solves the registry's vacuum-horizon problem without a lock service. One primitive should serve all of them.At a glance
asyncio.Queue, boot fan-outSKIP LOCKEDThe standard shape
A job table:
(id, kind, key, state, attempts, run_after, claimed_by, claimed_until, last_error, payload).UPDATE jobs SET claimed_by = :token, claimed_until = now() + :lease WHERE id = (SELECT id FROM jobs WHERE state = 'pending' AND run_after <= now() ORDER BY run_after, id FOR UPDATE SKIP LOCKED LIMIT 1) RETURNING *.state = 'done'where history matters. Failure:attempts + 1,run_after = now() + backoff(attempts),last_error; after N attemptsstate = 'failed', visible to operators and the API.claimed_untiland the next claimer takes it, so handlers are idempotent and at-least-once. Where the handler holds one transaction for the whole step (the segment store purge today), the transaction is the lease andclaimed_untilis unused. Where the step calls a remote system (vector purge, LLM), the lease is a timestamp and every write of a result carriesWHERE claimed_by = :token AND claimed_until > now(), which refuses a late result from an expired claim. The token is a per-claim uuid.run_afterandkindorder the work. A partial unique index on(kind, key) WHERE state IN ('pending', 'running')gives one job per key, so a repeated trigger resetsrun_afterinstead of inserting a duplicate (the replay job per session and subsystem indesign/server_redesign.md).This is the shape of the PostgreSQL-backed queues in general use (pg-boss, graphile-worker, Oban, Solid Queue, River):
SKIP LOCKEDclaims, a visibility timeout or a transaction-held lock, idempotent at-least-once handlers, backoff, a dead-letter state. A message broker buys nothing here: every job is tied to a row in the same database and would need an outbox anyway.Not part of this
A general lease lock service (#1722) unless a long-running indivisible job appears; leader election; coordination across databases.
Tracked
🤖 Written by Claude Code (Claude Fable 5.1) on behalf of @edwinyyyu.