diff --git a/server/api/jobs/run.get.ts b/server/api/jobs/run.get.ts index 9537e39..0d82c4a 100644 --- a/server/api/jobs/run.get.ts +++ b/server/api/jobs/run.get.ts @@ -1,6 +1,6 @@ import crypto from 'node:crypto' import { dispatch } from '#server/utils/job-handlers' -import { claim, complete, fail } from '#server/utils/queue' +import { claim, cleanupOldJobs, complete, fail } from '#server/utils/queue' const LEASE_MS = 5 * 60_000 // 5 min — generous for a sync job @@ -60,6 +60,11 @@ export default defineEventHandler(async event => { 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 @@ -96,5 +101,5 @@ export default defineEventHandler(async event => { await Promise.all(Array.from({ length: concurrency }, () => lane())) - return { ok: true, processed, failed, drained, concurrency } + return { ok: true, processed, failed, pruned, drained, concurrency } }) diff --git a/server/utils/queue.ts b/server/utils/queue.ts index cd60a63..651f708 100644 --- a/server/utils/queue.ts +++ b/server/utils/queue.ts @@ -104,6 +104,31 @@ export async function complete(id: number) { .where(sql`${job.id} = ${id}`) } +/** Terminal jobs older than this are pruned. `failed` rows are kept longer so + * there's a window to inspect them before they age out. */ +const DONE_RETENTION_MS = 24 * 60 * 60_000 // 1 day +const FAILED_RETENTION_MS = 7 * 24 * 60 * 60_000 // 7 days + +/** + * Delete old terminal jobs so the table stays small. Unbounded growth slows + * both the `claim()` scan and the dashboard's per-repo job aggregation (which + * has no index on the payload's repo id). Returns the number of rows removed. + * Safe to call opportunistically from the worker; it only touches rows past + * their retention window. + */ +export async function cleanupOldJobs(): Promise { + const db = useDb() + const doneCutoff = new Date(Date.now() - DONE_RETENTION_MS) + const failedCutoff = new Date(Date.now() - FAILED_RETENTION_MS) + const result = await db.execute(sql` + DELETE FROM ${job} + WHERE (${job.status} = 'done' AND ${job.updatedAt} < ${doneCutoff.toISOString()}) + OR (${job.status} = 'failed' AND ${job.updatedAt} < ${failedCutoff.toISOString()}) + `) + const affected = (result as { rowCount?: number, affectedRows?: number }) + return affected.rowCount ?? affected.affectedRows ?? 0 +} + /** * Record a failure. Re-queues with exponential backoff until `maxAttempts`, * after which the job is marked `failed` and stays put for inspection. diff --git a/test/unit/queue.spec.ts b/test/unit/queue.spec.ts index 6d93895..616e65f 100644 --- a/test/unit/queue.spec.ts +++ b/test/unit/queue.spec.ts @@ -2,7 +2,7 @@ import { sql } from 'drizzle-orm' import { afterEach, beforeEach, describe, expect, it } from 'vitest' import { job } from '../../server/db/schema' import { clearDb, setDb, useDb } from '../../server/utils/db' -import { claim, complete, enqueue, fail, MAX_ATTEMPTS } from '../../server/utils/queue' +import { claim, cleanupOldJobs, complete, enqueue, fail, MAX_ATTEMPTS } from '../../server/utils/queue' import { createTestDb } from '../utils/db' describe('queue', () => { @@ -123,3 +123,32 @@ describe('queue', () => { expect(rows[0]?.lastError).toBe('terminal') }) }) + +describe('cleanupOldJobs', () => { + beforeEach(async () => { + setDb(await createTestDb()) + }) + afterEach(() => clearDb()) + + it('deletes old done/failed jobs but keeps recent and active ones', async () => { + const db = useDb() + // Insert rows at controlled ages via raw updated_at. + await enqueue('github.push', { a: 1 }) // stays queued + const oldDone = await enqueue('github.push', { a: 2 }) + const recentDone = await enqueue('github.push', { a: 3 }) + const oldFailed = await enqueue('github.push', { a: 4 }) + + await db.execute(sql`UPDATE ${job} SET status='done', updated_at = now() - interval '2 days' WHERE ${job.id} = ${oldDone.id}`) + await db.execute(sql`UPDATE ${job} SET status='done', updated_at = now() - interval '1 hour' WHERE ${job.id} = ${recentDone.id}`) + await db.execute(sql`UPDATE ${job} SET status='failed', updated_at = now() - interval '8 days' WHERE ${job.id} = ${oldFailed.id}`) + + const removed = await cleanupOldJobs() + expect(removed).toBe(2) // oldDone + oldFailed + + const remaining = await db.select().from(job) + const ids = remaining.map(r => r.id) + expect(ids).not.toContain(oldDone.id) + expect(ids).not.toContain(oldFailed.id) + expect(ids).toContain(recentDone.id) + }) +})