From ed790e672b690c750ed7b33bdf06251d91b31fca Mon Sep 17 00:00:00 2001 From: Yun Chan Date: Mon, 28 Sep 2026 02:16:22 +0900 Subject: [PATCH] fix(storage): purge raw audio when its history entry or meeting is deleted --- .../functions/_shared/storage-objects.test.ts | 67 +++++ .../functions/_shared/storage-objects.ts | 82 +++++ .../functions/account-delete/index.ts | 46 +-- .../functions/audio-purge/handler.test.ts | 134 +++++++++ .../supabase/functions/audio-purge/handler.ts | 113 +++++++ .../supabase/functions/audio-purge/index.ts | 31 ++ .../supabase/functions/audio-purge/policy.ts | 45 +++ .../functions/audio-purge/queue-adapter.ts | 58 ++++ ...260929000004_audio_retention_on_delete.sql | 280 ++++++++++++++++++ .../audio-retention-on-delete.integration.sql | 256 ++++++++++++++++ 10 files changed, 1073 insertions(+), 39 deletions(-) create mode 100644 server/supabase/functions/_shared/storage-objects.test.ts create mode 100644 server/supabase/functions/_shared/storage-objects.ts create mode 100644 server/supabase/functions/audio-purge/handler.test.ts create mode 100644 server/supabase/functions/audio-purge/handler.ts create mode 100644 server/supabase/functions/audio-purge/index.ts create mode 100644 server/supabase/functions/audio-purge/policy.ts create mode 100644 server/supabase/functions/audio-purge/queue-adapter.ts create mode 100644 server/supabase/migrations/20260929000004_audio_retention_on_delete.sql create mode 100644 server/supabase/tests/audio-retention-on-delete.integration.sql diff --git a/server/supabase/functions/_shared/storage-objects.test.ts b/server/supabase/functions/_shared/storage-objects.test.ts new file mode 100644 index 0000000..b694cf9 --- /dev/null +++ b/server/supabase/functions/_shared/storage-objects.test.ts @@ -0,0 +1,67 @@ +import { + listStorageFiles, + removeStorageObjects, + STORAGE_BATCH_SIZE, + type StorageBucketPort, + type StorageEntry, +} from './storage-objects.ts' + +function assertEquals(actual: unknown, expected: unknown, message: string): void { + const a = JSON.stringify(actual) + const e = JSON.stringify(expected) + if (a !== e) throw new Error(`${message}: expected ${e}, got ${a}`) +} + +class FakeBucket implements StorageBucketPort { + readonly bucket = 'audio' + listCalls: Array<{ prefix: string; offset: number }> = [] + removed: string[][] = [] + constructor(private tree: Record, private failRemoveAt = -1) {} + list(prefix: string, offset: number, limit: number): Promise { + this.listCalls.push({ prefix, offset }) + return Promise.resolve((this.tree[prefix] ?? []).slice(offset, offset + limit)) + } + remove(paths: string[]): Promise { + if (this.removed.length === this.failRemoveAt) return Promise.reject(new Error('storage_delete_failed:audio')) + this.removed.push(paths) + return Promise.resolve() + } +} + +Deno.test('storage-objects: listStorageFiles walks folders and pages', async () => { + const many = Array.from({ length: STORAGE_BATCH_SIZE + 3 }, (_, i) => ({ name: `f${i}.wav`, id: `id${i}` })) + const bucket = new FakeBucket({ + u: [{ name: 'history', id: null }, { name: 'top.wav', id: 'top' }], + 'u/history': many, + }) + const files = await listStorageFiles(bucket, 'u') + + assertEquals(files.length, STORAGE_BATCH_SIZE + 4, 'all files found') + assertEquals(files[files.length - 1], 'u/top.wav', 'depth-first order') + assertEquals( + bucket.listCalls, + [{ prefix: 'u', offset: 0 }, { prefix: 'u/history', offset: 0 }, { prefix: 'u/history', offset: STORAGE_BATCH_SIZE }], + 'paged listing', + ) +}) + +Deno.test('storage-objects: removeStorageObjects batches and stops on the first failure', async () => { + const paths = Array.from({ length: STORAGE_BATCH_SIZE * 2 + 1 }, (_, i) => `u/${i}`) + const ok = new FakeBucket({}) + await removeStorageObjects(ok, paths) + assertEquals(ok.removed.map((b) => b.length), [STORAGE_BATCH_SIZE, STORAGE_BATCH_SIZE, 1], 'batch sizes') + + const failing = new FakeBucket({}, 1) + let message = '' + try { + await removeStorageObjects(failing, paths) + } catch (error) { + message = error instanceof Error ? error.message : '' + } + assertEquals(message, 'storage_delete_failed:audio', 'error propagates') + assertEquals(failing.removed.length, 1, 'later batches are not attempted') + + const empty = new FakeBucket({}) + await removeStorageObjects(empty, []) + assertEquals(empty.removed.length, 0, 'nothing to remove') +}) diff --git a/server/supabase/functions/_shared/storage-objects.ts b/server/supabase/functions/_shared/storage-objects.ts new file mode 100644 index 0000000..6821151 --- /dev/null +++ b/server/supabase/functions/_shared/storage-objects.ts @@ -0,0 +1,82 @@ +// server/supabase/functions/_shared/storage-objects.ts +// Storage bucket port + the listing/removal routines shared by account-delete +// and audio-purge. Callers depend on StorageBucketPort, not on supabase-js, so +// the routines are unit-testable with an in-memory fake. + +import type { SupabaseClient } from '@supabase/supabase-js' + +/** Storage API page and batch size (remove() accepts at most 100 paths well within limits). */ +export const STORAGE_BATCH_SIZE = 100 + +export interface StorageEntry { + name: string + /** null for a folder placeholder, the object id for a file. */ + id: string | null +} + +export interface StorageBucketPort { + readonly bucket: string + list(prefix: string, offset: number, limit: number): Promise + remove(paths: string[]): Promise +} + +/** Every object path under prefix, depth-first, in name order. */ +export async function listStorageFiles(storage: StorageBucketPort, prefix: string): Promise { + const result: string[] = [] + let offset = 0 + + while (true) { + const page = await storage.list(prefix, offset, STORAGE_BATCH_SIZE) + if (page.length === 0) break + + for (const entry of page) { + const childPath = `${prefix}/${entry.name}` + if (entry.id) { + result.push(childPath) + } else { + result.push(...await listStorageFiles(storage, childPath)) + } + } + + if (page.length < STORAGE_BATCH_SIZE) break + offset += page.length + } + + return result +} + +/** Remove paths in API-sized batches; the first failing batch aborts the rest. */ +export async function removeStorageObjects( + storage: Pick, + paths: readonly string[], +): Promise { + for (let index = 0; index < paths.length; index += STORAGE_BATCH_SIZE) { + await storage.remove(paths.slice(index, index + STORAGE_BATCH_SIZE)) + } +} + +type StorageClient = Pick + +/** supabase-js adapter. Errors keep the storage_list_failed / storage_delete_failed codes. */ +export function supabaseStorageBucket(client: StorageClient, bucket: string): StorageBucketPort { + return { + bucket, + async list(prefix, offset, limit) { + const { data, error } = await client.storage.from(bucket).list(prefix, { + limit, + offset, + sortBy: { column: 'name', order: 'asc' }, + }) + if (error) throw new Error(`storage_list_failed:${bucket}`) + return (data ?? []).map((entry) => ({ + name: entry.name, + id: typeof entry.id === 'string' && entry.id.length > 0 ? entry.id : null, + })) + }, + async remove(paths) { + if (paths.length === 0) return + const { error } = await client.storage.from(bucket).remove(paths) + if (error) throw new Error(`storage_delete_failed:${bucket}`) + }, + } +} diff --git a/server/supabase/functions/account-delete/index.ts b/server/supabase/functions/account-delete/index.ts index 95c2dc9..dd1b656 100644 --- a/server/supabase/functions/account-delete/index.ts +++ b/server/supabase/functions/account-delete/index.ts @@ -2,6 +2,11 @@ import { corsHeaders, handleCorsPreflightRequest } from '../_shared/cors.ts' import { requireUser, authErrorResponse, type AuthError } from '../_shared/auth.ts' import { createServiceRoleClient } from '../_shared/quota.ts' import { hasRecentAuthentication } from './recent-auth.ts' +import { + listStorageFiles, + removeStorageObjects, + supabaseStorageBucket, +} from '../_shared/storage-objects.ts' interface DeleteAccountRequest { confirmation: string @@ -17,39 +22,6 @@ function jsonResponse(body: Record, status = 200): Response { }) } -async function listStorageFiles( - serviceClient: ReturnType, - bucket: string, - prefix: string, -): Promise { - const result: string[] = [] - let offset = 0 - - while (true) { - const { data, error } = await serviceClient.storage.from(bucket).list(prefix, { - limit: 100, - offset, - sortBy: { column: 'name', order: 'asc' }, - }) - if (error) throw new Error(`storage_list_failed:${bucket}`) - if (!data || data.length === 0) break - - for (const entry of data) { - const childPath = `${prefix}/${entry.name}` - if (entry.id) { - result.push(childPath) - } else { - result.push(...await listStorageFiles(serviceClient, bucket, childPath)) - } - } - - if (data.length < 100) break - offset += data.length - } - - return result -} - Deno.serve(async (req: Request) => { const preflight = handleCorsPreflightRequest(req) if (preflight) return preflight @@ -110,12 +82,8 @@ Deno.serve(async (req: Request) => { } for (const bucket of STORAGE_BUCKETS) { - const files = await listStorageFiles(serviceClient, bucket, user.id) - for (let index = 0; index < files.length; index += 100) { - const batch = files.slice(index, index + 100) - const { error } = await serviceClient.storage.from(bucket).remove(batch) - if (error) throw new Error(`storage_delete_failed:${bucket}`) - } + const storage = supabaseStorageBucket(serviceClient, bucket) + await removeStorageObjects(storage, await listStorageFiles(storage, user.id)) } const { error: deleteError } = await serviceClient.auth.admin.deleteUser(user.id, false) diff --git a/server/supabase/functions/audio-purge/handler.test.ts b/server/supabase/functions/audio-purge/handler.test.ts new file mode 100644 index 0000000..4506d70 --- /dev/null +++ b/server/supabase/functions/audio-purge/handler.test.ts @@ -0,0 +1,134 @@ +import { drainAudioPurgeQueue, handleAudioPurgeRequest, type AudioPurgeQueuePort } from './handler.ts' +import { isPurgeableStorageKey, planPurge, type PurgeQueueEntry } from './policy.ts' +import { parseClaimedRows } from './queue-adapter.ts' + +function assert(condition: boolean, message: string): asserts condition { + if (!condition) throw new Error(message) +} + +function assertEquals(actual: unknown, expected: unknown, message: string): void { + const a = JSON.stringify(actual) + const e = JSON.stringify(expected) + if (a !== e) throw new Error(`${message}: expected ${e}, got ${a}`) +} + +const USER = '24000000-0000-4000-8000-000000000001' +const OTHER = '24000000-0000-4000-8000-000000000002' + +function entry(id: number, storageKey: string, userId = USER, leaseToken = 'lease-1'): PurgeQueueEntry { + return { id, userId, storageKey, leaseToken } +} + +class FakeQueue implements AudioPurgeQueuePort { + completed: Array<{ ids: number[]; leaseToken: string }> = [] + constructor(private batches: PurgeQueueEntry[][]) {} + claim(): Promise { + return Promise.resolve(this.batches.shift() ?? []) + } + complete(ids: number[], leaseToken: string): Promise { + this.completed.push({ ids, leaseToken }) + return Promise.resolve(ids.length) + } +} + +class FakeStorage { + removed: string[][] = [] + constructor(private failing = false) {} + remove(paths: string[]): Promise { + if (this.failing) return Promise.reject(new Error('storage_delete_failed:audio')) + this.removed.push(paths) + return Promise.resolve() + } +} + +Deno.test('redteam r2-24: keys inside the owner prefix are purgeable', () => { + assert(isPurgeableStorageKey(USER, `${USER}/history/abc.wav`), 'desktop history key') + assert(isPurgeableStorageKey(USER, `${USER}/${'a'.repeat(64)}/memo.m4a`), 'mobile content-addressed key') +}) + +Deno.test('redteam r2-24: keys that could reach another user or escape the prefix are rejected', () => { + assert(!isPurgeableStorageKey(USER, `${OTHER}/history/abc.wav`), 'other user prefix') + assert(!isPurgeableStorageKey(USER, USER), 'bare prefix (whole folder)') + assert(!isPurgeableStorageKey(USER, `${USER}/../${OTHER}/x.wav`), 'dot-dot segment') + assert(!isPurgeableStorageKey(USER, `${USER}/./x.wav`), 'dot segment') + assert(!isPurgeableStorageKey(USER, `${USER}//x.wav`), 'empty segment') + assert(!isPurgeableStorageKey(USER, `${USER}/x.wav/`), 'trailing slash') + assert(!isPurgeableStorageKey(USER, `${USER}\\..\\x.wav`), 'backslash') + assert(!isPurgeableStorageKey(USER, `${USER}/x\u0000.wav`), 'control character') + assert(!isPurgeableStorageKey('not-a-uuid', 'not-a-uuid/x.wav'), 'non-uuid owner') + assert(!isPurgeableStorageKey(USER, `${USER}/${'a'.repeat(1100)}`), 'overlong key') +}) + +Deno.test('redteam r2-24: planPurge splits removable from rejected', () => { + const plan = planPurge([entry(1, `${USER}/history/a.wav`), entry(2, `${OTHER}/history/b.wav`)]) + assertEquals(plan.removable.map((e) => e.id), [1], 'removable ids') + assertEquals(plan.rejected.map((e) => e.id), [2], 'rejected ids') +}) + +Deno.test('redteam r2-24: drain removes the objects, then completes the batch under its lease', async () => { + const queue = new FakeQueue([[ + entry(1, `${USER}/history/a.wav`), + entry(2, `${USER}/meetings/m.wav`), + entry(3, `${OTHER}/history/stolen.wav`), + ]]) + const storage = new FakeStorage() + const summary = await drainAudioPurgeQueue({ queue, storage }) + + assertEquals(storage.removed, [[`${USER}/history/a.wav`, `${USER}/meetings/m.wav`]], 'removed paths') + assertEquals(queue.completed, [{ ids: [3, 1, 2], leaseToken: 'lease-1' }], 'completed under lease') + assertEquals(summary, { batches: 1, claimed: 3, removed: 2, rejected: 1, failed: 0, completed: 3 }, 'summary') +}) + +Deno.test('redteam r2-24: a storage failure never completes the unremoved entries', async () => { + const queue = new FakeQueue([[entry(1, `${USER}/history/a.wav`), entry(2, `${OTHER}/x.wav`)]]) + const summary = await drainAudioPurgeQueue({ queue, storage: new FakeStorage(true) }) + + assertEquals(queue.completed, [{ ids: [2], leaseToken: 'lease-1' }], 'only the rejected entry completes') + assertEquals(summary.failed, 1, 'failed count') + assertEquals(summary.removed, 0, 'nothing removed') +}) + +Deno.test('redteam r2-24: drain keeps claiming full batches and stops on an empty claim', async () => { + const full = Array.from({ length: 100 }, (_, i) => entry(i + 1, `${USER}/history/${i}.wav`)) + const queue = new FakeQueue([full, [entry(101, `${USER}/history/last.wav`, USER, 'lease-2')]]) + const storage = new FakeStorage() + const summary = await drainAudioPurgeQueue({ queue, storage }) + + assertEquals(summary.batches, 2, 'two batches') + assertEquals(summary.completed, 101, 'all completed') + assertEquals(queue.completed.map((c) => c.leaseToken), ['lease-1', 'lease-2'], 'each batch under its own lease') + assertEquals(storage.removed.map((b) => b.length), [100, 1], 'API-sized removal batches') +}) + +Deno.test('redteam r2-24: parseClaimedRows accepts RPC rows and skips malformed ones', () => { + const rows = parseClaimedRows([ + { id: 7, user_id: USER, storage_key: `${USER}/a.wav`, lease_token: 'l' }, + { id: '8', user_id: USER, storage_key: `${USER}/b.wav`, lease_token: 'l' }, + { id: 0, user_id: USER, storage_key: `${USER}/c.wav`, lease_token: 'l' }, + { id: 9, user_id: USER, storage_key: null, lease_token: 'l' }, + 'junk', + ]) + assertEquals(rows.map((r) => r.id), [7, 8], 'parsed ids') + assertEquals(parseClaimedRows(null), [], 'null data') +}) + +Deno.test('redteam r2-24: only an authorized POST runs the purge', async () => { + let ran = 0 + const deps = (authorized: boolean) => ({ + authorize: () => Promise.resolve(authorized), + createDeps: () => { + ran += 1 + return { queue: new FakeQueue([]), storage: new FakeStorage() } + }, + }) + + const get = await handleAudioPurgeRequest(new Request('http://x/audio-purge'), deps(true)) + assertEquals(get.status, 405, 'GET rejected') + const denied = await handleAudioPurgeRequest(new Request('http://x/audio-purge', { method: 'POST' }), deps(false)) + assertEquals(denied.status, 401, 'unauthorized rejected') + assertEquals(ran, 0, 'no work before authorization') + + const ok = await handleAudioPurgeRequest(new Request('http://x/audio-purge', { method: 'POST' }), deps(true)) + assertEquals(ok.status, 200, 'authorized POST succeeds') + assertEquals(ran, 1, 'purge ran once') +}) diff --git a/server/supabase/functions/audio-purge/handler.ts b/server/supabase/functions/audio-purge/handler.ts new file mode 100644 index 0000000..dd744ac --- /dev/null +++ b/server/supabase/functions/audio-purge/handler.ts @@ -0,0 +1,113 @@ +// 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 }) + } +} diff --git a/server/supabase/functions/audio-purge/index.ts b/server/supabase/functions/audio-purge/index.ts new file mode 100644 index 0000000..4931b36 --- /dev/null +++ b/server/supabase/functions/audio-purge/index.ts @@ -0,0 +1,31 @@ +// server/supabase/functions/audio-purge/index.ts +// Service-only worker invoked by pg_cron (public.dispatch_audio_purge_v1) to +// remove audio objects whose history entry or meeting was deleted. + +import { isServiceRoleApiKey, isServiceRoleAuthorization } from '../_shared/push-contract.ts' +import { createServiceRoleClient } from '../_shared/quota.ts' +import { supabaseStorageBucket } from '../_shared/storage-objects.ts' +import { handleAudioPurgeRequest } from './handler.ts' +import { AUDIO_BUCKET } from './policy.ts' +import { supabaseAudioPurgeQueue } from './queue-adapter.ts' + +async function authorizeServiceCaller(req: Request): Promise { + const legacyServiceBearer = await isServiceRoleAuthorization( + req.headers.get('Authorization'), + Deno.env.get('SUPABASE_SERVICE_ROLE_KEY'), + ) + return legacyServiceBearer || await isServiceRoleApiKey(req.headers.get('apikey')) +} + +Deno.serve(async (req: Request) => { + return await handleAudioPurgeRequest(req, { + authorize: authorizeServiceCaller, + createDeps: () => { + const client = createServiceRoleClient() + return { + queue: supabaseAudioPurgeQueue(client), + storage: supabaseStorageBucket(client, AUDIO_BUCKET), + } + }, + }) +}) diff --git a/server/supabase/functions/audio-purge/policy.ts b/server/supabase/functions/audio-purge/policy.ts new file mode 100644 index 0000000..4c20286 --- /dev/null +++ b/server/supabase/functions/audio-purge/policy.ts @@ -0,0 +1,45 @@ +// server/supabase/functions/audio-purge/policy.ts +// Pure rules for which queued storage keys the service-role worker may remove. +// storage_key is client-written (RLS only pins its first segment), so the +// worker re-checks it before acting with service-role storage access. + +export const AUDIO_BUCKET = 'audio' +export const MAX_STORAGE_KEY_LENGTH = 1024 + +export interface PurgeQueueEntry { + id: number + userId: string + storageKey: string + leaseToken: string +} + +export interface PurgePlan { + /** Entries whose object may be removed from the audio bucket. */ + removable: PurgeQueueEntry[] + /** Entries dropped without touching storage (key outside the owner's prefix). */ + rejected: PurgeQueueEntry[] +} + +const UUID_PATTERN = /^[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}$/i +// deno-lint-ignore no-control-regex +const CONTROL_CHARACTERS = /[\u0000-\u001f\u007f]/ + +/** True when key is a plain object path inside the owner's own audio prefix. */ +export function isPurgeableStorageKey(userId: string, key: string): boolean { + if (!UUID_PATTERN.test(userId)) return false + if (key.length === 0 || key.length > MAX_STORAGE_KEY_LENGTH) return false + if (key.includes('\\') || CONTROL_CHARACTERS.test(key)) return false + const segments = key.split('/') + if (segments.length < 2 || segments[0] !== userId) return false + return segments.every((segment) => segment.length > 0 && segment !== '.' && segment !== '..') +} + +export function planPurge(entries: readonly PurgeQueueEntry[]): PurgePlan { + const removable: PurgeQueueEntry[] = [] + const rejected: PurgeQueueEntry[] = [] + for (const entry of entries) { + if (isPurgeableStorageKey(entry.userId, entry.storageKey)) removable.push(entry) + else rejected.push(entry) + } + return { removable, rejected } +} diff --git a/server/supabase/functions/audio-purge/queue-adapter.ts b/server/supabase/functions/audio-purge/queue-adapter.ts new file mode 100644 index 0000000..1512886 --- /dev/null +++ b/server/supabase/functions/audio-purge/queue-adapter.ts @@ -0,0 +1,58 @@ +// 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 + }, + } +} diff --git a/server/supabase/migrations/20260929000004_audio_retention_on_delete.sql b/server/supabase/migrations/20260929000004_audio_retention_on_delete.sql new file mode 100644 index 0000000..941ba7e --- /dev/null +++ b/server/supabase/migrations/20260929000004_audio_retention_on_delete.sql @@ -0,0 +1,280 @@ +-- Deleting a history entry or meeting must not leave its raw audio in storage. +-- +-- audio_files.history_id / meeting_id are ON DELETE SET NULL, so a parent +-- delete that does not purge first (every mobile delete, a desktop purge that +-- failed, the meeting re-record RPC) left an unreachable object in the audio +-- bucket until the whole account was deleted. +-- +-- 1. audio_purge_queue: service-only work queue of storage keys to remove. +-- 2. A BEFORE INSERT OR UPDATE trigger on audio_files: +-- - when a row loses its last owner link (the SET NULL of a parent delete +-- runs as an UPDATE, so this covers every delete path), the row becomes +-- upload_status = 'deleted' and its storage key is queued; +-- - when a live row claims a queued key again (a re-import of the same +-- content-addressed file), the queued purge is cancelled and the stale +-- 'deleted' rows for that key are dropped so the UNIQUE key is free. +-- While a worker holds the purge lease the write is refused (55006) so +-- the worker can never remove an object a live row just re-claimed. +-- 3. claim/complete RPCs (service_role only) used by the audio-purge edge +-- function, which removes objects through the Storage API (direct deletes +-- on storage.objects are blocked and would leave the bytes behind). +-- 4. pg_cron calls the edge function every 10 minutes through pg_net when the +-- queue has work. The URL and key come from Vault ('project_url', +-- 'service_role_key'); without them the job is a no-op. +-- 5. One-time backfill of audio already orphaned by earlier deletes. + +BEGIN; + +-- 1) Queue ----------------------------------------------------------------------- +-- user_id carries no FK: rows are enqueued inside the auth.users delete cascade +-- (history SET NULL -> audio_files UPDATE), where an FK to the vanishing user +-- would abort the account deletion. The worker tolerates missing objects. +CREATE TABLE IF NOT EXISTS public.audio_purge_queue ( + id bigint GENERATED ALWAYS AS IDENTITY PRIMARY KEY, + user_id uuid NOT NULL, + storage_key text NOT NULL CHECK (length(storage_key) BETWEEN 1 AND 1024), + attempts integer NOT NULL DEFAULT 0 CHECK (attempts >= 0), + lease_token uuid, + leased_until timestamptz, + enqueued_at timestamptz NOT NULL DEFAULT now(), + UNIQUE (user_id, storage_key) +); + +CREATE INDEX IF NOT EXISTS idx_audio_purge_queue_ready + ON public.audio_purge_queue(enqueued_at, id); + +ALTER TABLE public.audio_purge_queue ENABLE ROW LEVEL SECURITY; +REVOKE ALL ON TABLE public.audio_purge_queue FROM PUBLIC, anon, authenticated; +GRANT SELECT, INSERT, UPDATE, DELETE ON TABLE public.audio_purge_queue TO service_role; + +-- Key lookups from the trigger and the claim guard. +CREATE INDEX IF NOT EXISTS idx_audio_files_user_storage_key + ON public.audio_files(user_id, storage_key); + +-- 2) Orphan / reclaim trigger ---------------------------------------------------- +CREATE OR REPLACE FUNCTION public.audio_files_retention_v1() +RETURNS trigger +LANGUAGE plpgsql +SECURITY DEFINER +SET search_path = '' +AS $$ +BEGIN + -- The row lost its last owner: parent deleted (FK SET NULL) or detached. + IF TG_OP = 'UPDATE' + AND NEW.history_id IS NULL + AND NEW.meeting_id IS NULL + AND (OLD.history_id IS NOT NULL OR OLD.meeting_id IS NOT NULL) THEN + NEW.storage_key := OLD.storage_key; + NEW.upload_status := 'deleted'; + INSERT INTO public.audio_purge_queue (user_id, storage_key) + VALUES (OLD.user_id, OLD.storage_key) + ON CONFLICT (user_id, storage_key) DO NOTHING; + RETURN NEW; + END IF; + + -- A live row (re)claims a key: cancel any pending purge of that key. + IF NEW.upload_status <> 'deleted' + AND ( + TG_OP = 'INSERT' + OR OLD.upload_status = 'deleted' + OR OLD.storage_key IS DISTINCT FROM NEW.storage_key + OR OLD.user_id IS DISTINCT FROM NEW.user_id + ) THEN + DELETE FROM public.audio_purge_queue q + WHERE q.user_id = NEW.user_id + AND q.storage_key = NEW.storage_key + AND (q.leased_until IS NULL OR q.leased_until <= now()); + + IF EXISTS ( + SELECT 1 FROM public.audio_purge_queue q + WHERE q.user_id = NEW.user_id AND q.storage_key = NEW.storage_key + ) THEN + RAISE EXCEPTION 'audio purge in progress for this storage key' + USING ERRCODE = '55006'; + END IF; + + DELETE FROM public.audio_files a + WHERE a.user_id = NEW.user_id + AND a.storage_key = NEW.storage_key + AND a.upload_status = 'deleted' + AND a.id <> NEW.id; + END IF; + + RETURN NEW; +END; +$$; + +REVOKE ALL ON FUNCTION public.audio_files_retention_v1() FROM PUBLIC, anon, authenticated; + +DROP TRIGGER IF EXISTS audio_files_retention_v1 ON public.audio_files; +CREATE TRIGGER audio_files_retention_v1 + BEFORE INSERT OR UPDATE ON public.audio_files + FOR EACH ROW EXECUTE FUNCTION public.audio_files_retention_v1(); + +-- 3) Worker RPCs ----------------------------------------------------------------- +CREATE OR REPLACE FUNCTION public.claim_audio_purge_batch_v1( + p_limit integer DEFAULT 100, + p_lease_seconds integer DEFAULT 600 +) +RETURNS SETOF public.audio_purge_queue +LANGUAGE plpgsql +SECURITY DEFINER +SET search_path = '' +AS $$ +DECLARE + batch_size integer := least(greatest(coalesce(p_limit, 100), 1), 100); + lease_seconds integer := least(greatest(coalesce(p_lease_seconds, 600), 60), 3600); + token uuid := gen_random_uuid(); +BEGIN + -- A key that a live row still (or again) references must never be removed. + DELETE FROM public.audio_purge_queue q + WHERE (q.leased_until IS NULL OR q.leased_until <= now()) + AND EXISTS ( + SELECT 1 FROM public.audio_files a + WHERE a.user_id = q.user_id + AND a.storage_key = q.storage_key + AND a.upload_status <> 'deleted' + ); + + RETURN QUERY + WITH picked AS ( + SELECT q.id + FROM public.audio_purge_queue q + WHERE (q.leased_until IS NULL OR q.leased_until <= now()) + AND q.attempts < 20 + ORDER BY q.enqueued_at, q.id + LIMIT batch_size + FOR UPDATE SKIP LOCKED + ) + UPDATE public.audio_purge_queue q + SET lease_token = token, + leased_until = now() + make_interval(secs => lease_seconds), + attempts = q.attempts + 1 + FROM picked + WHERE q.id = picked.id + RETURNING q.*; +END; +$$; + +CREATE OR REPLACE FUNCTION public.complete_audio_purge_v1( + p_ids bigint[], + p_lease_token uuid +) +RETURNS integer +LANGUAGE plpgsql +SECURITY DEFINER +SET search_path = '' +AS $$ +DECLARE + completed integer; +BEGIN + IF p_lease_token IS NULL OR p_ids IS NULL OR cardinality(p_ids) = 0 THEN + RETURN 0; + END IF; + + WITH done AS ( + DELETE FROM public.audio_purge_queue q + WHERE q.id = ANY(p_ids) + AND q.lease_token = p_lease_token + RETURNING q.user_id, q.storage_key + ), dropped AS ( + DELETE FROM public.audio_files a + USING done + WHERE a.user_id = done.user_id + AND a.storage_key = done.storage_key + AND a.upload_status = 'deleted' + RETURNING a.id + ) + SELECT count(*)::integer INTO completed FROM done; + + RETURN completed; +END; +$$; + +REVOKE ALL ON FUNCTION public.claim_audio_purge_batch_v1(integer, integer) FROM PUBLIC, anon, authenticated; +REVOKE ALL ON FUNCTION public.complete_audio_purge_v1(bigint[], uuid) FROM PUBLIC, anon, authenticated; +GRANT EXECUTE ON FUNCTION public.claim_audio_purge_batch_v1(integer, integer) TO service_role; +GRANT EXECUTE ON FUNCTION public.complete_audio_purge_v1(bigint[], uuid) TO service_role; + +-- 4) Dispatcher called by pg_cron ------------------------------------------------ +CREATE OR REPLACE FUNCTION public.dispatch_audio_purge_v1() +RETURNS bigint +LANGUAGE plpgsql +SECURITY DEFINER +SET search_path = '' +AS $$ +DECLARE + project_url text; + service_key text; +BEGIN + IF NOT EXISTS ( + SELECT 1 FROM public.audio_purge_queue q + WHERE (q.leased_until IS NULL OR q.leased_until <= now()) + AND q.attempts < 20 + ) THEN + RETURN NULL; + END IF; + + BEGIN + SELECT s.decrypted_secret INTO project_url + FROM vault.decrypted_secrets s WHERE s.name = 'project_url'; + SELECT s.decrypted_secret INTO service_key + FROM vault.decrypted_secrets s WHERE s.name = 'service_role_key'; + EXCEPTION WHEN undefined_table OR invalid_schema_name OR insufficient_privilege THEN + RAISE NOTICE 'audio purge dispatch skipped: vault is unavailable'; + RETURN NULL; + END; + + IF project_url IS NULL OR service_key IS NULL THEN + RAISE NOTICE 'audio purge dispatch skipped: vault secrets project_url/service_role_key are not set'; + RETURN NULL; + END IF; + + RETURN net.http_post( + url := rtrim(project_url, '/') || '/functions/v1/audio-purge', + headers := jsonb_build_object( + 'Content-Type', 'application/json', + 'Authorization', 'Bearer ' || service_key, + 'apikey', service_key + ), + body := '{}'::jsonb, + timeout_milliseconds := 60000 + ); +END; +$$; + +REVOKE ALL ON FUNCTION public.dispatch_audio_purge_v1() FROM PUBLIC, anon, authenticated; + +-- 5) Backfill ---------------------------------------------------------------------- +-- Completed uploads that already lost their owner. Rows younger than a day or +-- still referenced by an active processing job may be mid-import, so they stay. +WITH orphaned AS ( + UPDATE public.audio_files a + SET upload_status = 'deleted' + WHERE a.history_id IS NULL + AND a.meeting_id IS NULL + AND a.upload_status = 'uploaded' + AND a.created_at < now() - interval '1 day' + AND NOT EXISTS ( + SELECT 1 FROM public.processing_jobs p + WHERE p.audio_file_id = a.id AND p.status IN ('queued', 'running') + ) + RETURNING a.user_id, a.storage_key +) +INSERT INTO public.audio_purge_queue (user_id, storage_key) +SELECT DISTINCT o.user_id, o.storage_key FROM orphaned o +ON CONFLICT (user_id, storage_key) DO NOTHING; + +COMMIT; + +-- 6) Schedule -------------------------------------------------------------------- +-- Outside the transaction like 20260927000035: extension creation and +-- cron.schedule (same job name replaces) are idempotent. +CREATE EXTENSION IF NOT EXISTS pg_net; +CREATE EXTENSION IF NOT EXISTS pg_cron WITH SCHEMA pg_catalog; + +SELECT cron.schedule( + 'purge-orphaned-audio', + '*/10 * * * *', + $$SELECT public.dispatch_audio_purge_v1()$$ +); diff --git a/server/supabase/tests/audio-retention-on-delete.integration.sql b/server/supabase/tests/audio-retention-on-delete.integration.sql new file mode 100644 index 0000000..7e4bba7 --- /dev/null +++ b/server/supabase/tests/audio-retention-on-delete.integration.sql @@ -0,0 +1,256 @@ +\set ON_ERROR_STOP on + +-- Regression (redteam r2-24): deleting a history entry or meeting (e.g. from the +-- phone, which deletes only the parent row) must queue its raw audio for removal +-- instead of leaving an unreachable object in the audio bucket forever. +-- Requires 20260929000004_audio_retention_on_delete.sql. + +BEGIN; + +CREATE OR REPLACE FUNCTION pg_temp.assert_true(condition boolean, message text) +RETURNS void +LANGUAGE plpgsql +AS $$ +BEGIN + IF condition IS NOT TRUE THEN + RAISE EXCEPTION 'assertion_failed: %', message; + END IF; +END; +$$; + +INSERT INTO auth.users ( + id, aud, role, email, encrypted_password, email_confirmed_at, + raw_app_meta_data, raw_user_meta_data, created_at, updated_at +) VALUES + ( + '24000000-0000-4000-8000-000000000001', 'authenticated', 'authenticated', + 'audio-retention-owner@example.invalid', crypt('fixture-password', gen_salt('bf')), now(), + '{"provider":"email","providers":["email"]}'::jsonb, '{}'::jsonb, now(), now() + ), + ( + '24000000-0000-4000-8000-000000000002', 'authenticated', 'authenticated', + 'audio-retention-leaver@example.invalid', crypt('fixture-password', gen_salt('bf')), now(), + '{"provider":"email","providers":["email"]}'::jsonb, '{}'::jsonb, now(), now() + ); + +INSERT INTO public.history (id, user_id, original_text, duration) +VALUES + ('24100000-0000-4000-8000-000000000001', '24000000-0000-4000-8000-000000000001', 'dictation one', 1), + ('24100000-0000-4000-8000-000000000002', '24000000-0000-4000-8000-000000000001', 'dictation two', 1), + ('24100000-0000-4000-8000-000000000003', '24000000-0000-4000-8000-000000000002', 'leaver dictation', 1); + +INSERT INTO public.meetings (id, user_id, title, status) +VALUES ('24200000-0000-4000-8000-000000000001', '24000000-0000-4000-8000-000000000001', 'Retention meeting', 'completed'); + +INSERT INTO public.audio_files ( + id, user_id, history_id, meeting_id, source, storage_key, mime_type, size_bytes, sha256, upload_status +) VALUES + ('24300000-0000-4000-8000-000000000001', '24000000-0000-4000-8000-000000000001', + '24100000-0000-4000-8000-000000000001', NULL, 'recording', + '24000000-0000-4000-8000-000000000001/history/one.wav', 'audio/wav', 10, repeat('a', 64), 'uploaded'), + ('24300000-0000-4000-8000-000000000002', '24000000-0000-4000-8000-000000000001', + NULL, '24200000-0000-4000-8000-000000000001', 'recording', + '24000000-0000-4000-8000-000000000001/meetings/m.wav', 'audio/wav', 10, repeat('b', 64), 'uploaded'), + ('24300000-0000-4000-8000-000000000003', '24000000-0000-4000-8000-000000000001', + '24100000-0000-4000-8000-000000000002', NULL, 'file-picker', + '24000000-0000-4000-8000-000000000001/' || repeat('c', 64) || '/memo.m4a', 'audio/mp4', 10, repeat('c', 64), 'uploaded'), + ('24300000-0000-4000-8000-000000000004', '24000000-0000-4000-8000-000000000002', + '24100000-0000-4000-8000-000000000003', NULL, 'recording', + '24000000-0000-4000-8000-000000000002/history/leaver.wav', 'audio/wav', 10, repeat('d', 64), 'uploaded'); + +-- 1) Clients cannot see or drive the purge queue. Checked through the catalog: +-- calling a function whose EXECUTE is denied segfaults the local +-- supabase/postgres 17.6.1.104 image, which would take the shared DB down. +SELECT pg_temp.assert_true( + NOT has_table_privilege('authenticated', 'public.audio_purge_queue', 'SELECT') + AND NOT has_table_privilege('anon', 'public.audio_purge_queue', 'SELECT') + AND NOT has_table_privilege('authenticated', 'public.audio_purge_queue', 'INSERT') + AND NOT has_function_privilege('authenticated', 'public.claim_audio_purge_batch_v1(integer, integer)', 'EXECUTE') + AND NOT has_function_privilege('anon', 'public.claim_audio_purge_batch_v1(integer, integer)', 'EXECUTE') + AND NOT has_function_privilege('authenticated', 'public.complete_audio_purge_v1(bigint[], uuid)', 'EXECUTE') + AND NOT has_function_privilege('authenticated', 'public.dispatch_audio_purge_v1()', 'EXECUTE') + AND has_function_privilege('service_role', 'public.claim_audio_purge_batch_v1(integer, integer)', 'EXECUTE') + AND has_function_privilege('service_role', 'public.complete_audio_purge_v1(bigint[], uuid)', 'EXECUTE'), + 'purge queue and worker RPCs are service-only' +); + +SET LOCAL ROLE authenticated; +SELECT set_config( + 'request.jwt.claims', + '{"sub":"24000000-0000-4000-8000-000000000001","role":"authenticated"}', + true +); + +-- 2) The phone deletes only the parent rows (RLS path, like deleteHistoryRevisionSafe). +DELETE FROM public.history WHERE id = '24100000-0000-4000-8000-000000000001'; +DELETE FROM public.meetings WHERE id = '24200000-0000-4000-8000-000000000001'; +DELETE FROM public.history WHERE id = '24100000-0000-4000-8000-000000000002'; + +RESET ROLE; + +SELECT pg_temp.assert_true( + (SELECT count(*) FROM public.audio_files + WHERE id IN ('24300000-0000-4000-8000-000000000001', + '24300000-0000-4000-8000-000000000002', + '24300000-0000-4000-8000-000000000003') + AND upload_status = 'deleted' + AND history_id IS NULL AND meeting_id IS NULL) = 3, + 'orphaned audio rows are marked deleted' +); + +SELECT pg_temp.assert_true( + (SELECT array_agg(storage_key ORDER BY storage_key) FROM public.audio_purge_queue + WHERE user_id = '24000000-0000-4000-8000-000000000001') + = ARRAY[ + '24000000-0000-4000-8000-000000000001/' || repeat('c', 64) || '/memo.m4a', + '24000000-0000-4000-8000-000000000001/history/one.wav', + '24000000-0000-4000-8000-000000000001/meetings/m.wav' + ], + 'history and meeting deletes queue their audio keys' +); + +-- 3) Re-importing the same content-addressed file reclaims the key: the pending +-- purge is cancelled and the stale deleted row no longer blocks the UNIQUE key. +SET LOCAL ROLE authenticated; +INSERT INTO public.audio_files ( + id, user_id, source, storage_key, mime_type, size_bytes, sha256, upload_status +) VALUES ( + '24300000-0000-4000-8000-000000000005', '24000000-0000-4000-8000-000000000001', 'file-picker', + '24000000-0000-4000-8000-000000000001/' || repeat('c', 64) || '/memo.m4a', 'audio/mp4', 10, repeat('c', 64), 'pending' +); +RESET ROLE; + +SELECT pg_temp.assert_true( + NOT EXISTS (SELECT 1 FROM public.audio_purge_queue + WHERE storage_key LIKE '%/memo.m4a'), + 'reclaimed key is no longer queued for purge' +); +SELECT pg_temp.assert_true( + NOT EXISTS (SELECT 1 FROM public.audio_files WHERE id = '24300000-0000-4000-8000-000000000003'), + 'stale deleted row for the reclaimed key is dropped' +); + +-- 4) A key referenced by a live row again is dropped at claim time, never removed. +INSERT INTO public.audio_purge_queue (user_id, storage_key) +VALUES ('24000000-0000-4000-8000-000000000001', + '24000000-0000-4000-8000-000000000001/' || repeat('c', 64) || '/memo.m4a'); + +SET LOCAL ROLE service_role; +CREATE TEMP TABLE claimed ON COMMIT DROP AS +SELECT * FROM public.claim_audio_purge_batch_v1(100, 600); +RESET ROLE; + +SELECT pg_temp.assert_true( + (SELECT count(*) FROM claimed) = 2 + AND NOT EXISTS (SELECT 1 FROM claimed WHERE storage_key LIKE '%/memo.m4a') + AND (SELECT count(DISTINCT lease_token) FROM claimed) = 1 + AND (SELECT bool_and(leased_until > now() AND attempts = 1) FROM claimed), + 'claim leases only keys with no live owner' +); + +-- 5) While a worker holds the lease, a live row cannot reclaim that key. +SET LOCAL ROLE authenticated; +DO $$ +BEGIN + BEGIN + INSERT INTO public.audio_files ( + user_id, meeting_id, source, storage_key, mime_type, size_bytes, sha256, upload_status + ) VALUES ( + '24000000-0000-4000-8000-000000000001', NULL, 'recording', + '24000000-0000-4000-8000-000000000001/meetings/m.wav', 'audio/wav', 10, repeat('e', 64), 'pending' + ); + RAISE EXCEPTION 'assertion_failed: leased key was reclaimed'; + EXCEPTION WHEN object_in_use THEN NULL; + END; +END; +$$; +RESET ROLE; + +-- 6) A second claim does not hand out leased work. +SET LOCAL ROLE service_role; +SELECT pg_temp.assert_true( + NOT EXISTS (SELECT 1 FROM public.claim_audio_purge_batch_v1(100, 600)), + 'leased entries are not claimed twice' +); + +-- 7) Completing with a wrong lease does nothing; the right lease clears queue and rows. +SELECT pg_temp.assert_true( + public.complete_audio_purge_v1( + (SELECT array_agg(id) FROM claimed), gen_random_uuid() + ) = 0, + 'foreign lease cannot complete' +); +SELECT pg_temp.assert_true( + public.complete_audio_purge_v1( + (SELECT array_agg(id) FROM claimed), (SELECT min(lease_token::text)::uuid FROM claimed) + ) = 2, + 'lease holder completes the batch' +); +RESET ROLE; + +SELECT pg_temp.assert_true( + NOT EXISTS (SELECT 1 FROM public.audio_purge_queue + WHERE user_id = '24000000-0000-4000-8000-000000000001'), + 'completed entries leave the queue' +); +SELECT pg_temp.assert_true( + NOT EXISTS (SELECT 1 FROM public.audio_files + WHERE id IN ('24300000-0000-4000-8000-000000000001', + '24300000-0000-4000-8000-000000000002')), + 'purged audio metadata is removed' +); +SELECT pg_temp.assert_true( + EXISTS (SELECT 1 FROM public.audio_files + WHERE id = '24300000-0000-4000-8000-000000000005' AND upload_status = 'pending'), + 'the reclaimed live row survives completion' +); + +-- 8) Desktop upsert (ON CONFLICT on the UNIQUE key) over a deleted row reclaims it. +INSERT INTO public.history (id, user_id, original_text, duration) +VALUES ('24100000-0000-4000-8000-000000000004', '24000000-0000-4000-8000-000000000001', 'desktop', 1); +INSERT INTO public.audio_files ( + id, user_id, history_id, source, storage_key, mime_type, size_bytes, sha256, upload_status +) VALUES ( + '24300000-0000-4000-8000-000000000006', '24000000-0000-4000-8000-000000000001', + '24100000-0000-4000-8000-000000000004', 'recording', + '24000000-0000-4000-8000-000000000001/history/desk.wav', 'audio/wav', 10, repeat('f', 64), 'uploaded' +); +DELETE FROM public.history WHERE id = '24100000-0000-4000-8000-000000000004'; +INSERT INTO public.history (id, user_id, original_text, duration) +VALUES ('24100000-0000-4000-8000-000000000004', '24000000-0000-4000-8000-000000000001', 'desktop again', 1); +INSERT INTO public.audio_files ( + user_id, history_id, source, storage_key, mime_type, size_bytes, sha256, upload_status +) VALUES ( + '24000000-0000-4000-8000-000000000001', '24100000-0000-4000-8000-000000000004', 'recording', + '24000000-0000-4000-8000-000000000001/history/desk.wav', 'audio/wav', 10, repeat('f', 64), 'uploaded' +) +ON CONFLICT (user_id, sha256, storage_key) DO UPDATE +SET history_id = EXCLUDED.history_id, upload_status = EXCLUDED.upload_status; + +SELECT pg_temp.assert_true( + (SELECT count(*) FROM public.audio_files + WHERE storage_key = '24000000-0000-4000-8000-000000000001/history/desk.wav') = 1 + AND EXISTS (SELECT 1 FROM public.audio_files + WHERE storage_key = '24000000-0000-4000-8000-000000000001/history/desk.wav' + AND upload_status = 'uploaded' + AND history_id = '24100000-0000-4000-8000-000000000004') + AND NOT EXISTS (SELECT 1 FROM public.audio_purge_queue WHERE storage_key LIKE '%/desk.wav'), + 'desktop upsert reclaims a deleted key' +); + +-- 9) Account deletion still cascades (the queue has no FK to auth.users). +DELETE FROM auth.users WHERE id = '24000000-0000-4000-8000-000000000002'; +SELECT pg_temp.assert_true( + NOT EXISTS (SELECT 1 FROM public.audio_files WHERE user_id = '24000000-0000-4000-8000-000000000002'), + 'account deletion removes the user audio rows' +); + +-- 10) The dispatcher is a no-op when there is no ready work. +SELECT pg_temp.assert_true( + public.dispatch_audio_purge_v1() IS NULL OR NOT EXISTS ( + SELECT 1 FROM public.audio_purge_queue WHERE user_id = '24000000-0000-4000-8000-000000000001' + ), + 'dispatcher does not fail' +); + +ROLLBACK;