// server/supabase/functions/audio-purge/queue-adapter.ts // AudioPurgeQueuePort over the service-role RPCs from // 20260929000004_audio_retention_on_delete.sql. import type { SupabaseClient } from '@supabase/supabase-js' import type { AudioPurgeQueuePort } from './handler.ts' import type { PurgeQueueEntry } from './policy.ts' function toId(value: unknown): number | null { if (typeof value === 'number' && Number.isSafeInteger(value) && value > 0) return value if (typeof value === 'string' && /^[1-9][0-9]{0,15}$/.test(value)) { const parsed = Number(value) return Number.isSafeInteger(parsed) ? parsed : null } return null } /** Rows returned by claim_audio_purge_batch_v1; malformed rows are skipped (their lease expires). */ export function parseClaimedRows(data: unknown): PurgeQueueEntry[] { if (!Array.isArray(data)) return [] const entries: PurgeQueueEntry[] = [] for (const row of data) { if (typeof row !== 'object' || row === null) continue const record = row as Record const id = toId(record.id) if ( id === null || typeof record.user_id !== 'string' || typeof record.storage_key !== 'string' || typeof record.lease_token !== 'string' ) continue entries.push({ id, userId: record.user_id, storageKey: record.storage_key, leaseToken: record.lease_token }) } return entries } type RpcClient = Pick export function supabaseAudioPurgeQueue(client: RpcClient): AudioPurgeQueuePort { return { async claim(limit, leaseSeconds) { const { data, error } = await client.rpc('claim_audio_purge_batch_v1', { p_limit: limit, p_lease_seconds: leaseSeconds, }) if (error) throw new Error('audio_purge_claim_failed') return parseClaimedRows(data) }, async complete(ids, leaseToken) { const { data, error } = await client.rpc('complete_audio_purge_v1', { p_ids: ids, p_lease_token: leaseToken, }) if (error) throw new Error('audio_purge_complete_failed') return typeof data === 'number' ? data : 0 }, } }