// server/supabase/functions/audio-purge/handler.ts // Drains public.audio_purge_queue: claim a leased batch, remove the objects // through the Storage API, then complete the batch (which also drops the // 'deleted' audio_files rows). Depends only on ports; index.ts wires Supabase. import type { StorageBucketPort } from '../_shared/storage-objects.ts' import { removeStorageObjects } from '../_shared/storage-objects.ts' import { planPurge, type PurgeQueueEntry } from './policy.ts' export const CLAIM_BATCH_SIZE = 100 export const LEASE_SECONDS = 600 export const MAX_BATCHES_PER_RUN = 20 export interface AudioPurgeQueuePort { claim(limit: number, leaseSeconds: number): Promise complete(ids: number[], leaseToken: string): Promise } export interface AudioPurgeDeps { queue: AudioPurgeQueuePort storage: Pick maxBatches?: number } export interface AudioPurgeSummary { batches: number claimed: number removed: number rejected: number failed: number completed: number } function groupByLease(entries: readonly PurgeQueueEntry[]): Map { const groups = new Map() for (const entry of entries) { const group = groups.get(entry.leaseToken) if (group) group.push(entry) else groups.set(entry.leaseToken, [entry]) } return groups } /** * Removes queued objects batch by batch. A storage failure stops the run and * leaves that batch leased, so it is retried once the lease expires; entries * are completed only after their objects are gone. */ export async function drainAudioPurgeQueue(deps: AudioPurgeDeps): Promise { const maxBatches = deps.maxBatches ?? MAX_BATCHES_PER_RUN const summary: AudioPurgeSummary = { batches: 0, claimed: 0, removed: 0, rejected: 0, failed: 0, completed: 0 } while (summary.batches < maxBatches) { const claimed = await deps.queue.claim(CLAIM_BATCH_SIZE, LEASE_SECONDS) if (claimed.length === 0) break summary.batches += 1 summary.claimed += claimed.length let storageFailed = false for (const [leaseToken, entries] of groupByLease(claimed)) { const plan = planPurge(entries) const done = [...plan.rejected] summary.rejected += plan.rejected.length if (!storageFailed && plan.removable.length > 0) { try { await removeStorageObjects(deps.storage, [...new Set(plan.removable.map((entry) => entry.storageKey))]) summary.removed += plan.removable.length done.push(...plan.removable) } catch { storageFailed = true summary.failed += plan.removable.length } } else if (storageFailed) { summary.failed += plan.removable.length } if (done.length > 0) { summary.completed += await deps.queue.complete(done.map((entry) => entry.id), leaseToken) } } if (storageFailed || claimed.length < CLAIM_BATCH_SIZE) break } return summary } function json(status: number, body: Record): Response { return new Response(JSON.stringify(body), { status, headers: { 'Content-Type': 'application/json', 'Cache-Control': 'no-store' }, }) } export interface AudioPurgeRequestDeps { /** True only for the service-role caller (pg_cron via pg_net). */ authorize(req: Request): Promise createDeps(): AudioPurgeDeps } export async function handleAudioPurgeRequest(req: Request, deps: AudioPurgeRequestDeps): Promise { if (req.method !== 'POST') return json(405, { error: 'Method not allowed', code: 'METHOD_NOT_ALLOWED' }) if (!await deps.authorize(req)) return json(401, { error: 'Unauthorized', code: 'UNAUTHORIZED' }) try { const summary = await drainAudioPurgeQueue(deps.createDeps()) return json(summary.failed > 0 ? 503 : 200, { ...summary }) } catch (error) { const code = error instanceof Error && /^[a-z_:]+$/.test(error.message) ? error.message : 'audio_purge_failed' return json(500, { error: 'Audio purge failed', code }) } }