Skip to content

feat(pipeline): a durable job queue - #19

Merged
ralyodio merged 1 commit into
mainfrom
feat/job-queue
Aug 13, 2026
Merged

ralyodio merged 1 commit into
mainfrom
feat/job-queue

Conversation

@ralyodio

Copy link
Copy Markdown
Contributor

Phase 1 of docs/url-first-pipeline.md, and the thing every later phase is blocked on.

JOB_KINDS has described this table since the first commit and nothing ever created it. Work either ran inline in the request — POST /prospects still runs the whole enrich → resolve → research → score → recommend → draft chain before responding — or the worker loop swept its own domain tables directly. Neither survives a batch: one URL fans out to a fetch, an extraction, an identity search per candidate person and a draft per recommendation.

What's here

  • migrations/0007_jobs.sql — the table, three indexes.
  • packages/pipeline/src/queue.ts — enqueue, claim, complete, fail, reclaim, drain, depth. It knows nothing about what a job does: drainQueue takes a handler, so the dispatcher lives in the server and tests pass a stub.
  • apps/server/src/index.ts — an explicit runJob dispatcher and a drain at the end of each tick.
  • packages/domain/src/ids.ts — a job id prefix; the domain had no kind for something the schema now stores.

Decisions worth reviewing

A table, not Redis. REDIS_URL sits empty in the vault and unreferenced in the code. Turso is already here and already backed up, and a queue that loses its contents on restart is not a queue.

The claim is atomic. One UPDATE … WHERE id = (SELECT …) RETURNING …, so two workers cannot both win — the second one's status = 'pending' predicate no longer matches. One replica runs today, so this is not yet load-bearing, but it costs nothing and a queue that quietly double-runs jobs is a bad thing to discover later.

Dedupe expires on its own. The unique index is partial, covering only pending and running rows, which makes it a "don't queue this twice right now" guard rather than "never do this again" — a site crawled last week is crawlable today. A duplicate is a reported outcome ({ queued: false }), not a thrown error: one URL pasted twice in a batch should not fail the batch.

Reclaim does not refund the attempt. Jobs abandoned by a dead worker return to the queue after a 15-minute lease, but keep their attempt count — otherwise one poisonous payload that kills its container retries forever and takes down the deployment.

Dead jobs are kept. A row that runs out of attempts becomes failed and holds its payload and last error. Deleting it makes "why did that URL never produce a card" unanswerable.

The drain runs serially, last in the tick. The work behind these jobs is rate-limited by other people's services — GitHub quota, crawl politeness, the model API — so concurrency here would only move the queue into someone else's 429s. limit (25) is the real throttle and bounds how long one tick holds the loop.

Verification

  • bun run check — format, typecheck, 410 tests, 0 fail (13 new).
  • Tests cover the atomic claim, dedupe suppression and its expiry after completion, run_after delay and ordering, retry-with-backoff, the dead-letter transition, reclaim with and without an expired lease, and that one throwing handler does not stop the rest of the tick.
  • Separately verified the incremental migration path the container actually takes: applied 0000–0006 first, then 0007 on top — 7 skipped, 1 applied, all three indexes created, and a re-run applies 0.

Not in this PR

No new job kinds are enqueued yet — nothing calls enqueue in production paths. The dispatcher handles rescore_prospect and process_deletion, and throws by name for anything else. Making POST /prospects asynchronous is Phase 3 work and a behaviour change to an existing endpoint, so it is deliberately separate.

🤖 Generated with Claude Code

Phase 1 of the URL-first plan, and the thing every later phase is
blocked on. `JOB_KINDS` has described this table since the first commit
and nothing ever created it: work either ran inline in the request —
`POST /prospects` still runs the whole enrich/resolve/research/score/
recommend/draft chain before responding — or the worker loop swept its
own domain tables directly. Neither survives a batch.

`migrations/0007_jobs.sql` adds the table. A table rather than Redis:
REDIS_URL sits empty in the vault and unreferenced in the code, Turso is
already here and already backed up, and a queue that loses its contents
on restart is not a queue.

`packages/pipeline/src/queue.ts` holds the state transitions and knows
nothing about what a job does — `drainQueue` takes a handler, so the
dispatcher lives in the server and the tests can pass a stub.

