fix(storage): purge raw audio when its history entry or meeting is deleted

This commit is contained in:
Yun Chan 2026-09-28 02:16:22 +09:00
parent 4807a5283d
commit ed790e672b
10 changed files with 1073 additions and 39 deletions

View file

@ -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<string, StorageEntry[]>, private failRemoveAt = -1) {}
list(prefix: string, offset: number, limit: number): Promise<StorageEntry[]> {
this.listCalls.push({ prefix, offset })
return Promise.resolve((this.tree[prefix] ?? []).slice(offset, offset + limit))
}
remove(paths: string[]): Promise<void> {
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')
})

View file

@ -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<StorageEntry[]>
remove(paths: string[]): Promise<void>
}
/** Every object path under prefix, depth-first, in name order. */
export async function listStorageFiles(storage: StorageBucketPort, prefix: string): Promise<string[]> {
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<StorageBucketPort, 'remove'>,
paths: readonly string[],
): Promise<void> {
for (let index = 0; index < paths.length; index += STORAGE_BATCH_SIZE) {
await storage.remove(paths.slice(index, index + STORAGE_BATCH_SIZE))
}
}
type StorageClient = Pick<SupabaseClient, 'storage'>
/** 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}`)
},
}
}

View file

@ -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<string, unknown>, status = 200): Response {
})
}
async function listStorageFiles(
serviceClient: ReturnType<typeof createServiceRoleClient>,
bucket: string,
prefix: string,
): Promise<string[]> {
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)

View file

@ -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<PurgeQueueEntry[]> {
return Promise.resolve(this.batches.shift() ?? [])
}
complete(ids: number[], leaseToken: string): Promise<number> {
this.completed.push({ ids, leaseToken })
return Promise.resolve(ids.length)
}
}
class FakeStorage {
removed: string[][] = []
constructor(private failing = false) {}
remove(paths: string[]): Promise<void> {
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')
})

View file

@ -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<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 })
}
}

View file

@ -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<boolean> {
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),
}
},
})
})

View file

@ -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 }
}

View file

@ -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<string, unknown>
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<SupabaseClient, 'rpc'>
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
},
}
}