Repository navigation
feat(pipeline): a durable job queue - #19
Merged
Merged
Conversation
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
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>
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Phase 1 of
docs/url-first-pipeline.md, and the thing every later phase is blocked on.JOB_KINDShas described this table since the first commit and nothing ever created it. Work either ran inline in the request —POST /prospectsstill 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:drainQueuetakes a handler, so the dispatcher lives in the server and tests pass a stub.apps/server/src/index.ts— an explicitrunJobdispatcher and a drain at the end of each tick.packages/domain/src/ids.ts— ajobid prefix; the domain had no kind for something the schema now stores.Decisions worth reviewing
A table, not Redis.
REDIS_URLsits 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'sstatus = '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
pendingandrunningrows, 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
failedand 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).run_afterdelay 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.0000–0006first, then0007on 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
enqueuein production paths. The dispatcher handlesrescore_prospectandprocess_deletion, and throws by name for anything else. MakingPOST /prospectsasynchronous is Phase 3 work and a behaviour change to an existing endpoint, so it is deliberately separate.🤖 Generated with Claude Code