Four decisions worth stating:

- The claim is one `UPDATE … WHERE id = (SELECT …)`, so two workers
  cannot both win. One replica runs today, but the guard costs nothing
  and a queue that quietly double-runs is a bad thing to find out later.
- Dedupe is a partial unique index over pending and running rows only.
  That makes it a "not twice right now" guard rather than "never again",
  so a site crawled last week is crawlable today. A duplicate is a
  reported outcome, not a thrown error — one URL pasted twice in a batch
  should not fail the batch.
- Abandoned jobs are reclaimed after a 15-minute lease, without
  refunding the attempt. A payload that kills its container must still
  run out of attempts instead of looping forever.
- A dead job keeps its row, its payload and its last error. Deleting it
  makes "why did that URL never produce a card" unanswerable.

The worker tick drains the queue last, after the existing bounded
sweeps, since this is the part that can take the whole tick.

Adds a `job` id prefix; the domain had no kind for something the schema
now stores.

13 new tests. Verified separately that 0007 applies incrementally on top
of an existing 7-migration database and that a re-run is a no-op, which
is what the container does at boot.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
@ralyodio
ralyodio merged commit 16f38ab into main Aug 13, 2026
4 checks passed
ralyodio added a commit that referenced this pull request Aug 13, 2026
Paste a company URL — or a hundred — and get approval cards. Phases 2,
3 and 4 of docs/url-first-pipeline.md, on top of the queue from #19.

Phase 2, the site provider (packages/providers/src/site/).

Hybrid extraction, as decided. The deterministic pass reads what a site
published about itself — JSON-LD, OpenGraph, rel=me, outbound links to
known networks — and the model runs only where that came back empty. A
page with decent markup never costs a token, and the grounding rule is
satisfied by construction: a value lifted from a site's own structured
data has evidence attached to it.

robots.txt is read and obeyed, the user agent names the product and
carries a contact URL, redirects are followed by hand so each hop is
re-checked against that origin's rules, and body size, redirect depth
and timeouts are all bounded. That is in the first commit rather than a
later hardening pass because the policy engine permits `website/observe`
as "permitted public web retrieval" — and what makes a retrieval
permitted is partly that the site said it was.

Phase 3, company-first intake.

`POST /prospects/by-url` takes one URL or up to a hundred and answers
202 with a batch id; `GET /batches/:id` reports each URL's state rather
than a bare count, because after a hundred URLs the question is which
ones failed and why. URLs dedupe by host, so example.com, www. and the
bare domain are one company. A URL already queued is reported as a
duplicate, never an error: pasting the same address twice is a normal
thing for a human to do and must not fail the other ninety-nine.

`runPipeline` splits: `runPipelineForCandidate` is the chain from a
candidate onwards, and the GitHub path is now one caller that produces
such a candidate. Provenance and identity source types follow the
provider that actually supplied the value instead of being hardcoded to
GitHub's — with a crawl in the mix, that would have labelled a scraped
name as an API fact and mis-classified what may be retained (PRD §35).
The old GitHub anchor rule generalises to `anchorNetwork`: a crawl has
no anchor, because nobody named on a company page is *proven* to be
that person.

Phase 4, fan-out — partial, deliberately.

`findIdentities` asks every configured provider what else it can vouch
for, then hands the lot to the existing resolver. It gathers claims and
merges nothing; keeping those apart is what stops a provider's
confidence becoming the product's. A provider is only asked about a
network we already hold a handle for — handing a display name to a
network with no verification and taking the first hit is how the wrong
human reaches an approval queue.

Bluesky ships because its AppView needs no key and no contract, so it
can be built and tested without a commercial decision attached. X and
the enrichment vendors are not here: both need credentials I cannot
test against, and an adapter that has never made a real call is code
pretending to be a feature.

62 new tests. Verified 0008 applies incrementally on an existing
database, as the container does at boot.

Known gap, stated rather than hidden: only GitHub produces signals, so a
person found by crawl reaches the queue with evidence but no activity.
The card still carries the prospect and its source, and the composer
refuses to invent a reason it cannot ground.

Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com>
@ralyodio
ralyodio deleted the feat/job-queue branch August 16, 2026 17:33
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant