diff --git a/server/api/jobs/run.get.ts b/server/api/jobs/run.get.ts index 0d82c4a..9770117 100644 --- a/server/api/jobs/run.get.ts +++ b/server/api/jobs/run.get.ts @@ -1,4 +1,5 @@ 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' @@ -91,8 +92,14 @@ export default defineEventHandler(async event => { } 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) + await fail(job.id, job.attempts, err, { retryForever }) } processed++ diff --git a/server/utils/queue.ts b/server/utils/queue.ts index 651f708..e581462 100644 --- a/server/utils/queue.ts +++ b/server/utils/queue.ts @@ -129,15 +129,30 @@ export async function cleanupOldJobs(): Promise { return affected.rowCount ?? affected.affectedRows ?? 0 } +export interface FailOptions { + maxAttempts?: number + /** + * When true, the job keeps re-queuing past `maxAttempts` instead of being + * retired to `failed`. For failures that are temporary by nature (a flaky + * knot SSH endpoint resetting connections), giving up permanently would + * abandon legitimate work: a push to a real branch must land eventually, not + * be dropped after 8 unlucky attempts. Backoff still caps at 1h, so an + * indefinitely-retrying job polls at most hourly rather than hammering. + */ + retryForever?: boolean +} + /** * Record a failure. Re-queues with exponential backoff until `maxAttempts`, - * after which the job is marked `failed` and stays put for inspection. + * after which the job is marked `failed` and stays put for inspection, unless + * `retryForever` is set (see `FailOptions`). */ -export async function fail(id: number, attempts: number, err: unknown, maxAttempts = MAX_ATTEMPTS) { +export async function fail(id: number, attempts: number, err: unknown, opts: FailOptions = {}) { + const { maxAttempts = MAX_ATTEMPTS, retryForever = false } = opts const db = useDb() const message = err instanceof Error ? err.message : String(err) - if (attempts >= maxAttempts) { + if (!retryForever && attempts >= maxAttempts) { await db.update(job) .set({ status: 'failed', lastError: message, lockedBy: null, lockedUntil: null, updatedAt: new Date() }) .where(sql`${job.id} = ${id}`) diff --git a/test/unit/queue.spec.ts b/test/unit/queue.spec.ts index 616e65f..c468cfc 100644 --- a/test/unit/queue.spec.ts +++ b/test/unit/queue.spec.ts @@ -115,13 +115,25 @@ describe('queue', () => { it('marks failed once attempts >= maxAttempts', async () => { await enqueue('github.push', {}) const claimed = await claim('worker-1', 60_000) - await fail(claimed!.id, 5, new Error('terminal'), 5) + await fail(claimed!.id, 5, new Error('terminal'), { maxAttempts: 5 }) const db = useDb() const rows = await db.select().from(job).where(sql`${job.id} = ${claimed!.id}`) expect(rows[0]?.status).toBe('failed') expect(rows[0]?.lastError).toBe('terminal') }) + + it('keeps re-queuing past maxAttempts when retryForever is set', async () => { + await enqueue('github.push', {}) + const claimed = await claim('worker-1', 60_000) + // Well past the ceiling, but a transient transport failure must not retire. + await fail(claimed!.id, 99, new Error('ssh error: read ECONNRESET'), { retryForever: true }) + + const db = useDb() + const rows = await db.select().from(job).where(sql`${job.id} = ${claimed!.id}`) + expect(rows[0]?.status).toBe('queued') + expect(new Date(rows[0]?.runAfter ?? 0).getTime()).toBeGreaterThan(Date.now()) + }) }) describe('cleanupOldJobs', () => {