d3ro-voice/server/supabase/migrations/20260929000004_audio_retention_on_delete.sql

280 lines
9.9 KiB
PL/PgSQL

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