d3ro-voice/server/supabase/functions/audio-purge/handler.ts

113 lines
4 KiB
TypeScript

// 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<PurgeQueueEntry[]>
complete(ids: number[], leaseToken: string): Promise<number>
}
export interface AudioPurgeDeps {
queue: AudioPurgeQueuePort
storage: Pick<StorageBucketPort, 'remove'>
maxBatches?: number
}
export interface AudioPurgeSummary {
batches: number
claimed: number
removed: number
rejected: number
failed: number
completed: number
}
function groupByLease(entries: readonly PurgeQueueEntry[]): Map<string, PurgeQueueEntry[]> {
const groups = new Map<string, PurgeQueueEntry[]>()
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<AudioPurgeSummary> {
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<string, unknown>): 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<boolean>
createDeps(): AudioPurgeDeps
}
export async function handleAudioPurgeRequest(req: Request, deps: AudioPurgeRequestDeps): Promise<Response> {
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 })
}
}