Nakodo turns a brand's website into a list of YouTube creators who fit, with contact details and a drafted first message. The landing page describes that as four steps: brief, survey, identify, make contact. Behind those four words is a background pipeline that spends a rationed third-party API, calls a language model, and has to survive being killed halfway through, because it runs on scheduled serverless invocations.
It is one Postgres table and about 200 lines. This post is the queue, not the analysis: what each stage does with its data is the product, and that part stays indoors.
The table
export const jobs = pgTable("jobs", {
id: uuid().primaryKey().defaultRandom(),
type: text({ enum: JOB_TYPES }).notNull(),
campaignId: uuid().references(() => campaigns.id, { onDelete: "cascade" }),
payload: jsonb().$type<JobPayload>().notNull(),
dedupeKey: text(),
status: text({ enum: ["pending", "running", "done", "failed"] }).notNull().default("pending"),
attempts: integer().notNull().default(0),
runAfter: timestamp({ withTimezone: true }).notNull().defaultNow(),
lockedAt: timestamp({ withTimezone: true }),
lastError: text(),
runId: uuid().references(() => runs.id, { onDelete: "set null" }),
createdAt: createdAt(),
finishedAt: timestamp({ withTimezone: true }),
}, (t) => [
index("jobs_claim_idx").on(t.status, t.runAfter),
uniqueIndex("jobs_open_dedupe_uq").on(t.dedupeKey).where(sql`${t.status} in ('pending', 'running')`),
]);
The partial unique index is my favourite line in the schema. Enqueuing is a plain insert with onConflictDoNothing(), and the index decides what a conflict means: only one open job per dedupe key. Ask twice for the same work while it is queued and the second ask disappears. Ask again next week, after the first one finished, and it is allowed, because done rows are not in the index. No application-level "is this already queued" query, which would be a race anyway.
runAfter is the other load-bearing column. It is how backoff, scheduled re-runs and "wait until the API quota resets tomorrow" are all the same mechanism.
Claiming work without two invocations fighting
The cron runs every five minutes and an invocation can outlive its slot, so two runs overlap more often than you would like. Claiming is one transaction: pick ids with FOR UPDATE SKIP LOCKED, then flip them to running.
const ids = await tx
.with(served)
.select({ id: jobs.id })
.from(jobs)
.leftJoin(campaigns, eq(campaigns.id, jobs.campaignId))
.leftJoin(subscriptions, eq(subscriptions.userId, campaigns.userId))
.leftJoin(served, eq(served.campaignId, jobs.campaignId))
.where(and(
eq(jobs.type, type),
lte(jobs.runAfter, sql`now()`),
or(
eq(jobs.status, "pending"),
and(eq(jobs.status, "running"), lt(jobs.lockedAt, STALE_LOCK), lt(jobs.attempts, MAX_ATTEMPTS)),
),
or(isNull(jobs.campaignId), eq(campaigns.status, "active")),
))
.orderBy(sql`${PRIORITY} desc`, sql`coalesce(${served.n}, 0)`, jobs.createdAt)
.limit(limit)
.for("update", { of: jobs, skipLocked: true });
return tx.update(jobs)
.set({ status: "running", lockedAt: new Date(), attempts: sql`${jobs.attempts} + 1`, runId })
.where(inArray(jobs.id, ids.map((r) => r.id)))
.returning();
skipLocked is what makes the overlap harmless: the second invocation walks past rows the first is holding instead of blocking on them. of: jobs matters too, because the query joins campaigns and subscriptions and you do not want to lock a user's subscription row to claim a job.
The or in the status clause is crash recovery. A serverless invocation that is killed leaves rows stuck in running forever, so a row whose lock is older than fifteen minutes is claimable again:
const STALE_LOCK = sql`now() - interval '15 minutes'`;
with attempts < MAX_ATTEMPTS so this recovery cannot loop forever on a job that keeps killing its invocation. And attempts + 1 happens at claim time, not at failure time, which is the only version that counts a crash. If you increment when a handler throws, a job that reliably kills the process is retried for eternity.
Fairness is three ORDER BY terms
Everything interesting about this queue is in that orderBy.
First term: the plan. Paid plans go first when there is a backlog. The priority number lives with the plans, not in SQL, so the CASE is generated from the same object the pricing page renders:
const PRIORITY = sql.raw(
`case when campaigns.user_id is null then ${PLANS.pro.limits.priority} ` +
PLAN_ORDER.map((id) => `when subscriptions.plan = '${id}' then ${PLANS[id].limits.priority}`).join(" ") +
" else 0 end",
);
On the pricing comparison table this is the "Search queue" row: Standard, Ahead of Free, First. That row is this CASE expression, written for customers. I like that you can read the policy in both places and check they agree.
The first when is the quietly important one: a job with no owner, such as a maintenance refresh, is treated as the middle plan rather than as priority zero. Housekeeping that always loses is housekeeping that never happens.
Second term: the campaign that has had the least of this work lately. Priority alone means one large campaign on the top plan starves every other campaign on the top plan. So a CTE counts how many jobs of this type each campaign has had served in the last day, and the campaign with the smallest count goes first:
const served = tx.$with("served").as(
tx.select({ campaignId: jobs.campaignId, n: sql<number>`count(*)::int`.as("n") })
.from(jobs)
.where(and(
eq(jobs.type, type),
inArray(jobs.status, ["done", "running"]),
gte(sql`coalesce(${jobs.finishedAt}, ${jobs.lockedAt})`, sql`now() - interval '1 day'`),
))
.groupBy(jobs.campaignId),
);
coalesce(finished_at, locked_at) counts in-flight work as served, otherwise a campaign currently being worked on looks idle and gets more.
Third term: oldest first. Only as a tie-break. FIFO is the fallback, not the policy.
A job that cannot run yet is not a job that failed
This distinction took two rewrites to get right, and it is the thing I would most want someone to copy.
A failure means the work was attempted and went wrong. Retry with backoff, then give up:
export async function failJob(job: Job, error: unknown): Promise<"retry" | "failed"> {
const giveUp = job.attempts >= MAX_ATTEMPTS;
await db.update(jobs).set(
giveUp
? { status: "failed", finishedAt: new Date(), lockedAt: null, lastError: message }
: { status: "pending", lockedAt: null, lastError: message, runAfter: new Date(Date.now() + 2 ** job.attempts * 60_000) },
).where(and(eq(jobs.id, job.id), eq(jobs.status, "running")));
return giveUp ? "failed" : "retry";
}
Two minutes, then four, then it is failed with the error text kept on the row. Note the and(eq(jobs.id, ...), eq(jobs.status, "running")): the update only applies to a job this run still holds, so a reclaimed job is not stomped by the invocation that lost it.
A deferral means nothing went wrong, the work just cannot happen now: the run is out of time, or the account's daily search allowance is spent, or the shared API quota is gone until tomorrow. Those go back without burning an attempt:
export class DeferJob extends Error {
constructor(readonly until: Date = new Date()) { super("Deferred"); this.name = "DeferJob"; }
}
export async function deferJobs(ids: string[], until: Date): Promise<void> {
await db.update(jobs)
.set({ status: "pending", lockedAt: null, runAfter: until, attempts: sql`greatest(${jobs.attempts} - 1, 0)` })
.where(and(inArray(jobs.id, ids), eq(jobs.status, "running")));
}
The attempts - 1 undoes the increment that claiming added. Without it, three quiet days of API pressure would retire jobs that were never actually tried, and the user would see work silently vanish. A custom error type carrying a wake-up time is a small thing that made the handlers much easier to write: a handler that notices it cannot proceed throws new DeferJob(resetTime) and stops thinking about queue mechanics.
There is also the third case, which is a job whose invocation died on its last attempt. Nobody will ever reclaim it, because reclaiming requires attempts < MAX_ATTEMPTS:
// Jobs whose invocation died on their last attempt: never reclaimed, so mark
// them failed. Returns them so paid work can be refunded.
export async function failStuckJobs(): Promise<Job[]>
It returns the rows rather than just counting them, because some of that work was paid for in plan credits, and a daily job hands those back:
async function gaveUp(job: Job): Promise<void> {
if (job.type === "draft") await refundDraftJob(job);
}
A metered product that charges for work it never delivered is a product that generates refund email. Cheaper to write the ten lines.
The loop restarts from the top on purpose
The run loop walks a list of stages, claims one batch from the first stage that has work, and then goes back to the beginning rather than continuing down the list.
// After each batch the loop starts again from the top, so earlier stages win.
const STAGES = [
{ type: "enrich", batch: 25, run: runEnrichJobs },
{ type: "contacts", batch: 5, run: each(5, runContactsJob) },
{ type: "draft", batch: 5, run: each(5, runDraftJob) },
{ type: "rationale", batch: 8, run: each(8, runRationaleJob) },
{ type: "comments", batch: 8, run: each(8, runCommentsJob) },
{ type: "score", batch: 1, run: each(1, runScoreJob) },
{ type: "search", batch: 1, run: each(1, runSearchJob) },
];
The expensive, work-generating stage is last, and finishing work beats starting more of it. Drain-the-first-stage-first is the difference between a user watching finished entries appear one by one and a user watching a progress bar that never resolves because every pass found more to do.
The run has a budget, and the budget is smaller than the timeout
// Runs every 5 minutes. Works for about 3.5 minutes, then stops claiming jobs.
export const maxDuration = 300;
export async function GET(request: Request) {
if (!isCronRequest(request)) return new Response("Unauthorized", { status: 401 });
const result = await runPipeline({ trigger: "cron", timeBudgetMs: 220_000 });
return Response.json({ runId: result.runId, stopReason: result.stopReason, jobsDone: result.jobsDone, jobsFailed: result.jobsFailed, units: result.units });
}
The platform limit is 300 seconds and the run stops claiming at 220. The gap is for finishing what is already in flight, so a run ends by completing its last batch rather than by being executed mid-write. Anything still queued is somebody else's turn in five minutes.
Every invocation also writes a row of its own:
stopReason: text({ enum: ["no_work", "time_budget", "quota_reserve", "quota_exhausted", "error"] }),
jobsDone, jobsFailed, quotaUnits, summary, error, startedAt, finishedAt
Five stop reasons, and four of them are healthy. "Ran out of time with work left" and "ran out of third-party quota" are normal operating states for this system, not incidents, and a run log that cannot tell them apart from error is a run log that cries wolf. The summary counts searches, channels enriched, contacts found, drafts written and tokens spent, which means the question "what did it actually do at 3am" has an answer in one row.
Would I reach for a queue service instead?
Not here. The work is lumpy, latency-tolerant, and already next to its data: every handler reads and writes the same Postgres, so a separate broker would add a second place for truth to live and a second thing to be down. The whole mechanism is a table, one index, one claim query and two status transitions, and select * from jobs where status = 'failed' is the entire debugging story.
The moment I would change my mind is when claiming becomes the bottleneck, or when work needs to start within seconds of being asked for. Neither is true of a pipeline that spends a daily API budget.
If you want to see the output rather than the plumbing, nakodo.app runs the free plan with no card, and the pricing page's search queue row is this ORDER BY written for customers.
Top comments (0)