diff --git a/.gitignore b/.gitignore index 02aa7473a5..95ea5a0360 100644 --- a/.gitignore +++ b/.gitignore @@ -114,3 +114,4 @@ cli-helper/dist/ cli-helper/npm/*/CapgoKeychainHelper.app cli-helper/npm/*/CapgoAscKeyHelper.app .tinbase/ +.context/ diff --git a/read_replicate/schema_replicate.catalog.json b/read_replicate/schema_replicate.catalog.json index 4857f5c5bb..f7b40ad9f5 100644 --- a/read_replicate/schema_replicate.catalog.json +++ b/read_replicate/schema_replicate.catalog.json @@ -2052,6 +2052,20 @@ "table": "app_versions", "valid": true }, + { + "constraintOwned": false, + "definition": "CREATE INDEX idx_app_versions_deleted_at_id ON public.app_versions USING btree (deleted_at, id) WHERE (deleted = true)", + "name": "idx_app_versions_deleted_at_id", + "table": "app_versions", + "valid": true + }, + { + "constraintOwned": false, + "definition": "CREATE INDEX idx_app_versions_deleted_with_manifest ON public.app_versions USING btree (id) WHERE ((deleted = true) AND (manifest_count > 0))", + "name": "idx_app_versions_deleted_with_manifest", + "table": "app_versions", + "valid": true + }, { "constraintOwned": false, "definition": "CREATE INDEX idx_app_versions_id ON public.app_versions USING btree (id)", diff --git a/read_replicate/schema_replicate.sql b/read_replicate/schema_replicate.sql index 94e69f1dc1..ec28434c82 100644 --- a/read_replicate/schema_replicate.sql +++ b/read_replicate/schema_replicate.sql @@ -768,6 +768,20 @@ CREATE INDEX idx_app_versions_deleted ON public.app_versions USING btree (delete CREATE INDEX idx_app_versions_deleted_at ON public.app_versions USING btree (deleted_at) WHERE (deleted_at IS NOT NULL); +-- +-- Name: idx_app_versions_deleted_at_id; Type: INDEX; Schema: public; Owner: - +-- + +CREATE INDEX idx_app_versions_deleted_at_id ON public.app_versions USING btree (deleted_at, id) WHERE (deleted = true); + + +-- +-- Name: idx_app_versions_deleted_with_manifest; Type: INDEX; Schema: public; Owner: - +-- + +CREATE INDEX idx_app_versions_deleted_with_manifest ON public.app_versions USING btree (id) WHERE ((deleted = true) AND (manifest_count > 0)); + + -- -- Name: idx_app_versions_id; Type: INDEX; Schema: public; Owner: - -- diff --git a/scripts/reclaim_deleted_version_manifests.ts b/scripts/reclaim_deleted_version_manifests.ts new file mode 100644 index 0000000000..682b3e5991 --- /dev/null +++ b/scripts/reclaim_deleted_version_manifests.ts @@ -0,0 +1,240 @@ +/** + * Reclaim ALL soft-deleted app_versions that still have public.manifest rows. + * + * For each file: + * 1) if R2 object exists → move to deleted-after-7-days/ + * 2) if missing → continue + * 3) only then delete the DB row + * + * Shared delta files referenced by other versions stay in R2. + * + * Usage: + * bun scripts/reclaim_deleted_version_manifests.ts + */ +import { CopyObjectCommand, DeleteObjectCommand, HeadObjectCommand, S3Client } from '@aws-sdk/client-s3' +import pg from 'pg' + +const ENV_FILE = './internal/cloudflare/.env.prod' +const TRASH_PREFIX = 'deleted-after-7-days/' +const CONCURRENCY = 50 +const VERSION_PAGE = 100 +const DB_URL_ENV_KEYS = [ + 'MAIN_SUPABASE_DB_URL', + 'DATABASE_URL', + 'POSTGRES_URL', + 'SUPABASE_DB_URL', + 'SUPABASE_DB_DIRECT_URL', + 'DIRECT_URL', +] + +function loadEnv(filePath: string) { + return Bun.file(filePath).text().then((text) => { + const env: Record = {} + for (const line of text.split('\n')) { + const trimmed = line.trim() + if (!trimmed || trimmed.startsWith('#')) + continue + const idx = trimmed.indexOf('=') + if (idx <= 0) + continue + env[trimmed.slice(0, idx)] = trimmed.slice(idx + 1) + } + return env + }) +} + +function requireDbUrl(env: Record) { + for (const key of DB_URL_ENV_KEYS) { + if (env[key]) + return env[key] + } + throw new Error(`Missing DB URL. Set one of: ${DB_URL_ENV_KEYS.join(', ')}`) +} + +async function mapPool(items: T[], concurrency: number, fn: (item: T) => Promise) { + let idx = 0 + const workers = Array.from({ length: Math.min(concurrency, Math.max(items.length, 1)) }, async () => { + while (idx < items.length) { + const current = items[idx++] + await fn(current) + } + }) + await Promise.all(workers) +} + +async function objectExists(s3: S3Client, bucket: string, key: string) { + try { + await s3.send(new HeadObjectCommand({ Bucket: bucket, Key: key })) + return true + } + catch (error: any) { + const status = error?.$metadata?.httpStatusCode ?? error?.statusCode ?? error?.status + const code = error?.name ?? error?.Code ?? error?.code + if (status === 404 || code === 'NotFound' || code === 'NoSuchKey') + return false + throw error + } +} + +async function moveToTrash(s3: S3Client, bucket: string, key: string) { + if (key.startsWith(TRASH_PREFIX)) + return + + const exists = await objectExists(s3, bucket, key) + if (!exists) + return + + const encodedKey = key.split('/').map(segment => encodeURIComponent(segment)).join('/') + await s3.send(new CopyObjectCommand({ + Bucket: bucket, + CopySource: `${bucket}/${encodedKey}`, + Key: `${TRASH_PREFIX}${key}`, + })) + await s3.send(new DeleteObjectCommand({ + Bucket: bucket, + Key: key, + })) +} + +function shouldAllowSelfSignedPgCertificate(env: Record, databaseUrl: string) { + const rejectUnauthorized = env.PG_SSL_REJECT_UNAUTHORIZED?.trim() + if (rejectUnauthorized === '0') + return true + if (rejectUnauthorized === '1') + return false + // Local/docker URLs keep default Node TLS verification off only when explicitly local. + return databaseUrl.includes('localhost') || databaseUrl.includes('127.0.0.1') +} + +async function main() { + const env = await loadEnv(ENV_FILE) + const databaseUrl = requireDbUrl(env) + const usesLocalDatabase = databaseUrl.includes('localhost') || databaseUrl.includes('127.0.0.1') + const db = new pg.Client({ + connectionString: databaseUrl, + ssl: usesLocalDatabase + ? false + : { rejectUnauthorized: !shouldAllowSelfSignedPgCertificate(env, databaseUrl) }, + }) + await db.connect() + + const s3 = new S3Client({ + credentials: { + accessKeyId: env.S3_ACCESS_KEY_ID, + secretAccessKey: env.S3_SECRET_ACCESS_KEY, + }, + endpoint: `https://${env.S3_ENDPOINT}`, + region: env.S3_REGION || 'auto', + forcePathStyle: true, + }) + const bucket = env.S3_BUCKET || 'capgo' + + console.log('Listing soft-deleted versions with leftover manifests...') + const versionsRes = await db.query<{ id: number, app_id: string, manifest_count: number }>(` + SELECT av.id, av.app_id, av.manifest_count + FROM public.app_versions AS av + WHERE av.deleted = true + AND ( + av.manifest_count > 0 + OR EXISTS ( + SELECT 1 FROM public.manifest AS m WHERE m.app_version_id = av.id + ) + ) + AND av.app_id NOT LIKE 'com.capdemo%' + ORDER BY av.id + `) + const versions = versionsRes.rows + console.log(`Found ${versions.length} versions`) + + let done = 0 + let totalTrashed = 0 + let totalDeleted = 0 + + for (let i = 0; i < versions.length; i += VERSION_PAGE) { + const page = versions.slice(i, i + VERSION_PAGE) + for (const version of page) { + const entriesRes = await db.query<{ + id: number + file_hash: string + file_name: string + s3_path: string | null + }>( + `SELECT id, file_hash, file_name, s3_path + FROM public.manifest + WHERE app_version_id = $1 + ORDER BY id`, + [version.id], + ) + const entries = entriesRes.rows + let trashed = 0 + let deletedRows = 0 + + await mapPool(entries, CONCURRENCY, async (entry) => { + if (entry.s3_path) { + const ref = await db.query( + `SELECT 1 AS ok + FROM public.manifest + WHERE file_hash = $1 + AND file_name = $2 + AND app_version_id <> $3 + LIMIT 1`, + [entry.file_hash, entry.file_name, version.id], + ) + + if (ref.rows.length === 0) { + await moveToTrash(s3, bucket, entry.s3_path) + trashed += 1 + } + } + + await db.query(`DELETE FROM public.manifest WHERE id = $1`, [entry.id]) + deletedRows += 1 + }) + + const remaining = await db.query( + `SELECT COUNT(*)::int AS count FROM public.manifest WHERE app_version_id = $1`, + [version.id], + ) + if (Number(remaining.rows[0]?.count ?? 0) > 0) + throw new Error(`version ${version.id} still has manifest rows`) + + await db.query('BEGIN') + try { + await db.query( + `UPDATE public.app_versions + SET manifest_count = 0, manifest = NULL + WHERE id = $1`, + [version.id], + ) + if (deletedRows > 0 || (version.manifest_count ?? 0) > 0) { + await db.query( + `UPDATE public.apps + SET manifest_bundle_count = GREATEST(manifest_bundle_count - 1, 0), + updated_at = now() + WHERE app_id = $1`, + [version.app_id], + ) + } + await db.query('COMMIT') + } + catch (error) { + await db.query('ROLLBACK') + throw error + } + + done += 1 + totalTrashed += trashed + totalDeleted += deletedRows + process.stdout.write(`\rCleaned ${done}/${versions.length} versions (rows=${totalDeleted}, r2=${totalTrashed})`) + } + } + + await db.end() + process.stdout.write('\n') + console.log('Done.') + console.log(`Versions cleaned: ${done}`) + console.log(`Manifest rows deleted: ${totalDeleted}`) + console.log(`R2 objects moved to ${TRASH_PREFIX}: ${totalTrashed}`) +} + +await main() diff --git a/supabase/functions/_backend/triggers/on_version_update.ts b/supabase/functions/_backend/triggers/on_version_update.ts index 93214db3f8..5859fbb987 100644 --- a/supabase/functions/_backend/triggers/on_version_update.ts +++ b/supabase/functions/_backend/triggers/on_version_update.ts @@ -316,111 +316,187 @@ async function updateIt(c: Context, record: Database['public']['Tables']['app_ve return c.json(BRES) } +const MANIFEST_TRASH_CONCURRENCY = 10 + +type ManifestCleanupEntry = { + id: number + file_hash: string + file_name: string + s3_path: string | null +} + /** - * Deletes manifest rows and moves orphaned S3 assets to the R2 trash prefix. + * Trash unreferenced R2 objects first (exist → move to deleted-after-7-days/, + * missing → ok), then delete that DB row. Never drop DB tracking before R2 is handled. + * Incomplete work throws so the queue retries; already-trashed paths are idempotent. */ async function deleteManifest(c: Context, record: Database['public']['Tables']['app_versions']['Row']) { - // Delete manifest entries - first get them to delete from S3 - const pgClient = getPgClient(c, true) // READ-ONLY: deletes use SDK, not Drizzle - const drizzleClient = getDrizzleClient(pgClient) + const readPgClient = getPgClient(c, true) + const drizzleClient = getDrizzleClient(readPgClient) + let manifestEntries: ManifestCleanupEntry[] = [] try { - const manifestEntries = await drizzleClient - .select() + manifestEntries = await drizzleClient + .select({ + id: manifest.id, + file_hash: manifest.file_hash, + file_name: manifest.file_name, + s3_path: manifest.s3_path, + }) .from(manifest) .where(eq(manifest.app_version_id, record.id)) + } + finally { + await closeClient(c, readPgClient) + } - if (manifestEntries && manifestEntries.length > 0) { - const manifestCount = manifestEntries.length - - // Move each unreferenced file to the R2 trash prefix. - const promisesMoveToTrash = [] - for (const entry of manifestEntries) { - if (entry.s3_path) { - promisesMoveToTrash.push( - // First delete the manifest row from database - supabaseAdmin(c) - .from('manifest') - .delete() - .eq('id', entry.id) - .then(({ error: deleteError }) => { - if (deleteError) { - cloudlog({ requestId: c.get('requestId'), message: 'error deleting manifest row', id: entry.id, error: deleteError }) - return null // Signal to skip S3 cleanup - } - // After deleting, check if any other rows still reference this file - // This avoids race condition where concurrent deletes both skip S3 cleanup - return supabaseAdmin(c) - .from('manifest') - .select('id') - .eq('file_hash', entry.file_hash) - .eq('file_name', entry.file_name) - .limit(1) - .maybeSingle() - }) - .then((v) => { - if (!v) - return // Delete failed, skip S3 cleanup - if (v.error) { - cloudlog({ requestId: c.get('requestId'), message: 'error checking manifest references', error: v.error }) - return // Don't delete S3 if we can't confirm no other references - } - if (v.data) { - // Other versions still use this file, S3 cleanup not needed - return - } - // No other versions use this file, move it to the R2 trash prefix. - cloudlog({ requestId: c.get('requestId'), message: 'moving manifest file to R2 trash', s3_path: entry.s3_path }) - return s3.moveObjectToTrash(c, entry.s3_path) - .then((moved) => { - if (!moved) { - throw simpleError('cannot_move_manifest_s3_to_trash', 'Cannot move S3 object for deleted manifest file to trash', { id: entry.id, s3_path: entry.s3_path }) - } - }) - }), + const startedWithRows = manifestEntries.length > 0 + + if (startedWithRows) { + for (let i = 0; i < manifestEntries.length; i += MANIFEST_TRASH_CONCURRENCY) { + const batch = manifestEntries.slice(i, i + MANIFEST_TRASH_CONCURRENCY) + await Promise.all(batch.map(async (entry) => { + const entryPg = getPgClient(c, false) + try { + await entryPg.query('BEGIN') + // Serialize shared-hash cleanup across concurrent deleted versions. + await entryPg.query( + `SELECT pg_advisory_xact_lock(hashtextextended($1 || chr(0) || $2, 0))`, + [entry.file_hash, entry.file_name], ) - } - } - await Promise.all(promisesMoveToTrash) - - // After deleting manifest entries, update manifest_count and decrement manifest_bundle_count - const updatePgClient = getPgClient(c, false) - try { - await updatePgClient.query( - `UPDATE app_versions SET manifest_count = 0, manifest = NULL WHERE id = $1`, - [record.id], - ) - - // Only decrement if this version had manifests - if (manifestCount > 0) { - await updatePgClient.query( - `UPDATE apps - SET manifest_bundle_count = GREATEST(manifest_bundle_count - 1, 0), - updated_at = now() - WHERE app_id = $1`, - [record.app_id], + + if (entry.s3_path) { + const refs = await entryPg.query( + `SELECT 1 AS ok + FROM public.manifest + WHERE file_hash = $1 + AND file_name = $2 + AND app_version_id <> $3 + LIMIT 1`, + [entry.file_hash, entry.file_name, record.id], + ) + + if (refs.rows.length === 0) { + const moved = await s3.moveObjectToTrash(c, entry.s3_path) + if (!moved) { + throw simpleError('cannot_move_manifest_s3_to_trash', 'Cannot move S3 object for deleted manifest file to trash', { + id: entry.id, + s3_path: entry.s3_path, + }) + } + } + } + + // Only delete the DB row after R2 is handled (or shared and kept). + await entryPg.query( + `DELETE FROM public.manifest WHERE id = $1`, + [entry.id], ) + await entryPg.query('COMMIT') } + catch (error) { + try { + await entryPg.query('ROLLBACK') + } + catch { + // ignore rollback errors + } + throw error + } + finally { + await closeClient(c, entryPg) + } + })) + } + } + + const writePgClient = getPgClient(c, false) + try { + await writePgClient.query('BEGIN') + try { + const remaining = await writePgClient.query( + `SELECT COUNT(*)::int AS count FROM public.manifest WHERE app_version_id = $1`, + [record.id], + ) + const remainingCount = Number(remaining.rows[0]?.count ?? 0) + if (remainingCount > 0) { + throw simpleError('manifest_cleanup_incomplete', 'Manifest rows still present after trash/delete pass', { + id: record.id, + remainingCount, + }) } - catch (error) { - cloudlog({ requestId: c.get('requestId'), message: 'error update counters on delete', error }) - } - finally { - await closeClient(c, updatePgClient) - } + + await writePgClient.query( + `WITH prev AS ( + SELECT id, app_id, manifest_count, (manifest IS NOT NULL) AS has_json + FROM public.app_versions + WHERE id = $1 + FOR UPDATE + ), + upd AS ( + UPDATE public.app_versions AS av + SET manifest_count = 0, + manifest = NULL + FROM prev + WHERE av.id = prev.id + AND (prev.manifest_count > 0 OR prev.has_json OR $2::boolean) + RETURNING prev.app_id, prev.manifest_count AS prev_count + ) + UPDATE public.apps AS a + SET manifest_bundle_count = GREATEST(a.manifest_bundle_count - 1, 0), + updated_at = now() + FROM upd + WHERE a.app_id = upd.app_id + AND (upd.prev_count > 0 OR $2::boolean)`, + [record.id, startedWithRows], + ) + + await writePgClient.query('COMMIT') + } + catch (error) { + await writePgClient.query('ROLLBACK') + throw error } } catch (error) { - cloudlog({ requestId: c.get('requestId'), message: 'error deleting manifest entries', error }) + cloudlog({ requestId: c.get('requestId'), message: 'error finalizing manifest cleanup', error, id: record.id }) + throw error } finally { - await closeClient(c, pgClient) + await closeClient(c, writePgClient) } } export async function deleteIt(c: Context, record: Database['public']['Tables']['app_versions']['Row']) { cloudlog({ requestId: c.get('requestId'), message: 'Delete', r2_path: record.r2_path }) + // Manifest files: trash R2 first, then drop DB rows. Must finish before ACK. + await deleteManifest(c, record) + + const { data, error: dbError } = await supabaseAdmin(c) + .from('app_versions_meta') + .select() + .eq('id', record.id) + .single() + if (dbError || !data) { + cloudlog({ requestId: c.get('requestId'), message: 'Cannot find version meta', id: record.id }) + } + else { + const { error: errorCreateStatsMeta } = await createStatsMeta(c, record.app_id, record.id, -data.size) + if (errorCreateStatsMeta) + cloudlog({ requestId: c.get('requestId'), message: 'error createStatsMeta', error: errorCreateStatsMeta }) + + const { error: errorUpdate } = await supabaseAdmin(c) + .from('app_versions_meta') + .update({ size: 0 }) + .eq('id', record.id) + if (errorUpdate) { + cloudlog({ requestId: c.get('requestId'), message: 'error', error: errorUpdate }) + throw simpleError('cannot_update_version_meta', 'Cannot update version metadata for deleted version', { id: record.id }, errorUpdate) + } + } + + // Bundle zip: move to lifecycle trash. Retry via queue if this fails; manifests already cleared. if (record.r2_path) { let moved = false try { @@ -439,30 +515,6 @@ export async function deleteIt(c: Context, record: Database['public']['Tables'][ cloudlog({ requestId: c.get('requestId'), message: 'No r2 path for deleted version', id: record.id }) } - const { data, error: dbError } = await supabaseAdmin(c) - .from('app_versions_meta') - .select() - .eq('id', record.id) - .single() - if (dbError || !data) { - cloudlog({ requestId: c.get('requestId'), message: 'Cannot find version meta', id: record.id }) - return c.json(BRES) - } - const { error: errorCreateStatsMeta } = await createStatsMeta(c, record.app_id, record.id, -data.size) - if (errorCreateStatsMeta) - cloudlog({ requestId: c.get('requestId'), message: 'error createStatsMeta', error: errorCreateStatsMeta }) - // set app_versions_meta versionSize = 0 - const { error: errorUpdate } = await supabaseAdmin(c) - .from('app_versions_meta') - .update({ size: 0 }) - .eq('id', record.id) - if (errorUpdate) { - cloudlog({ requestId: c.get('requestId'), message: 'error', error: errorUpdate }) - throw simpleError('cannot_update_version_meta', 'Cannot update version metadata for deleted version', { id: record.id }, errorUpdate) - } - - await deleteManifest(c, record) - return c.json(BRES) } @@ -505,4 +557,5 @@ app.post('/', middlewareAPISecret, triggerValidator('app_versions', 'UPDATE'), a export const onVersionUpdateTestUtils = { getDeletedVersionAction, + deleteManifest, } diff --git a/supabase/functions/_backend/triggers/queue_consumer.ts b/supabase/functions/_backend/triggers/queue_consumer.ts index e73b025bda..4f5c9c5633 100644 --- a/supabase/functions/_backend/triggers/queue_consumer.ts +++ b/supabase/functions/_backend/triggers/queue_consumer.ts @@ -23,9 +23,10 @@ const DEFAULT_QUEUE_VISIBILITY_TIMEOUT_SECONDS = 120 const VERSION_QUEUE_VISIBILITY_TIMEOUT_SECONDS = 900 const MANIFEST_QUEUE_VISIBILITY_TIMEOUT_SECONDS = 900 const QUEUE_HTTP_TIMEOUT_MS = 15_000 -const VERSION_QUEUE_HTTP_TIMEOUT_MS = 60_000 +const VERSION_QUEUE_HTTP_TIMEOUT_MS = 300_000 // large deleted manifests: trash then DB delete const HEALTHCHECK_HTTP_TIMEOUT_MS = 8_000 export const MAX_QUEUE_READS = 5 +const VERSION_QUEUE_MAX_READS = 30 // deleted manifests can need many partial trash/delete passes const DISCORD_IGNORED_ERROR_CODES = new Set(['version_not_found', 'no_channel']) const integerLikeSchema = type('number.integer').or(type('string.numeric.parse |> number.integer')) @@ -137,9 +138,9 @@ function getQueueMessageTrace(functionName: string, body: Record { - if (detail.read_count < MAX_QUEUE_READS) + if (detail.read_count < retryBudget) return false return !detail.error_code || !DISCORD_IGNORED_ERROR_CODES.has(detail.error_code) }) @@ -281,6 +282,12 @@ function getQueueHttpTimeoutMs(functionName: string): number { return QUEUE_HTTP_TIMEOUT_MS } +function getQueueMaxReads(queueName: string): number { + if (isVersionQueueFunction(queueName)) + return VERSION_QUEUE_MAX_READS + return MAX_QUEUE_READS +} + function resolveFunctionUrl(c: Context, function_name: string, function_type: string | null | undefined): string { const cfPpUrl = getEnv(c, 'CLOUDFLARE_PP_FUNCTION_URL') const cfUrl = getEnv(c, 'CLOUDFLARE_FUNCTION_URL') @@ -675,6 +682,7 @@ async function reportQueueFailures(c: Context, queueName: string, messagesFailed cloudlog({ requestId: c.get('requestId'), message: `[${queueName}] Failed to process ${messagesFailed.length} messages.` }) + const retryBudget = getQueueMaxReads(queueName) const timestamp = new Date().toISOString() const failureDetails = messagesFailed.map(msg => ({ function_name: msg.message?.function_name ?? 'unknown', @@ -692,7 +700,7 @@ async function reportQueueFailures(c: Context, queueName: string, messagesFailed target_url: msg.targetUrl ?? undefined, })) - const actionableFailures = getActionableQueueFailures(failureDetails) + const actionableFailures = getActionableQueueFailures(failureDetails, retryBudget) const groupedByFunction = actionableFailures.reduce((acc, detail) => { const key = detail.function_name acc[key] ??= [] @@ -727,7 +735,7 @@ async function reportQueueFailures(c: Context, queueName: string, messagesFailed const messageInfo = detail.error_message ? ` | ${truncateDiscordField(detail.error_message.replace(/\s+/g, ' ').trim(), 180)}` : '' const durationInfo = typeof detail.duration_ms === 'number' ? ` | ${detail.duration_ms}ms` : '' const targetInfo = detail.target_url ? ` | Target: ${truncateDiscordField(detail.target_url, 120)}` : '' - return `**${detail.function_name}** | Status: ${detail.status} | Read: ${detail.read_count}/${MAX_QUEUE_READS}${durationInfo}${errorInfo}${messageInfo}${targetInfo} | [CF Logs](${cfLogUrl})` + return `**${detail.function_name}** | Status: ${detail.status} | Read: ${detail.read_count}/${retryBudget}${durationInfo}${errorInfo}${messageInfo}${targetInfo} | [CF Logs](${cfLogUrl})` }).join('\n')), inline: false, }, @@ -770,7 +778,7 @@ async function reportQueueFailures(c: Context, queueName: string, messagesFailed cloudlog({ requestId: c.get('requestId'), message: `[${queueName}] Suppressed Discord alert for retryable or ignored queue failures.`, - retryingFailures: failureDetails.filter(detail => detail.read_count < MAX_QUEUE_READS).length, + retryingFailures: failureDetails.filter(detail => detail.read_count < retryBudget).length, ignoredErrors: Array.from(DISCORD_IGNORED_ERROR_CODES), }) } @@ -799,8 +807,9 @@ async function processQueue(c: Context, db: ReturnType, queu } } + const retryBudget = getQueueMaxReads(queueName) const [messagesToProcess, messagesToSkip] = messages.reduce((acc, message) => { - acc[message.read_ct <= MAX_QUEUE_READS ? 0 : 1].push(message) + acc[message.read_ct <= retryBudget ? 0 : 1].push(message) return acc }, [[], []] as [typeof messages, typeof messages]) @@ -812,7 +821,7 @@ async function processQueue(c: Context, db: ReturnType, queu processingCount: messagesToProcess.length, skippedCount: messagesToSkip.length, concurrency: processConcurrency, - retryBudget: MAX_QUEUE_READS, + retryBudget, }) // Archive messages after the configured retry budget is exhausted. @@ -822,7 +831,7 @@ async function processQueue(c: Context, db: ReturnType, queu message: `[${queueName}] Archiving messages that exceeded the retry budget.`, queueName, archiveCount: messagesToSkip.length, - retryBudget: MAX_QUEUE_READS, + retryBudget, }) await archive_queue_messages(c, db, queueName, messagesToSkip.map(msg => msg.msg_id)) } @@ -978,7 +987,7 @@ export async function http_post_helper( headers['x-capgo-queue-name'] = metadata.queueName headers['x-capgo-queue-msg-id'] = String(metadata.msgId) headers['x-capgo-queue-read-count'] = String(metadata.readCount) - headers['x-capgo-queue-max-reads'] = String(MAX_QUEUE_READS) + headers['x-capgo-queue-max-reads'] = String(metadata ? getQueueMaxReads(metadata.queueName) : MAX_QUEUE_READS) } if (waitForCompletion) headers[WAIT_FOR_COMPLETION_HEADER] = 'true' @@ -1243,6 +1252,7 @@ export const __queueConsumerTestUtils__ = { getQueueAckChunkSize, getQueueHttpConcurrency, getQueueHttpTimeoutMs, + getQueueMaxReads, httpExceptionToQueueResponse, getQueueVisibilityTimeout, normalizeQueueFunctionType, diff --git a/supabase/functions/_backend/utils/s3.ts b/supabase/functions/_backend/utils/s3.ts index 42dda690b8..e2d966f853 100644 --- a/supabase/functions/_backend/utils/s3.ts +++ b/supabase/functions/_backend/utils/s3.ts @@ -141,10 +141,64 @@ function shouldUseSizeRangeFallback(size: number, headError: unknown): boolean { return !size && !isMissingObjectError(headError) } +type ObjectPresence = 'present' | 'absent' | 'unknown' + +async function getObjectPresence(c: Context, fileId: string | null): Promise { + if (!fileId) + return 'absent' + + try { + const client = initS3(c) + const url = await client.getPresignedUrl('HEAD', fileId) + const response = await fetch(url, { + method: 'HEAD', + signal: AbortSignal.timeout(10_000), + }) + await response.body?.cancel() + + if (response.status === 404) + return 'absent' + if (response.status === 200) + return 'present' + + cloudlogErr({ + requestId: c.get('requestId'), + message: 'getObjectPresence unexpected HEAD status', + fileId, + status: response.status, + statusText: response.statusText, + }) + return 'unknown' + } + catch (error) { + if (isMissingObjectError(error)) + return 'absent' + cloudlogErr({ + requestId: c.get('requestId'), + message: 'getObjectPresence failed', + fileId, + error: serializeStorageError(error), + }) + return 'unknown' + } +} + async function moveObjectToTrash(c: Context, fileId: string) { if (fileId.startsWith(R2_TRASH_PREFIX)) return true + // Only skip copy on a definitive absent object. Unknown HEAD must fail closed + // so callers keep DB tracking until trash succeeds. + const presence = await getObjectPresence(c, fileId) + if (presence === 'absent') { + cloudlog({ requestId: c.get('requestId'), message: 'R2 object missing before trash move, skip copy', fileId }) + return true + } + if (presence === 'unknown') { + cloudlogErr({ requestId: c.get('requestId'), message: 'R2 object presence unknown, refuse trash skip', fileId }) + return false + } + const client = initS3(c) const trashPath = getTrashPath(fileId) try { @@ -155,7 +209,7 @@ async function moveObjectToTrash(c: Context, fileId: string) { } catch (error) { if (isMissingObjectError(error)) { - cloudlog({ requestId: c.get('requestId'), message: 'R2 object already missing before trash move', fileId, error: serializeStorageError(error) }) + cloudlog({ requestId: c.get('requestId'), message: 'R2 object disappeared during trash move', fileId, error: serializeStorageError(error) }) return true } diff --git a/supabase/migrations/20260722095157_sweep_deleted_version_manifests.sql b/supabase/migrations/20260722095157_sweep_deleted_version_manifests.sql new file mode 100644 index 0000000000..69119503e7 --- /dev/null +++ b/supabase/migrations/20260722095157_sweep_deleted_version_manifests.sql @@ -0,0 +1,146 @@ +-- Sweep soft-deleted app_versions + +CREATE INDEX IF NOT EXISTS idx_app_versions_deleted_with_manifest + ON public.app_versions (id) + WHERE deleted = true AND manifest_count > 0; + +CREATE INDEX IF NOT EXISTS idx_app_versions_deleted_at_id + ON public.app_versions (deleted_at, id) + WHERE deleted = true; + +-- Sweeps soft-deleted app_versions that still have manifest rows or stale counters. +-- Touches a bounded batch so on_version_update re-runs cleanup_manifest. +-- Also zeros stale manifest_count when no rows remain. + +CREATE OR REPLACE FUNCTION "public"."sweep_deleted_version_manifests"("p_batch_size" integer DEFAULT 100) +RETURNS bigint +LANGUAGE plpgsql +SECURITY DEFINER +SET search_path = '' +AS $$ +DECLARE + stale_fixed bigint := 0; + requeued bigint := 0; +BEGIN + IF p_batch_size IS NULL OR p_batch_size < 1 THEN + p_batch_size := 100; + END IF; + + -- Fix stale counters: deleted versions with manifest_count > 0 but no rows. + WITH stale AS ( + SELECT av.id, av.app_id + FROM public.app_versions AS av + WHERE av.deleted = true + AND av.manifest_count > 0 + AND NOT EXISTS ( + SELECT 1 + FROM public.manifest AS m + WHERE m.app_version_id = av.id + ) + ORDER BY av.deleted_at NULLS LAST, av.id + LIMIT p_batch_size + ), + cleared AS ( + UPDATE public.app_versions AS av + SET manifest_count = 0, + manifest = NULL, + updated_at = now() + FROM stale + WHERE av.id = stale.id + RETURNING stale.app_id + ), + app_counts AS ( + SELECT app_id, COUNT(*)::int AS cleared_count + FROM cleared + GROUP BY app_id + ) + UPDATE public.apps AS a + SET manifest_bundle_count = GREATEST(a.manifest_bundle_count - app_counts.cleared_count, 0), + updated_at = now() + FROM app_counts + WHERE a.app_id = app_counts.app_id; + + GET DIAGNOSTICS stale_fixed = ROW_COUNT; + + -- Re-queue deleted versions that still have manifest rows by touching them. + -- on_version_update trigger enqueues cleanup when deleted_at is unchanged and + -- manifest_count > 0. + -- Start from deleted versions (bounded) and probe manifest via app_version_id index. + WITH candidates AS ( + SELECT av.id + FROM public.app_versions AS av + WHERE av.deleted = true + AND EXISTS ( + SELECT 1 + FROM public.manifest AS m + WHERE m.app_version_id = av.id + ) + ORDER BY av.deleted_at NULLS LAST, av.id + LIMIT p_batch_size + ) + UPDATE public.app_versions AS av + SET manifest_count = GREATEST(av.manifest_count, 1), + updated_at = now() + FROM candidates + WHERE av.id = candidates.id; + + GET DIAGNOSTICS requeued = ROW_COUNT; + + IF stale_fixed > 0 OR requeued > 0 THEN + RAISE NOTICE 'sweep_deleted_version_manifests: stale_counters=% requeued=%', stale_fixed, requeued; + END IF; + + RETURN requeued; +END; +$$; + +ALTER FUNCTION public.sweep_deleted_version_manifests(integer) OWNER TO postgres; +REVOKE ALL ON FUNCTION public.sweep_deleted_version_manifests(integer) FROM PUBLIC; +GRANT ALL ON FUNCTION public.sweep_deleted_version_manifests(integer) TO service_role; + +COMMENT ON FUNCTION public.sweep_deleted_version_manifests(integer) IS + 'Bounded sweeper for soft-deleted versions with leftover manifest rows or stale manifest_count. Re-touches rows so on_version_update runs cleanup_manifest (DB delete + R2 trash).'; + +INSERT INTO public.cron_tasks ( + name, + description, + task_type, + target, + batch_size, + payload, + second_interval, + minute_interval, + hour_interval, + run_at_hour, + run_at_minute, + run_at_second, + run_on_dow, + run_on_day, + enabled, + healthcheck_url +) VALUES ( + 'sweep_deleted_version_manifests', + 'Re-queue soft-deleted versions that still have manifest rows; zero stale manifest_count', + 'function', + 'public.sweep_deleted_version_manifests(100)', + NULL, + NULL, + NULL, + 15, + NULL, + NULL, + NULL, + 0, + NULL, + NULL, + true, + NULL +) +ON CONFLICT (name) DO UPDATE +SET + description = EXCLUDED.description, + task_type = EXCLUDED.task_type, + target = EXCLUDED.target, + minute_interval = EXCLUDED.minute_interval, + enabled = EXCLUDED.enabled, + updated_at = now(); diff --git a/tests/on-version-update-cleanup.unit.test.ts b/tests/on-version-update-cleanup.unit.test.ts index 7b3a5aa299..3de90ada57 100644 --- a/tests/on-version-update-cleanup.unit.test.ts +++ b/tests/on-version-update-cleanup.unit.test.ts @@ -4,29 +4,22 @@ const { appVersionsMetaSelectEq, appVersionsMetaUpdate, appVersionsMetaUpdateEq, + callOrder, closeClient, createStatsMeta, deleteObject, getDrizzleClient, getPgClient, - manifestDeleteEq, - manifestReferenceMaybeSingle, manifestSelectWhere, moveObjectToTrash, pgQuery, supabaseAdmin, } = vi.hoisted(() => { + const callOrder: string[] = [] const appVersionsMetaSelectEq = vi.fn() const appVersionsMetaSelect = vi.fn(() => ({ eq: appVersionsMetaSelectEq })) const appVersionsMetaUpdateEq = vi.fn() const appVersionsMetaUpdate = vi.fn(() => ({ eq: appVersionsMetaUpdateEq })) - const manifestDeleteEq = vi.fn() - const manifestDelete = vi.fn(() => ({ eq: manifestDeleteEq })) - const manifestReferenceMaybeSingle = vi.fn() - const manifestReferenceLimit = vi.fn(() => ({ maybeSingle: manifestReferenceMaybeSingle })) - const manifestReferenceEqFileName = vi.fn(() => ({ limit: manifestReferenceLimit })) - const manifestReferenceEqFileHash = vi.fn(() => ({ eq: manifestReferenceEqFileName })) - const manifestSelect = vi.fn(() => ({ eq: manifestReferenceEqFileHash })) const supabaseFrom = vi.fn((table: string) => { if (table === 'app_versions_meta') { return { @@ -34,21 +27,40 @@ const { update: appVersionsMetaUpdate, } } - if (table === 'manifest') { - return { - delete: manifestDelete, - select: manifestSelect, - } - } return {} }) const manifestSelectWhere = vi.fn(async (): Promise => []) - const pgQuery = vi.fn(async () => ({ rows: [] })) + const pgQuery = vi.fn(async (sql: string, params?: any[]) => { + if (sql === 'BEGIN') + callOrder.push('begin') + if (sql.includes('pg_advisory_xact_lock')) + callOrder.push('lock') + if (sql.includes('SELECT 1 AS ok')) + return { rows: [], rowCount: 0 } + if (sql.includes('DELETE FROM public.manifest WHERE id')) { + callOrder.push(`db_delete_row:${params?.[0]}`) + return { rows: [], rowCount: 1 } + } + if (sql === 'COMMIT') + callOrder.push('commit_entry') + if (sql === 'ROLLBACK') + callOrder.push('rollback_entry') + if (sql.includes('SELECT COUNT(*)')) + return { rows: [{ count: 0 }], rowCount: 1 } + if (sql.includes('WITH prev AS')) + return { rows: [], rowCount: 1 } + return { rows: [], rowCount: 0 } + }) + const moveObjectToTrash = vi.fn(async (..._args: any[]) => { + callOrder.push('r2_trash') + return true + }) return { appVersionsMetaSelectEq, appVersionsMetaUpdate, appVersionsMetaUpdateEq, + callOrder, closeClient: vi.fn(), createStatsMeta: vi.fn(), deleteObject: vi.fn(), @@ -60,13 +72,10 @@ const { })), })), getPgClient: vi.fn(() => ({ query: pgQuery })), - manifestDeleteEq, - manifestReferenceMaybeSingle, manifestSelectWhere, - moveObjectToTrash: vi.fn(), + moveObjectToTrash, pgQuery, supabaseAdmin: vi.fn(() => ({ from: supabaseFrom })), - supabaseFrom, } }) @@ -97,7 +106,7 @@ vi.mock('../supabase/functions/_backend/utils/logging.ts', () => ({ cloudlogErr: vi.fn(), })) -const { deleteIt } = await import('../supabase/functions/_backend/triggers/on_version_update.ts') +const { deleteIt, onVersionUpdateTestUtils } = await import('../supabase/functions/_backend/triggers/on_version_update.ts') function createContext() { return { @@ -111,6 +120,7 @@ function createVersion(overrides: Record = {}) { app_id: 'com.cleanup.test', id: 123, manifest: null, + manifest_count: 0, name: '1.0.0', r2_path: 'orgs/org-1/apps/com.cleanup.test/1.0.0.zip', storage_provider: 'r2', @@ -118,16 +128,47 @@ function createVersion(overrides: Record = {}) { } as any } +function makeEntries(count: number) { + return Array.from({ length: count }, (_, i) => ({ + id: 1000 + i, + file_hash: `hash-${i}`, + file_name: `file-${i}.js`, + s3_path: `orgs/org-1/apps/com.cleanup.test/delta/file-${i}.js`, + })) +} + describe('on_version_update deleted version cleanup', () => { beforeEach(() => { vi.clearAllMocks() + callOrder.length = 0 deleteObject.mockResolvedValue(true) - moveObjectToTrash.mockResolvedValue(true) + moveObjectToTrash.mockImplementation(async () => { + callOrder.push('r2_trash') + return true + }) createStatsMeta.mockResolvedValue({ error: null }) manifestSelectWhere.mockResolvedValue([]) - manifestDeleteEq.mockResolvedValue({ error: null }) - manifestReferenceMaybeSingle.mockResolvedValue({ data: null, error: null }) - pgQuery.mockResolvedValue({ rows: [] }) + pgQuery.mockImplementation(async (sql: string, params?: any[]) => { + if (sql === 'BEGIN') + callOrder.push('begin') + if (sql.includes('pg_advisory_xact_lock')) + callOrder.push('lock') + if (sql.includes('SELECT 1 AS ok')) + return { rows: [], rowCount: 0 } + if (sql.includes('DELETE FROM public.manifest WHERE id')) { + callOrder.push(`db_delete_row:${params?.[0]}`) + return { rows: [], rowCount: 1 } + } + if (sql === 'COMMIT') + callOrder.push('commit_entry') + if (sql === 'ROLLBACK') + callOrder.push('rollback_entry') + if (sql.includes('SELECT COUNT(*)')) + return { rows: [{ count: 0 }], rowCount: 1 } + if (sql.includes('WITH prev AS')) + return { rows: [], rowCount: 1 } + return { rows: [], rowCount: 0 } + }) appVersionsMetaSelectEq.mockReturnValue({ single: vi.fn(async () => ({ data: { size: 1234 }, error: null })), }) @@ -136,47 +177,181 @@ describe('on_version_update deleted version cleanup', () => { it('moves the bundle to trash and clears stored size for soft-deleted versions', async () => { const response = await deleteIt(createContext(), createVersion()) - expect(response.status).toBe(200) expect(moveObjectToTrash).toHaveBeenCalledWith(expect.anything(), 'orgs/org-1/apps/com.cleanup.test/1.0.0.zip') - expect(deleteObject).not.toHaveBeenCalled() expect(appVersionsMetaUpdate).toHaveBeenCalledWith({ size: 0 }) - expect(appVersionsMetaUpdateEq).toHaveBeenCalledWith('id', 123) - expect(createStatsMeta).toHaveBeenCalledWith(expect.anything(), 'com.cleanup.test', 123, -1234) }) - it('still clears stale metadata when the deleted version has no bundle path', async () => { - const response = await deleteIt(createContext(), createVersion({ r2_path: null })) + it('locks, trashes R2, then deletes each manifest DB row', async () => { + manifestSelectWhere.mockResolvedValue(makeEntries(1)) + + await deleteIt(createContext(), createVersion({ r2_path: null, manifest_count: 1 })) + + expect(callOrder).toContain('lock') + expect(callOrder.indexOf('r2_trash')).toBeGreaterThan(callOrder.indexOf('lock')) + expect(callOrder.indexOf('db_delete_row:1000')).toBeGreaterThan(callOrder.indexOf('r2_trash')) + expect(pgQuery).toHaveBeenCalledWith(expect.stringContaining('WITH prev AS'), expect.any(Array)) + }) + + it('does not delete DB rows when R2 trash fails', async () => { + manifestSelectWhere.mockResolvedValue(makeEntries(1)) + moveObjectToTrash.mockImplementation(async () => { + callOrder.push('r2_trash') + return false + }) + + await expect(deleteIt(createContext(), createVersion({ r2_path: null, manifest_count: 1 }))).rejects.toThrow( + 'Cannot move S3 object for deleted manifest file to trash', + ) + expect(callOrder).toContain('r2_trash') + expect(callOrder.some(v => v.startsWith('db_delete_row:'))).toBe(false) + expect(callOrder).toContain('rollback_entry') + }) + + it('skips R2 trash when another version still references the file, then deletes the row', async () => { + manifestSelectWhere.mockResolvedValue(makeEntries(1)) + pgQuery.mockImplementation((async (sql: string, params?: any[]) => { + if (sql === 'BEGIN') + callOrder.push('begin') + if (sql.includes('pg_advisory_xact_lock')) + callOrder.push('lock') + if (sql.includes('SELECT 1 AS ok')) + return { rows: [{ ok: 1 }], rowCount: 1 } + if (sql.includes('DELETE FROM public.manifest WHERE id')) { + callOrder.push(`db_delete_row:${params?.[0]}`) + return { rows: [], rowCount: 1 } + } + if (sql === 'COMMIT') + callOrder.push('commit_entry') + if (sql.includes('SELECT COUNT(*)')) + return { rows: [{ count: 0 }], rowCount: 1 } + if (sql.includes('WITH prev AS')) + return { rows: [], rowCount: 1 } + return { rows: [], rowCount: 0 } + }) as any) + + await deleteIt(createContext(), createVersion({ r2_path: null, manifest_count: 1 })) - expect(response.status).toBe(200) expect(moveObjectToTrash).not.toHaveBeenCalled() - expect(appVersionsMetaUpdate).toHaveBeenCalledWith({ size: 0 }) - expect(createStatsMeta).toHaveBeenCalledWith(expect.anything(), 'com.cleanup.test', 123, -1234) + expect(callOrder).toContain('db_delete_row:1000') }) - it('moves unreferenced manifest files to trash instead of hard deleting them', async () => { - manifestSelectWhere.mockResolvedValue([{ - app_version_id: 123, - file_hash: 'manifest-hash', - file_name: 'index.js', - id: 456, - s3_path: 'orgs/org-1/apps/com.cleanup.test/manifest/index.js', - }]) + it('still clears manifests when version meta is missing', async () => { + appVersionsMetaSelectEq.mockReturnValue({ + single: vi.fn(async () => ({ data: null, error: { message: 'not found' } })), + }) + manifestSelectWhere.mockResolvedValue(makeEntries(1)) - const response = await deleteIt(createContext(), createVersion({ r2_path: null })) + const response = await deleteIt(createContext(), createVersion({ manifest_count: 1 })) expect(response.status).toBe(200) - expect(moveObjectToTrash).toHaveBeenCalledWith(expect.anything(), 'orgs/org-1/apps/com.cleanup.test/manifest/index.js') - expect(deleteObject).not.toHaveBeenCalled() - expect(pgQuery).toHaveBeenCalledWith('UPDATE app_versions SET manifest_count = 0, manifest = NULL WHERE id = $1', [123]) - expect(pgQuery).toHaveBeenCalledWith(expect.stringContaining('manifest_bundle_count = GREATEST(manifest_bundle_count - 1, 0)'), ['com.cleanup.test']) + expect(callOrder.indexOf('r2_trash')).toBeLessThan(callOrder.indexOf('db_delete_row:1000')) + expect(moveObjectToTrash).toHaveBeenCalledWith(expect.anything(), 'orgs/org-1/apps/com.cleanup.test/1.0.0.zip') }) - it('keeps the queue retryable when moving the bundle to trash fails', async () => { - moveObjectToTrash.mockResolvedValue(false) + it('keeps the queue retryable when moving the bundle to trash fails after manifest cleanup', async () => { + manifestSelectWhere.mockResolvedValue(makeEntries(1)) + moveObjectToTrash.mockImplementation(async (_c: unknown, path: string) => { + callOrder.push(path.includes('.zip') ? 'bundle_trash' : 'r2_trash') + return !path.includes('.zip') + }) + + await expect(deleteIt(createContext(), createVersion({ manifest_count: 1 }))).rejects.toThrow( + 'Cannot move S3 object for deleted version to trash', + ) + expect(callOrder).toContain('r2_trash') + expect(callOrder).toContain('db_delete_row:1000') + }) + + it('throws when rows remain after the trash/delete pass', async () => { + manifestSelectWhere.mockResolvedValue(makeEntries(1)) + pgQuery.mockImplementation(async (sql: string, params?: any[]) => { + if (sql === 'BEGIN') + callOrder.push('begin') + if (sql.includes('pg_advisory_xact_lock')) + callOrder.push('lock') + if (sql.includes('SELECT 1 AS ok')) + return { rows: [], rowCount: 0 } + if (sql.includes('DELETE FROM public.manifest WHERE id')) { + callOrder.push(`db_delete_row:${params?.[0]}`) + return { rows: [], rowCount: 1 } + } + if (sql === 'COMMIT') + callOrder.push('commit_entry') + if (sql === 'ROLLBACK') + callOrder.push('rollback_entry') + if (sql.includes('SELECT COUNT(*)')) + return { rows: [{ count: 2 }], rowCount: 1 } + return { rows: [], rowCount: 0 } + }) + + await expect(deleteIt(createContext(), createVersion({ r2_path: null, manifest_count: 1 }))).rejects.toThrow( + 'Manifest rows still present after trash/delete pass', + ) + expect(callOrder).toContain('rollback_entry') + }) + + it('routes already-deleted versions with leftover counts to cleanup_manifest', () => { + expect(onVersionUpdateTestUtils.getDeletedVersionAction( + createVersion({ deleted_at: '2026-01-01T00:00:00Z', manifest_count: 3 }), + createVersion({ deleted_at: '2026-01-01T00:00:00Z', manifest_count: 3 }), + )).toBe('cleanup_manifest') + }) +}) - await expect(deleteIt(createContext(), createVersion())).rejects.toThrow('Cannot move S3 object for deleted version to trash') - expect(appVersionsMetaUpdate).not.toHaveBeenCalled() - expect(createStatsMeta).not.toHaveBeenCalled() +describe('on_version_update manifest cleanup load', () => { + beforeEach(() => { + vi.clearAllMocks() + callOrder.length = 0 + createStatsMeta.mockResolvedValue({ error: null }) + appVersionsMetaSelectEq.mockReturnValue({ + single: vi.fn(async () => ({ data: { size: 0 }, error: null })), + }) + appVersionsMetaUpdateEq.mockResolvedValue({ error: null }) + moveObjectToTrash.mockImplementation(async () => { + callOrder.push('r2_trash') + return true + }) + pgQuery.mockImplementation(async (sql: string, params?: any[]) => { + if (sql.includes('SELECT 1 AS ok')) + return { rows: [], rowCount: 0 } + if (sql.includes('DELETE FROM public.manifest WHERE id')) { + callOrder.push(`db_delete_row:${params?.[0]}`) + return { rows: [], rowCount: 1 } + } + if (sql.includes('SELECT COUNT(*)')) + return { rows: [{ count: 0 }], rowCount: 1 } + if (sql.includes('WITH prev AS')) + return { rows: [], rowCount: 1 } + return { rows: [], rowCount: 0 } + }) }) + + it('handles 5000-file manifests with R2 before every DB delete', async () => { + manifestSelectWhere.mockResolvedValue(makeEntries(5000)) + + const response = await deleteIt(createContext(), createVersion({ r2_path: null, manifest_count: 5000 })) + + expect(response.status).toBe(200) + expect(moveObjectToTrash).toHaveBeenCalledTimes(5000) + expect(callOrder.filter(v => v.startsWith('db_delete_row:'))).toHaveLength(5000) + expect(pgQuery).toHaveBeenCalledWith(expect.stringContaining('WITH prev AS'), expect.any(Array)) + }, 60_000) + + it('keeps remaining rows retryable when one file in a large batch fails trash', async () => { + manifestSelectWhere.mockResolvedValue(makeEntries(200)) + moveObjectToTrash.mockImplementation(async (_c: unknown, path: string) => { + callOrder.push('r2_trash') + if (path.endsWith('file-150.js')) + return false + return true + }) + + await expect(deleteIt(createContext(), createVersion({ r2_path: null, manifest_count: 200 }))).rejects.toThrow( + 'Cannot move S3 object for deleted manifest file to trash', + ) + + const deletedIds = callOrder.filter(v => v.startsWith('db_delete_row:')).map(v => Number(v.split(':')[1])) + expect(deletedIds).not.toContain(1150) + }, 30_000) }) diff --git a/tests/queue-consumer-message-shape.unit.test.ts b/tests/queue-consumer-message-shape.unit.test.ts index 51cd2b27b1..352bd09135 100644 --- a/tests/queue-consumer-message-shape.unit.test.ts +++ b/tests/queue-consumer-message-shape.unit.test.ts @@ -158,7 +158,9 @@ describe('queue_consumer legacy message compatibility', () => { expect(__queueConsumerTestUtils__.getQueueVisibilityTimeout('on_manifest_create')).toBe(900) expect(__queueConsumerTestUtils__.getQueueVisibilityTimeout('cron_email')).toBe(120) expect(__queueConsumerTestUtils__.getQueueVisibilityTimeout('on_version_update')).toBe(900) - expect(__queueConsumerTestUtils__.getQueueHttpTimeoutMs('on_version_update')).toBe(60_000) + expect(__queueConsumerTestUtils__.getQueueHttpTimeoutMs('on_version_update')).toBe(300_000) + expect(__queueConsumerTestUtils__.getQueueMaxReads('on_version_update')).toBe(30) + expect(__queueConsumerTestUtils__.getQueueMaxReads('on_manifest_create')).toBe(5) expect(__queueConsumerTestUtils__.getQueueHttpTimeoutMs('cron_email')).toBe(15_000) expect(__queueConsumerTestUtils__.shouldRunQueueSyncInBackground('on_manifest_create')).toBe(false) expect(__queueConsumerTestUtils__.shouldRunQueueSyncInBackground('cron_email')).toBe(true) @@ -235,6 +237,31 @@ describe('queue_consumer legacy message compatibility', () => { )).toBe('continue') }) + it.concurrent('uses the version queue retry budget for Discord failure alerts', () => { + const versionRetryBudget = __queueConsumerTestUtils__.getQueueMaxReads('on_version_update') + const midRetry = { + cf_id: 'cf-version-mid', + error_code: 'manifest_cleanup_incomplete', + function_name: 'on_version_update', + function_type: 'supabase', + msg_id: 2, + payload_size: 10, + read_count: MAX_QUEUE_READS, + status: 500, + status_text: 'Internal Server Error', + } + const exhausted = { + ...midRetry, + cf_id: 'cf-version-done', + msg_id: 3, + read_count: versionRetryBudget, + } + + expect(versionRetryBudget).toBe(30) + expect(__queueConsumerTestUtils__.getActionableQueueFailures([midRetry], versionRetryBudget)).toEqual([]) + expect(__queueConsumerTestUtils__.getActionableQueueFailures([exhausted], versionRetryBudget)).toEqual([exhausted]) + }) + it.concurrent('alerts Discord after retry budget is exhausted', () => { const failure = { cf_id: 'cf-1',