From 9be51373830744edd2889d3b63d542f984894c9e Mon Sep 17 00:00:00 2001 From: Theodore Li Date: Thu, 1 Oct 2026 12:08:41 -0700 Subject: [PATCH 1/6] improvement(cron): run stale execution cleanup as a background task --- .../cleanup-stale-executions/route.test.ts | 482 +---------- .../cron/cleanup-stale-executions/route.ts | 804 +----------------- .../cleanup-stale-executions.test.ts | 451 ++++++++++ .../background/cleanup-stale-executions.ts | 785 +++++++++++++++++ .../core/async-jobs/backends/trigger-dev.ts | 1 + apps/sim/lib/core/async-jobs/types.ts | 1 + 6 files changed, 1302 insertions(+), 1222 deletions(-) create mode 100644 apps/sim/background/cleanup-stale-executions.test.ts create mode 100644 apps/sim/background/cleanup-stale-executions.ts diff --git a/apps/sim/app/api/cron/cleanup-stale-executions/route.test.ts b/apps/sim/app/api/cron/cleanup-stale-executions/route.test.ts index 08476db1683..059a2ea899a 100644 --- a/apps/sim/app/api/cron/cleanup-stale-executions/route.test.ts +++ b/apps/sim/app/api/cron/cleanup-stale-executions/route.test.ts @@ -1,62 +1,18 @@ -import { asyncJobs, tableJobs, workflowExecutionLogs } from '@sim/db/schema' -import { createLogger } from '@sim/logger' -import { createMockRequest, dbChainMockFns, queueTableRows, resetDbChainMock } from '@sim/testing' +import { createMockRequest } from '@sim/testing' +import { asyncJobsMock, asyncJobsMockFns } from '@sim/testing/mocks/async-jobs.mock' import { authInternalMock, authInternalMockFns } from '@sim/testing/mocks/auth-internal.mock' -import { storageServiceMock, storageServiceMockFns } from '@sim/testing/mocks/storage-service.mock' -import { beforeEach, describe, expect, it, vi } from 'vitest' -import { JOB_RETENTION_HOURS } from '@/lib/core/async-jobs' -import { - SCHEDULE_CARRIER_IRRECOVERABLE_METADATA_KEY, - SCHEDULE_CARRIER_IRRECOVERABLE_RETENTION_HOURS, - SCHEDULE_CARRIER_RECONCILED_METADATA_KEY, -} from '@/lib/workflows/schedules/carrier-metadata' +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' vi.mock('@/lib/auth/internal', () => authInternalMock) -vi.mock('@/lib/uploads/core/storage-service', () => storageServiceMock) +vi.mock('@/lib/core/async-jobs', () => asyncJobsMock) import { GET } from '@/app/api/cron/cleanup-stale-executions/route' -const mockVerifyCronAuth = authInternalMockFns.mockVerifyCronAuth -const mockDeleteFile = storageServiceMockFns.mockDeleteFile -mockDeleteFile.mockResolvedValue(undefined) +const { mockVerifyCronAuth } = authInternalMockFns -const cleanupLogger = - vi.mocked(createLogger).mock.results[ - vi.mocked(createLogger).mock.calls.findIndex(([name]) => name === 'CleanupStaleExecutions') - ].value +const mockEnqueue = asyncJobsMockFns.mockJobQueue.enqueue -interface MockCondition { - type?: string - conditions?: unknown[] - left?: unknown - right?: unknown - column?: unknown - values?: unknown - toSQL?: () => { sql: string; params: unknown[] } -} - -function flattenConditions(condition: unknown): MockCondition[] { - if (!condition || typeof condition !== 'object') return [] - const node = condition as MockCondition - return [node, ...(node.conditions?.flatMap((child) => flattenConditions(child)) ?? [])] -} - -function hasToSQL(value: unknown): value is { toSQL: () => { sql: string; params: unknown[] } } { - return typeof value === 'object' && value !== null && 'toSQL' in value -} - -/** - * Collects the leaves of a nested `sql` expression. The duration expression is - * built by a shared helper, so the values it binds sit one level below the - * fragment this route assembles rather than directly in its own params. - */ -function flattenSqlParams(expression: { sql: string; params: unknown[] }): unknown[] { - return expression.params.flatMap((param) => - hasToSQL(param) ? flattenSqlParams(param.toSQL()) : [param] - ) -} - -function createRequest() { +function request() { return createMockRequest( 'GET', undefined, @@ -65,419 +21,53 @@ function createRequest() { ) } -/** Recursively renders a mocked drizzle fragment, expanding nested fragments. */ -function renderSql(fragment: unknown): string { - if (fragment === null || fragment === undefined) return '' - if (typeof fragment !== 'object') return String(fragment) - const candidate = fragment as { rawSql?: string; strings?: string[]; values?: unknown[] } - if (typeof candidate.rawSql === 'string') return candidate.rawSql - if (!candidate.strings) return '' - return candidate.strings - .map((part, index) => - index < (candidate.values?.length ?? 0) - ? `${part}${renderSql(candidate.values?.[index])}` - : part - ) - .join('') -} - -/** Recursively collects the non-fragment bind values of a mocked fragment. */ -function collectSqlParams(fragment: unknown): unknown[] { - if (!fragment || typeof fragment !== 'object') return [] - const candidate = fragment as { values?: unknown[] } - if (!candidate.values) return [] - return candidate.values.flatMap((value) => { - const nested = collectSqlParams(value) - return nested.length > 0 ? nested : [value] - }) -} - -describe('stale execution cleanup deadline grace', () => { +describe('stale execution cleanup route', () => { beforeEach(() => { - resetDbChainMock() - cleanupLogger.info.mockReset() - mockVerifyCronAuth.mockReturnValue(null) - }) - - it('waits five minutes past a workflow execution deadline in both cleanup predicates', async () => { vi.useFakeTimers() - vi.setSystemTime(new Date('2026-08-03T12:10:00.000Z')) - queueTableRows(workflowExecutionLogs, [{ id: 'log-1' }]) - dbChainMockFns.returning.mockResolvedValueOnce([{ id: 'log-1' }]) - - try { - const response = await GET(createRequest()) - - expect(response.status).toBe(200) - const expectedThreshold = new Date('2026-08-03T12:05:00.000Z') - const deadlineComparisons = dbChainMockFns.where.mock.calls - .flatMap(([condition]) => flattenConditions(condition)) - .filter( - (condition) => - condition.type === 'lt' && - condition.right instanceof Date && - condition.right.getTime() === expectedThreshold.getTime() - ) - - expect(deadlineComparisons).toHaveLength(2) - expect(deadlineComparisons.map(({ right }) => right)).toEqual([ - expectedThreshold, - expectedThreshold, - ]) - - const executionUpdateIndex = dbChainMockFns.update.mock.calls.findIndex( - ([table]) => table === workflowExecutionLogs - ) - const update = dbChainMockFns.set.mock.calls[executionUpdateIndex]?.[0] as { - endedAt: Date - totalDurationMs: { toSQL: () => { sql: string; params: unknown[] } } - executionData: { toSQL: () => { sql: string; params: unknown[] } } - } - const totalDurationLeaves = flattenSqlParams(update.totalDurationMs.toSQL()) - const renderedError = renderSql(update.executionData) - const errorLeaves = collectSqlParams(update.executionData) - - expect(renderedError).toContain('CASE') - expect(renderedError).toContain('IS NOT NULL') - expect(renderedError).toContain('ROUND') - expect(renderedError).toContain('EXTRACT(EPOCH') - expect(errorLeaves).toContain(workflowExecutionLogs.executionDeadlineAt) - expect(errorLeaves).toContain('Execution timed out') - expect(errorLeaves).toContain('Execution terminated: worker timeout or crash after ') - expect(errorLeaves).toContain(workflowExecutionLogs.startedAt) - expect(totalDurationLeaves).toContain(2_147_483_647) - expect(totalDurationLeaves).toContain(workflowExecutionLogs.startedAt) - expect(totalDurationLeaves).toContainEqual(new Date('2026-08-03T12:10:00.000Z')) - expect(update.endedAt).toEqual(new Date('2026-08-03T12:10:00.000Z')) - } finally { - vi.useRealTimers() - } - }) - - it('sweeps redacting logs on the generic window, never the execution deadline', async () => { - const response = await GET(createRequest()) - - expect(response.status).toBe(200) - const redactingPredicates = dbChainMockFns.where.mock.calls - .map(([condition]) => flattenConditions(condition)) - .filter((conditions) => - conditions.some( - (condition) => - condition.type === 'eq' && - condition.left === workflowExecutionLogs.status && - condition.right === 'redacting' - ) - ) - - expect(redactingPredicates.length).toBeGreaterThan(0) - for (const conditions of redactingPredicates) { - expect( - conditions.some((condition) => condition.left === workflowExecutionLogs.executionDeadlineAt) - ).toBe(false) - expect( - conditions.some( - (condition) => - condition.type === 'lt' && condition.left === workflowExecutionLogs.startedAt - ) - ).toBe(true) - } - }) - - it('leaves pending and processing schedule jobs to schedule recovery', async () => { - const response = await GET(createRequest()) - - expect(response.status).toBe(200) - const activeAsyncPredicates = dbChainMockFns.where.mock.calls - .map(([condition]) => flattenConditions(condition)) - .filter((conditions) => - conditions.some( - (condition) => - condition.type === 'ne' && - condition.left === asyncJobs.type && - condition.right === 'schedule-execution' - ) - ) - - expect( - activeAsyncPredicates.some((conditions) => - conditions.some( - (condition) => - condition.type === 'eq' && - condition.left === asyncJobs.status && - condition.right === 'processing' - ) - ) - ).toBe(true) - expect( - activeAsyncPredicates.some((conditions) => - conditions.some( - (condition) => - condition.type === 'eq' && - condition.left === asyncJobs.status && - condition.right === 'pending' - ) - ) - ).toBe(true) - }) - - it('retains terminal schedule carriers until reconciliation is recorded', async () => { - const response = await GET(createRequest()) - - expect(response.status).toBe(200) - const retentionConditions = dbChainMockFns.where.mock.calls.flatMap(([condition]) => - flattenConditions(condition) - ) - const reconciliationMarker = retentionConditions.find((condition) => - renderSql(condition).includes(SCHEDULE_CARRIER_RECONCILED_METADATA_KEY) - ) - - expect(collectSqlParams(reconciliationMarker)).toContain(asyncJobs.metadata) - expect( - retentionConditions.some( - (condition) => - condition.type === 'ne' && - condition.left === asyncJobs.type && - condition.right === 'schedule-execution' - ) - ).toBe(true) + vi.setSystemTime(new Date('2026-10-01T17:31:00Z')) + mockVerifyCronAuth.mockReturnValue(null) + mockEnqueue.mockReset() + mockEnqueue.mockResolvedValue('job-stale-1') }) - it('spells carrier metadata keys as SQL literals so the partial index matches', async () => { - const response = await GET(createRequest()) - - expect(response.status).toBe(200) - const reconciliationMarker = dbChainMockFns.where.mock.calls - .flatMap(([condition]) => flattenConditions(condition)) - .find((condition) => renderSql(condition).includes(SCHEDULE_CARRIER_RECONCILED_METADATA_KEY)) - - expect(renderSql(reconciliationMarker)).toContain( - `'${SCHEDULE_CARRIER_RECONCILED_METADATA_KEY}'` - ) - expect(collectSqlParams(reconciliationMarker)).not.toContain( - SCHEDULE_CARRIER_RECONCILED_METADATA_KEY - ) + afterEach(() => { + vi.useRealTimers() }) - it('deletes irrecoverable schedule carrier tombstones once their longer window lapses', async () => { - const response = await GET(createRequest()) + it('hands the cleanup to the job queue and answers without waiting for it', async () => { + const response = await GET(request()) expect(response.status).toBe(200) - const retentionConditions = dbChainMockFns.where.mock.calls.flatMap(([condition]) => - flattenConditions(condition) + await expect(response.json()).resolves.toEqual({ triggered: true, jobId: 'job-stale-1' }) + expect(mockEnqueue).toHaveBeenCalledWith( + 'cleanup-stale-executions', + {}, + expect.objectContaining({ maxAttempts: 1, concurrencyLimit: 1 }) ) - const irrecoverableExclusion = retentionConditions.find((condition) => - renderSql(condition).includes(SCHEDULE_CARRIER_IRRECOVERABLE_METADATA_KEY) - ) - - expect(renderSql(irrecoverableExclusion)).toContain("<> 'true'") - expect(collectSqlParams(irrecoverableExclusion)).toContain(asyncJobs.metadata) - - const tombstoneWindow = retentionConditions.filter( - (condition) => - condition.type === 'lt' && - condition.left === asyncJobs.completedAt && - condition.right instanceof Date - ) - const oldest = Math.min(...tombstoneWindow.map(({ right }) => (right as Date).getTime())) - const newest = Math.max(...tombstoneWindow.map(({ right }) => (right as Date).getTime())) - expect(newest - oldest).toBe( - (SCHEDULE_CARRIER_IRRECOVERABLE_RETENTION_HOURS - JOB_RETENTION_HOURS) * 60 * 60 * 1000 - ) - }) - - it('keeps table-job heartbeat cleanup independent from workflow timeout policy', async () => { - vi.useFakeTimers() - vi.setSystemTime(new Date('2026-08-03T12:00:00.000Z')) - queueTableRows(tableJobs, [{ id: 'table-job-1' }]) - dbChainMockFns.returning.mockResolvedValueOnce([{ id: 'table-job-1' }]) - - try { - const response = await GET(createRequest()) - - expect(response.status).toBe(200) - const expectedThreshold = new Date('2026-08-03T10:25:00.000Z') - const tableJobComparisons = dbChainMockFns.where.mock.calls - .flatMap(([condition]) => flattenConditions(condition)) - .filter( - (condition) => - condition.type === 'lt' && - condition.left === tableJobs.updatedAt && - condition.right instanceof Date && - condition.right.getTime() === expectedThreshold.getTime() - ) - - expect(tableJobComparisons).toHaveLength(2) - expect(tableJobComparisons.map(({ right }) => right)).toEqual([ - expectedThreshold, - expectedThreshold, - ]) - - const tableJobUpdateIndex = dbChainMockFns.update.mock.calls.findIndex( - ([table]) => table === tableJobs - ) - const update = dbChainMockFns.set.mock.calls[tableJobUpdateIndex]?.[0] as { - error: string - } - expect(update.error).toBe( - 'Job terminated: no progress for more than 95 minutes (worker timeout or crash)' - ) - } finally { - vi.useRealTimers() - } }) - it('claims every cleanup page without overlapping concurrent workers', async () => { - const response = await GET(createRequest()) - - expect(response.status).toBe(200) - // Nine batched arms: the connector sync-log retention pass is the newest. - expect(dbChainMockFns.transaction).toHaveBeenCalledTimes(9) - expect(dbChainMockFns.for).toHaveBeenCalledTimes(9) - for (const [strength, options] of dbChainMockFns.for.mock.calls) { - expect(strength).toBe('update') - expect(options).toEqual({ skipLocked: true }) - } - }) - - it('caps every bulk mutation and returns only scalar export cleanup fields', async () => { - const stateBatch = Array.from({ length: 1000 }, (_, index) => ({ id: `state-${index}` })) - const retentionBatch = Array.from({ length: 2000 }, (_, index) => ({ - id: `retention-${index}`, - })) - const exportBatch = Array.from({ length: 100 }, (_, index) => ({ - type: 'export', - resultKey: `workspace/workspace-1/exports/table-1/job-${index}/export.csv`, - })) - const exportCandidates = Array.from({ length: 100 }, (_, index) => ({ - id: `export-${index}`, - })) - - for (let batch = 0; batch < 10; batch++) { - const workflowBatch = Array.from({ length: 100 }, (_, index) => ({ - id: `workflow-state-${batch}-${index}`, - })) - queueTableRows(workflowExecutionLogs, workflowBatch) - dbChainMockFns.returning.mockResolvedValueOnce(workflowBatch) - } - for (let batch = 0; batch < 10; batch++) { - queueTableRows(asyncJobs, stateBatch) - dbChainMockFns.returning.mockResolvedValueOnce(stateBatch) - } - for (let batch = 0; batch < 10; batch++) { - queueTableRows(tableJobs, stateBatch) - dbChainMockFns.returning.mockResolvedValueOnce(stateBatch) - } - for (let batch = 0; batch < 10; batch++) { - queueTableRows(tableJobs, exportCandidates) - dbChainMockFns.returning.mockResolvedValueOnce(exportBatch) - } - for (let batch = 0; batch < 10; batch++) { - queueTableRows(asyncJobs, stateBatch) - dbChainMockFns.returning.mockResolvedValueOnce(stateBatch) - } - for (let batch = 0; batch < 10; batch++) { - queueTableRows(asyncJobs, retentionBatch) - dbChainMockFns.returning.mockResolvedValueOnce(retentionBatch) - } - dbChainMockFns.returning.mockResolvedValueOnce([]) - - const response = await GET(createRequest()) + it('deduplicates retries within the same thirty-minute schedule window', async () => { + await GET(request()) + vi.advanceTimersByTime(28 * 60 * 1000) + await GET(request()) - expect(response.status).toBe(200) - await expect(response.json()).resolves.toMatchObject({ - executions: { - found: 1000, - cleaned: 1000, - failed: 0, - }, - asyncJobs: { - staleProcessingMarkedFailed: 10_000, - stalePendingMarkedFailed: 10_000, - oldDeleted: 20_000, - }, - tableJobs: { - staleMarkedFailed: 10_000, - }, - }) - expect(mockDeleteFile).toHaveBeenCalledTimes(1000) - - const limits = dbChainMockFns.limit.mock.calls.map(([limit]) => limit) - expect(limits.filter((limit) => limit === 100)).toHaveLength(20) - expect(limits.filter((limit) => limit === 1000)).toHaveLength(30) - expect(limits.filter((limit) => limit === 2000)).toHaveLength(12) - - const workflowUpdates = dbChainMockFns.update.mock.calls.filter( - ([table]) => table === workflowExecutionLogs - ) - expect(workflowUpdates).toHaveLength(10) - - const returningShapes = dbChainMockFns.returning.mock.calls - .map(([shape]) => shape) - .filter((shape): shape is Record => Boolean(shape)) - expect(returningShapes.some((shape) => 'payload' in shape)).toBe(false) - expect(returningShapes.some((shape) => 'type' in shape && 'resultKey' in shape)).toBe(true) - - const claimedIds = dbChainMockFns.where.mock.calls - .flatMap(([condition]) => flattenConditions(condition)) - .filter((condition) => condition.type === 'inArray') - .map((condition) => condition.values) - expect(claimedIds.length).toBeGreaterThan(0) - expect(claimedIds.every((ids) => Array.isArray(ids))).toBe(true) + expect(mockEnqueue.mock.calls[0]?.[2]?.jobId).toBe(mockEnqueue.mock.calls[1]?.[2]?.jobId) }) - it('preserves committed workflow cleanup counts when a later batch fails', async () => { - const firstBatch = Array.from({ length: 100 }, (_, index) => ({ - id: `execution-${index}`, - })) - const failedBatch = Array.from({ length: 37 }, (_, index) => ({ - id: `failed-execution-${index}`, - })) - queueTableRows(workflowExecutionLogs, firstBatch) - queueTableRows(workflowExecutionLogs, failedBatch) - queueTableRows(asyncJobs, [{ id: 'async-job-1' }]) - dbChainMockFns.returning - .mockResolvedValueOnce(firstBatch) - .mockRejectedValueOnce(new Error('database unavailable')) - .mockResolvedValueOnce([{ id: 'async-job-1' }]) - - const response = await GET(createRequest()) + it('uses a new id immediately after the next thirty-minute window begins', async () => { + vi.setSystemTime(new Date('2026-10-01T17:59:59.999Z')) + await GET(request()) + vi.setSystemTime(new Date('2026-10-01T18:00:00.000Z')) + await GET(request()) - expect(response.status).toBe(200) - await expect(response.json()).resolves.toMatchObject({ - executions: { - found: 137, - cleaned: 100, - failed: 37, - }, - }) - expect( - dbChainMockFns.update.mock.calls.filter(([table]) => table === workflowExecutionLogs) - ).toHaveLength(2) - expect(dbChainMockFns.update.mock.calls.some(([table]) => table === asyncJobs)).toBe(true) + expect(mockEnqueue.mock.calls[0]?.[2]?.jobId).not.toBe(mockEnqueue.mock.calls[1]?.[2]?.jobId) }) - it('continues draining when an atomic race updates fewer rows than were selected', async () => { - const firstCandidates = Array.from({ length: 100 }, (_, index) => ({ - id: `execution-${index}`, - })) - const firstUpdated = firstCandidates.slice(0, 99) - const secondBatch = [{ id: 'execution-100' }] - queueTableRows(workflowExecutionLogs, firstCandidates) - queueTableRows(workflowExecutionLogs, secondBatch) - dbChainMockFns.returning.mockResolvedValueOnce(firstUpdated).mockResolvedValueOnce(secondBatch) + it('fails the cron invocation when the job cannot be enqueued', async () => { + mockEnqueue.mockRejectedValue(new Error('queue unavailable')) - const response = await GET(createRequest()) + const response = await GET(request()) - expect(response.status).toBe(200) - await expect(response.json()).resolves.toMatchObject({ - executions: { - found: 101, - cleaned: 100, - failed: 0, - }, - }) - expect( - dbChainMockFns.update.mock.calls.filter(([table]) => table === workflowExecutionLogs) - ).toHaveLength(2) + expect(response.status).toBe(500) }) }) diff --git a/apps/sim/app/api/cron/cleanup-stale-executions/route.ts b/apps/sim/app/api/cron/cleanup-stale-executions/route.ts index 0eba10da8cc..cf514015843 100644 --- a/apps/sim/app/api/cron/cleanup-stale-executions/route.ts +++ b/apps/sim/app/api/cron/cleanup-stale-executions/route.ts @@ -1,794 +1,46 @@ -import { db } from '@sim/db' -import { - asyncJobs, - knowledgeConnectorSyncLog, - tableJobs, - workflowDeploymentOperation, - workflowExecutionLogs, -} from '@sim/db/schema' import { createLogger } from '@sim/logger' -import { toError } from '@sim/utils/errors' -import { and, eq, exists, gt, inArray, isNull, lt, ne, or, sql } from 'drizzle-orm' -import { alias } from 'drizzle-orm/pg-core' import { type NextRequest, NextResponse } from 'next/server' import { verifyCronAuth } from '@/lib/auth/internal' -import { - JOB_PENDING_RETENTION_HOURS, - JOB_RETENTION_HOURS, - JOB_STATUS, - MAX_JOB_DURATION_SECONDS, - MIN_JOB_DURATION_SECONDS, - TERMINAL_JOB_STATUSES, -} from '@/lib/core/async-jobs' -import { - getExecutionReservationTtlMs, - getTimeoutErrorMessage, - RESERVATION_TTL_BUFFER_MS, -} from '@/lib/core/execution-limits' +import { getJobQueue } from '@/lib/core/async-jobs' import { withRouteHandler } from '@/lib/core/utils/with-route-handler' -import type { DbTransaction } from '@/lib/db/types' -import { elapsedDurationMsSql } from '@/lib/logs/execution/duration' -import { - STALE_SWEEPABLE_EXECUTION_STATUSES, - type StaleSweepableExecutionStatus, -} from '@/lib/logs/types' -import { sweepOrphanedRuns } from '@/lib/mothership/async-runs/orphaned-runs' -import { cancelStaleDispatches } from '@/lib/table/dispatcher' -import { deleteFile } from '@/lib/uploads/core/storage-service' -import { - carrierNotIrrecoverableSql, - carrierReconciledSql, - SCHEDULE_CARRIER_IRRECOVERABLE_RETENTION_HOURS, -} from '@/lib/workflows/schedules/carrier-metadata' -import { SCHEDULE_EXECUTION_QUEUE_NAME } from '@/lib/workflows/schedules/execution-limits' -const logger = createLogger('CleanupStaleExecutions') +export const dynamic = 'force-dynamic' -const STALE_THRESHOLD_MS = getExecutionReservationTtlMs() -const STALE_THRESHOLD_MINUTES = Math.ceil(STALE_THRESHOLD_MS / 60000) -const GENERIC_STALE_PROCESSING_ERROR = `Job terminated: stuck in processing for more than ${STALE_THRESHOLD_MINUTES} minutes` -const EXECUTION_DEADLINE_ERROR = getTimeoutErrorMessage(undefined) -/** - * Table jobs run as detached workers with progress heartbeats, independently of workflow timeout - * policy. Preserve their historical 90-minute task window plus five-minute cleanup grace. - */ -const TABLE_JOB_STALE_THRESHOLD_MINUTES = 95 -/** Terminal table-jobs older than this are pruned; only the latest job per table is ever read. */ -const TABLE_JOB_RETENTION_HOURS = 24 -/** - * A table run dispatch whose holder has not made progress for this long is - * treated as dead. Same shape and window as the table-job threshold above: the - * 90-minute Trigger.dev task ceiling (`maxDuration` in `trigger.config.ts`) plus - * five minutes of cleanup grace, measured from the dispatcher's own per-window - * heartbeat rather than from when the run was requested. - */ -const TABLE_DISPATCH_STALE_THRESHOLD_MINUTES = 95 -/** Per-run ceiling on reaped dispatches, so one tick cannot fan out unbounded SSE. */ -const TABLE_DISPATCH_MAX_PER_RUN = 200 -/** - * Terminal deployment operations older than this are pruned. Every reader of - * this table is latest-generation-only, and idempotency keys only need to - * survive a client retry window, so 30 days is generous. - */ -const DEPLOYMENT_OPERATION_RETENTION_DAYS = 30 -/** - * Terminal connector sync logs older than this are pruned. Nothing pruned them - * before, so the table grew by one row per sync run forever — a connector on a - * fifteen-minute interval writes about 35,000 rows a year on its own. That cost - * lands on `loadPreviousListingObservation`, which reads the newest `completed` - * row per connector through an index covering `connector_id` alone, so every - * retained row makes the sort behind the deletion-safety corroboration slower. - */ -const CONNECTOR_SYNC_LOG_RETENTION_DAYS = 30 -const CONNECTOR_SYNC_LOG_PRUNE_BATCH_SIZE = 2000 -const CONNECTOR_SYNC_LOG_MAX_ROWS_PER_RUN = 20_000 -const DEPLOYMENT_OPERATION_PRUNE_BATCH_SIZE = 2000 -const DEPLOYMENT_OPERATION_PRUNE_MAX_BATCHES = 10 -const WORKFLOW_EXECUTION_MUTATION_BATCH_SIZE = 100 -const WORKFLOW_EXECUTION_MAX_ROWS_PER_RUN = 1000 -const STATE_MUTATION_BATCH_SIZE = 1000 -const STATE_MUTATION_MAX_ROWS_PER_RUN = 10_000 -const RETENTION_DELETE_BATCH_SIZE = 2000 -const RETENTION_DELETE_MAX_ROWS_PER_RUN = 20_000 -const TABLE_JOB_PRUNE_BATCH_SIZE = 100 -const TABLE_JOB_PRUNE_MAX_ROWS_PER_RUN = 1000 - -interface BatchedMutationResult { - affected: number - reachedLimit: boolean -} - -interface RunBatchedMutationOptions { - batchSize: number - maxRowsPerRun: number - claim: (tx: DbTransaction, limit: number) => Promise - mutation: (tx: DbTransaction, candidateIds: string[]) => Promise - onBatch?: (rows: TRow[]) => Promise -} - -/** - * Runs a mutation in bounded, atomic pages. Each page claims explicit rows with - * `FOR UPDATE SKIP LOCKED` and mutates those same rows before committing, so - * concurrent cleanup workers cannot overlap and `LIMIT` is evaluated exactly once. - */ -async function runBatchedMutation({ - batchSize, - maxRowsPerRun, - claim, - mutation, - onBatch, -}: RunBatchedMutationOptions): Promise { - let affected = 0 - - while (affected < maxRowsPerRun) { - const limit = Math.min(batchSize, maxRowsPerRun - affected) - const { candidates, rows } = await db.transaction(async (tx) => { - const candidates = await claim(tx, limit) - if (candidates.length > limit) { - throw new Error(`Cleanup claimed ${candidates.length} rows for a ${limit}-row batch`) - } - if (candidates.length === 0) return { candidates, rows: [] as TRow[] } - - const rows = await mutation( - tx, - candidates.map(({ id }) => id) - ) - if (rows.length !== candidates.length) { - throw new Error( - `Cleanup mutation returned ${rows.length} rows for ${candidates.length} claimed rows` - ) - } - return { candidates, rows } - }) - if (candidates.length === 0) break - - affected += rows.length - if (onBatch) await onBatch(rows) - if (candidates.length < limit) break - } - - return { affected, reachedLimit: affected >= maxRowsPerRun } -} +const logger = createLogger('CleanupStaleExecutionsApi') +const STALE_EXECUTION_CLEANUP_INTERVAL_MS = 30 * 60 * 1000 export const GET = withRouteHandler(async (request: NextRequest) => { try { const authError = verifyCronAuth(request, 'Stale execution cleanup') - if (authError) { - return authError - } - - logger.info('Starting stale execution cleanup job') - - const now = new Date() - const staleDeadlineThreshold = new Date(now.getTime() - RESERVATION_TTL_BUFFER_MS) - const staleThreshold = new Date(now.getTime() - STALE_THRESHOLD_MINUTES * 60 * 1000) - const stalePendingThreshold = new Date( - now.getTime() - JOB_PENDING_RETENTION_HOURS * 60 * 60 * 1000 - ) - const staleTableJobThreshold = new Date( - now.getTime() - TABLE_JOB_STALE_THRESHOLD_MINUTES * 60 * 1000 - ) - const staleDispatchThreshold = new Date( - now.getTime() - TABLE_DISPATCH_STALE_THRESHOLD_MINUTES * 60 * 1000 - ) - - let staleExecutionsFound = 0 - let cleaned = 0 - let failed = 0 - let currentWorkflowBatchSize = 0 - - try { - /** - * `running` is swept on its execution deadline plus the cleanup grace, or - * on the generic stale window when it has no deadline. - * - * `redacting` gets the generic window only. That status is set *after* the - * run finished, while a live worker masks the payload, so the execution - * deadline is already in the past by the time redaction starts — the - * deadline rule would fail a masking pass that is merely slow, and - * schedule recovery would read that as a failed occurrence even though - * the worker is about to persist `completed`. A redaction that genuinely - * crashed still clears within the generic window. - */ - const staleExecutionTimePredicate = (status: StaleSweepableExecutionStatus) => - status === 'redacting' - ? lt(workflowExecutionLogs.startedAt, staleThreshold) - : or( - lt(workflowExecutionLogs.executionDeadlineAt, staleDeadlineThreshold), - and( - isNull(workflowExecutionLogs.executionDeadlineAt), - lt(workflowExecutionLogs.startedAt, staleThreshold) - ) - ) - const cleanupTimestamp = sql.param(now, workflowExecutionLogs.startedAt) - const staleDurationMinutes = sql`ROUND( - EXTRACT(EPOCH FROM (${cleanupTimestamp} - ${workflowExecutionLogs.startedAt})) / 60 - )::integer` - const totalDurationMs = elapsedDurationMsSql(now) - const staleDurationError = sql`${'Execution terminated: worker timeout or crash after '}::text - || ${staleDurationMinutes}::text - || ' minutes'` - /** - * A `redacting` row is never swept by the deadline rule, so the deadline - * message cannot apply to it. - */ - const staleExecutionError = (status: StaleSweepableExecutionStatus) => - status === 'redacting' - ? staleDurationError - : sql`CASE - WHEN ${workflowExecutionLogs.executionDeadlineAt} IS NOT NULL - THEN ${EXECUTION_DEADLINE_ERROR}::text - ELSE ${staleDurationError} - END` - /** - * Swept one status at a time so each pass stays on its own partial index; - * `status IN (...)` would match neither. The row budget is shared across - * both passes so the per-run cap keeps meaning what its name says. - */ - let workflowRowsConsidered = 0 - for (const executionStatus of STALE_SWEEPABLE_EXECUTION_STATUSES) { - const staleExecutionPredicate = and( - eq(workflowExecutionLogs.status, executionStatus), - staleExecutionTimePredicate(executionStatus) - ) - while (workflowRowsConsidered < WORKFLOW_EXECUTION_MAX_ROWS_PER_RUN) { - const limit = Math.min( - WORKFLOW_EXECUTION_MUTATION_BATCH_SIZE, - WORKFLOW_EXECUTION_MAX_ROWS_PER_RUN - workflowRowsConsidered + if (authError) return authError + + const queue = await getJobQueue() + const scheduleWindow = Math.floor(Date.now() / STALE_EXECUTION_CLEANUP_INTERVAL_MS) + const jobId = await queue.enqueue( + 'cleanup-stale-executions', + {}, + { + maxAttempts: 1, + jobId: `cleanup-stale-executions:${scheduleWindow}`, + name: 'Stale execution cleanup', + concurrencyKey: 'cleanup:stale-executions', + concurrencyLimit: 1, + runner: async () => { + const { runCleanupStaleExecutions } = await import( + '@/background/cleanup-stale-executions' ) - currentWorkflowBatchSize = 0 - const { candidates, updatedExecutions } = await db.transaction(async (tx) => { - const candidates = await tx - .select({ id: workflowExecutionLogs.id }) - .from(workflowExecutionLogs) - .where(staleExecutionPredicate) - .limit(limit) - .for('update', { skipLocked: true }) - currentWorkflowBatchSize = candidates.length - if (candidates.length === 0) return { candidates, updatedExecutions: [] } - - const updatedExecutions = await tx - .update(workflowExecutionLogs) - .set({ - status: 'failed', - endedAt: now, - executionDeadlineAt: null, - totalDurationMs, - executionData: sql`jsonb_set( - COALESCE(execution_data, '{}'::jsonb), - ARRAY['error'], - to_jsonb(${staleExecutionError(executionStatus)}) - )`, - }) - .where( - and( - staleExecutionPredicate, - inArray( - workflowExecutionLogs.id, - candidates.map(({ id }) => id) - ) - ) - ) - .returning({ id: workflowExecutionLogs.id }) - - return { candidates, updatedExecutions } - }) - currentWorkflowBatchSize = 0 - staleExecutionsFound += candidates.length - if (candidates.length === 0) break - - cleaned += updatedExecutions.length - workflowRowsConsidered += candidates.length - if (candidates.length < limit) break - } - } - - if (workflowRowsConsidered >= WORKFLOW_EXECUTION_MAX_ROWS_PER_RUN) { - logger.info('Deferred remaining stale workflow executions after reaching the per-run cap', { - maxRowsPerRun: WORKFLOW_EXECUTION_MAX_ROWS_PER_RUN, - }) - } - } catch (error) { - logger.error('Failed to clean up stale workflow executions:', { - error: toError(error).message, - }) - staleExecutionsFound += currentWorkflowBatchSize - failed += currentWorkflowBatchSize - } - - logger.info(`Stale execution cleanup completed. Cleaned: ${cleaned}, Failed: ${failed}`) - - // Clean up stale async jobs (stuck in processing) - let asyncJobsMarkedFailed = 0 - - try { - const hasPositiveMaxDuration = sql`CASE - WHEN jsonb_typeof(${asyncJobs.metadata}->'maxDurationSeconds') = 'number' - THEN (${asyncJobs.metadata}->>'maxDurationSeconds')::numeric >= ${MIN_JOB_DURATION_SECONDS} - AND (${asyncJobs.metadata}->>'maxDurationSeconds')::numeric <= ${MAX_JOB_DURATION_SECONDS} - AND trunc((${asyncJobs.metadata}->>'maxDurationSeconds')::numeric) - = (${asyncJobs.metadata}->>'maxDurationSeconds')::numeric - ELSE FALSE - END` - // A bare `Date` in a raw template reaches the driver unserialized; `sql.param` - // binds it through the column encoder, as `lt(column, date)` does elsewhere here. - const staleProcessingCutoff = sql.param(now, asyncJobs.startedAt) - const staleProcessingFallbackCutoff = sql.param(staleThreshold, asyncJobs.startedAt) - const staleProcessingDurationPredicate = sql`CASE - WHEN ${hasPositiveMaxDuration} - THEN ${asyncJobs.startedAt} + ((${asyncJobs.metadata}->>'maxDurationSeconds')::double precision * interval '1 second') < ${staleProcessingCutoff} - ELSE ${asyncJobs.startedAt} < ${staleProcessingFallbackCutoff} - END` - const staleProcessingPredicate = and( - eq(asyncJobs.status, JOB_STATUS.PROCESSING), - ne(asyncJobs.type, SCHEDULE_EXECUTION_QUEUE_NAME), - staleProcessingDurationPredicate - ) - const staleProcessingResult = await runBatchedMutation({ - batchSize: STATE_MUTATION_BATCH_SIZE, - maxRowsPerRun: STATE_MUTATION_MAX_ROWS_PER_RUN, - claim: (tx, limit) => - tx - .select({ id: asyncJobs.id }) - .from(asyncJobs) - .where(staleProcessingPredicate) - .limit(limit) - .for('update', { skipLocked: true }), - mutation: (tx, candidateIds) => - tx - .update(asyncJobs) - .set({ - status: JOB_STATUS.FAILED, - completedAt: new Date(), - error: sql`CASE - WHEN ${hasPositiveMaxDuration} - THEN 'Job terminated: stuck in processing for more than ' - || (${asyncJobs.metadata}->>'maxDurationSeconds') - || ' seconds (worker cleanup deadline)' - ELSE ${GENERIC_STALE_PROCESSING_ERROR} - END`, - updatedAt: new Date(), - }) - .where(and(staleProcessingPredicate, inArray(asyncJobs.id, candidateIds))) - .returning({ id: asyncJobs.id }), - }) - - asyncJobsMarkedFailed = staleProcessingResult.affected - if (asyncJobsMarkedFailed > 0) { - logger.info(`Marked ${asyncJobsMarkedFailed} stale async jobs as failed`) - } - if (staleProcessingResult.reachedLimit) { - logger.info('Deferred remaining stale async jobs after reaching the per-run cap', { - maxRowsPerRun: STATE_MUTATION_MAX_ROWS_PER_RUN, - }) - } - } catch (error) { - logger.error('Failed to clean up stale async jobs:', { - error: toError(error).message, - }) - } - - // Mark stale table jobs (import, export, or delete) as failed. Jobs run detached on the web container - // and are lost if the pod is killed mid-run. `updated_at` is bumped by progress updates, so a - // `running` job with no recent update has stalled (not merely slow). Committed work is left in - // place (no rollback); the user retries. Also prune long-settled terminal jobs so the table - // doesn't grow unbounded (the latest job per table is what list/detail reads surface). - let staleTableJobsMarkedFailed = 0 - try { - const staleTableJobPredicate = and( - eq(tableJobs.status, 'running'), - lt(tableJobs.updatedAt, staleTableJobThreshold) - ) - const staleTableJobResult = await runBatchedMutation({ - batchSize: STATE_MUTATION_BATCH_SIZE, - maxRowsPerRun: STATE_MUTATION_MAX_ROWS_PER_RUN, - claim: (tx, limit) => - tx - .select({ id: tableJobs.id }) - .from(tableJobs) - .where(staleTableJobPredicate) - .limit(limit) - .for('update', { skipLocked: true }), - mutation: (tx, candidateIds) => - tx - .update(tableJobs) - .set({ - status: 'failed', - error: `Job terminated: no progress for more than ${TABLE_JOB_STALE_THRESHOLD_MINUTES} minutes (worker timeout or crash)`, - completedAt: now, - updatedAt: now, - }) - .where(and(staleTableJobPredicate, inArray(tableJobs.id, candidateIds))) - .returning({ id: tableJobs.id }), - }) - - staleTableJobsMarkedFailed = staleTableJobResult.affected - if (staleTableJobsMarkedFailed > 0) { - logger.info(`Marked ${staleTableJobsMarkedFailed} stale table jobs as failed`) - } - - const terminalRetention = new Date(Date.now() - TABLE_JOB_RETENTION_HOURS * 60 * 60 * 1000) - const terminalTableJobPredicate = and( - inArray(tableJobs.status, ['ready', 'failed', 'canceled']), - lt(tableJobs.updatedAt, terminalRetention) - ) - const terminalTableJobResult = await runBatchedMutation({ - batchSize: TABLE_JOB_PRUNE_BATCH_SIZE, - maxRowsPerRun: TABLE_JOB_PRUNE_MAX_ROWS_PER_RUN, - claim: (tx, limit) => - tx - .select({ id: tableJobs.id }) - .from(tableJobs) - .where(terminalTableJobPredicate) - .limit(limit) - .for('update', { skipLocked: true }), - mutation: (tx, candidateIds) => - tx - .delete(tableJobs) - .where(and(terminalTableJobPredicate, inArray(tableJobs.id, candidateIds))) - .returning({ - type: tableJobs.type, - resultKey: sql`${tableJobs.payload}->>'resultKey'`, - }), - onBatch: async (jobs) => { - /** - * Pruned export jobs carry the generated file's storage key. The scalar - * key is returned instead of the full JSON payload, and cleanup stays - * sequential so storage concurrency is bounded at one request. - */ - for (const { type, resultKey } of jobs) { - if (type !== 'export' || !resultKey) continue - await deleteFile({ key: resultKey, context: 'workspace' }).catch((err) => { - logger.warn('Failed to delete pruned export file', { - resultKey, - error: toError(err).message, - }) - }) - } + return runCleanupStaleExecutions() }, - }) - if (terminalTableJobResult.reachedLimit) { - logger.info('Deferred remaining terminal table jobs after reaching the per-run cap', { - maxRowsPerRun: TABLE_JOB_PRUNE_MAX_ROWS_PER_RUN, - }) - } - } catch (error) { - logger.error('Failed to clean up stale table jobs:', { - error: toError(error).message, - }) - } - - // Clean up stale pending jobs (never started, e.g., due to server crash before startJob()) - let stalePendingJobsMarkedFailed = 0 - - try { - const stalePendingPredicate = and( - eq(asyncJobs.status, JOB_STATUS.PENDING), - ne(asyncJobs.type, SCHEDULE_EXECUTION_QUEUE_NAME), - lt(asyncJobs.createdAt, stalePendingThreshold) - ) - const stalePendingResult = await runBatchedMutation({ - batchSize: STATE_MUTATION_BATCH_SIZE, - maxRowsPerRun: STATE_MUTATION_MAX_ROWS_PER_RUN, - claim: (tx, limit) => - tx - .select({ id: asyncJobs.id }) - .from(asyncJobs) - .where(stalePendingPredicate) - .limit(limit) - .for('update', { skipLocked: true }), - mutation: (tx, candidateIds) => - tx - .update(asyncJobs) - .set({ - status: JOB_STATUS.FAILED, - completedAt: new Date(), - error: `Job terminated: stuck in pending state for more than ${JOB_PENDING_RETENTION_HOURS} hours (never started)`, - updatedAt: new Date(), - }) - .where(and(stalePendingPredicate, inArray(asyncJobs.id, candidateIds))) - .returning({ id: asyncJobs.id }), - }) - - stalePendingJobsMarkedFailed = stalePendingResult.affected - if (stalePendingJobsMarkedFailed > 0) { - logger.info(`Marked ${stalePendingJobsMarkedFailed} stale pending jobs as failed`) } - if (stalePendingResult.reachedLimit) { - logger.info('Deferred remaining stale pending jobs after reaching the per-run cap', { - maxRowsPerRun: STATE_MUTATION_MAX_ROWS_PER_RUN, - }) - } - } catch (error) { - logger.error('Failed to clean up stale pending jobs:', { - error: toError(error).message, - }) - } - - const retentionNow = Date.now() - const retentionThreshold = new Date(retentionNow - JOB_RETENTION_HOURS * 60 * 60 * 1000) - const irrecoverableCarrierRetentionThreshold = new Date( - retentionNow - SCHEDULE_CARRIER_IRRECOVERABLE_RETENTION_HOURS * 60 * 60 * 1000 ) - let asyncJobsDeleted = 0 - try { - const retainedJobPredicate = and( - inArray(asyncJobs.status, TERMINAL_JOB_STATUSES), - or( - ne(asyncJobs.type, SCHEDULE_EXECUTION_QUEUE_NAME), - /** - * Schedule recovery owns a carrier until it stamps the reconciled - * marker, so retention waits for it rather than deleting an - * occurrence that has not been accounted for yet. - */ - and( - carrierReconciledSql(asyncJobs.metadata), - or( - carrierNotIrrecoverableSql(asyncJobs.metadata), - lt(asyncJobs.completedAt, irrecoverableCarrierRetentionThreshold) - ) - ) - ), - lt(asyncJobs.completedAt, retentionThreshold) - ) - const retainedJobResult = await runBatchedMutation({ - batchSize: RETENTION_DELETE_BATCH_SIZE, - maxRowsPerRun: RETENTION_DELETE_MAX_ROWS_PER_RUN, - claim: (tx, limit) => - tx - .select({ id: asyncJobs.id }) - .from(asyncJobs) - .where(retainedJobPredicate) - .limit(limit) - .for('update', { skipLocked: true }), - mutation: (tx, candidateIds) => - tx - .delete(asyncJobs) - .where(and(retainedJobPredicate, inArray(asyncJobs.id, candidateIds))) - .returning({ id: asyncJobs.id }), - }) - - asyncJobsDeleted = retainedJobResult.affected - if (asyncJobsDeleted > 0) { - logger.info( - `Deleted ${asyncJobsDeleted} old async jobs (retention: ${JOB_RETENTION_HOURS}h)` - ) - } - if (retainedJobResult.reachedLimit) { - logger.info('Deferred remaining retained async jobs after reaching the per-run cap', { - maxRowsPerRun: RETENTION_DELETE_MAX_ROWS_PER_RUN, - }) - } - } catch (error) { - logger.error('Failed to delete old async jobs:', { - error: toError(error).message, - }) - } - - /** - * Prune terminal connector sync logs past retention. - * - * HARD INVARIANT: the newest row per connector must survive, and so must the - * newest `completed` row. `loadPreviousListingObservation` reconstructs the - * previous listing from the latest `completed` log, and that reconstruction - * decides whether a suspect listing is corroborated — i.e. whether - * reconciliation may delete documents. Pruning the last `completed` row - * would silently change deletion behaviour, so both `exists` guards below - * are load-bearing rather than defensive. - * - * `started` rows are never eligible: they are either in flight or waiting on - * the scheduler's own sweep to close them. - */ - let connectorSyncLogsPruned = 0 - try { - const syncLogRetention = new Date( - Date.now() - CONNECTOR_SYNC_LOG_RETENTION_DAYS * 24 * 60 * 60 * 1000 - ) - const newerSyncLog = alias(knowledgeConnectorSyncLog, 'newer_sync_log') - const newerCompletedSyncLog = alias(knowledgeConnectorSyncLog, 'newer_completed_sync_log') - const syncLogPredicate = and( - inArray(knowledgeConnectorSyncLog.status, ['completed', 'failed']), - lt(knowledgeConnectorSyncLog.startedAt, syncLogRetention), - exists( - db - .select({ id: newerSyncLog.id }) - .from(newerSyncLog) - .where( - and( - eq(newerSyncLog.connectorId, knowledgeConnectorSyncLog.connectorId), - gt(newerSyncLog.startedAt, knowledgeConnectorSyncLog.startedAt) - ) - ) - ), - or( - ne(knowledgeConnectorSyncLog.status, 'completed'), - exists( - db - .select({ id: newerCompletedSyncLog.id }) - .from(newerCompletedSyncLog) - .where( - and( - eq(newerCompletedSyncLog.connectorId, knowledgeConnectorSyncLog.connectorId), - eq(newerCompletedSyncLog.status, 'completed'), - gt(newerCompletedSyncLog.startedAt, knowledgeConnectorSyncLog.startedAt) - ) - ) - ) - ) - ) - const syncLogResult = await runBatchedMutation({ - batchSize: CONNECTOR_SYNC_LOG_PRUNE_BATCH_SIZE, - maxRowsPerRun: CONNECTOR_SYNC_LOG_MAX_ROWS_PER_RUN, - claim: (tx, limit) => - tx - .select({ id: knowledgeConnectorSyncLog.id }) - .from(knowledgeConnectorSyncLog) - .where(syncLogPredicate) - .limit(limit) - .for('update', { skipLocked: true }), - mutation: (tx, candidateIds) => - tx - .delete(knowledgeConnectorSyncLog) - .where(inArray(knowledgeConnectorSyncLog.id, candidateIds)) - .returning({ id: knowledgeConnectorSyncLog.id }), - }) - connectorSyncLogsPruned = syncLogResult.affected - if (connectorSyncLogsPruned > 0) { - logger.info( - `Pruned ${connectorSyncLogsPruned} old connector sync logs (retention: ${CONNECTOR_SYNC_LOG_RETENTION_DAYS}d)` - ) - } - if (syncLogResult.reachedLimit) { - logger.info('Deferred remaining connector sync logs after reaching the per-run cap', { - maxRowsPerRun: CONNECTOR_SYNC_LOG_MAX_ROWS_PER_RUN, - }) - } - } catch (error) { - logger.error('Failed to prune old connector sync logs:', { - error: toError(error).message, - }) - } - - /** - * Prune terminal deployment operations past retention. HARD INVARIANT: - * the newest-generation row per workflow must always survive — the next - * deploy computes `generation = MAX(generation) + 1`, and the webhook - * registration store fences rows with lt/gt comparisons against stored - * generations, so generation reuse after a full wipe would permanently - * wedge that workflow's deployments. The `exists(newer)` predicate - * guarantees the max-generation row is never eligible; the status filter - * keeps in-flight rows (an outbox worker may still hold their fence). - */ - let deploymentOperationsPruned = 0 - try { - const deploymentOpRetention = new Date( - Date.now() - DEPLOYMENT_OPERATION_RETENTION_DAYS * 24 * 60 * 60 * 1000 - ) - const newerOperation = alias(workflowDeploymentOperation, 'newer_operation') - const deploymentOpPredicate = and( - inArray(workflowDeploymentOperation.status, ['active', 'failed', 'superseded']), - lt(workflowDeploymentOperation.completedAt, deploymentOpRetention), - exists( - db - .select({ id: newerOperation.id }) - .from(newerOperation) - .where( - and( - eq(newerOperation.workflowId, workflowDeploymentOperation.workflowId), - gt(newerOperation.generation, workflowDeploymentOperation.generation) - ) - ) - ) - ) - const deploymentOpResult = await runBatchedMutation({ - batchSize: DEPLOYMENT_OPERATION_PRUNE_BATCH_SIZE, - maxRowsPerRun: - DEPLOYMENT_OPERATION_PRUNE_BATCH_SIZE * DEPLOYMENT_OPERATION_PRUNE_MAX_BATCHES, - claim: (tx, limit) => - tx - .select({ id: workflowDeploymentOperation.id }) - .from(workflowDeploymentOperation) - .where(deploymentOpPredicate) - .limit(limit) - .for('update', { skipLocked: true }), - mutation: (tx, candidateIds) => - tx - .delete(workflowDeploymentOperation) - .where(inArray(workflowDeploymentOperation.id, candidateIds)) - .returning({ id: workflowDeploymentOperation.id }), - }) - deploymentOperationsPruned = deploymentOpResult.affected - if (deploymentOperationsPruned > 0) { - logger.info( - `Pruned ${deploymentOperationsPruned} old deployment operations (retention: ${DEPLOYMENT_OPERATION_RETENTION_DAYS}d)` - ) - } - if (deploymentOpResult.reachedLimit) { - logger.info('Deferred remaining deployment operations after reaching the per-run cap', { - maxRowsPerRun: - DEPLOYMENT_OPERATION_PRUNE_BATCH_SIZE * DEPLOYMENT_OPERATION_PRUNE_MAX_BATCHES, - }) - } - } catch (error) { - logger.error('Failed to prune old deployment operations:', { - error: toError(error).message, - }) - } - - /** - * Cancel table run dispatches abandoned by a dead dispatcher. Nothing else - * reclaims them — every other terminal transition is user- or flow-initiated - * — so a dispatcher killed mid-loop left the row `dispatching` forever and - * the client's "X running" overlay with it. Ages from the dispatcher's - * per-window heartbeat, so a slow-but-live dispatch is spared. - */ - let staleDispatchesCancelled = 0 - try { - staleDispatchesCancelled = ( - await cancelStaleDispatches(staleDispatchThreshold, TABLE_DISPATCH_MAX_PER_RUN) - ).length - if (staleDispatchesCancelled > 0) { - logger.warn(`Cancelled ${staleDispatchesCancelled} abandoned table run dispatches`, { - thresholdMinutes: TABLE_DISPATCH_STALE_THRESHOLD_MINUTES, - }) - } - } catch (error) { - logger.error('Failed to cancel abandoned table run dispatches:', { - error: toError(error).message, - }) - } - - /** - * Settle Chat runs no controller will finish: their process died, their - * controller was superseded without a successor, or Stop found none. Without - * this they stay unfinished forever and keep their chat marked as busy. - */ - let orphanedRunsSettled = 0 - try { - orphanedRunsSettled = (await sweepOrphanedRuns()).settledRunIds.length - } catch (error) { - logger.error('Failed to settle orphaned Chat runs:', { - error: toError(error).message, - }) - } - - return NextResponse.json({ - success: true, - executions: { - found: staleExecutionsFound, - cleaned, - failed, - thresholdMinutes: STALE_THRESHOLD_MINUTES, - }, - asyncJobs: { - staleProcessingMarkedFailed: asyncJobsMarkedFailed, - stalePendingMarkedFailed: stalePendingJobsMarkedFailed, - oldDeleted: asyncJobsDeleted, - staleThresholdMinutes: STALE_THRESHOLD_MINUTES, - retentionHours: JOB_RETENTION_HOURS, - }, - tableJobs: { - staleMarkedFailed: staleTableJobsMarkedFailed, - }, - connectorSyncLogs: { - pruned: connectorSyncLogsPruned, - retentionDays: CONNECTOR_SYNC_LOG_RETENTION_DAYS, - }, - tableRunDispatches: { - staleCancelled: staleDispatchesCancelled, - thresholdMinutes: TABLE_DISPATCH_STALE_THRESHOLD_MINUTES, - }, - deploymentOperations: { - pruned: deploymentOperationsPruned, - retentionDays: DEPLOYMENT_OPERATION_RETENTION_DAYS, - }, - chatRuns: { - orphanedSettled: orphanedRunsSettled, - }, - }) + logger.info('Stale execution cleanup dispatched', { jobId }) + return NextResponse.json({ triggered: true, jobId }) } catch (error) { - logger.error('Error in stale execution cleanup job:', error) - return NextResponse.json({ error: 'Internal server error' }, { status: 500 }) + logger.error('Failed to dispatch stale execution cleanup', { error }) + return NextResponse.json( + { error: 'Failed to dispatch stale execution cleanup' }, + { status: 500 } + ) } }) diff --git a/apps/sim/background/cleanup-stale-executions.test.ts b/apps/sim/background/cleanup-stale-executions.test.ts new file mode 100644 index 00000000000..6f552046988 --- /dev/null +++ b/apps/sim/background/cleanup-stale-executions.test.ts @@ -0,0 +1,451 @@ +import { asyncJobs, tableJobs, workflowExecutionLogs } from '@sim/db/schema' +import { createLogger } from '@sim/logger' +import { dbChainMockFns, queueTableRows, resetDbChainMock } from '@sim/testing' +import { storageServiceMock, storageServiceMockFns } from '@sim/testing/mocks/storage-service.mock' +import { beforeEach, describe, expect, it, vi } from 'vitest' +import { JOB_RETENTION_HOURS } from '@/lib/core/async-jobs' +import { + SCHEDULE_CARRIER_IRRECOVERABLE_METADATA_KEY, + SCHEDULE_CARRIER_IRRECOVERABLE_RETENTION_HOURS, + SCHEDULE_CARRIER_RECONCILED_METADATA_KEY, +} from '@/lib/workflows/schedules/carrier-metadata' + +vi.mock('@/lib/uploads/core/storage-service', () => storageServiceMock) + +import { runCleanupStaleExecutions } from '@/background/cleanup-stale-executions' + +const mockDeleteFile = storageServiceMockFns.mockDeleteFile +mockDeleteFile.mockResolvedValue(undefined) + +const cleanupLogger = + vi.mocked(createLogger).mock.results[ + vi.mocked(createLogger).mock.calls.findIndex(([name]) => name === 'CleanupStaleExecutions') + ].value + +interface MockCondition { + type?: string + conditions?: unknown[] + left?: unknown + right?: unknown + column?: unknown + values?: unknown + toSQL?: () => { sql: string; params: unknown[] } +} + +function flattenConditions(condition: unknown): MockCondition[] { + if (!condition || typeof condition !== 'object') return [] + const node = condition as MockCondition + return [node, ...(node.conditions?.flatMap((child) => flattenConditions(child)) ?? [])] +} + +function hasToSQL(value: unknown): value is { toSQL: () => { sql: string; params: unknown[] } } { + return typeof value === 'object' && value !== null && 'toSQL' in value +} + +/** + * Collects the leaves of a nested `sql` expression. The duration expression is + * built by a shared helper, so the values it binds sit one level below the + * fragment this route assembles rather than directly in its own params. + */ +function flattenSqlParams(expression: { sql: string; params: unknown[] }): unknown[] { + return expression.params.flatMap((param) => + hasToSQL(param) ? flattenSqlParams(param.toSQL()) : [param] + ) +} + +/** Recursively renders a mocked drizzle fragment, expanding nested fragments. */ +function renderSql(fragment: unknown): string { + if (fragment === null || fragment === undefined) return '' + if (typeof fragment !== 'object') return String(fragment) + const candidate = fragment as { rawSql?: string; strings?: string[]; values?: unknown[] } + if (typeof candidate.rawSql === 'string') return candidate.rawSql + if (!candidate.strings) return '' + return candidate.strings + .map((part, index) => + index < (candidate.values?.length ?? 0) + ? `${part}${renderSql(candidate.values?.[index])}` + : part + ) + .join('') +} + +/** Recursively collects the non-fragment bind values of a mocked fragment. */ +function collectSqlParams(fragment: unknown): unknown[] { + if (!fragment || typeof fragment !== 'object') return [] + const candidate = fragment as { values?: unknown[] } + if (!candidate.values) return [] + return candidate.values.flatMap((value) => { + const nested = collectSqlParams(value) + return nested.length > 0 ? nested : [value] + }) +} + +describe('stale execution cleanup deadline grace', () => { + beforeEach(() => { + resetDbChainMock() + cleanupLogger.info.mockReset() + }) + + it('waits five minutes past a workflow execution deadline in both cleanup predicates', async () => { + vi.useFakeTimers() + vi.setSystemTime(new Date('2026-08-03T12:10:00.000Z')) + queueTableRows(workflowExecutionLogs, [{ id: 'log-1' }]) + dbChainMockFns.returning.mockResolvedValueOnce([{ id: 'log-1' }]) + + try { + await runCleanupStaleExecutions() + const expectedThreshold = new Date('2026-08-03T12:05:00.000Z') + const deadlineComparisons = dbChainMockFns.where.mock.calls + .flatMap(([condition]) => flattenConditions(condition)) + .filter( + (condition) => + condition.type === 'lt' && + condition.right instanceof Date && + condition.right.getTime() === expectedThreshold.getTime() + ) + + expect(deadlineComparisons).toHaveLength(2) + expect(deadlineComparisons.map(({ right }) => right)).toEqual([ + expectedThreshold, + expectedThreshold, + ]) + + const executionUpdateIndex = dbChainMockFns.update.mock.calls.findIndex( + ([table]) => table === workflowExecutionLogs + ) + const update = dbChainMockFns.set.mock.calls[executionUpdateIndex]?.[0] as { + endedAt: Date + totalDurationMs: { toSQL: () => { sql: string; params: unknown[] } } + executionData: { toSQL: () => { sql: string; params: unknown[] } } + } + const totalDurationLeaves = flattenSqlParams(update.totalDurationMs.toSQL()) + const renderedError = renderSql(update.executionData) + const errorLeaves = collectSqlParams(update.executionData) + + expect(renderedError).toContain('CASE') + expect(renderedError).toContain('IS NOT NULL') + expect(renderedError).toContain('ROUND') + expect(renderedError).toContain('EXTRACT(EPOCH') + expect(errorLeaves).toContain(workflowExecutionLogs.executionDeadlineAt) + expect(errorLeaves).toContain('Execution timed out') + expect(errorLeaves).toContain('Execution terminated: worker timeout or crash after ') + expect(errorLeaves).toContain(workflowExecutionLogs.startedAt) + expect(totalDurationLeaves).toContain(2_147_483_647) + expect(totalDurationLeaves).toContain(workflowExecutionLogs.startedAt) + expect(totalDurationLeaves).toContainEqual(new Date('2026-08-03T12:10:00.000Z')) + expect(update.endedAt).toEqual(new Date('2026-08-03T12:10:00.000Z')) + } finally { + vi.useRealTimers() + } + }) + + it('sweeps redacting logs on the generic window, never the execution deadline', async () => { + await runCleanupStaleExecutions() + const redactingPredicates = dbChainMockFns.where.mock.calls + .map(([condition]) => flattenConditions(condition)) + .filter((conditions) => + conditions.some( + (condition) => + condition.type === 'eq' && + condition.left === workflowExecutionLogs.status && + condition.right === 'redacting' + ) + ) + + expect(redactingPredicates.length).toBeGreaterThan(0) + for (const conditions of redactingPredicates) { + expect( + conditions.some((condition) => condition.left === workflowExecutionLogs.executionDeadlineAt) + ).toBe(false) + expect( + conditions.some( + (condition) => + condition.type === 'lt' && condition.left === workflowExecutionLogs.startedAt + ) + ).toBe(true) + } + }) + + it('leaves pending and processing schedule jobs to schedule recovery', async () => { + await runCleanupStaleExecutions() + const activeAsyncPredicates = dbChainMockFns.where.mock.calls + .map(([condition]) => flattenConditions(condition)) + .filter((conditions) => + conditions.some( + (condition) => + condition.type === 'ne' && + condition.left === asyncJobs.type && + condition.right === 'schedule-execution' + ) + ) + + expect( + activeAsyncPredicates.some((conditions) => + conditions.some( + (condition) => + condition.type === 'eq' && + condition.left === asyncJobs.status && + condition.right === 'processing' + ) + ) + ).toBe(true) + expect( + activeAsyncPredicates.some((conditions) => + conditions.some( + (condition) => + condition.type === 'eq' && + condition.left === asyncJobs.status && + condition.right === 'pending' + ) + ) + ).toBe(true) + }) + + it('retains terminal schedule carriers until reconciliation is recorded', async () => { + await runCleanupStaleExecutions() + const retentionConditions = dbChainMockFns.where.mock.calls.flatMap(([condition]) => + flattenConditions(condition) + ) + const reconciliationMarker = retentionConditions.find((condition) => + renderSql(condition).includes(SCHEDULE_CARRIER_RECONCILED_METADATA_KEY) + ) + + expect(collectSqlParams(reconciliationMarker)).toContain(asyncJobs.metadata) + expect( + retentionConditions.some( + (condition) => + condition.type === 'ne' && + condition.left === asyncJobs.type && + condition.right === 'schedule-execution' + ) + ).toBe(true) + }) + + it('spells carrier metadata keys as SQL literals so the partial index matches', async () => { + await runCleanupStaleExecutions() + const reconciliationMarker = dbChainMockFns.where.mock.calls + .flatMap(([condition]) => flattenConditions(condition)) + .find((condition) => renderSql(condition).includes(SCHEDULE_CARRIER_RECONCILED_METADATA_KEY)) + + expect(renderSql(reconciliationMarker)).toContain( + `'${SCHEDULE_CARRIER_RECONCILED_METADATA_KEY}'` + ) + expect(collectSqlParams(reconciliationMarker)).not.toContain( + SCHEDULE_CARRIER_RECONCILED_METADATA_KEY + ) + }) + + it('deletes irrecoverable schedule carrier tombstones once their longer window lapses', async () => { + await runCleanupStaleExecutions() + const retentionConditions = dbChainMockFns.where.mock.calls.flatMap(([condition]) => + flattenConditions(condition) + ) + const irrecoverableExclusion = retentionConditions.find((condition) => + renderSql(condition).includes(SCHEDULE_CARRIER_IRRECOVERABLE_METADATA_KEY) + ) + + expect(renderSql(irrecoverableExclusion)).toContain("<> 'true'") + expect(collectSqlParams(irrecoverableExclusion)).toContain(asyncJobs.metadata) + + const tombstoneWindow = retentionConditions.filter( + (condition) => + condition.type === 'lt' && + condition.left === asyncJobs.completedAt && + condition.right instanceof Date + ) + const oldest = Math.min(...tombstoneWindow.map(({ right }) => (right as Date).getTime())) + const newest = Math.max(...tombstoneWindow.map(({ right }) => (right as Date).getTime())) + expect(newest - oldest).toBe( + (SCHEDULE_CARRIER_IRRECOVERABLE_RETENTION_HOURS - JOB_RETENTION_HOURS) * 60 * 60 * 1000 + ) + }) + + it('keeps table-job heartbeat cleanup independent from workflow timeout policy', async () => { + vi.useFakeTimers() + vi.setSystemTime(new Date('2026-08-03T12:00:00.000Z')) + queueTableRows(tableJobs, [{ id: 'table-job-1' }]) + dbChainMockFns.returning.mockResolvedValueOnce([{ id: 'table-job-1' }]) + + try { + await runCleanupStaleExecutions() + const expectedThreshold = new Date('2026-08-03T10:25:00.000Z') + const tableJobComparisons = dbChainMockFns.where.mock.calls + .flatMap(([condition]) => flattenConditions(condition)) + .filter( + (condition) => + condition.type === 'lt' && + condition.left === tableJobs.updatedAt && + condition.right instanceof Date && + condition.right.getTime() === expectedThreshold.getTime() + ) + + expect(tableJobComparisons).toHaveLength(2) + expect(tableJobComparisons.map(({ right }) => right)).toEqual([ + expectedThreshold, + expectedThreshold, + ]) + + const tableJobUpdateIndex = dbChainMockFns.update.mock.calls.findIndex( + ([table]) => table === tableJobs + ) + const update = dbChainMockFns.set.mock.calls[tableJobUpdateIndex]?.[0] as { + error: string + } + expect(update.error).toBe( + 'Job terminated: no progress for more than 95 minutes (worker timeout or crash)' + ) + } finally { + vi.useRealTimers() + } + }) + + it('claims every cleanup page without overlapping concurrent workers', async () => { + await runCleanupStaleExecutions() + // Nine batched arms: the connector sync-log retention pass is the newest. + expect(dbChainMockFns.transaction).toHaveBeenCalledTimes(9) + expect(dbChainMockFns.for).toHaveBeenCalledTimes(9) + for (const [strength, options] of dbChainMockFns.for.mock.calls) { + expect(strength).toBe('update') + expect(options).toEqual({ skipLocked: true }) + } + }) + + it('caps every bulk mutation and returns only scalar export cleanup fields', async () => { + const stateBatch = Array.from({ length: 1000 }, (_, index) => ({ id: `state-${index}` })) + const retentionBatch = Array.from({ length: 2000 }, (_, index) => ({ + id: `retention-${index}`, + })) + const exportBatch = Array.from({ length: 100 }, (_, index) => ({ + type: 'export', + resultKey: `workspace/workspace-1/exports/table-1/job-${index}/export.csv`, + })) + const exportCandidates = Array.from({ length: 100 }, (_, index) => ({ + id: `export-${index}`, + })) + + for (let batch = 0; batch < 10; batch++) { + const workflowBatch = Array.from({ length: 100 }, (_, index) => ({ + id: `workflow-state-${batch}-${index}`, + })) + queueTableRows(workflowExecutionLogs, workflowBatch) + dbChainMockFns.returning.mockResolvedValueOnce(workflowBatch) + } + for (let batch = 0; batch < 10; batch++) { + queueTableRows(asyncJobs, stateBatch) + dbChainMockFns.returning.mockResolvedValueOnce(stateBatch) + } + for (let batch = 0; batch < 10; batch++) { + queueTableRows(tableJobs, stateBatch) + dbChainMockFns.returning.mockResolvedValueOnce(stateBatch) + } + for (let batch = 0; batch < 10; batch++) { + queueTableRows(tableJobs, exportCandidates) + dbChainMockFns.returning.mockResolvedValueOnce(exportBatch) + } + for (let batch = 0; batch < 10; batch++) { + queueTableRows(asyncJobs, stateBatch) + dbChainMockFns.returning.mockResolvedValueOnce(stateBatch) + } + for (let batch = 0; batch < 10; batch++) { + queueTableRows(asyncJobs, retentionBatch) + dbChainMockFns.returning.mockResolvedValueOnce(retentionBatch) + } + dbChainMockFns.returning.mockResolvedValueOnce([]) + + const result = await runCleanupStaleExecutions() + + expect(result).toMatchObject({ + executions: { + found: 1000, + cleaned: 1000, + failed: 0, + }, + asyncJobs: { + staleProcessingMarkedFailed: 10_000, + stalePendingMarkedFailed: 10_000, + oldDeleted: 20_000, + }, + tableJobs: { + staleMarkedFailed: 10_000, + }, + }) + expect(mockDeleteFile).toHaveBeenCalledTimes(1000) + + const limits = dbChainMockFns.limit.mock.calls.map(([limit]) => limit) + expect(limits.filter((limit) => limit === 100)).toHaveLength(20) + expect(limits.filter((limit) => limit === 1000)).toHaveLength(30) + expect(limits.filter((limit) => limit === 2000)).toHaveLength(12) + + const workflowUpdates = dbChainMockFns.update.mock.calls.filter( + ([table]) => table === workflowExecutionLogs + ) + expect(workflowUpdates).toHaveLength(10) + + const returningShapes = dbChainMockFns.returning.mock.calls + .map(([shape]) => shape) + .filter((shape): shape is Record => Boolean(shape)) + expect(returningShapes.some((shape) => 'payload' in shape)).toBe(false) + expect(returningShapes.some((shape) => 'type' in shape && 'resultKey' in shape)).toBe(true) + + const claimedIds = dbChainMockFns.where.mock.calls + .flatMap(([condition]) => flattenConditions(condition)) + .filter((condition) => condition.type === 'inArray') + .map((condition) => condition.values) + expect(claimedIds.length).toBeGreaterThan(0) + expect(claimedIds.every((ids) => Array.isArray(ids))).toBe(true) + }) + + it('preserves committed workflow cleanup counts when a later batch fails', async () => { + const firstBatch = Array.from({ length: 100 }, (_, index) => ({ + id: `execution-${index}`, + })) + const failedBatch = Array.from({ length: 37 }, (_, index) => ({ + id: `failed-execution-${index}`, + })) + queueTableRows(workflowExecutionLogs, firstBatch) + queueTableRows(workflowExecutionLogs, failedBatch) + queueTableRows(asyncJobs, [{ id: 'async-job-1' }]) + dbChainMockFns.returning + .mockResolvedValueOnce(firstBatch) + .mockRejectedValueOnce(new Error('database unavailable')) + .mockResolvedValueOnce([{ id: 'async-job-1' }]) + + const result = await runCleanupStaleExecutions() + + expect(result).toMatchObject({ + executions: { + found: 137, + cleaned: 100, + failed: 37, + }, + }) + expect( + dbChainMockFns.update.mock.calls.filter(([table]) => table === workflowExecutionLogs) + ).toHaveLength(2) + expect(dbChainMockFns.update.mock.calls.some(([table]) => table === asyncJobs)).toBe(true) + }) + + it('continues draining when an atomic race updates fewer rows than were selected', async () => { + const firstCandidates = Array.from({ length: 100 }, (_, index) => ({ + id: `execution-${index}`, + })) + const firstUpdated = firstCandidates.slice(0, 99) + const secondBatch = [{ id: 'execution-100' }] + queueTableRows(workflowExecutionLogs, firstCandidates) + queueTableRows(workflowExecutionLogs, secondBatch) + dbChainMockFns.returning.mockResolvedValueOnce(firstUpdated).mockResolvedValueOnce(secondBatch) + + const result = await runCleanupStaleExecutions() + + expect(result).toMatchObject({ + executions: { + found: 101, + cleaned: 100, + failed: 0, + }, + }) + expect( + dbChainMockFns.update.mock.calls.filter(([table]) => table === workflowExecutionLogs) + ).toHaveLength(2) + }) +}) diff --git a/apps/sim/background/cleanup-stale-executions.ts b/apps/sim/background/cleanup-stale-executions.ts new file mode 100644 index 00000000000..3c0d9010475 --- /dev/null +++ b/apps/sim/background/cleanup-stale-executions.ts @@ -0,0 +1,785 @@ +import { db } from '@sim/db' +import { + asyncJobs, + knowledgeConnectorSyncLog, + tableJobs, + workflowDeploymentOperation, + workflowExecutionLogs, +} from '@sim/db/schema' +import { createLogger } from '@sim/logger' +import { toError } from '@sim/utils/errors' +import { task } from '@trigger.dev/sdk' +import { and, eq, exists, gt, inArray, isNull, lt, ne, or, sql } from 'drizzle-orm' +import { alias } from 'drizzle-orm/pg-core' +import { + JOB_PENDING_RETENTION_HOURS, + JOB_RETENTION_HOURS, + JOB_STATUS, + MAX_JOB_DURATION_SECONDS, + MIN_JOB_DURATION_SECONDS, + TERMINAL_JOB_STATUSES, +} from '@/lib/core/async-jobs' +import { + getExecutionReservationTtlMs, + getTimeoutErrorMessage, + RESERVATION_TTL_BUFFER_MS, +} from '@/lib/core/execution-limits' +import type { DbTransaction } from '@/lib/db/types' +import { elapsedDurationMsSql } from '@/lib/logs/execution/duration' +import { + STALE_SWEEPABLE_EXECUTION_STATUSES, + type StaleSweepableExecutionStatus, +} from '@/lib/logs/types' +import { sweepOrphanedRuns } from '@/lib/mothership/async-runs/orphaned-runs' +import { cancelStaleDispatches } from '@/lib/table/dispatcher' +import { deleteFile } from '@/lib/uploads/core/storage-service' +import { + carrierNotIrrecoverableSql, + carrierReconciledSql, + SCHEDULE_CARRIER_IRRECOVERABLE_RETENTION_HOURS, +} from '@/lib/workflows/schedules/carrier-metadata' +import { SCHEDULE_EXECUTION_QUEUE_NAME } from '@/lib/workflows/schedules/execution-limits' + +const logger = createLogger('CleanupStaleExecutions') + +const STALE_THRESHOLD_MS = getExecutionReservationTtlMs() +const STALE_THRESHOLD_MINUTES = Math.ceil(STALE_THRESHOLD_MS / 60000) +const GENERIC_STALE_PROCESSING_ERROR = `Job terminated: stuck in processing for more than ${STALE_THRESHOLD_MINUTES} minutes` +const EXECUTION_DEADLINE_ERROR = getTimeoutErrorMessage(undefined) +/** + * Table jobs run as detached workers with progress heartbeats, independently of workflow timeout + * policy. Preserve their historical 90-minute task window plus five-minute cleanup grace. + */ +const TABLE_JOB_STALE_THRESHOLD_MINUTES = 95 +/** Terminal table-jobs older than this are pruned; only the latest job per table is ever read. */ +const TABLE_JOB_RETENTION_HOURS = 24 +/** + * A table run dispatch whose holder has not made progress for this long is + * treated as dead. Same shape and window as the table-job threshold above: the + * 90-minute Trigger.dev task ceiling (`maxDuration` in `trigger.config.ts`) plus + * five minutes of cleanup grace, measured from the dispatcher's own per-window + * heartbeat rather than from when the run was requested. + */ +const TABLE_DISPATCH_STALE_THRESHOLD_MINUTES = 95 +/** Per-run ceiling on reaped dispatches, so one tick cannot fan out unbounded SSE. */ +const TABLE_DISPATCH_MAX_PER_RUN = 200 +/** + * Terminal deployment operations older than this are pruned. Every reader of + * this table is latest-generation-only, and idempotency keys only need to + * survive a client retry window, so 30 days is generous. + */ +const DEPLOYMENT_OPERATION_RETENTION_DAYS = 30 +/** + * Terminal connector sync logs older than this are pruned. Nothing pruned them + * before, so the table grew by one row per sync run forever — a connector on a + * fifteen-minute interval writes about 35,000 rows a year on its own. That cost + * lands on `loadPreviousListingObservation`, which reads the newest `completed` + * row per connector through an index covering `connector_id` alone, so every + * retained row makes the sort behind the deletion-safety corroboration slower. + */ +const CONNECTOR_SYNC_LOG_RETENTION_DAYS = 30 +const CONNECTOR_SYNC_LOG_PRUNE_BATCH_SIZE = 2000 +const CONNECTOR_SYNC_LOG_MAX_ROWS_PER_RUN = 20_000 +const DEPLOYMENT_OPERATION_PRUNE_BATCH_SIZE = 2000 +const DEPLOYMENT_OPERATION_PRUNE_MAX_BATCHES = 10 +const WORKFLOW_EXECUTION_MUTATION_BATCH_SIZE = 100 +const WORKFLOW_EXECUTION_MAX_ROWS_PER_RUN = 1000 +const STATE_MUTATION_BATCH_SIZE = 1000 +const STATE_MUTATION_MAX_ROWS_PER_RUN = 10_000 +const RETENTION_DELETE_BATCH_SIZE = 2000 +const RETENTION_DELETE_MAX_ROWS_PER_RUN = 20_000 +const TABLE_JOB_PRUNE_BATCH_SIZE = 100 +const TABLE_JOB_PRUNE_MAX_ROWS_PER_RUN = 1000 + +interface BatchedMutationResult { + affected: number + reachedLimit: boolean +} + +interface RunBatchedMutationOptions { + batchSize: number + maxRowsPerRun: number + claim: (tx: DbTransaction, limit: number) => Promise + mutation: (tx: DbTransaction, candidateIds: string[]) => Promise + onBatch?: (rows: TRow[]) => Promise +} + +/** + * Runs a mutation in bounded, atomic pages. Each page claims explicit rows with + * `FOR UPDATE SKIP LOCKED` and mutates those same rows before committing, so + * concurrent cleanup workers cannot overlap and `LIMIT` is evaluated exactly once. + */ +async function runBatchedMutation({ + batchSize, + maxRowsPerRun, + claim, + mutation, + onBatch, +}: RunBatchedMutationOptions): Promise { + let affected = 0 + + while (affected < maxRowsPerRun) { + const limit = Math.min(batchSize, maxRowsPerRun - affected) + const { candidates, rows } = await db.transaction(async (tx) => { + const candidates = await claim(tx, limit) + if (candidates.length > limit) { + throw new Error(`Cleanup claimed ${candidates.length} rows for a ${limit}-row batch`) + } + if (candidates.length === 0) return { candidates, rows: [] as TRow[] } + + const rows = await mutation( + tx, + candidates.map(({ id }) => id) + ) + if (rows.length !== candidates.length) { + throw new Error( + `Cleanup mutation returned ${rows.length} rows for ${candidates.length} claimed rows` + ) + } + return { candidates, rows } + }) + if (candidates.length === 0) break + + affected += rows.length + if (onBatch) await onBatch(rows) + if (candidates.length < limit) break + } + + return { affected, reachedLimit: affected >= maxRowsPerRun } +} + +/** Settles everything a dead worker left unfinished and prunes expired bookkeeping rows. */ +export async function runCleanupStaleExecutions() { + logger.info('Starting stale execution cleanup job') + + const now = new Date() + const staleDeadlineThreshold = new Date(now.getTime() - RESERVATION_TTL_BUFFER_MS) + const staleThreshold = new Date(now.getTime() - STALE_THRESHOLD_MINUTES * 60 * 1000) + const stalePendingThreshold = new Date( + now.getTime() - JOB_PENDING_RETENTION_HOURS * 60 * 60 * 1000 + ) + const staleTableJobThreshold = new Date( + now.getTime() - TABLE_JOB_STALE_THRESHOLD_MINUTES * 60 * 1000 + ) + const staleDispatchThreshold = new Date( + now.getTime() - TABLE_DISPATCH_STALE_THRESHOLD_MINUTES * 60 * 1000 + ) + + let staleExecutionsFound = 0 + let cleaned = 0 + let failed = 0 + let currentWorkflowBatchSize = 0 + + try { + /** + * `running` is swept on its execution deadline plus the cleanup grace, or + * on the generic stale window when it has no deadline. + * + * `redacting` gets the generic window only. That status is set *after* the + * run finished, while a live worker masks the payload, so the execution + * deadline is already in the past by the time redaction starts — the + * deadline rule would fail a masking pass that is merely slow, and + * schedule recovery would read that as a failed occurrence even though + * the worker is about to persist `completed`. A redaction that genuinely + * crashed still clears within the generic window. + */ + const staleExecutionTimePredicate = (status: StaleSweepableExecutionStatus) => + status === 'redacting' + ? lt(workflowExecutionLogs.startedAt, staleThreshold) + : or( + lt(workflowExecutionLogs.executionDeadlineAt, staleDeadlineThreshold), + and( + isNull(workflowExecutionLogs.executionDeadlineAt), + lt(workflowExecutionLogs.startedAt, staleThreshold) + ) + ) + const cleanupTimestamp = sql.param(now, workflowExecutionLogs.startedAt) + const staleDurationMinutes = sql`ROUND( + EXTRACT(EPOCH FROM (${cleanupTimestamp} - ${workflowExecutionLogs.startedAt})) / 60 + )::integer` + const totalDurationMs = elapsedDurationMsSql(now) + const staleDurationError = sql`${'Execution terminated: worker timeout or crash after '}::text + || ${staleDurationMinutes}::text + || ' minutes'` + /** + * A `redacting` row is never swept by the deadline rule, so the deadline + * message cannot apply to it. + */ + const staleExecutionError = (status: StaleSweepableExecutionStatus) => + status === 'redacting' + ? staleDurationError + : sql`CASE + WHEN ${workflowExecutionLogs.executionDeadlineAt} IS NOT NULL + THEN ${EXECUTION_DEADLINE_ERROR}::text + ELSE ${staleDurationError} + END` + /** + * Swept one status at a time so each pass stays on its own partial index; + * `status IN (...)` would match neither. The row budget is shared across + * both passes so the per-run cap keeps meaning what its name says. + */ + let workflowRowsConsidered = 0 + for (const executionStatus of STALE_SWEEPABLE_EXECUTION_STATUSES) { + const staleExecutionPredicate = and( + eq(workflowExecutionLogs.status, executionStatus), + staleExecutionTimePredicate(executionStatus) + ) + while (workflowRowsConsidered < WORKFLOW_EXECUTION_MAX_ROWS_PER_RUN) { + const limit = Math.min( + WORKFLOW_EXECUTION_MUTATION_BATCH_SIZE, + WORKFLOW_EXECUTION_MAX_ROWS_PER_RUN - workflowRowsConsidered + ) + currentWorkflowBatchSize = 0 + const { candidates, updatedExecutions } = await db.transaction(async (tx) => { + const candidates = await tx + .select({ id: workflowExecutionLogs.id }) + .from(workflowExecutionLogs) + .where(staleExecutionPredicate) + .limit(limit) + .for('update', { skipLocked: true }) + currentWorkflowBatchSize = candidates.length + if (candidates.length === 0) return { candidates, updatedExecutions: [] } + + const updatedExecutions = await tx + .update(workflowExecutionLogs) + .set({ + status: 'failed', + endedAt: now, + executionDeadlineAt: null, + totalDurationMs, + executionData: sql`jsonb_set( + COALESCE(execution_data, '{}'::jsonb), + ARRAY['error'], + to_jsonb(${staleExecutionError(executionStatus)}) + )`, + }) + .where( + and( + staleExecutionPredicate, + inArray( + workflowExecutionLogs.id, + candidates.map(({ id }) => id) + ) + ) + ) + .returning({ id: workflowExecutionLogs.id }) + + return { candidates, updatedExecutions } + }) + currentWorkflowBatchSize = 0 + staleExecutionsFound += candidates.length + if (candidates.length === 0) break + + cleaned += updatedExecutions.length + workflowRowsConsidered += candidates.length + if (candidates.length < limit) break + } + } + + if (workflowRowsConsidered >= WORKFLOW_EXECUTION_MAX_ROWS_PER_RUN) { + logger.info('Deferred remaining stale workflow executions after reaching the per-run cap', { + maxRowsPerRun: WORKFLOW_EXECUTION_MAX_ROWS_PER_RUN, + }) + } + } catch (error) { + logger.error('Failed to clean up stale workflow executions:', { + error: toError(error).message, + }) + staleExecutionsFound += currentWorkflowBatchSize + failed += currentWorkflowBatchSize + } + + logger.info(`Stale execution cleanup completed. Cleaned: ${cleaned}, Failed: ${failed}`) + + // Clean up stale async jobs (stuck in processing) + let asyncJobsMarkedFailed = 0 + + try { + const hasPositiveMaxDuration = sql`CASE + WHEN jsonb_typeof(${asyncJobs.metadata}->'maxDurationSeconds') = 'number' + THEN (${asyncJobs.metadata}->>'maxDurationSeconds')::numeric >= ${MIN_JOB_DURATION_SECONDS} + AND (${asyncJobs.metadata}->>'maxDurationSeconds')::numeric <= ${MAX_JOB_DURATION_SECONDS} + AND trunc((${asyncJobs.metadata}->>'maxDurationSeconds')::numeric) + = (${asyncJobs.metadata}->>'maxDurationSeconds')::numeric + ELSE FALSE + END` + // A bare `Date` in a raw template reaches the driver unserialized; `sql.param` + // binds it through the column encoder, as `lt(column, date)` does elsewhere here. + const staleProcessingCutoff = sql.param(now, asyncJobs.startedAt) + const staleProcessingFallbackCutoff = sql.param(staleThreshold, asyncJobs.startedAt) + const staleProcessingDurationPredicate = sql`CASE + WHEN ${hasPositiveMaxDuration} + THEN ${asyncJobs.startedAt} + ((${asyncJobs.metadata}->>'maxDurationSeconds')::double precision * interval '1 second') < ${staleProcessingCutoff} + ELSE ${asyncJobs.startedAt} < ${staleProcessingFallbackCutoff} + END` + const staleProcessingPredicate = and( + eq(asyncJobs.status, JOB_STATUS.PROCESSING), + ne(asyncJobs.type, SCHEDULE_EXECUTION_QUEUE_NAME), + staleProcessingDurationPredicate + ) + const staleProcessingResult = await runBatchedMutation({ + batchSize: STATE_MUTATION_BATCH_SIZE, + maxRowsPerRun: STATE_MUTATION_MAX_ROWS_PER_RUN, + claim: (tx, limit) => + tx + .select({ id: asyncJobs.id }) + .from(asyncJobs) + .where(staleProcessingPredicate) + .limit(limit) + .for('update', { skipLocked: true }), + mutation: (tx, candidateIds) => + tx + .update(asyncJobs) + .set({ + status: JOB_STATUS.FAILED, + completedAt: new Date(), + error: sql`CASE + WHEN ${hasPositiveMaxDuration} + THEN 'Job terminated: stuck in processing for more than ' + || (${asyncJobs.metadata}->>'maxDurationSeconds') + || ' seconds (worker cleanup deadline)' + ELSE ${GENERIC_STALE_PROCESSING_ERROR} + END`, + updatedAt: new Date(), + }) + .where(and(staleProcessingPredicate, inArray(asyncJobs.id, candidateIds))) + .returning({ id: asyncJobs.id }), + }) + + asyncJobsMarkedFailed = staleProcessingResult.affected + if (asyncJobsMarkedFailed > 0) { + logger.info(`Marked ${asyncJobsMarkedFailed} stale async jobs as failed`) + } + if (staleProcessingResult.reachedLimit) { + logger.info('Deferred remaining stale async jobs after reaching the per-run cap', { + maxRowsPerRun: STATE_MUTATION_MAX_ROWS_PER_RUN, + }) + } + } catch (error) { + logger.error('Failed to clean up stale async jobs:', { + error: toError(error).message, + }) + } + + // Mark stale table jobs (import, export, or delete) as failed. Jobs run detached on the web container + // and are lost if the pod is killed mid-run. `updated_at` is bumped by progress updates, so a + // `running` job with no recent update has stalled (not merely slow). Committed work is left in + // place (no rollback); the user retries. Also prune long-settled terminal jobs so the table + // doesn't grow unbounded (the latest job per table is what list/detail reads surface). + let staleTableJobsMarkedFailed = 0 + try { + const staleTableJobPredicate = and( + eq(tableJobs.status, 'running'), + lt(tableJobs.updatedAt, staleTableJobThreshold) + ) + const staleTableJobResult = await runBatchedMutation({ + batchSize: STATE_MUTATION_BATCH_SIZE, + maxRowsPerRun: STATE_MUTATION_MAX_ROWS_PER_RUN, + claim: (tx, limit) => + tx + .select({ id: tableJobs.id }) + .from(tableJobs) + .where(staleTableJobPredicate) + .limit(limit) + .for('update', { skipLocked: true }), + mutation: (tx, candidateIds) => + tx + .update(tableJobs) + .set({ + status: 'failed', + error: `Job terminated: no progress for more than ${TABLE_JOB_STALE_THRESHOLD_MINUTES} minutes (worker timeout or crash)`, + completedAt: now, + updatedAt: now, + }) + .where(and(staleTableJobPredicate, inArray(tableJobs.id, candidateIds))) + .returning({ id: tableJobs.id }), + }) + + staleTableJobsMarkedFailed = staleTableJobResult.affected + if (staleTableJobsMarkedFailed > 0) { + logger.info(`Marked ${staleTableJobsMarkedFailed} stale table jobs as failed`) + } + + const terminalRetention = new Date(Date.now() - TABLE_JOB_RETENTION_HOURS * 60 * 60 * 1000) + const terminalTableJobPredicate = and( + inArray(tableJobs.status, ['ready', 'failed', 'canceled']), + lt(tableJobs.updatedAt, terminalRetention) + ) + const terminalTableJobResult = await runBatchedMutation({ + batchSize: TABLE_JOB_PRUNE_BATCH_SIZE, + maxRowsPerRun: TABLE_JOB_PRUNE_MAX_ROWS_PER_RUN, + claim: (tx, limit) => + tx + .select({ id: tableJobs.id }) + .from(tableJobs) + .where(terminalTableJobPredicate) + .limit(limit) + .for('update', { skipLocked: true }), + mutation: (tx, candidateIds) => + tx + .delete(tableJobs) + .where(and(terminalTableJobPredicate, inArray(tableJobs.id, candidateIds))) + .returning({ + type: tableJobs.type, + resultKey: sql`${tableJobs.payload}->>'resultKey'`, + }), + onBatch: async (jobs) => { + /** + * Pruned export jobs carry the generated file's storage key. The scalar + * key is returned instead of the full JSON payload, and cleanup stays + * sequential so storage concurrency is bounded at one request. + */ + for (const { type, resultKey } of jobs) { + if (type !== 'export' || !resultKey) continue + await deleteFile({ key: resultKey, context: 'workspace' }).catch((err) => { + logger.warn('Failed to delete pruned export file', { + resultKey, + error: toError(err).message, + }) + }) + } + }, + }) + if (terminalTableJobResult.reachedLimit) { + logger.info('Deferred remaining terminal table jobs after reaching the per-run cap', { + maxRowsPerRun: TABLE_JOB_PRUNE_MAX_ROWS_PER_RUN, + }) + } + } catch (error) { + logger.error('Failed to clean up stale table jobs:', { + error: toError(error).message, + }) + } + + // Clean up stale pending jobs (never started, e.g., due to server crash before startJob()) + let stalePendingJobsMarkedFailed = 0 + + try { + const stalePendingPredicate = and( + eq(asyncJobs.status, JOB_STATUS.PENDING), + ne(asyncJobs.type, SCHEDULE_EXECUTION_QUEUE_NAME), + lt(asyncJobs.createdAt, stalePendingThreshold) + ) + const stalePendingResult = await runBatchedMutation({ + batchSize: STATE_MUTATION_BATCH_SIZE, + maxRowsPerRun: STATE_MUTATION_MAX_ROWS_PER_RUN, + claim: (tx, limit) => + tx + .select({ id: asyncJobs.id }) + .from(asyncJobs) + .where(stalePendingPredicate) + .limit(limit) + .for('update', { skipLocked: true }), + mutation: (tx, candidateIds) => + tx + .update(asyncJobs) + .set({ + status: JOB_STATUS.FAILED, + completedAt: new Date(), + error: `Job terminated: stuck in pending state for more than ${JOB_PENDING_RETENTION_HOURS} hours (never started)`, + updatedAt: new Date(), + }) + .where(and(stalePendingPredicate, inArray(asyncJobs.id, candidateIds))) + .returning({ id: asyncJobs.id }), + }) + + stalePendingJobsMarkedFailed = stalePendingResult.affected + if (stalePendingJobsMarkedFailed > 0) { + logger.info(`Marked ${stalePendingJobsMarkedFailed} stale pending jobs as failed`) + } + if (stalePendingResult.reachedLimit) { + logger.info('Deferred remaining stale pending jobs after reaching the per-run cap', { + maxRowsPerRun: STATE_MUTATION_MAX_ROWS_PER_RUN, + }) + } + } catch (error) { + logger.error('Failed to clean up stale pending jobs:', { + error: toError(error).message, + }) + } + + const retentionNow = Date.now() + const retentionThreshold = new Date(retentionNow - JOB_RETENTION_HOURS * 60 * 60 * 1000) + const irrecoverableCarrierRetentionThreshold = new Date( + retentionNow - SCHEDULE_CARRIER_IRRECOVERABLE_RETENTION_HOURS * 60 * 60 * 1000 + ) + let asyncJobsDeleted = 0 + + try { + const retainedJobPredicate = and( + inArray(asyncJobs.status, TERMINAL_JOB_STATUSES), + or( + ne(asyncJobs.type, SCHEDULE_EXECUTION_QUEUE_NAME), + /** + * Schedule recovery owns a carrier until it stamps the reconciled + * marker, so retention waits for it rather than deleting an + * occurrence that has not been accounted for yet. + */ + and( + carrierReconciledSql(asyncJobs.metadata), + or( + carrierNotIrrecoverableSql(asyncJobs.metadata), + lt(asyncJobs.completedAt, irrecoverableCarrierRetentionThreshold) + ) + ) + ), + lt(asyncJobs.completedAt, retentionThreshold) + ) + const retainedJobResult = await runBatchedMutation({ + batchSize: RETENTION_DELETE_BATCH_SIZE, + maxRowsPerRun: RETENTION_DELETE_MAX_ROWS_PER_RUN, + claim: (tx, limit) => + tx + .select({ id: asyncJobs.id }) + .from(asyncJobs) + .where(retainedJobPredicate) + .limit(limit) + .for('update', { skipLocked: true }), + mutation: (tx, candidateIds) => + tx + .delete(asyncJobs) + .where(and(retainedJobPredicate, inArray(asyncJobs.id, candidateIds))) + .returning({ id: asyncJobs.id }), + }) + + asyncJobsDeleted = retainedJobResult.affected + if (asyncJobsDeleted > 0) { + logger.info(`Deleted ${asyncJobsDeleted} old async jobs (retention: ${JOB_RETENTION_HOURS}h)`) + } + if (retainedJobResult.reachedLimit) { + logger.info('Deferred remaining retained async jobs after reaching the per-run cap', { + maxRowsPerRun: RETENTION_DELETE_MAX_ROWS_PER_RUN, + }) + } + } catch (error) { + logger.error('Failed to delete old async jobs:', { + error: toError(error).message, + }) + } + + /** + * Prune terminal connector sync logs past retention. + * + * HARD INVARIANT: the newest row per connector must survive, and so must the + * newest `completed` row. `loadPreviousListingObservation` reconstructs the + * previous listing from the latest `completed` log, and that reconstruction + * decides whether a suspect listing is corroborated — i.e. whether + * reconciliation may delete documents. Pruning the last `completed` row + * would silently change deletion behaviour, so both `exists` guards below + * are load-bearing rather than defensive. + * + * `started` rows are never eligible: they are either in flight or waiting on + * the scheduler's own sweep to close them. + */ + let connectorSyncLogsPruned = 0 + try { + const syncLogRetention = new Date( + Date.now() - CONNECTOR_SYNC_LOG_RETENTION_DAYS * 24 * 60 * 60 * 1000 + ) + const newerSyncLog = alias(knowledgeConnectorSyncLog, 'newer_sync_log') + const newerCompletedSyncLog = alias(knowledgeConnectorSyncLog, 'newer_completed_sync_log') + const syncLogPredicate = and( + inArray(knowledgeConnectorSyncLog.status, ['completed', 'failed']), + lt(knowledgeConnectorSyncLog.startedAt, syncLogRetention), + exists( + db + .select({ id: newerSyncLog.id }) + .from(newerSyncLog) + .where( + and( + eq(newerSyncLog.connectorId, knowledgeConnectorSyncLog.connectorId), + gt(newerSyncLog.startedAt, knowledgeConnectorSyncLog.startedAt) + ) + ) + ), + or( + ne(knowledgeConnectorSyncLog.status, 'completed'), + exists( + db + .select({ id: newerCompletedSyncLog.id }) + .from(newerCompletedSyncLog) + .where( + and( + eq(newerCompletedSyncLog.connectorId, knowledgeConnectorSyncLog.connectorId), + eq(newerCompletedSyncLog.status, 'completed'), + gt(newerCompletedSyncLog.startedAt, knowledgeConnectorSyncLog.startedAt) + ) + ) + ) + ) + ) + const syncLogResult = await runBatchedMutation({ + batchSize: CONNECTOR_SYNC_LOG_PRUNE_BATCH_SIZE, + maxRowsPerRun: CONNECTOR_SYNC_LOG_MAX_ROWS_PER_RUN, + claim: (tx, limit) => + tx + .select({ id: knowledgeConnectorSyncLog.id }) + .from(knowledgeConnectorSyncLog) + .where(syncLogPredicate) + .limit(limit) + .for('update', { skipLocked: true }), + mutation: (tx, candidateIds) => + tx + .delete(knowledgeConnectorSyncLog) + .where(inArray(knowledgeConnectorSyncLog.id, candidateIds)) + .returning({ id: knowledgeConnectorSyncLog.id }), + }) + connectorSyncLogsPruned = syncLogResult.affected + if (connectorSyncLogsPruned > 0) { + logger.info( + `Pruned ${connectorSyncLogsPruned} old connector sync logs (retention: ${CONNECTOR_SYNC_LOG_RETENTION_DAYS}d)` + ) + } + if (syncLogResult.reachedLimit) { + logger.info('Deferred remaining connector sync logs after reaching the per-run cap', { + maxRowsPerRun: CONNECTOR_SYNC_LOG_MAX_ROWS_PER_RUN, + }) + } + } catch (error) { + logger.error('Failed to prune old connector sync logs:', { + error: toError(error).message, + }) + } + + /** + * Prune terminal deployment operations past retention. HARD INVARIANT: + * the newest-generation row per workflow must always survive — the next + * deploy computes `generation = MAX(generation) + 1`, and the webhook + * registration store fences rows with lt/gt comparisons against stored + * generations, so generation reuse after a full wipe would permanently + * wedge that workflow's deployments. The `exists(newer)` predicate + * guarantees the max-generation row is never eligible; the status filter + * keeps in-flight rows (an outbox worker may still hold their fence). + */ + let deploymentOperationsPruned = 0 + try { + const deploymentOpRetention = new Date( + Date.now() - DEPLOYMENT_OPERATION_RETENTION_DAYS * 24 * 60 * 60 * 1000 + ) + const newerOperation = alias(workflowDeploymentOperation, 'newer_operation') + const deploymentOpPredicate = and( + inArray(workflowDeploymentOperation.status, ['active', 'failed', 'superseded']), + lt(workflowDeploymentOperation.completedAt, deploymentOpRetention), + exists( + db + .select({ id: newerOperation.id }) + .from(newerOperation) + .where( + and( + eq(newerOperation.workflowId, workflowDeploymentOperation.workflowId), + gt(newerOperation.generation, workflowDeploymentOperation.generation) + ) + ) + ) + ) + const deploymentOpResult = await runBatchedMutation({ + batchSize: DEPLOYMENT_OPERATION_PRUNE_BATCH_SIZE, + maxRowsPerRun: DEPLOYMENT_OPERATION_PRUNE_BATCH_SIZE * DEPLOYMENT_OPERATION_PRUNE_MAX_BATCHES, + claim: (tx, limit) => + tx + .select({ id: workflowDeploymentOperation.id }) + .from(workflowDeploymentOperation) + .where(deploymentOpPredicate) + .limit(limit) + .for('update', { skipLocked: true }), + mutation: (tx, candidateIds) => + tx + .delete(workflowDeploymentOperation) + .where(inArray(workflowDeploymentOperation.id, candidateIds)) + .returning({ id: workflowDeploymentOperation.id }), + }) + deploymentOperationsPruned = deploymentOpResult.affected + if (deploymentOperationsPruned > 0) { + logger.info( + `Pruned ${deploymentOperationsPruned} old deployment operations (retention: ${DEPLOYMENT_OPERATION_RETENTION_DAYS}d)` + ) + } + if (deploymentOpResult.reachedLimit) { + logger.info('Deferred remaining deployment operations after reaching the per-run cap', { + maxRowsPerRun: + DEPLOYMENT_OPERATION_PRUNE_BATCH_SIZE * DEPLOYMENT_OPERATION_PRUNE_MAX_BATCHES, + }) + } + } catch (error) { + logger.error('Failed to prune old deployment operations:', { + error: toError(error).message, + }) + } + + /** + * Cancel table run dispatches abandoned by a dead dispatcher. Nothing else + * reclaims them — every other terminal transition is user- or flow-initiated + * — so a dispatcher killed mid-loop left the row `dispatching` forever and + * the client's "X running" overlay with it. Ages from the dispatcher's + * per-window heartbeat, so a slow-but-live dispatch is spared. + */ + let staleDispatchesCancelled = 0 + try { + staleDispatchesCancelled = ( + await cancelStaleDispatches(staleDispatchThreshold, TABLE_DISPATCH_MAX_PER_RUN) + ).length + if (staleDispatchesCancelled > 0) { + logger.warn(`Cancelled ${staleDispatchesCancelled} abandoned table run dispatches`, { + thresholdMinutes: TABLE_DISPATCH_STALE_THRESHOLD_MINUTES, + }) + } + } catch (error) { + logger.error('Failed to cancel abandoned table run dispatches:', { + error: toError(error).message, + }) + } + + /** + * Settle Chat runs no controller will finish: their process died, their + * controller was superseded without a successor, or Stop found none. Without + * this they stay unfinished forever and keep their chat marked as busy. + */ + let orphanedRunsSettled = 0 + try { + orphanedRunsSettled = (await sweepOrphanedRuns()).settledRunIds.length + } catch (error) { + logger.error('Failed to settle orphaned Chat runs:', { + error: toError(error).message, + }) + } + + return { + executions: { + found: staleExecutionsFound, + cleaned, + failed, + thresholdMinutes: STALE_THRESHOLD_MINUTES, + }, + asyncJobs: { + staleProcessingMarkedFailed: asyncJobsMarkedFailed, + stalePendingMarkedFailed: stalePendingJobsMarkedFailed, + oldDeleted: asyncJobsDeleted, + staleThresholdMinutes: STALE_THRESHOLD_MINUTES, + retentionHours: JOB_RETENTION_HOURS, + }, + tableJobs: { + staleMarkedFailed: staleTableJobsMarkedFailed, + }, + connectorSyncLogs: { + pruned: connectorSyncLogsPruned, + retentionDays: CONNECTOR_SYNC_LOG_RETENTION_DAYS, + }, + tableRunDispatches: { + staleCancelled: staleDispatchesCancelled, + thresholdMinutes: TABLE_DISPATCH_STALE_THRESHOLD_MINUTES, + }, + deploymentOperations: { + pruned: deploymentOperationsPruned, + retentionDays: DEPLOYMENT_OPERATION_RETENTION_DAYS, + }, + chatRuns: { + orphanedSettled: orphanedRunsSettled, + }, + } +} + +export const cleanupStaleExecutionsTask = task({ + id: 'cleanup-stale-executions', + queue: { concurrencyLimit: 1 }, + run: () => runCleanupStaleExecutions(), +}) diff --git a/apps/sim/lib/core/async-jobs/backends/trigger-dev.ts b/apps/sim/lib/core/async-jobs/backends/trigger-dev.ts index e2f13aba5c9..ecce4f97894 100644 --- a/apps/sim/lib/core/async-jobs/backends/trigger-dev.ts +++ b/apps/sim/lib/core/async-jobs/backends/trigger-dev.ts @@ -253,6 +253,7 @@ const JOB_TYPE_TO_TASK_ID: Record = { 'cleanup-logs': 'cleanup-logs', 'cleanup-soft-deletes': 'cleanup-soft-deletes', 'cleanup-table-row-ttl': 'cleanup-table-row-ttl', + 'cleanup-stale-executions': 'cleanup-stale-executions', 'cleanup-tasks': 'cleanup-tasks', 'cleanup-file-versions': 'cleanup-file-versions', 'run-data-drain': 'run-data-drain', diff --git a/apps/sim/lib/core/async-jobs/types.ts b/apps/sim/lib/core/async-jobs/types.ts index d3c893cd54a..cf5eb54e2d8 100644 --- a/apps/sim/lib/core/async-jobs/types.ts +++ b/apps/sim/lib/core/async-jobs/types.ts @@ -47,6 +47,7 @@ export type JobType = | 'cleanup-logs' | 'cleanup-soft-deletes' | 'cleanup-table-row-ttl' + | 'cleanup-stale-executions' | 'cleanup-tasks' | 'cleanup-file-versions' | 'run-data-drain' From 10e03739c0998944155665dce96db58e2b7fdc33 Mon Sep 17 00:00:00 2001 From: Theodore Li Date: Thu, 1 Oct 2026 12:26:37 -0700 Subject: [PATCH 2/6] fix(cron): run sync-log retention test against the cleanup task --- .../connectors/sync-log-retention.test.ts | 14 +++----------- 1 file changed, 3 insertions(+), 11 deletions(-) diff --git a/apps/sim/lib/knowledge/connectors/sync-log-retention.test.ts b/apps/sim/lib/knowledge/connectors/sync-log-retention.test.ts index fc206e7f02a..e22f9072163 100644 --- a/apps/sim/lib/knowledge/connectors/sync-log-retention.test.ts +++ b/apps/sim/lib/knowledge/connectors/sync-log-retention.test.ts @@ -1,19 +1,11 @@ import { knowledgeConnectorSyncLog } from '@sim/db/schema' -import { - createMockRequest, - dbChainMockFns, - queueTableRows, - resetDbChainMock, - schemaMock, -} from '@sim/testing' -import { authInternalMock } from '@sim/testing/mocks/auth-internal.mock' +import { dbChainMockFns, queueTableRows, resetDbChainMock, schemaMock } from '@sim/testing' import { storageServiceMock } from '@sim/testing/mocks/storage-service.mock' import { beforeEach, describe, expect, it, vi } from 'vitest' -vi.mock('@/lib/auth/internal', () => authInternalMock) vi.mock('@/lib/uploads/core/storage-service', () => storageServiceMock) -import { GET } from '@/app/api/cron/cleanup-stale-executions/route' +import { runCleanupStaleExecutions } from '@/background/cleanup-stale-executions' /** Flattens a drizzle condition tree into the values it compares against. */ function collectValues(value: unknown, out: unknown[] = []): unknown[] { @@ -71,7 +63,7 @@ describe('connector sync log retention', () => { * behaviour, so the two `exists` guards are load-bearing. */ it('claims only terminal rows that still have a newer sibling', async () => { - await GET(createMockRequest('GET') as never) + await runCleanupStaleExecutions() const predicate = syncLogClaimPredicate() // The arm has to actually run, or this asserts nothing. From e092a9ad173bec97601f8ce8af0cb62093ac7bce Mon Sep 17 00:00:00 2001 From: Theodore Li Date: Thu, 1 Oct 2026 12:47:02 -0700 Subject: [PATCH 3/6] fix(cron): assert stale cleanup route by its response only --- .../cleanup-stale-executions/route.test.ts | 35 +++---------------- 1 file changed, 4 insertions(+), 31 deletions(-) diff --git a/apps/sim/app/api/cron/cleanup-stale-executions/route.test.ts b/apps/sim/app/api/cron/cleanup-stale-executions/route.test.ts index 059a2ea899a..705e42d9ab5 100644 --- a/apps/sim/app/api/cron/cleanup-stale-executions/route.test.ts +++ b/apps/sim/app/api/cron/cleanup-stale-executions/route.test.ts @@ -1,7 +1,7 @@ import { createMockRequest } from '@sim/testing' import { asyncJobsMock, asyncJobsMockFns } from '@sim/testing/mocks/async-jobs.mock' import { authInternalMock, authInternalMockFns } from '@sim/testing/mocks/auth-internal.mock' -import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' +import { beforeEach, describe, expect, it, vi } from 'vitest' vi.mock('@/lib/auth/internal', () => authInternalMock) vi.mock('@/lib/core/async-jobs', () => asyncJobsMock) @@ -23,48 +23,21 @@ function request() { describe('stale execution cleanup route', () => { beforeEach(() => { - vi.useFakeTimers() - vi.setSystemTime(new Date('2026-10-01T17:31:00Z')) mockVerifyCronAuth.mockReturnValue(null) mockEnqueue.mockReset() - mockEnqueue.mockResolvedValue('job-stale-1') }) - afterEach(() => { - vi.useRealTimers() - }) + it('answers with the dispatched job once the cleanup is enqueued', async () => { + mockEnqueue.mockResolvedValueOnce('job-stale-1') - it('hands the cleanup to the job queue and answers without waiting for it', async () => { const response = await GET(request()) expect(response.status).toBe(200) await expect(response.json()).resolves.toEqual({ triggered: true, jobId: 'job-stale-1' }) - expect(mockEnqueue).toHaveBeenCalledWith( - 'cleanup-stale-executions', - {}, - expect.objectContaining({ maxAttempts: 1, concurrencyLimit: 1 }) - ) - }) - - it('deduplicates retries within the same thirty-minute schedule window', async () => { - await GET(request()) - vi.advanceTimersByTime(28 * 60 * 1000) - await GET(request()) - - expect(mockEnqueue.mock.calls[0]?.[2]?.jobId).toBe(mockEnqueue.mock.calls[1]?.[2]?.jobId) - }) - - it('uses a new id immediately after the next thirty-minute window begins', async () => { - vi.setSystemTime(new Date('2026-10-01T17:59:59.999Z')) - await GET(request()) - vi.setSystemTime(new Date('2026-10-01T18:00:00.000Z')) - await GET(request()) - - expect(mockEnqueue.mock.calls[0]?.[2]?.jobId).not.toBe(mockEnqueue.mock.calls[1]?.[2]?.jobId) }) it('fails the cron invocation when the job cannot be enqueued', async () => { - mockEnqueue.mockRejectedValue(new Error('queue unavailable')) + mockEnqueue.mockRejectedValueOnce(new Error('queue unavailable')) const response = await GET(request()) From 8cd3f0d773feb16643826d1c79ce22f86c47dad7 Mon Sep 17 00:00:00 2001 From: Theodore Li Date: Thu, 1 Oct 2026 13:23:51 -0700 Subject: [PATCH 4/6] fix(cron): cover stale cleanup dispatch window through the route response --- .../cleanup-stale-executions/route.test.ts | 38 ++++++++++++++++++- 1 file changed, 37 insertions(+), 1 deletion(-) diff --git a/apps/sim/app/api/cron/cleanup-stale-executions/route.test.ts b/apps/sim/app/api/cron/cleanup-stale-executions/route.test.ts index 705e42d9ab5..7d7c8bfb72a 100644 --- a/apps/sim/app/api/cron/cleanup-stale-executions/route.test.ts +++ b/apps/sim/app/api/cron/cleanup-stale-executions/route.test.ts @@ -1,7 +1,7 @@ import { createMockRequest } from '@sim/testing' import { asyncJobsMock, asyncJobsMockFns } from '@sim/testing/mocks/async-jobs.mock' import { authInternalMock, authInternalMockFns } from '@sim/testing/mocks/auth-internal.mock' -import { beforeEach, describe, expect, it, vi } from 'vitest' +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' vi.mock('@/lib/auth/internal', () => authInternalMock) vi.mock('@/lib/core/async-jobs', () => asyncJobsMock) @@ -21,12 +21,26 @@ function request() { ) } +/** The job id a dispatch was keyed to, as the queue reports it back. */ +async function dispatchedJobId(): Promise { + const response = await GET(request()) + expect(response.status).toBe(200) + const body = (await response.json()) as { jobId: string } + return body.jobId +} + describe('stale execution cleanup route', () => { beforeEach(() => { + vi.useFakeTimers() + vi.setSystemTime(new Date('2026-10-01T17:31:00Z')) mockVerifyCronAuth.mockReturnValue(null) mockEnqueue.mockReset() }) + afterEach(() => { + vi.useRealTimers() + }) + it('answers with the dispatched job once the cleanup is enqueued', async () => { mockEnqueue.mockResolvedValueOnce('job-stale-1') @@ -43,4 +57,26 @@ describe('stale execution cleanup route', () => { expect(response.status).toBe(500) }) + + describe('with a queue that keys each job by the id it is given', () => { + beforeEach(() => { + mockEnqueue.mockImplementation( + async (_type: string, _payload: unknown, options: { jobId: string }) => options.jobId + ) + }) + + it('dispatches a retry inside the same thirty-minute window as the same job', async () => { + const first = await dispatchedJobId() + vi.setSystemTime(new Date('2026-10-01T17:59:59.999Z')) + + expect(await dispatchedJobId()).toBe(first) + }) + + it('dispatches a new job once the next thirty-minute window begins', async () => { + const first = await dispatchedJobId() + vi.setSystemTime(new Date('2026-10-01T18:00:00.000Z')) + + expect(await dispatchedJobId()).not.toBe(first) + }) + }) }) From 4d0974dbce030411c9a13f75512711378c937e1b Mon Sep 17 00:00:00 2001 From: Theodore Li Date: Thu, 1 Oct 2026 15:04:11 -0700 Subject: [PATCH 5/6] improvement(cron): dispatch file version cleanup from a background task --- .../cron/cleanup-file-versions/route.test.ts | 73 +++++++++++++++++++ .../api/cron/cleanup-file-versions/route.ts | 26 +++++-- apps/sim/background/cleanup-dispatch.ts | 16 ++++ .../core/async-jobs/backends/trigger-dev.ts | 1 + apps/sim/lib/core/async-jobs/types.ts | 1 + 5 files changed, 112 insertions(+), 5 deletions(-) create mode 100644 apps/sim/app/api/cron/cleanup-file-versions/route.test.ts create mode 100644 apps/sim/background/cleanup-dispatch.ts diff --git a/apps/sim/app/api/cron/cleanup-file-versions/route.test.ts b/apps/sim/app/api/cron/cleanup-file-versions/route.test.ts new file mode 100644 index 00000000000..4b3707a81e0 --- /dev/null +++ b/apps/sim/app/api/cron/cleanup-file-versions/route.test.ts @@ -0,0 +1,73 @@ +import { createMockRequest } from '@sim/testing' +import { asyncJobsMock, asyncJobsMockFns } from '@sim/testing/mocks/async-jobs.mock' +import { authInternalMock, authInternalMockFns } from '@sim/testing/mocks/auth-internal.mock' +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' + +vi.mock('@/lib/auth/internal', () => authInternalMock) +vi.mock('@/lib/core/async-jobs', () => asyncJobsMock) + +import { GET } from '@/app/api/cron/cleanup-file-versions/route' + +const { mockVerifyCronAuth } = authInternalMockFns + +const mockEnqueue = asyncJobsMockFns.mockJobQueue.enqueue + +function request() { + return createMockRequest( + 'GET', + undefined, + {}, + 'http://localhost:3000/api/cron/cleanup-file-versions' + ) +} + +/** The job id a dispatch was keyed to, as the queue reports it back. */ +async function dispatchedJobId(): Promise { + const response = await GET(request()) + expect(response.status).toBe(200) + const body = (await response.json()) as { jobId: string } + return body.jobId +} + +describe('file version cleanup route', () => { + beforeEach(() => { + vi.useFakeTimers() + vi.setSystemTime(new Date('2026-10-01T03:30:00Z')) + mockVerifyCronAuth.mockReturnValue(null) + mockEnqueue.mockReset() + }) + + afterEach(() => { + vi.useRealTimers() + }) + + it('fails the cron invocation when the dispatch cannot be enqueued', async () => { + mockEnqueue.mockRejectedValueOnce(new Error('queue unavailable')) + + const response = await GET(request()) + + expect(response.status).toBe(500) + }) + + describe('with a queue that keys each job by the id it is given', () => { + beforeEach(() => { + mockEnqueue.mockImplementation( + async (_type: string, _payload: unknown, options: { jobId: string }) => options.jobId + ) + }) + + it('dispatches a retry on the same day as the same job', async () => { + const first = await dispatchedJobId() + vi.setSystemTime(new Date('2026-10-01T23:59:59.999Z')) + + expect(await dispatchedJobId()).toBe(first) + }) + + it('dispatches a new job on the next day', async () => { + const first = await dispatchedJobId() + vi.setSystemTime(new Date('2026-10-02T00:00:00.000Z')) + + expect(await dispatchedJobId()).not.toBe(first) + }) + }) +}) diff --git a/apps/sim/app/api/cron/cleanup-file-versions/route.ts b/apps/sim/app/api/cron/cleanup-file-versions/route.ts index a4b312347ce..5f0a3cc24c1 100644 --- a/apps/sim/app/api/cron/cleanup-file-versions/route.ts +++ b/apps/sim/app/api/cron/cleanup-file-versions/route.ts @@ -1,12 +1,13 @@ import { createLogger } from '@sim/logger' import { type NextRequest, NextResponse } from 'next/server' import { verifyCronAuth } from '@/lib/auth/internal' -import { dispatchCleanupJobs } from '@/lib/billing/cleanup-dispatcher' +import { getJobQueue } from '@/lib/core/async-jobs' import { withRouteHandler } from '@/lib/core/utils/with-route-handler' export const dynamic = 'force-dynamic' const logger = createLogger('FileVersionCleanupAPI') +const FILE_VERSION_CLEANUP_INTERVAL_MS = 24 * 60 * 60 * 1000 /** GET /api/cron/cleanup-file-versions — dispatch retention for superseded workspace file versions. */ export const GET = withRouteHandler(async (request: NextRequest) => { @@ -14,11 +15,26 @@ export const GET = withRouteHandler(async (request: NextRequest) => { const authError = verifyCronAuth(request, 'file version cleanup') if (authError) return authError - const result = await dispatchCleanupJobs('cleanup-file-versions') + const queue = await getJobQueue() + const scheduleWindow = Math.floor(Date.now() / FILE_VERSION_CLEANUP_INTERVAL_MS) + const jobId = await queue.enqueue( + 'cleanup-dispatch', + { jobType: 'cleanup-file-versions' }, + { + maxAttempts: 1, + jobId: `cleanup-dispatch:cleanup-file-versions:${scheduleWindow}`, + name: 'File version cleanup dispatch', + concurrencyKey: 'cleanup-dispatch:cleanup-file-versions', + concurrencyLimit: 1, + runner: async () => { + const { dispatchCleanupJobs } = await import('@/lib/billing/cleanup-dispatcher') + return dispatchCleanupJobs('cleanup-file-versions') + }, + } + ) - logger.info('File version cleanup jobs dispatched', result) - - return NextResponse.json({ triggered: true, ...result }) + logger.info('File version cleanup dispatch enqueued', { jobId }) + return NextResponse.json({ triggered: true, jobId }) } catch (error) { logger.error('Failed to dispatch file version cleanup jobs:', { error }) return NextResponse.json({ error: 'Failed to dispatch file version cleanup' }, { status: 500 }) diff --git a/apps/sim/background/cleanup-dispatch.ts b/apps/sim/background/cleanup-dispatch.ts new file mode 100644 index 00000000000..b119004aaf2 --- /dev/null +++ b/apps/sim/background/cleanup-dispatch.ts @@ -0,0 +1,16 @@ +import { task } from '@trigger.dev/sdk' +import { type CleanupJobType, dispatchCleanupJobs } from '@/lib/billing/cleanup-dispatcher' + +export interface CleanupDispatchPayload { + jobType: CleanupJobType +} + +/** + * Resolves a retention job's workspace scope and fans out its per-chunk cleanup runs. The scan + * covers every active workspace, so it runs here rather than inside the cron request. + */ +export const cleanupDispatchTask = task({ + id: 'cleanup-dispatch', + retry: { maxAttempts: 1 }, + run: ({ jobType }: CleanupDispatchPayload) => dispatchCleanupJobs(jobType), +}) diff --git a/apps/sim/lib/core/async-jobs/backends/trigger-dev.ts b/apps/sim/lib/core/async-jobs/backends/trigger-dev.ts index ecce4f97894..59dd758ef32 100644 --- a/apps/sim/lib/core/async-jobs/backends/trigger-dev.ts +++ b/apps/sim/lib/core/async-jobs/backends/trigger-dev.ts @@ -256,6 +256,7 @@ const JOB_TYPE_TO_TASK_ID: Record = { 'cleanup-stale-executions': 'cleanup-stale-executions', 'cleanup-tasks': 'cleanup-tasks', 'cleanup-file-versions': 'cleanup-file-versions', + 'cleanup-dispatch': 'cleanup-dispatch', 'run-data-drain': 'run-data-drain', 'knowledge-connector-directory-sync': 'knowledge-connector-directory-sync', } diff --git a/apps/sim/lib/core/async-jobs/types.ts b/apps/sim/lib/core/async-jobs/types.ts index cf5eb54e2d8..c65b5df56fd 100644 --- a/apps/sim/lib/core/async-jobs/types.ts +++ b/apps/sim/lib/core/async-jobs/types.ts @@ -50,6 +50,7 @@ export type JobType = | 'cleanup-stale-executions' | 'cleanup-tasks' | 'cleanup-file-versions' + | 'cleanup-dispatch' | 'run-data-drain' | 'knowledge-connector-directory-sync' From 078e466f6d4d0d6098acda3ee88c46d5ba293622 Mon Sep 17 00:00:00 2001 From: Theodore Li Date: Thu, 1 Oct 2026 15:15:23 -0700 Subject: [PATCH 6/6] fix(cron): retry and serialize cleanup dispatch on the task definition --- apps/sim/app/api/cron/cleanup-file-versions/route.ts | 3 ++- apps/sim/background/cleanup-dispatch.ts | 9 +++++++-- 2 files changed, 9 insertions(+), 3 deletions(-) diff --git a/apps/sim/app/api/cron/cleanup-file-versions/route.ts b/apps/sim/app/api/cron/cleanup-file-versions/route.ts index 5f0a3cc24c1..13f7565b599 100644 --- a/apps/sim/app/api/cron/cleanup-file-versions/route.ts +++ b/apps/sim/app/api/cron/cleanup-file-versions/route.ts @@ -3,6 +3,7 @@ import { type NextRequest, NextResponse } from 'next/server' import { verifyCronAuth } from '@/lib/auth/internal' import { getJobQueue } from '@/lib/core/async-jobs' import { withRouteHandler } from '@/lib/core/utils/with-route-handler' +import { CLEANUP_DISPATCH_MAX_ATTEMPTS } from '@/background/cleanup-dispatch' export const dynamic = 'force-dynamic' @@ -21,7 +22,7 @@ export const GET = withRouteHandler(async (request: NextRequest) => { 'cleanup-dispatch', { jobType: 'cleanup-file-versions' }, { - maxAttempts: 1, + maxAttempts: CLEANUP_DISPATCH_MAX_ATTEMPTS, jobId: `cleanup-dispatch:cleanup-file-versions:${scheduleWindow}`, name: 'File version cleanup dispatch', concurrencyKey: 'cleanup-dispatch:cleanup-file-versions', diff --git a/apps/sim/background/cleanup-dispatch.ts b/apps/sim/background/cleanup-dispatch.ts index b119004aaf2..d0e1f94f5d2 100644 --- a/apps/sim/background/cleanup-dispatch.ts +++ b/apps/sim/background/cleanup-dispatch.ts @@ -5,12 +5,17 @@ export interface CleanupDispatchPayload { jobType: CleanupJobType } +/** Attempts per dispatch. A rerun only re-triggers chunks that delete already-expired data. */ +export const CLEANUP_DISPATCH_MAX_ATTEMPTS = 3 + /** * Resolves a retention job's workspace scope and fans out its per-chunk cleanup runs. The scan - * covers every active workspace, so it runs here rather than inside the cron request. + * covers every active workspace, so it runs here rather than inside the cron request. One + * dispatch at a time per job type: the route keys runs by `concurrencyKey`. */ export const cleanupDispatchTask = task({ id: 'cleanup-dispatch', - retry: { maxAttempts: 1 }, + queue: { concurrencyLimit: 1 }, + retry: { maxAttempts: CLEANUP_DISPATCH_MAX_ATTEMPTS }, run: ({ jobType }: CleanupDispatchPayload) => dispatchCleanupJobs(jobType), })