-- 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()$$ );