import crypto from 'node:crypto' import { isTransientTransport } from '#server/utils/git-wire/errors' import { dispatch } from '#server/utils/job-handlers' import { claim, cleanupOldJobs, complete, fail } from '#server/utils/queue' const LEASE_MS = 5 * 60_000 // 5 min — generous for a sync job // Use most of Vercel's 300s `maxDuration` (see nuxt.config `functions`). The // old 25s budget drained only a handful of jobs per minute, so the queue grew // unboundedly under real push volume. Leave headroom for the response and for // in-flight jobs to settle. const DEFAULT_BUDGET_MS = 270_000 // How many jobs to run at once inside one invocation. Each job is mostly // network wait (SSH to the knot, HTTPS to GitHub), so concurrency buys real // throughput without CPU contention. A single slow/hung job no longer blocks // the others for the whole budget. const DEFAULT_CONCURRENCY = 6 // Hard cap on a single job's wall-clock time. A knot SSH connection can hang // past its `readyTimeout`; without this the slot stays busy until the budget // ends. On timeout the job is recorded as a failure so backoff and the attempt // ceiling apply, and the slot is freed for the next job. const JOB_TIMEOUT_MS = 45_000 class JobTimeoutError extends Error { constructor(ms: number) { super(`job exceeded ${ms}ms wall-clock cap`) this.name = 'JobTimeoutError' } } async function withTimeout(p: Promise, ms: number): Promise { let timer: ReturnType | undefined try { return await Promise.race([ p, new Promise((_, reject) => { timer = setTimeout(() => reject(new JobTimeoutError(ms)), ms) }), ]) } finally { if (timer) clearTimeout(timer) } } export default defineEventHandler(async event => { const cronSecret = process.env.CRON_SECRET if (!cronSecret) { throw createError({ statusCode: 500, statusMessage: 'cron secret not configured' }) } const auth = getRequestHeader(event, 'authorization') if (auth !== `Bearer ${cronSecret}`) { throw createError({ statusCode: 401, statusMessage: 'unauthorized' }) } const config = useRuntimeConfig() const budgetMs = Number(config.workerBudgetMs) || DEFAULT_BUDGET_MS const concurrency = Number(config.workerConcurrency) || DEFAULT_CONCURRENCY const deadline = Date.now() + budgetMs // Opportunistic pruning of old terminal jobs. Cheap, and keeps the table // from growing without bound (which would slow claim() and the dashboard's // per-repo job aggregation). Failures here must not block draining. const pruned = await cleanupOldJobs().catch(() => 0) let processed = 0 let failed = 0 let drained = false // Each lane runs the claim -> dispatch -> record loop independently until the // queue drains or the budget runs out. `claim()` is an atomic // `FOR UPDATE SKIP LOCKED`, so lanes never race for the same row; concurrency // comes from lanes waiting on different jobs' network I/O at the same time. async function lane() { const workerId = `${process.env.VERCEL_DEPLOYMENT_ID ?? 'local'}:${crypto.randomUUID()}` while (Date.now() < deadline) { // eslint-disable-next-line no-await-in-loop -- each iteration processes one job to completion const job = await claim(workerId, LEASE_MS) if (!job) { drained = true return } try { // eslint-disable-next-line no-await-in-loop await withTimeout(dispatch(job), JOB_TIMEOUT_MS) // eslint-disable-next-line no-await-in-loop await complete(job.id) } catch (err) { failed++ // A flaky knot resets connections for minutes at a time; a push to a // real branch must survive that and land once the knot recovers, not // be retired to `failed` after 8 unlucky attempts. Keep transient // transport failures (and the per-job wall-clock cap, which usually // means "slow right now") retrying indefinitely on the capped backoff. const retryForever = isTransientTransport(err) || err instanceof JobTimeoutError // eslint-disable-next-line no-await-in-loop await fail(job.id, job.attempts, err, { retryForever }) } processed++ } } await Promise.all(Array.from({ length: concurrency }, () => lane())) return { ok: true, processed, failed, pruned, drained, concurrency } })