From 31fe40b978b5d01c77722856ae20b061ef9647ed Mon Sep 17 00:00:00 2001 From: Martin Donadieu Date: Wed, 22 Jul 2026 11:47:02 +0300 Subject: [PATCH 01/16] fix(db): reclaim Capgo-EU swap pressure without upgrading Batched queue/audit cleanups, hourly http_response truncate, slim audit payloads, and allow nulling migrated dual-storage manifest arrays that were blocked by bundle_already_ready. Co-authored-by: Cursor --- scripts/ops/reclaim_supabase_swap.sql | 167 +++++ scripts/ops/verify_supabase_swap.sql | 63 ++ ...0260722082019_fix_supabase_swap_memory.sql | 654 ++++++++++++++++++ tests/audit-logs.test.ts | 3 + tests/cleanup_swap_memory.test.ts | 168 +++++ 5 files changed, 1055 insertions(+) create mode 100644 scripts/ops/reclaim_supabase_swap.sql create mode 100644 scripts/ops/verify_supabase_swap.sql create mode 100644 supabase/migrations/20260722082019_fix_supabase_swap_memory.sql create mode 100644 tests/cleanup_swap_memory.test.ts diff --git a/scripts/ops/reclaim_supabase_swap.sql b/scripts/ops/reclaim_supabase_swap.sql new file mode 100644 index 0000000000..51463c7fb1 --- /dev/null +++ b/scripts/ops/reclaim_supabase_swap.sql @@ -0,0 +1,167 @@ +-- Capgo-EU Phase A reclaim (run manually in a maintenance window). +-- Safe order: truncate empty bloat -> batched archive deletes -> null dual manifests -> trim audit. +-- Do NOT wrap the whole file in one transaction. VACUUM cannot run inside a transaction block. +-- Prefer Supabase SQL editor / psql as postgres. Re-run sections until counts hit zero. + +\timing on +\set ON_ERROR_STOP on + +-- --------------------------------------------------------------------------- +-- 0) Baseline sizes +-- --------------------------------------------------------------------------- +SELECT pg_size_pretty(pg_database_size(current_database())::bigint) AS db_size; + +SELECT + relname, + n_live_tup, + pg_size_pretty(pg_total_relation_size(format('%I.%I', schemaname, relname)::regclass)::bigint) AS total +FROM pg_stat_user_tables +WHERE (schemaname, relname) IN ( + ('net', '_http_response'), + ('public', 'audit_logs'), + ('public', 'app_versions'), + ('public', 'manifest'), + ('pgmq', 'a_on_version_update'), + ('pgmq', 'a_on_manifest_create'), + ('pgmq', 'a_webhook_dispatcher'), + ('pgmq', 'a_on_channel_update') +) +ORDER BY pg_total_relation_size(format('%I.%I', schemaname, relname)::regclass) DESC; + +-- --------------------------------------------------------------------------- +-- 1) Truncate pg_net response bloat (~5GB empty table in prod) +-- --------------------------------------------------------------------------- +TRUNCATE TABLE net._http_response; + +-- --------------------------------------------------------------------------- +-- 2) Purge pgmq archives older than 2 days (batched). Repeat until deleted=0. +-- --------------------------------------------------------------------------- +DO $$ +DECLARE + queue_name text; + cutoff timestamptz := now() - interval '2 days'; + batch_size integer := 10000; + deleted_batch integer; + deleted_total bigint; +BEGIN + FOREACH queue_name IN ARRAY ARRAY[ + 'on_version_update', + 'on_manifest_create', + 'webhook_dispatcher', + 'on_channel_update' + ] + LOOP + deleted_total := 0; + LOOP + EXECUTE format( + 'WITH doomed AS ( + SELECT ctid FROM pgmq.a_%I WHERE archived_at < $1 LIMIT $2 + ) + DELETE FROM pgmq.a_%I AS archive + USING doomed + WHERE archive.ctid = doomed.ctid', + queue_name, queue_name + ) USING cutoff, batch_size; + GET DIAGNOSTICS deleted_batch = ROW_COUNT; + deleted_total := deleted_total + deleted_batch; + EXIT WHEN deleted_batch = 0; + END LOOP; + RAISE NOTICE 'purged % rows from pgmq.a_%', deleted_total, queue_name; + END LOOP; +END $$; + +VACUUM (VERBOSE) pgmq.a_on_version_update; +VACUUM (VERBOSE) pgmq.a_on_manifest_create; +VACUUM (VERBOSE) pgmq.a_webhook_dispatcher; +VACUUM (VERBOSE) pgmq.a_on_channel_update; + +-- Optional hard reclaim if VACUUM leaves a lot of empty pages (takes stronger locks): +-- VACUUM (FULL, VERBOSE) pgmq.a_on_version_update; +-- VACUUM (FULL, VERBOSE) pgmq.a_on_manifest_create; +-- VACUUM (FULL, VERBOSE) pgmq.a_webhook_dispatcher; +-- VACUUM (FULL, VERBOSE) pgmq.a_on_channel_update; + +-- --------------------------------------------------------------------------- +-- 3) Null migrated app_versions.manifest arrays (dual storage leftovers) +-- --------------------------------------------------------------------------- +DO $$ +DECLARE + batch_size integer := 200; + updated_batch integer; + updated_total bigint := 0; +BEGIN + LOOP + WITH doomed AS ( + SELECT av.id + FROM public.app_versions AS av + WHERE av.manifest IS NOT NULL + AND cardinality(av.manifest) > 0 + AND EXISTS ( + SELECT 1 FROM public.manifest AS m WHERE m.app_version_id = av.id + ) + ORDER BY av.id + LIMIT batch_size + ) + UPDATE public.app_versions AS av + SET manifest = NULL + FROM doomed + WHERE av.id = doomed.id; + GET DIAGNOSTICS updated_batch = ROW_COUNT; + updated_total := updated_total + updated_batch; + EXIT WHEN updated_batch = 0; + END LOOP; + RAISE NOTICE 'nulled manifest arrays on % versions', updated_total; +END $$; + +VACUUM (VERBOSE) public.app_versions; + +-- --------------------------------------------------------------------------- +-- 4) Trim audit_logs older than 30 days (batched) +-- --------------------------------------------------------------------------- +DO $$ +DECLARE + cutoff timestamptz := now() - interval '30 days'; + batch_size integer := 5000; + deleted_batch integer; + deleted_total bigint := 0; +BEGIN + LOOP + WITH doomed AS ( + SELECT ctid + FROM public.audit_logs + WHERE created_at < cutoff + LIMIT batch_size + ) + DELETE FROM public.audit_logs AS audit_logs + USING doomed + WHERE audit_logs.ctid = doomed.ctid; + GET DIAGNOSTICS deleted_batch = ROW_COUNT; + deleted_total := deleted_total + deleted_batch; + EXIT WHEN deleted_batch = 0; + END LOOP; + RAISE NOTICE 'deleted % audit_logs rows older than 30 days', deleted_total; +END $$; + +VACUUM (VERBOSE) public.audit_logs; + +-- --------------------------------------------------------------------------- +-- 5) Final sizes +-- --------------------------------------------------------------------------- +SELECT pg_size_pretty(pg_database_size(current_database())::bigint) AS db_size_after; + +SELECT + relname, + n_live_tup, + pg_size_pretty(pg_total_relation_size(format('%I.%I', schemaname, relname)::regclass)::bigint) AS total +FROM pg_stat_user_tables +WHERE (schemaname, relname) IN ( + ('net', '_http_response'), + ('public', 'audit_logs'), + ('public', 'app_versions'), + ('public', 'manifest'), + ('pgmq', 'a_on_version_update'), + ('pgmq', 'a_on_manifest_create'), + ('pgmq', 'a_webhook_dispatcher'), + ('pgmq', 'a_on_channel_update') +) +ORDER BY pg_total_relation_size(format('%I.%I', schemaname, relname)::regclass) DESC; diff --git a/scripts/ops/verify_supabase_swap.sql b/scripts/ops/verify_supabase_swap.sql new file mode 100644 index 0000000000..6b976c92f1 --- /dev/null +++ b/scripts/ops/verify_supabase_swap.sql @@ -0,0 +1,63 @@ +-- Post-deploy / post-reclaim verification for Capgo-EU swap pressure. + +SELECT pg_size_pretty(pg_database_size(current_database())::bigint) AS db_size; + +SELECT + name, + setting +FROM pg_settings +WHERE name IN ('shared_buffers', 'work_mem', 'max_connections'); + +SELECT + relname, + n_live_tup, + pg_size_pretty(pg_total_relation_size(format('%I.%I', schemaname, relname)::regclass)::bigint) AS total +FROM pg_stat_user_tables +WHERE (schemaname, relname) IN ( + ('net', '_http_response'), + ('public', 'audit_logs'), + ('public', 'app_versions'), + ('public', 'manifest'), + ('pgmq', 'a_on_version_update'), + ('pgmq', 'a_on_manifest_create'), + ('pgmq', 'a_webhook_dispatcher'), + ('pgmq', 'a_on_channel_update') +) +ORDER BY pg_total_relation_size(format('%I.%I', schemaname, relname)::regclass) DESC; + +SELECT count(*) AS versions_with_array_manifest +FROM public.app_versions +WHERE manifest IS NOT NULL AND cardinality(manifest) > 0; + +SELECT name, enabled, hour_interval, run_at_hour, run_at_minute, target +FROM public.cron_tasks +WHERE name IN ( + 'cleanup_queue_messages', + 'cleanup_net_http_response', + 'cleanup_old_audit_logs', + 'null_migrated_app_version_manifests' +) +ORDER BY name; + +SELECT status, return_message, count(*) AS n +FROM cron.job_run_details +WHERE start_time > now() - interval '24 hours' + OR (start_time IS NULL AND status IN ('failed', 'connecting')) +GROUP BY status, return_message +ORDER BY n DESC +LIMIT 20; + +SELECT + count(*) FILTER (WHERE archived_at < now() - interval '2 days') AS archives_older_than_2d, + count(*) AS archives_total +FROM pgmq.a_on_manifest_create; + +SELECT + 'index hit rate' AS name, + ROUND((sum(idx_blks_hit) / nullif(sum(idx_blks_hit + idx_blks_read), 0) * 100)::numeric, 2) AS ratio +FROM pg_statio_user_indexes +UNION ALL +SELECT + 'table hit rate', + ROUND((sum(heap_blks_hit) / nullif(sum(heap_blks_hit) + sum(heap_blks_read), 0) * 100)::numeric, 2) +FROM pg_statio_user_tables; diff --git a/supabase/migrations/20260722082019_fix_supabase_swap_memory.sql b/supabase/migrations/20260722082019_fix_supabase_swap_memory.sql new file mode 100644 index 0000000000..bfc1f21404 --- /dev/null +++ b/supabase/migrations/20260722082019_fix_supabase_swap_memory.sql @@ -0,0 +1,654 @@ +-- Reduce Capgo-EU swap pressure without upgrading compute: +-- 1) batched pgmq archive cleanup (2-day retention, hourly) +-- 2) reclaim net._http_response (truncate hourly) +-- 3) slim audit payloads + batched 30-day audit cleanup +-- 4) null leftover app_versions.manifest arrays after table migration +-- 5) also omit native_packages from app_versions queue payloads + +-- --------------------------------------------------------------------------- +-- Batched pgmq archive / stuck-message cleanup +-- --------------------------------------------------------------------------- +CREATE OR REPLACE FUNCTION "public"."cleanup_queue_messages"() RETURNS "void" + LANGUAGE "plpgsql" SECURITY DEFINER + SET "search_path" TO '' + AS $$ +DECLARE + queue_name text; + cutoff timestamptz := pg_catalog.now() - INTERVAL '2 days'; + batch_size integer := 10000; + deleted_batch integer; + deleted_total bigint; +BEGIN + FOR queue_name IN ( + SELECT q.queue_name FROM pgmq.list_queues() q + ) LOOP + deleted_total := 0; + + LOOP + EXECUTE format( + 'WITH doomed AS ( + SELECT ctid + FROM pgmq.a_%I + WHERE archived_at < $1 + LIMIT $2 + ) + DELETE FROM pgmq.a_%I AS archive + USING doomed + WHERE archive.ctid = doomed.ctid', + queue_name, + queue_name + ) + USING cutoff, batch_size; + + GET DIAGNOSTICS deleted_batch = ROW_COUNT; + deleted_total := deleted_total + deleted_batch; + EXIT WHEN deleted_batch = 0; + END LOOP; + + IF deleted_total > 0 THEN + RAISE NOTICE 'cleanup_queue_messages: deleted % archived rows from a_%', deleted_total, queue_name; + END IF; + + deleted_total := 0; + LOOP + EXECUTE format( + 'WITH doomed AS ( + SELECT ctid + FROM pgmq.q_%I + WHERE read_ct > 5 + LIMIT $1 + ) + DELETE FROM pgmq.q_%I AS queue + USING doomed + WHERE queue.ctid = doomed.ctid', + queue_name, + queue_name + ) + USING batch_size; + + GET DIAGNOSTICS deleted_batch = ROW_COUNT; + deleted_total := deleted_total + deleted_batch; + EXIT WHEN deleted_batch = 0; + END LOOP; + + IF deleted_total > 0 THEN + RAISE NOTICE 'cleanup_queue_messages: deleted % stuck rows from q_%', deleted_total, queue_name; + END IF; + END LOOP; +END; +$$; + +ALTER FUNCTION "public"."cleanup_queue_messages"() OWNER TO "postgres"; +REVOKE ALL ON FUNCTION "public"."cleanup_queue_messages"() FROM PUBLIC; +GRANT ALL ON FUNCTION "public"."cleanup_queue_messages"() TO "service_role"; + +-- --------------------------------------------------------------------------- +-- Reclaim pg_net response bloat (DELETE alone never shrinks the table) +-- --------------------------------------------------------------------------- +CREATE OR REPLACE FUNCTION "public"."cleanup_net_http_response"() RETURNS "void" + LANGUAGE "plpgsql" SECURITY DEFINER + SET "search_path" TO '' + AS $$ +BEGIN + -- Responses are only used for short-lived async HTTP debugging. + -- Truncate reclaims disk; row deletes do not. + TRUNCATE TABLE net._http_response; + RAISE NOTICE 'cleanup_net_http_response: truncated net._http_response'; +END; +$$; + +ALTER FUNCTION "public"."cleanup_net_http_response"() OWNER TO "postgres"; +REVOKE ALL ON FUNCTION "public"."cleanup_net_http_response"() FROM PUBLIC; +GRANT ALL ON FUNCTION "public"."cleanup_net_http_response"() TO "service_role"; + +-- --------------------------------------------------------------------------- +-- Batched audit log cleanup (30-day retention) +-- --------------------------------------------------------------------------- +CREATE OR REPLACE FUNCTION "public"."cleanup_old_audit_logs"() RETURNS "void" + LANGUAGE "plpgsql" SECURITY DEFINER + SET "search_path" TO '' + AS $$ +DECLARE + cutoff timestamptz := pg_catalog.now() - INTERVAL '30 days'; + batch_size integer := 5000; + deleted_batch integer; + deleted_total bigint := 0; +BEGIN + LOOP + WITH doomed AS ( + SELECT ctid + FROM public.audit_logs + WHERE created_at < cutoff + LIMIT batch_size + ) + DELETE FROM public.audit_logs AS audit_logs + USING doomed + WHERE audit_logs.ctid = doomed.ctid; + + GET DIAGNOSTICS deleted_batch = ROW_COUNT; + deleted_total := deleted_total + deleted_batch; + EXIT WHEN deleted_batch = 0; + END LOOP; + + IF deleted_total > 0 THEN + RAISE NOTICE 'cleanup_old_audit_logs: deleted % rows older than 30 days', deleted_total; + END IF; +END; +$$; + +ALTER FUNCTION "public"."cleanup_old_audit_logs"() OWNER TO "postgres"; +REVOKE ALL ON FUNCTION "public"."cleanup_old_audit_logs"() FROM PUBLIC; +GRANT ALL ON FUNCTION "public"."cleanup_old_audit_logs"() TO "service_role"; + +-- --------------------------------------------------------------------------- +-- Null leftover dual-storage app_versions.manifest arrays +-- --------------------------------------------------------------------------- +CREATE OR REPLACE FUNCTION "public"."null_migrated_app_version_manifests"() RETURNS "void" + LANGUAGE "plpgsql" SECURITY DEFINER + SET "search_path" TO '' + AS $$ +DECLARE + batch_size integer := 200; + updated_batch integer; + updated_total bigint := 0; +BEGIN + LOOP + WITH doomed AS ( + SELECT av.id + FROM public.app_versions AS av + WHERE av.manifest IS NOT NULL + AND pg_catalog.cardinality(av.manifest) > 0 + AND EXISTS ( + SELECT 1 + FROM public.manifest AS m + WHERE m.app_version_id = av.id + ) + ORDER BY av.id + LIMIT batch_size + ) + UPDATE public.app_versions AS av + SET manifest = NULL + FROM doomed + WHERE av.id = doomed.id; + + GET DIAGNOSTICS updated_batch = ROW_COUNT; + updated_total := updated_total + updated_batch; + EXIT WHEN updated_batch = 0; + END LOOP; + + IF updated_total > 0 THEN + RAISE NOTICE 'null_migrated_app_version_manifests: nulled manifest arrays on % versions', updated_total; + END IF; +END; +$$; + +ALTER FUNCTION "public"."null_migrated_app_version_manifests"() OWNER TO "postgres"; +REVOKE ALL ON FUNCTION "public"."null_migrated_app_version_manifests"() FROM PUBLIC; +GRANT ALL ON FUNCTION "public"."null_migrated_app_version_manifests"() TO "service_role"; + +-- --------------------------------------------------------------------------- +-- Slim audit payloads (app_versions fat columns) +-- --------------------------------------------------------------------------- +CREATE OR REPLACE FUNCTION "public"."audit_log_trigger"() RETURNS "trigger" + LANGUAGE "plpgsql" SECURITY DEFINER + SET "search_path" TO '' + AS $$ +DECLARE + v_old_record jsonb; + v_new_record jsonb; + v_changed_fields text[]; + v_org_id uuid; + v_record_id text; + v_user_id uuid; + v_key text; + v_api_key_text text; + v_api_key public.apikeys%ROWTYPE; + v_actor_type text := 'system'; + v_actor_user_id uuid; + v_actor_user_email text; + v_actor_apikey_id bigint; + v_actor_apikey_name text; + v_stats_refresh_fields constant text[] := ARRAY['stats_refresh_requested_at', 'stats_updated_at', 'updated_at']; + v_background_counter_fields constant text[] := ARRAY['channel_device_count', 'manifest_bundle_count', 'updated_at']; + v_fat_app_version_fields constant text[] := ARRAY['manifest', 'native_packages']; +BEGIN + SELECT auth.uid() INTO v_actor_user_id; + + IF v_actor_user_id IS NOT NULL THEN + v_actor_type := 'user'; + ELSE + SELECT public.get_apikey_header() INTO v_api_key_text; + + IF v_api_key_text IS NOT NULL THEN + SELECT * + INTO v_api_key + FROM public.find_apikey_by_value(v_api_key_text) + LIMIT 1; + + -- Attribute only valid, write-capable API keys; a read-only key present on + -- a request must not be recorded as the actor of a mutation. + IF v_api_key.id IS NOT NULL + AND NOT public.is_apikey_expired(v_api_key.expires_at) + AND ( + public.is_allowed_capgkey(v_api_key_text, '{upload}'::text[]) + OR public.is_allowed_capgkey(v_api_key_text, '{write}'::text[]) + OR public.is_allowed_capgkey(v_api_key_text, '{all}'::text[]) + ) THEN + v_actor_type := 'apikey'; + v_actor_user_id := v_api_key.user_id; + v_actor_apikey_id := v_api_key.id; + v_actor_apikey_name := v_api_key.name; + END IF; + END IF; + END IF; + + IF v_actor_user_id IS NOT NULL THEN + SELECT users.email + INTO v_actor_user_email + FROM public.users AS users + WHERE users.id = v_actor_user_id; + END IF; + + v_user_id := v_actor_user_id; + + IF TG_OP = 'DELETE' THEN + v_old_record := pg_catalog.to_jsonb(OLD); + v_new_record := NULL; + ELSIF TG_OP = 'INSERT' THEN + v_old_record := NULL; + v_new_record := pg_catalog.to_jsonb(NEW); + ELSE + v_old_record := pg_catalog.to_jsonb(OLD); + v_new_record := pg_catalog.to_jsonb(NEW); + + FOR v_key IN SELECT pg_catalog.jsonb_object_keys(v_new_record) + LOOP + IF v_old_record->v_key IS DISTINCT FROM v_new_record->v_key THEN + v_changed_fields := pg_catalog.array_append(v_changed_fields, v_key); + END IF; + END LOOP; + + IF TG_TABLE_NAME = ANY(ARRAY['apps', 'orgs']) + AND v_changed_fields && ARRAY['stats_refresh_requested_at', 'stats_updated_at'] + AND NOT EXISTS ( + SELECT 1 + FROM pg_catalog.unnest(v_changed_fields) AS changed_field(field_name) + WHERE changed_field.field_name <> ALL(v_stats_refresh_fields) + ) THEN + RETURN NEW; + END IF; + + IF v_actor_type = 'system' + AND TG_TABLE_NAME = 'apps' + AND v_changed_fields && ARRAY['channel_device_count', 'manifest_bundle_count'] + AND NOT EXISTS ( + SELECT 1 + FROM pg_catalog.unnest(v_changed_fields) AS changed_field(field_name) + WHERE changed_field.field_name <> ALL(v_background_counter_fields) + ) THEN + RETURN NEW; + END IF; + END IF; + + -- Never persist multi-MB array/json columns in audit TOAST. + IF TG_TABLE_NAME = 'app_versions' THEN + IF v_old_record IS NOT NULL THEN + v_old_record := v_old_record - v_fat_app_version_fields; + END IF; + IF v_new_record IS NOT NULL THEN + v_new_record := v_new_record - v_fat_app_version_fields; + END IF; + IF v_changed_fields IS NOT NULL THEN + SELECT pg_catalog.array_agg(field_name) + INTO v_changed_fields + FROM pg_catalog.unnest(v_changed_fields) AS changed_field(field_name) + WHERE changed_field.field_name <> ALL(v_fat_app_version_fields); + END IF; + + -- Skip updates that only touched stripped fat columns (e.g. manifest nulling). + IF TG_OP = 'UPDATE' + AND (v_changed_fields IS NULL OR pg_catalog.cardinality(v_changed_fields) = 0) THEN + RETURN NEW; + END IF; + END IF; + + CASE TG_TABLE_NAME + WHEN 'orgs' THEN + v_org_id := COALESCE(NEW.id, OLD.id); + v_record_id := COALESCE(NEW.id, OLD.id)::text; + WHEN 'apps' THEN + v_org_id := COALESCE(NEW.owner_org, OLD.owner_org); + v_record_id := COALESCE(NEW.app_id, OLD.app_id)::text; + WHEN 'channels' THEN + v_org_id := COALESCE(NEW.owner_org, OLD.owner_org); + v_record_id := COALESCE(NEW.id, OLD.id)::text; + WHEN 'app_versions' THEN + v_org_id := COALESCE(NEW.owner_org, OLD.owner_org); + v_record_id := COALESCE(NEW.id, OLD.id)::text; + WHEN 'org_users' THEN + v_org_id := COALESCE(NEW.org_id, OLD.org_id); + v_record_id := COALESCE(NEW.id, OLD.id)::text; + ELSE + v_org_id := NULL; + v_record_id := NULL; + END CASE; + + IF v_org_id IS NOT NULL THEN + INSERT INTO public.audit_logs ( + table_name, + record_id, + operation, + user_id, + org_id, + old_record, + new_record, + changed_fields, + actor_type, + actor_user_id, + actor_user_email, + actor_apikey_id, + actor_apikey_name + ) VALUES ( + TG_TABLE_NAME, + v_record_id, + TG_OP, + v_user_id, + v_org_id, + v_old_record, + v_new_record, + v_changed_fields, + v_actor_type, + v_actor_user_id, + v_actor_user_email, + v_actor_apikey_id, + v_actor_apikey_name + ); + END IF; + + RETURN COALESCE(NEW, OLD); +END; +$$; + +ALTER FUNCTION "public"."audit_log_trigger"() OWNER TO "postgres"; + +-- --------------------------------------------------------------------------- +-- Also omit native_packages from app_versions queue payloads +-- --------------------------------------------------------------------------- +CREATE OR REPLACE FUNCTION "public"."trigger_http_queue_post_to_function"() RETURNS "trigger" + LANGUAGE "plpgsql" SECURITY DEFINER + SET "search_path" TO '' + AS $$ +DECLARE + payload jsonb; + record_payload jsonb; + old_record_payload jsonb; + function_type text; +BEGIN + function_type := CASE + WHEN NULLIF(TG_ARGV[1], '') IS NULL THEN 'cloudflare' + WHEN lower(TG_ARGV[1]) = 'supabase' THEN 'cloudflare' + ELSE TG_ARGV[1] + END; + + record_payload := to_jsonb(NEW); + old_record_payload := to_jsonb(OLD); + + -- app_versions fat columns can be multi-MB. Never enqueue them; handlers reload when needed. + IF TG_TABLE_NAME = 'app_versions' THEN + IF record_payload IS NOT NULL THEN + record_payload := record_payload - 'manifest' - 'native_packages'; + END IF; + IF old_record_payload IS NOT NULL THEN + old_record_payload := old_record_payload - 'manifest' - 'native_packages'; + END IF; + END IF; + + payload := jsonb_build_object( + 'function_name', TG_ARGV[0], + 'function_type', function_type, + 'payload', jsonb_build_object( + 'old_record', old_record_payload, + 'record', record_payload, + 'type', TG_OP, + 'table', TG_TABLE_NAME, + 'schema', TG_TABLE_SCHEMA + ) + ); + + IF TG_ARGV[0] IS NOT NULL THEN + PERFORM "pgmq"."send"(TG_ARGV[0], payload); + END IF; + RETURN NEW; +END; +$$; + +ALTER FUNCTION "public"."trigger_http_queue_post_to_function"() OWNER TO "postgres"; + +-- --------------------------------------------------------------------------- +-- Cron schedules: make reclaim jobs hourly / reliable +-- --------------------------------------------------------------------------- +UPDATE public.cron_tasks +SET + second_interval = NULL, + minute_interval = NULL, + hour_interval = 1, + run_at_hour = NULL, + run_at_minute = 0, + run_at_second = NULL, + updated_at = pg_catalog.now() +WHERE name = 'cleanup_queue_messages'; + +UPDATE public.cron_tasks +SET + second_interval = NULL, + minute_interval = NULL, + hour_interval = NULL, + run_at_hour = 3, + run_at_minute = 0, + run_at_second = NULL, + updated_at = pg_catalog.now() +WHERE name = 'cleanup_old_audit_logs'; + +INSERT INTO public.cron_tasks ( + name, + description, + task_type, + target, + batch_size, + payload, + second_interval, + minute_interval, + hour_interval, + run_at_hour, + run_at_minute, + run_at_second, + run_on_dow, + run_on_day, + enabled +) +VALUES + ( + 'cleanup_net_http_response', + 'Truncate net._http_response so pg_net response history cannot bloat disk/RAM', + 'function', + 'public.cleanup_net_http_response()', + NULL, + NULL, + NULL, + NULL, + 1, + NULL, + 5, + NULL, + NULL, + NULL, + true + ), + ( + 'null_migrated_app_version_manifests', + 'Null leftover app_versions.manifest arrays after rows exist in public.manifest', + 'function', + 'public.null_migrated_app_version_manifests()', + NULL, + NULL, + NULL, + NULL, + 1, + NULL, + 15, + NULL, + NULL, + NULL, + true + ) +ON CONFLICT (name) DO UPDATE +SET + description = EXCLUDED.description, + task_type = EXCLUDED.task_type, + target = EXCLUDED.target, + hour_interval = EXCLUDED.hour_interval, + run_at_minute = EXCLUDED.run_at_minute, + run_at_hour = NULL, + second_interval = NULL, + minute_interval = NULL, + enabled = true, + updated_at = pg_catalog.now(); + +-- --------------------------------------------------------------------------- +-- Allow clearing dual-storage fat columns +-- after upload (null only). Unblocks on_version_update + reclaim jobs. +-- --------------------------------------------------------------------------- +CREATE OR REPLACE FUNCTION "public"."check_encrypted_bundle_on_insert"() RETURNS "trigger" + LANGUAGE "plpgsql" SECURITY DEFINER + SET "search_path" TO '' + AS $$ +DECLARE + org_id uuid; + org_enforcing boolean; + org_required_key varchar(21); + bundle_is_encrypted boolean; + bundle_key_id varchar(20); + bundle_was_ready boolean; +BEGIN + IF TG_OP = 'UPDATE' THEN + bundle_was_ready := OLD.storage_provider IS DISTINCT FROM 'r2-direct'; + + -- Nulling migrated dual-storage columns is allowed after upload completes. + -- Rewriting non-null manifest/native_packages content stays locked. + IF bundle_was_ready + AND ( + NEW.name IS DISTINCT FROM OLD.name + OR NEW.app_id IS DISTINCT FROM OLD.app_id + OR NEW.session_key IS DISTINCT FROM OLD.session_key + OR NEW.key_id IS DISTINCT FROM OLD.key_id + OR NEW.storage_provider IS DISTINCT FROM OLD.storage_provider + OR NEW.r2_path IS DISTINCT FROM OLD.r2_path + OR NEW.external_url IS DISTINCT FROM OLD.external_url + OR NEW.checksum IS DISTINCT FROM OLD.checksum + OR (NEW.manifest IS DISTINCT FROM OLD.manifest AND NEW.manifest IS NOT NULL) + OR (NEW.native_packages IS DISTINCT FROM OLD.native_packages AND NEW.native_packages IS NOT NULL) + ) + THEN + PERFORM public.pg_log('deny: BUNDLE_CONTENT_LOCKED_TRIGGER', + jsonb_build_object( + 'org_id', OLD.owner_org, + 'app_id', OLD.app_id, + 'version_name', OLD.name, + 'user_id', OLD.user_id, + 'old_storage_provider', OLD.storage_provider, + 'new_storage_provider', NEW.storage_provider, + 'reason', 'bundle_ready' + )); + RAISE EXCEPTION '%', + 'bundle_already_ready: Bundle content cannot be changed ' + || 'after upload is complete. Upload a new bundle instead.'; + END IF; + END IF; + + -- Derive org_id from NEW.app_id first because + -- force_valid_owner_org_app_versions runs after this trigger. + SELECT apps.owner_org INTO org_id + FROM public.apps + WHERE apps.app_id = NEW.app_id; + + IF org_id IS NULL THEN + org_id := NEW.owner_org; + END IF; + + -- If org not found, allow the existing foreign-key/owner checks to fail. + IF org_id IS NULL THEN + RETURN NEW; + END IF; + + SELECT enforce_encrypted_bundles, required_encryption_key + INTO org_enforcing, org_required_key + FROM public.orgs + WHERE id = org_id; + + IF org_enforcing IS NULL OR org_enforcing = false THEN + RETURN NEW; + END IF; + + bundle_is_encrypted := public.is_bundle_encrypted(NEW.session_key); + bundle_key_id := NULLIF(btrim(NEW.key_id), '')::varchar(20); + + IF NOT bundle_is_encrypted THEN + PERFORM public.pg_log('deny: ORG_REQUIRES_ENCRYPTED_BUNDLES_TRIGGER', + jsonb_build_object( + 'org_id', org_id, + 'app_id', NEW.app_id, + 'version_name', NEW.name, + 'user_id', NEW.user_id, + 'reason', 'not_encrypted' + )); + RAISE EXCEPTION '%', + 'encryption_required: This organization requires all bundles to be ' + || 'encrypted. Please upload an encrypted bundle with a session_key.'; + END IF; + + IF org_required_key IS NOT NULL AND org_required_key <> '' THEN + IF bundle_key_id IS NULL THEN + PERFORM public.pg_log('deny: ORG_REQUIRES_SPECIFIC_ENCRYPTION_KEY_TRIGGER', + jsonb_build_object( + 'org_id', org_id, + 'app_id', NEW.app_id, + 'version_name', NEW.name, + 'user_id', NEW.user_id, + 'required_key', org_required_key, + 'bundle_key_id', bundle_key_id, + 'reason', 'missing_key_id' + )); + RAISE EXCEPTION '%', + 'encryption_key_required: This organization requires bundles to be ' + || 'encrypted with a specific key. The uploaded bundle does not have ' + || 'a key_id.'; + END IF; + + -- key_id is 20 chars and required_encryption_key may be 20 or 21 chars. + IF NOT ( + bundle_key_id = LEFT(org_required_key, 20) + OR LEFT(bundle_key_id, LENGTH(org_required_key)) = org_required_key + ) THEN + PERFORM public.pg_log('deny: ORG_REQUIRES_SPECIFIC_ENCRYPTION_KEY_TRIGGER', + jsonb_build_object( + 'org_id', org_id, + 'app_id', NEW.app_id, + 'version_name', NEW.name, + 'user_id', NEW.user_id, + 'required_key', org_required_key, + 'bundle_key_id', bundle_key_id, + 'reason', 'key_mismatch' + )); + RAISE EXCEPTION '%', + 'encryption_key_mismatch: This organization requires bundles to be ' + || 'encrypted with a specific key. The uploaded bundle was encrypted ' + || 'with a different key.'; + END IF; + END IF; + + RETURN NEW; +END; +$$; + + +ALTER FUNCTION "public"."check_encrypted_bundle_on_insert"() OWNER TO "postgres"; diff --git a/tests/audit-logs.test.ts b/tests/audit-logs.test.ts index 4f8fe68956..7ae360b020 100644 --- a/tests/audit-logs.test.ts +++ b/tests/audit-logs.test.ts @@ -558,6 +558,9 @@ describe('audit logs for app_versions via API key', () => { expect(versionAuditLog.new_record).toBeTruthy() if (versionAuditLog.new_record && typeof versionAuditLog.new_record === 'object') { expect((versionAuditLog.new_record as Record).name).toBe(testVersionName) + // Fat columns must stay out of audit TOAST to protect primary DB memory. + expect((versionAuditLog.new_record as Record).manifest).toBeUndefined() + expect((versionAuditLog.new_record as Record).native_packages).toBeUndefined() } } } diff --git a/tests/cleanup_swap_memory.test.ts b/tests/cleanup_swap_memory.test.ts new file mode 100644 index 0000000000..80eb2d2424 --- /dev/null +++ b/tests/cleanup_swap_memory.test.ts @@ -0,0 +1,168 @@ +import { randomUUID } from 'node:crypto' +import { afterAll, describe, expect, it } from 'vitest' +import { cleanupPostgresClient, executeSQL } from './test-utils.ts' + +describe('swap memory cleanup functions', () => { + afterAll(async () => { + await cleanupPostgresClient() + }) + + it('cleanup_queue_messages deletes archived rows older than 2 days in batches', async () => { + const marker = `swap-cleanup-${randomUUID()}` + const baseMsgId = BigInt(Date.now()) * 1000n + + await executeSQL( + `INSERT INTO pgmq.a_on_version_update (msg_id, read_ct, enqueued_at, archived_at, vt, message) + VALUES + ($1, 0, now() - interval '10 days', now() - interval '10 days', now(), $3::jsonb), + ($2, 0, now() - interval '1 hour', now() - interval '1 hour', now(), $4::jsonb)`, + [ + (baseMsgId + 1n).toString(), + (baseMsgId + 2n).toString(), + JSON.stringify({ marker, age: 'old' }), + JSON.stringify({ marker, age: 'fresh' }), + ], + ) + + await executeSQL(`SELECT public.cleanup_queue_messages()`) + + const rows = await executeSQL( + `SELECT message->>'age' AS age + FROM pgmq.a_on_version_update + WHERE message->>'marker' = $1 + ORDER BY age`, + [marker], + ) + + expect(rows).toHaveLength(1) + expect(rows[0]?.age).toBe('fresh') + + await executeSQL( + `DELETE FROM pgmq.a_on_version_update WHERE message->>'marker' = $1`, + [marker], + ) + }) + + it('cleanup_net_http_response truncates net._http_response', async () => { + const id = BigInt(Date.now()) * 1000n + 7n + await executeSQL( + `INSERT INTO net._http_response (id, status_code, content, created) + VALUES ($1, 200, 'swap-cleanup-test', now())`, + [id.toString()], + ) + + await executeSQL(`SELECT public.cleanup_net_http_response()`) + + const rows = await executeSQL( + `SELECT count(*)::int AS n FROM net._http_response`, + ) + expect(rows[0]?.n).toBe(0) + }) + + it('null_migrated_app_version_manifests clears dual-storage arrays', async () => { + const appId = `com.swap.nullmanifest.${randomUUID().slice(0, 8)}` + const orgRows = await executeSQL( + `SELECT id FROM public.orgs ORDER BY created_at LIMIT 1`, + ) + const orgId = orgRows[0]?.id as string + expect(orgId).toBeTruthy() + + await executeSQL( + `INSERT INTO public.apps (app_id, name, icon_url, owner_org) + VALUES ($1, 'swap-null-manifest', '', $2::uuid)`, + [appId, orgId], + ) + + const versionRows = await executeSQL( + `INSERT INTO public.app_versions (app_id, name, owner_org, storage_provider, manifest, manifest_count) + VALUES ( + $1, + $2, + $3::uuid, + 'r2', + ARRAY[ROW('index.html', 'apps/test/index.html', 'abc123')::public.manifest_entry], + 1 + ) + RETURNING id`, + [appId, `1.0.0-${randomUUID().slice(0, 8)}`, orgId], + ) + const versionId = versionRows[0]?.id as number + expect(versionId).toBeTruthy() + + await executeSQL( + `INSERT INTO public.manifest (app_version_id, file_name, s3_path, file_hash) + VALUES ($1, 'index.html', 'apps/test/index.html', 'abc123')`, + [versionId], + ) + + await executeSQL(`SELECT public.null_migrated_app_version_manifests()`) + + const after = await executeSQL( + `SELECT manifest IS NULL AS is_null, manifest_count + FROM public.app_versions + WHERE id = $1`, + [versionId], + ) + expect(after[0]?.is_null).toBe(true) + expect(after[0]?.manifest_count).toBe(1) + + await executeSQL(`DELETE FROM public.manifest WHERE app_version_id = $1`, [versionId]) + await executeSQL(`DELETE FROM public.app_versions WHERE id = $1`, [versionId]) + await executeSQL(`DELETE FROM public.apps WHERE app_id = $1`, [appId]) + }) + + it('audit_log_trigger strips fat app_versions fields', async () => { + const appId = `com.swap.auditfat.${randomUUID().slice(0, 8)}` + const orgRows = await executeSQL( + `SELECT id FROM public.orgs ORDER BY created_at LIMIT 1`, + ) + const orgId = orgRows[0]?.id as string + + await executeSQL( + `INSERT INTO public.apps (app_id, name, icon_url, owner_org) + VALUES ($1, 'swap-audit', '', $2::uuid)`, + [appId, orgId], + ) + + const versionRows = await executeSQL( + `INSERT INTO public.app_versions (app_id, name, owner_org, storage_provider, comment) + VALUES ($1, $2, $3::uuid, 'r2-direct', 'before') + RETURNING id`, + [appId, `1.0.0-${randomUUID().slice(0, 8)}`, orgId], + ) + const versionId = versionRows[0]?.id as number + + await executeSQL( + `UPDATE public.app_versions + SET + comment = 'after', + manifest = ARRAY[ROW('a.js', 'apps/a.js', 'hash')::public.manifest_entry], + native_packages = ARRAY['{"name":"cordova-plugin"}'::jsonb] + WHERE id = $1`, + [versionId], + ) + + const logs = await executeSQL( + `SELECT new_record, changed_fields + FROM public.audit_logs + WHERE table_name = 'app_versions' + AND record_id = $1 + AND operation = 'UPDATE' + ORDER BY id DESC + LIMIT 1`, + [String(versionId)], + ) + + expect(logs).toHaveLength(1) + expect(logs[0]?.new_record?.manifest).toBeUndefined() + expect(logs[0]?.new_record?.native_packages).toBeUndefined() + expect(logs[0]?.new_record?.comment).toBe('after') + expect(logs[0]?.changed_fields).toContain('comment') + expect(logs[0]?.changed_fields ?? []).not.toContain('manifest') + expect(logs[0]?.changed_fields ?? []).not.toContain('native_packages') + + await executeSQL(`DELETE FROM public.audit_logs WHERE record_id = $1 AND table_name = 'app_versions'`, [String(versionId)]) + await executeSQL(`DELETE FROM public.app_versions WHERE id = $1`, [versionId]) + await executeSQL(`DELETE FROM public.apps WHERE app_id = $1`, [appId]) + }) +}) From c63d4dbc9ca81da79fff0370214ce12574ad41ec Mon Sep 17 00:00:00 2001 From: Martin Donadieu Date: Wed, 22 Jul 2026 11:51:42 +0300 Subject: [PATCH 02/16] test(db): align audit cleanup cron test with 30-day retention Co-authored-by: Cursor --- .../20260722082019_fix_supabase_swap_memory.sql | 2 +- supabase/tests/55_test_audit_log_cleanup_cron.sql | 8 ++++---- 2 files changed, 5 insertions(+), 5 deletions(-) diff --git a/supabase/migrations/20260722082019_fix_supabase_swap_memory.sql b/supabase/migrations/20260722082019_fix_supabase_swap_memory.sql index bfc1f21404..921d54e07d 100644 --- a/supabase/migrations/20260722082019_fix_supabase_swap_memory.sql +++ b/supabase/migrations/20260722082019_fix_supabase_swap_memory.sql @@ -445,7 +445,7 @@ SET hour_interval = NULL, run_at_hour = 3, run_at_minute = 0, - run_at_second = NULL, + run_at_second = 0, updated_at = pg_catalog.now() WHERE name = 'cleanup_old_audit_logs'; diff --git a/supabase/tests/55_test_audit_log_cleanup_cron.sql b/supabase/tests/55_test_audit_log_cleanup_cron.sql index 26d785990e..8ebc57e567 100644 --- a/supabase/tests/55_test_audit_log_cleanup_cron.sql +++ b/supabase/tests/55_test_audit_log_cleanup_cron.sql @@ -68,7 +68,7 @@ INSERT INTO public.audit_logs ( ) VALUES ( - now() - interval '91 days', + now() - interval '31 days', 'audit_log_retention_test', 'audit-log-retention-old', 'INSERT', @@ -79,7 +79,7 @@ VALUES ARRAY['retention_probe']::text [] ), ( - now() - interval '89 days', + now() - interval '29 days', 'audit_log_retention_test', 'audit-log-retention-fresh', 'INSERT', @@ -101,7 +101,7 @@ SELECT is( AND table_name = 'audit_log_retention_test' ), 0, - 'cleanup_old_audit_logs deletes rows older than 90 days' + 'cleanup_old_audit_logs deletes rows older than 30 days' ); SELECT is( @@ -113,7 +113,7 @@ SELECT is( AND table_name = 'audit_log_retention_test' ), 1, - 'cleanup_old_audit_logs keeps rows newer than 90 days' + 'cleanup_old_audit_logs keeps rows newer than 30 days' ); SELECT tests.clear_authentication(); From ad18df675c65dfe68e96dcfa770a56395ddad4dd Mon Sep 17 00:00:00 2001 From: Martin Donadieu Date: Wed, 22 Jul 2026 12:01:01 +0300 Subject: [PATCH 03/16] fix(db): make queue cleanup SQL plpgsql_check-safe Also fill created_by_apikey_rbac_id on builtin version stubs so frontend typecheck stays green. Co-authored-by: Cursor --- src/services/versions.ts | 1 + ...0260722082019_fix_supabase_swap_memory.sql | 24 +++++++------------ 2 files changed, 10 insertions(+), 15 deletions(-) diff --git a/src/services/versions.ts b/src/services/versions.ts index a04188e85f..3bed6f3b60 100644 --- a/src/services/versions.ts +++ b/src/services/versions.ts @@ -26,6 +26,7 @@ export function createBuiltinChannelVersion(channel: { cli_version: null, comment: null, created_at: channel.created_at, + created_by_apikey_rbac_id: null, deleted: false, deleted_at: null, external_url: null, diff --git a/supabase/migrations/20260722082019_fix_supabase_swap_memory.sql b/supabase/migrations/20260722082019_fix_supabase_swap_memory.sql index 921d54e07d..5b9b5b9ca2 100644 --- a/supabase/migrations/20260722082019_fix_supabase_swap_memory.sql +++ b/supabase/migrations/20260722082019_fix_supabase_swap_memory.sql @@ -26,15 +26,13 @@ BEGIN LOOP EXECUTE format( - 'WITH doomed AS ( + 'DELETE FROM pgmq.a_%I + WHERE ctid IN ( SELECT ctid FROM pgmq.a_%I WHERE archived_at < $1 LIMIT $2 - ) - DELETE FROM pgmq.a_%I AS archive - USING doomed - WHERE archive.ctid = doomed.ctid', + )', queue_name, queue_name ) @@ -52,15 +50,13 @@ BEGIN deleted_total := 0; LOOP EXECUTE format( - 'WITH doomed AS ( + 'DELETE FROM pgmq.q_%I + WHERE ctid IN ( SELECT ctid FROM pgmq.q_%I WHERE read_ct > 5 LIMIT $1 - ) - DELETE FROM pgmq.q_%I AS queue - USING doomed - WHERE queue.ctid = doomed.ctid', + )', queue_name, queue_name ) @@ -115,15 +111,13 @@ DECLARE deleted_total bigint := 0; BEGIN LOOP - WITH doomed AS ( + DELETE FROM public.audit_logs + WHERE ctid IN ( SELECT ctid FROM public.audit_logs WHERE created_at < cutoff LIMIT batch_size - ) - DELETE FROM public.audit_logs AS audit_logs - USING doomed - WHERE audit_logs.ctid = doomed.ctid; + ); GET DIAGNOSTICS deleted_batch = ROW_COUNT; deleted_total := deleted_total + deleted_batch; From 0347360504b3e71289e3a19f3c19169ec8734b31 Mon Sep 17 00:00:00 2001 From: Martin Donadieu Date: Wed, 22 Jul 2026 12:30:06 +0300 Subject: [PATCH 04/16] fix(db): skip encryption checks when nulling dual-storage columns Legacy unencrypted bundles in encryption-enforced orgs must still be reclaimable by null_migrated_app_version_manifests. Co-authored-by: Cursor --- ...0260722082019_fix_supabase_swap_memory.sql | 18 ++++++ tests/cleanup_swap_memory.test.ts | 61 +++++++++++++++++++ 2 files changed, 79 insertions(+) diff --git a/supabase/migrations/20260722082019_fix_supabase_swap_memory.sql b/supabase/migrations/20260722082019_fix_supabase_swap_memory.sql index 5b9b5b9ca2..36929afc80 100644 --- a/supabase/migrations/20260722082019_fix_supabase_swap_memory.sql +++ b/supabase/migrations/20260722082019_fix_supabase_swap_memory.sql @@ -559,6 +559,24 @@ BEGIN END IF; END IF; + -- Manifest/native_packages nulling must not re-run encryption enforcement. + -- Legacy rows can predate org encryption requirements; reclaim only clears + -- dual-storage columns and must not abort on those orgs. + IF TG_OP = 'UPDATE' + AND NEW.session_key IS NOT DISTINCT FROM OLD.session_key + AND NEW.key_id IS NOT DISTINCT FROM OLD.key_id + AND NEW.name IS NOT DISTINCT FROM OLD.name + AND NEW.app_id IS NOT DISTINCT FROM OLD.app_id + AND NEW.storage_provider IS NOT DISTINCT FROM OLD.storage_provider + AND NEW.r2_path IS NOT DISTINCT FROM OLD.r2_path + AND NEW.external_url IS NOT DISTINCT FROM OLD.external_url + AND NEW.checksum IS NOT DISTINCT FROM OLD.checksum + AND (NEW.manifest IS NULL OR NEW.manifest IS NOT DISTINCT FROM OLD.manifest) + AND (NEW.native_packages IS NULL OR NEW.native_packages IS NOT DISTINCT FROM OLD.native_packages) + THEN + RETURN NEW; + END IF; + -- Derive org_id from NEW.app_id first because -- force_valid_owner_org_app_versions runs after this trigger. SELECT apps.owner_org INTO org_id diff --git a/tests/cleanup_swap_memory.test.ts b/tests/cleanup_swap_memory.test.ts index 80eb2d2424..b020e00de6 100644 --- a/tests/cleanup_swap_memory.test.ts +++ b/tests/cleanup_swap_memory.test.ts @@ -165,4 +165,65 @@ describe('swap memory cleanup functions', () => { await executeSQL(`DELETE FROM public.app_versions WHERE id = $1`, [versionId]) await executeSQL(`DELETE FROM public.apps WHERE app_id = $1`, [appId]) }) + it('null_migrated_app_version_manifests works when org requires encryption', async () => { + const appId = `com.swap.encnull.${randomUUID().slice(0, 8)}` + const orgRows = await executeSQL( + `SELECT id FROM public.orgs ORDER BY created_at LIMIT 1`, + ) + const orgId = orgRows[0]?.id as string + + await executeSQL( + `INSERT INTO public.apps (app_id, name, icon_url, owner_org) + VALUES ($1, 'swap-enc-null', '', $2::uuid)`, + [appId, orgId], + ) + + // Create an unencrypted ready bundle while enforcement is off. + const versionRows = await executeSQL( + `INSERT INTO public.app_versions (app_id, name, owner_org, storage_provider, session_key, manifest, manifest_count) + VALUES ( + $1, + $2, + $3::uuid, + 'r2', + NULL, + ARRAY[ROW('index.html', 'apps/test/index.html', 'abc123')::public.manifest_entry], + 1 + ) + RETURNING id`, + [appId, `1.0.0-${randomUUID().slice(0, 8)}`, orgId], + ) + const versionId = versionRows[0]?.id as number + + await executeSQL( + `INSERT INTO public.manifest (app_version_id, file_name, s3_path, file_hash) + VALUES ($1, 'index.html', 'apps/test/index.html', 'abc123')`, + [versionId], + ) + + try { + await executeSQL( + `UPDATE public.orgs SET enforce_encrypted_bundles = true WHERE id = $1::uuid`, + [orgId], + ) + + await executeSQL(`SELECT public.null_migrated_app_version_manifests()`) + + const after = await executeSQL( + `SELECT manifest IS NULL AS is_null FROM public.app_versions WHERE id = $1`, + [versionId], + ) + expect(after[0]?.is_null).toBe(true) + } + finally { + await executeSQL( + `UPDATE public.orgs SET enforce_encrypted_bundles = false WHERE id = $1::uuid`, + [orgId], + ) + await executeSQL(`DELETE FROM public.manifest WHERE app_version_id = $1`, [versionId]) + await executeSQL(`DELETE FROM public.app_versions WHERE id = $1`, [versionId]) + await executeSQL(`DELETE FROM public.apps WHERE app_id = $1`, [appId]) + } + }) + }) From d731d6df99f6aacd071713f76352bddc4d82bed5 Mon Sep 17 00:00:00 2001 From: Martin Donadieu Date: Wed, 22 Jul 2026 12:37:11 +0300 Subject: [PATCH 05/16] fix(db): harden swap reclaim batching and dual-storage null gate Bound cron cleanup work per invocation, require full manifest migration before nulling arrays, and keep native_packages locked on ready bundles. Co-authored-by: Cursor --- scripts/ops/reclaim_supabase_swap.sql | 109 +++--------------- scripts/ops/verify_supabase_swap.sql | 63 ++++++++-- ...0260722082019_fix_supabase_swap_memory.sql | 67 ++++++++--- tests/cleanup_swap_memory.test.ts | 57 ++++++++- 4 files changed, 178 insertions(+), 118 deletions(-) diff --git a/scripts/ops/reclaim_supabase_swap.sql b/scripts/ops/reclaim_supabase_swap.sql index 51463c7fb1..592507b487 100644 --- a/scripts/ops/reclaim_supabase_swap.sql +++ b/scripts/ops/reclaim_supabase_swap.sql @@ -1,10 +1,9 @@ -- Capgo-EU Phase A reclaim (run manually in a maintenance window). +-- Prefer psql (VACUUM cannot run inside a transaction / SQL-editor DO block). +-- Example: +-- PGPASSWORD=... psql "postgresql://..." -v ON_ERROR_STOP=1 -f scripts/ops/reclaim_supabase_swap.sql -- Safe order: truncate empty bloat -> batched archive deletes -> null dual manifests -> trim audit. --- Do NOT wrap the whole file in one transaction. VACUUM cannot run inside a transaction block. --- Prefer Supabase SQL editor / psql as postgres. Re-run sections until counts hit zero. - -\timing on -\set ON_ERROR_STOP on +-- Each batch commits (separate statements). Re-run until notices show 0 deleted/updated. -- --------------------------------------------------------------------------- -- 0) Baseline sizes @@ -29,118 +28,40 @@ WHERE (schemaname, relname) IN ( ORDER BY pg_total_relation_size(format('%I.%I', schemaname, relname)::regclass) DESC; -- --------------------------------------------------------------------------- --- 1) Truncate pg_net response bloat (~5GB empty table in prod) +-- 1) Truncate pg_net response bloat -- --------------------------------------------------------------------------- TRUNCATE TABLE net._http_response; -- --------------------------------------------------------------------------- --- 2) Purge pgmq archives older than 2 days (batched). Repeat until deleted=0. +-- 2) Purge pgmq archives older than 2 days (one committed batch per statement). +-- Re-run this section until deleted totals are 0. -- --------------------------------------------------------------------------- -DO $$ -DECLARE - queue_name text; - cutoff timestamptz := now() - interval '2 days'; - batch_size integer := 10000; - deleted_batch integer; - deleted_total bigint; -BEGIN - FOREACH queue_name IN ARRAY ARRAY[ - 'on_version_update', - 'on_manifest_create', - 'webhook_dispatcher', - 'on_channel_update' - ] - LOOP - deleted_total := 0; - LOOP - EXECUTE format( - 'WITH doomed AS ( - SELECT ctid FROM pgmq.a_%I WHERE archived_at < $1 LIMIT $2 - ) - DELETE FROM pgmq.a_%I AS archive - USING doomed - WHERE archive.ctid = doomed.ctid', - queue_name, queue_name - ) USING cutoff, batch_size; - GET DIAGNOSTICS deleted_batch = ROW_COUNT; - deleted_total := deleted_total + deleted_batch; - EXIT WHEN deleted_batch = 0; - END LOOP; - RAISE NOTICE 'purged % rows from pgmq.a_%', deleted_total, queue_name; - END LOOP; -END $$; +SELECT public.cleanup_queue_messages(); VACUUM (VERBOSE) pgmq.a_on_version_update; VACUUM (VERBOSE) pgmq.a_on_manifest_create; VACUUM (VERBOSE) pgmq.a_webhook_dispatcher; VACUUM (VERBOSE) pgmq.a_on_channel_update; --- Optional hard reclaim if VACUUM leaves a lot of empty pages (takes stronger locks): +-- Optional hard reclaim if VACUUM leaves empty pages (stronger locks): -- VACUUM (FULL, VERBOSE) pgmq.a_on_version_update; -- VACUUM (FULL, VERBOSE) pgmq.a_on_manifest_create; -- VACUUM (FULL, VERBOSE) pgmq.a_webhook_dispatcher; -- VACUUM (FULL, VERBOSE) pgmq.a_on_channel_update; -- --------------------------------------------------------------------------- --- 3) Null migrated app_versions.manifest arrays (dual storage leftovers) +-- 3) Null fully migrated app_versions.manifest arrays +-- Requires every expected legacy entry to exist in public.manifest. +-- Re-run until notice shows 0. -- --------------------------------------------------------------------------- -DO $$ -DECLARE - batch_size integer := 200; - updated_batch integer; - updated_total bigint := 0; -BEGIN - LOOP - WITH doomed AS ( - SELECT av.id - FROM public.app_versions AS av - WHERE av.manifest IS NOT NULL - AND cardinality(av.manifest) > 0 - AND EXISTS ( - SELECT 1 FROM public.manifest AS m WHERE m.app_version_id = av.id - ) - ORDER BY av.id - LIMIT batch_size - ) - UPDATE public.app_versions AS av - SET manifest = NULL - FROM doomed - WHERE av.id = doomed.id; - GET DIAGNOSTICS updated_batch = ROW_COUNT; - updated_total := updated_total + updated_batch; - EXIT WHEN updated_batch = 0; - END LOOP; - RAISE NOTICE 'nulled manifest arrays on % versions', updated_total; -END $$; +SELECT public.null_migrated_app_version_manifests(); VACUUM (VERBOSE) public.app_versions; -- --------------------------------------------------------------------------- --- 4) Trim audit_logs older than 30 days (batched) +-- 4) Trim audit_logs older than 30 days (bounded batches). Re-run until 0. -- --------------------------------------------------------------------------- -DO $$ -DECLARE - cutoff timestamptz := now() - interval '30 days'; - batch_size integer := 5000; - deleted_batch integer; - deleted_total bigint := 0; -BEGIN - LOOP - WITH doomed AS ( - SELECT ctid - FROM public.audit_logs - WHERE created_at < cutoff - LIMIT batch_size - ) - DELETE FROM public.audit_logs AS audit_logs - USING doomed - WHERE audit_logs.ctid = doomed.ctid; - GET DIAGNOSTICS deleted_batch = ROW_COUNT; - deleted_total := deleted_total + deleted_batch; - EXIT WHEN deleted_batch = 0; - END LOOP; - RAISE NOTICE 'deleted % audit_logs rows older than 30 days', deleted_total; -END $$; +SELECT public.cleanup_old_audit_logs(); VACUUM (VERBOSE) public.audit_logs; diff --git a/scripts/ops/verify_supabase_swap.sql b/scripts/ops/verify_supabase_swap.sql index 6b976c92f1..49d46e0bed 100644 --- a/scripts/ops/verify_supabase_swap.sql +++ b/scripts/ops/verify_supabase_swap.sql @@ -1,10 +1,9 @@ -- Post-deploy / post-reclaim verification for Capgo-EU swap pressure. +-- Prefer psql. Avoids unbounded whole-table counts where possible. SELECT pg_size_pretty(pg_database_size(current_database())::bigint) AS db_size; -SELECT - name, - setting +SELECT name, setting FROM pg_settings WHERE name IN ('shared_buffers', 'work_mem', 'max_connections'); @@ -25,9 +24,21 @@ WHERE (schemaname, relname) IN ( ) ORDER BY pg_total_relation_size(format('%I.%I', schemaname, relname)::regclass) DESC; -SELECT count(*) AS versions_with_array_manifest -FROM public.app_versions -WHERE manifest IS NOT NULL AND cardinality(manifest) > 0; +-- Sample of dual-storage leftovers that are eligible for nulling (bounded). +SELECT count(*) AS eligible_dual_storage_sample +FROM ( + SELECT av.id + FROM public.app_versions AS av + WHERE av.manifest IS NOT NULL + AND cardinality(av.manifest) > 0 + AND ( + SELECT count(*)::integer + FROM public.manifest AS m + WHERE m.app_version_id = av.id + ) >= GREATEST(COALESCE(av.manifest_count, 0), cardinality(av.manifest)) + ORDER BY av.id + LIMIT 1000 +) AS sample; SELECT name, enabled, hour_interval, run_at_hour, run_at_minute, target FROM public.cron_tasks @@ -47,10 +58,42 @@ GROUP BY status, return_message ORDER BY n DESC LIMIT 20; -SELECT - count(*) FILTER (WHERE archived_at < now() - interval '2 days') AS archives_older_than_2d, - count(*) AS archives_total -FROM pgmq.a_on_manifest_create; +-- Bounded existence checks per archive queue (no full-table aggregates). +SELECT queue_name, has_rows_older_than_2d +FROM ( + SELECT 'a_on_manifest_create' AS queue_name, + EXISTS ( + SELECT 1 + FROM pgmq.a_on_manifest_create + WHERE archived_at < now() - interval '2 days' + LIMIT 1 + ) AS has_rows_older_than_2d + UNION ALL + SELECT 'a_on_version_update', + EXISTS ( + SELECT 1 + FROM pgmq.a_on_version_update + WHERE archived_at < now() - interval '2 days' + LIMIT 1 + ) + UNION ALL + SELECT 'a_webhook_dispatcher', + EXISTS ( + SELECT 1 + FROM pgmq.a_webhook_dispatcher + WHERE archived_at < now() - interval '2 days' + LIMIT 1 + ) + UNION ALL + SELECT 'a_on_channel_update', + EXISTS ( + SELECT 1 + FROM pgmq.a_on_channel_update + WHERE archived_at < now() - interval '2 days' + LIMIT 1 + ) +) AS archives +ORDER BY queue_name; SELECT 'index hit rate' AS name, diff --git a/supabase/migrations/20260722082019_fix_supabase_swap_memory.sql b/supabase/migrations/20260722082019_fix_supabase_swap_memory.sql index 36929afc80..dec644e755 100644 --- a/supabase/migrations/20260722082019_fix_supabase_swap_memory.sql +++ b/supabase/migrations/20260722082019_fix_supabase_swap_memory.sql @@ -16,6 +16,8 @@ DECLARE queue_name text; cutoff timestamptz := pg_catalog.now() - INTERVAL '2 days'; batch_size integer := 10000; + max_batches integer := 20; + batch_no integer; deleted_batch integer; deleted_total bigint; BEGIN @@ -23,9 +25,13 @@ BEGIN SELECT q.queue_name FROM pgmq.list_queues() q ) LOOP deleted_total := 0; + batch_no := 0; LOOP - EXECUTE format( + batch_no := batch_no + 1; + EXIT WHEN batch_no > max_batches; + + EXECUTE pg_catalog.format( 'DELETE FROM pgmq.a_%I WHERE ctid IN ( SELECT ctid @@ -48,8 +54,12 @@ BEGIN END IF; deleted_total := 0; + batch_no := 0; LOOP - EXECUTE format( + batch_no := batch_no + 1; + EXIT WHEN batch_no > max_batches; + + EXECUTE pg_catalog.format( 'DELETE FROM pgmq.q_%I WHERE ctid IN ( SELECT ctid @@ -107,10 +117,15 @@ CREATE OR REPLACE FUNCTION "public"."cleanup_old_audit_logs"() RETURNS "void" DECLARE cutoff timestamptz := pg_catalog.now() - INTERVAL '30 days'; batch_size integer := 5000; + max_batches integer := 40; + batch_no integer := 0; deleted_batch integer; deleted_total bigint := 0; BEGIN LOOP + batch_no := batch_no + 1; + EXIT WHEN batch_no > max_batches; + DELETE FROM public.audit_logs WHERE ctid IN ( SELECT ctid @@ -143,19 +158,27 @@ CREATE OR REPLACE FUNCTION "public"."null_migrated_app_version_manifests"() RETU AS $$ DECLARE batch_size integer := 200; + max_batches integer := 50; + batch_no integer := 0; updated_batch integer; updated_total bigint := 0; BEGIN LOOP + batch_no := batch_no + 1; + EXIT WHEN batch_no > max_batches; + WITH doomed AS ( SELECT av.id FROM public.app_versions AS av WHERE av.manifest IS NOT NULL AND pg_catalog.cardinality(av.manifest) > 0 - AND EXISTS ( - SELECT 1 + AND ( + SELECT count(*)::integer FROM public.manifest AS m WHERE m.app_version_id = av.id + ) >= pg_catalog.GREATEST( + COALESCE(av.manifest_count, 0), + pg_catalog.cardinality(av.manifest) ) ORDER BY av.id LIMIT batch_size @@ -205,6 +228,7 @@ DECLARE v_stats_refresh_fields constant text[] := ARRAY['stats_refresh_requested_at', 'stats_updated_at', 'updated_at']; v_background_counter_fields constant text[] := ARRAY['channel_device_count', 'manifest_bundle_count', 'updated_at']; v_fat_app_version_fields constant text[] := ARRAY['manifest', 'native_packages']; + v_noise_app_version_fields constant text[] := ARRAY['manifest', 'native_packages', 'updated_at']; BEGIN SELECT auth.uid() INTO v_actor_user_id; @@ -299,9 +323,13 @@ BEGIN WHERE changed_field.field_name <> ALL(v_fat_app_version_fields); END IF; - -- Skip updates that only touched stripped fat columns (e.g. manifest nulling). + -- Skip updates that only touched stripped fat columns / auto timestamps. IF TG_OP = 'UPDATE' - AND (v_changed_fields IS NULL OR pg_catalog.cardinality(v_changed_fields) = 0) THEN + AND NOT EXISTS ( + SELECT 1 + FROM pg_catalog.unnest(COALESCE(v_changed_fields, ARRAY[]::text[])) AS changed_field(field_name) + WHERE changed_field.field_name <> ALL(v_noise_app_version_fields) + ) THEN RETURN NEW; END IF; END IF; @@ -380,12 +408,12 @@ DECLARE BEGIN function_type := CASE WHEN NULLIF(TG_ARGV[1], '') IS NULL THEN 'cloudflare' - WHEN lower(TG_ARGV[1]) = 'supabase' THEN 'cloudflare' + WHEN pg_catalog.lower(TG_ARGV[1]) = 'supabase' THEN 'cloudflare' ELSE TG_ARGV[1] END; - record_payload := to_jsonb(NEW); - old_record_payload := to_jsonb(OLD); + record_payload := pg_catalog.to_jsonb(NEW); + old_record_payload := pg_catalog.to_jsonb(OLD); -- app_versions fat columns can be multi-MB. Never enqueue them; handlers reload when needed. IF TG_TABLE_NAME = 'app_versions' THEN @@ -397,10 +425,10 @@ BEGIN END IF; END IF; - payload := jsonb_build_object( + payload := pg_catalog.jsonb_build_object( 'function_name', TG_ARGV[0], 'function_type', function_type, - 'payload', jsonb_build_object( + 'payload', pg_catalog.jsonb_build_object( 'old_record', old_record_payload, 'record', record_payload, 'type', TG_OP, @@ -540,7 +568,20 @@ BEGIN OR NEW.external_url IS DISTINCT FROM OLD.external_url OR NEW.checksum IS DISTINCT FROM OLD.checksum OR (NEW.manifest IS DISTINCT FROM OLD.manifest AND NEW.manifest IS NOT NULL) - OR (NEW.native_packages IS DISTINCT FROM OLD.native_packages AND NEW.native_packages IS NOT NULL) + -- Nulling is allowed only when public.manifest has every expected entry. + OR ( + NEW.manifest IS NULL + AND OLD.manifest IS NOT NULL + AND ( + SELECT count(*)::integer + FROM public.manifest AS m + WHERE m.app_version_id = OLD.id + ) < pg_catalog.GREATEST( + COALESCE(OLD.manifest_count, 0), + COALESCE(pg_catalog.cardinality(OLD.manifest), 0) + ) + ) + OR NEW.native_packages IS DISTINCT FROM OLD.native_packages ) THEN PERFORM public.pg_log('deny: BUNDLE_CONTENT_LOCKED_TRIGGER', @@ -571,8 +612,8 @@ BEGIN AND NEW.r2_path IS NOT DISTINCT FROM OLD.r2_path AND NEW.external_url IS NOT DISTINCT FROM OLD.external_url AND NEW.checksum IS NOT DISTINCT FROM OLD.checksum + AND NEW.native_packages IS NOT DISTINCT FROM OLD.native_packages AND (NEW.manifest IS NULL OR NEW.manifest IS NOT DISTINCT FROM OLD.manifest) - AND (NEW.native_packages IS NULL OR NEW.native_packages IS NOT DISTINCT FROM OLD.native_packages) THEN RETURN NEW; END IF; diff --git a/tests/cleanup_swap_memory.test.ts b/tests/cleanup_swap_memory.test.ts index b020e00de6..8e7043e17b 100644 --- a/tests/cleanup_swap_memory.test.ts +++ b/tests/cleanup_swap_memory.test.ts @@ -59,7 +59,7 @@ describe('swap memory cleanup functions', () => { expect(rows[0]?.n).toBe(0) }) - it('null_migrated_app_version_manifests clears dual-storage arrays', async () => { + it('null_migrated_app_version_manifests clears fully migrated dual-storage arrays', async () => { const appId = `com.swap.nullmanifest.${randomUUID().slice(0, 8)}` const orgRows = await executeSQL( `SELECT id FROM public.orgs ORDER BY created_at LIMIT 1`, @@ -165,6 +165,61 @@ describe('swap memory cleanup functions', () => { await executeSQL(`DELETE FROM public.app_versions WHERE id = $1`, [versionId]) await executeSQL(`DELETE FROM public.apps WHERE app_id = $1`, [appId]) }) + + it('null_migrated_app_version_manifests skips partially migrated arrays', async () => { + const appId = `com.swap.partialmanifest.${randomUUID().slice(0, 8)}` + const orgRows = await executeSQL( + `SELECT id FROM public.orgs ORDER BY created_at LIMIT 1`, + ) + const orgId = orgRows[0]?.id as string + + await executeSQL( + `INSERT INTO public.apps (app_id, name, icon_url, owner_org) + VALUES ($1, 'swap-partial-manifest', '', $2::uuid)`, + [appId, orgId], + ) + + const versionRows = await executeSQL( + `INSERT INTO public.app_versions (app_id, name, owner_org, storage_provider, manifest, manifest_count) + VALUES ( + $1, + $2, + $3::uuid, + 'r2', + ARRAY[ + ROW('index.html', 'apps/test/index.html', 'abc123')::public.manifest_entry, + ROW('main.js', 'apps/test/main.js', 'def456')::public.manifest_entry + ], + 2 + ) + RETURNING id`, + [appId, `1.0.0-${randomUUID().slice(0, 8)}`, orgId], + ) + const versionId = versionRows[0]?.id as number + + // Only one of two files migrated. + await executeSQL( + `INSERT INTO public.manifest (app_version_id, file_name, s3_path, file_hash) + VALUES ($1, 'index.html', 'apps/test/index.html', 'abc123')`, + [versionId], + ) + + await executeSQL(`SELECT public.null_migrated_app_version_manifests()`) + + const after = await executeSQL( + `SELECT manifest IS NULL AS is_null, cardinality(manifest) AS n + FROM public.app_versions + WHERE id = $1`, + [versionId], + ) + expect(after[0]?.is_null).toBe(false) + expect(after[0]?.n).toBe(2) + + await executeSQL(`DELETE FROM public.manifest WHERE app_version_id = $1`, [versionId]) + await executeSQL(`DELETE FROM public.app_versions WHERE id = $1`, [versionId]) + await executeSQL(`DELETE FROM public.apps WHERE app_id = $1`, [appId]) + }) + it('null_migrated_app_version_manifests works when org requires encryption', async () => { const appId = `com.swap.encnull.${randomUUID().slice(0, 8)}` const orgRows = await executeSQL( From 169b939f09bed7db49ee2f10677f97d33853836e Mon Sep 17 00:00:00 2001 From: Martin Donadieu Date: Wed, 22 Jul 2026 12:51:10 +0300 Subject: [PATCH 06/16] fix(db): avoid pg_catalog.GREATEST in dual-storage null gate GREATEST is SQL syntax, not a pg_catalog function under empty search_path, so the ready-bundle trigger was failing with 42883 instead of locking. Co-authored-by: Cursor --- scripts/ops/verify_supabase_swap.sql | 8 +++++++- ...20260722082019_fix_supabase_swap_memory.sql | 18 ++++++++++++------ 2 files changed, 19 insertions(+), 7 deletions(-) diff --git a/scripts/ops/verify_supabase_swap.sql b/scripts/ops/verify_supabase_swap.sql index 49d46e0bed..1eb91759c5 100644 --- a/scripts/ops/verify_supabase_swap.sql +++ b/scripts/ops/verify_supabase_swap.sql @@ -35,7 +35,13 @@ FROM ( SELECT count(*)::integer FROM public.manifest AS m WHERE m.app_version_id = av.id - ) >= GREATEST(COALESCE(av.manifest_count, 0), cardinality(av.manifest)) + ) >= ( + CASE + WHEN COALESCE(av.manifest_count, 0) >= cardinality(av.manifest) + THEN COALESCE(av.manifest_count, 0) + ELSE cardinality(av.manifest) + END + ) ORDER BY av.id LIMIT 1000 ) AS sample; diff --git a/supabase/migrations/20260722082019_fix_supabase_swap_memory.sql b/supabase/migrations/20260722082019_fix_supabase_swap_memory.sql index dec644e755..df9f184184 100644 --- a/supabase/migrations/20260722082019_fix_supabase_swap_memory.sql +++ b/supabase/migrations/20260722082019_fix_supabase_swap_memory.sql @@ -176,9 +176,12 @@ BEGIN SELECT count(*)::integer FROM public.manifest AS m WHERE m.app_version_id = av.id - ) >= pg_catalog.GREATEST( - COALESCE(av.manifest_count, 0), - pg_catalog.cardinality(av.manifest) + ) >= ( + CASE + WHEN COALESCE(av.manifest_count, 0) >= pg_catalog.cardinality(av.manifest) + THEN COALESCE(av.manifest_count, 0) + ELSE pg_catalog.cardinality(av.manifest) + END ) ORDER BY av.id LIMIT batch_size @@ -576,9 +579,12 @@ BEGIN SELECT count(*)::integer FROM public.manifest AS m WHERE m.app_version_id = OLD.id - ) < pg_catalog.GREATEST( - COALESCE(OLD.manifest_count, 0), - COALESCE(pg_catalog.cardinality(OLD.manifest), 0) + ) < ( + CASE + WHEN COALESCE(OLD.manifest_count, 0) >= COALESCE(pg_catalog.cardinality(OLD.manifest), 0) + THEN COALESCE(OLD.manifest_count, 0) + ELSE COALESCE(pg_catalog.cardinality(OLD.manifest), 0) + END ) ) OR NEW.native_packages IS DISTINCT FROM OLD.native_packages From 8b41f3ff4618a8e5b5e7f7e471655c98674ca1f7 Mon Sep 17 00:00:00 2001 From: Martin Donadieu Date: Wed, 22 Jul 2026 13:23:52 +0300 Subject: [PATCH 07/16] fix(db): bound queue cleanup globally and skip reclaim queue fan-out Cap cleanup_queue_messages with an invocation-wide batch budget, skip on_version_update enqueue for fat-only nulling, and add a partial index for leftover manifest arrays. Co-authored-by: Cursor --- scripts/ops/reclaim_supabase_swap.sql | 30 ++++----- scripts/ops/verify_supabase_swap.sql | 42 +++++++------ ...0260722082019_fix_supabase_swap_memory.sql | 61 ++++++++++++------- 3 files changed, 76 insertions(+), 57 deletions(-) diff --git a/scripts/ops/reclaim_supabase_swap.sql b/scripts/ops/reclaim_supabase_swap.sql index 592507b487..3fab978bf2 100644 --- a/scripts/ops/reclaim_supabase_swap.sql +++ b/scripts/ops/reclaim_supabase_swap.sql @@ -3,7 +3,8 @@ -- Example: -- PGPASSWORD=... psql "postgresql://..." -v ON_ERROR_STOP=1 -f scripts/ops/reclaim_supabase_swap.sql -- Safe order: truncate empty bloat -> batched archive deletes -> null dual manifests -> trim audit. --- Each batch commits (separate statements). Re-run until notices show 0 deleted/updated. +-- Each statement commits separately. Re-run until cleanup notices report deleted/updated = 0 +-- (functions always emit a notice, including zero totals). -- --------------------------------------------------------------------------- -- 0) Baseline sizes @@ -33,33 +34,34 @@ ORDER BY pg_total_relation_size(format('%I.%I', schemaname, relname)::regclass) TRUNCATE TABLE net._http_response; -- --------------------------------------------------------------------------- --- 2) Purge pgmq archives older than 2 days (one committed batch per statement). --- Re-run this section until deleted totals are 0. +-- 2) Purge pgmq archives/stuck messages (global batch budget per call). +-- Re-run this SELECT until the notice shows archived_deleted=0 and stuck_deleted=0. -- --------------------------------------------------------------------------- SELECT public.cleanup_queue_messages(); -VACUUM (VERBOSE) pgmq.a_on_version_update; -VACUUM (VERBOSE) pgmq.a_on_manifest_create; -VACUUM (VERBOSE) pgmq.a_webhook_dispatcher; -VACUUM (VERBOSE) pgmq.a_on_channel_update; +-- Vacuum every pgmq archive + queue table (psql \gexec; VACUUM cannot run in DO/tx). +SELECT format('VACUUM (VERBOSE) pgmq.a_%I;', queue_name) +FROM pgmq.list_queues() +\gexec +SELECT format('VACUUM (VERBOSE) pgmq.q_%I;', queue_name) +FROM pgmq.list_queues() +\gexec --- Optional hard reclaim if VACUUM leaves empty pages (stronger locks): --- VACUUM (FULL, VERBOSE) pgmq.a_on_version_update; --- VACUUM (FULL, VERBOSE) pgmq.a_on_manifest_create; --- VACUUM (FULL, VERBOSE) pgmq.a_webhook_dispatcher; --- VACUUM (FULL, VERBOSE) pgmq.a_on_channel_update; +-- Optional hard reclaim (stronger locks): +-- SELECT format('VACUUM (FULL, VERBOSE) pgmq.a_%I;', queue_name) FROM pgmq.list_queues() \gexec +-- SELECT format('VACUUM (FULL, VERBOSE) pgmq.q_%I;', queue_name) FROM pgmq.list_queues() \gexec -- --------------------------------------------------------------------------- -- 3) Null fully migrated app_versions.manifest arrays -- Requires every expected legacy entry to exist in public.manifest. --- Re-run until notice shows 0. +-- Re-run until notice shows updated=0. -- --------------------------------------------------------------------------- SELECT public.null_migrated_app_version_manifests(); VACUUM (VERBOSE) public.app_versions; -- --------------------------------------------------------------------------- --- 4) Trim audit_logs older than 30 days (bounded batches). Re-run until 0. +-- 4) Trim audit_logs older than 30 days (bounded batches). Re-run until deleted=0. -- --------------------------------------------------------------------------- SELECT public.cleanup_old_audit_logs(); diff --git a/scripts/ops/verify_supabase_swap.sql b/scripts/ops/verify_supabase_swap.sql index 1eb91759c5..e4bafd6d42 100644 --- a/scripts/ops/verify_supabase_swap.sql +++ b/scripts/ops/verify_supabase_swap.sql @@ -24,29 +24,32 @@ WHERE (schemaname, relname) IN ( ) ORDER BY pg_total_relation_size(format('%I.%I', schemaname, relname)::regclass) DESC; --- Sample of dual-storage leftovers that are eligible for nulling (bounded). +-- Bound candidate discovery first, then evaluate eligibility inside the sample. SELECT count(*) AS eligible_dual_storage_sample FROM ( - SELECT av.id - FROM public.app_versions AS av - WHERE av.manifest IS NOT NULL - AND cardinality(av.manifest) > 0 + SELECT sample.id + FROM ( + SELECT av.id, av.manifest, av.manifest_count + FROM public.app_versions AS av + WHERE av.manifest IS NOT NULL + ORDER BY av.id + LIMIT 1000 + ) AS sample + WHERE cardinality(sample.manifest) > 0 AND ( SELECT count(*)::integer FROM public.manifest AS m - WHERE m.app_version_id = av.id + WHERE m.app_version_id = sample.id ) >= ( CASE - WHEN COALESCE(av.manifest_count, 0) >= cardinality(av.manifest) - THEN COALESCE(av.manifest_count, 0) - ELSE cardinality(av.manifest) + WHEN COALESCE(sample.manifest_count, 0) >= cardinality(sample.manifest) + THEN COALESCE(sample.manifest_count, 0) + ELSE cardinality(sample.manifest) END ) - ORDER BY av.id - LIMIT 1000 -) AS sample; +) AS eligible; -SELECT name, enabled, hour_interval, run_at_hour, run_at_minute, target +SELECT name, enabled, hour_interval, run_at_hour, run_at_minute, target, updated_at FROM public.cron_tasks WHERE name IN ( 'cleanup_queue_messages', @@ -56,13 +59,12 @@ WHERE name IN ( ) ORDER BY name; -SELECT status, return_message, count(*) AS n -FROM cron.job_run_details -WHERE start_time > now() - interval '24 hours' - OR (start_time IS NULL AND status IN ('failed', 'connecting')) -GROUP BY status, return_message -ORDER BY n DESC -LIMIT 20; +-- process_all_cron_tasks() swallows per-task errors; cron.job_run_details only +-- reflects the outer job. Prefer Postgres logs / healthchecks for task failures. +SELECT indexname +FROM pg_indexes +WHERE schemaname = 'public' + AND indexname = 'app_versions_manifest_present_idx'; -- Bounded existence checks per archive queue (no full-table aggregates). SELECT queue_name, has_rows_older_than_2d diff --git a/supabase/migrations/20260722082019_fix_supabase_swap_memory.sql b/supabase/migrations/20260722082019_fix_supabase_swap_memory.sql index df9f184184..a91afb826b 100644 --- a/supabase/migrations/20260722082019_fix_supabase_swap_memory.sql +++ b/supabase/migrations/20260722082019_fix_supabase_swap_memory.sql @@ -16,20 +16,23 @@ DECLARE queue_name text; cutoff timestamptz := pg_catalog.now() - INTERVAL '2 days'; batch_size integer := 10000; - max_batches integer := 20; - batch_no integer; + -- Hard cap across ALL queues/phases for one cron/reclaim invocation. + max_batches_total integer := 40; + batches_used integer := 0; deleted_batch integer; deleted_total bigint; + deleted_archived_total bigint := 0; + deleted_stuck_total bigint := 0; BEGIN FOR queue_name IN ( SELECT q.queue_name FROM pgmq.list_queues() q ) LOOP - deleted_total := 0; - batch_no := 0; + EXIT WHEN batches_used >= max_batches_total; + deleted_total := 0; LOOP - batch_no := batch_no + 1; - EXIT WHEN batch_no > max_batches; + EXIT WHEN batches_used >= max_batches_total; + batches_used := batches_used + 1; EXECUTE pg_catalog.format( 'DELETE FROM pgmq.a_%I @@ -46,18 +49,14 @@ BEGIN GET DIAGNOSTICS deleted_batch = ROW_COUNT; deleted_total := deleted_total + deleted_batch; + deleted_archived_total := deleted_archived_total + deleted_batch; EXIT WHEN deleted_batch = 0; END LOOP; - IF deleted_total > 0 THEN - RAISE NOTICE 'cleanup_queue_messages: deleted % archived rows from a_%', deleted_total, queue_name; - END IF; - deleted_total := 0; - batch_no := 0; LOOP - batch_no := batch_no + 1; - EXIT WHEN batch_no > max_batches; + EXIT WHEN batches_used >= max_batches_total; + batches_used := batches_used + 1; EXECUTE pg_catalog.format( 'DELETE FROM pgmq.q_%I @@ -74,13 +73,17 @@ BEGIN GET DIAGNOSTICS deleted_batch = ROW_COUNT; deleted_total := deleted_total + deleted_batch; + deleted_stuck_total := deleted_stuck_total + deleted_batch; EXIT WHEN deleted_batch = 0; END LOOP; - - IF deleted_total > 0 THEN - RAISE NOTICE 'cleanup_queue_messages: deleted % stuck rows from q_%', deleted_total, queue_name; - END IF; END LOOP; + + RAISE NOTICE + 'cleanup_queue_messages: archived_deleted=% stuck_deleted=% batches_used=%/%', + deleted_archived_total, + deleted_stuck_total, + batches_used, + max_batches_total; END; $$; @@ -139,9 +142,7 @@ BEGIN EXIT WHEN deleted_batch = 0; END LOOP; - IF deleted_total > 0 THEN - RAISE NOTICE 'cleanup_old_audit_logs: deleted % rows older than 30 days', deleted_total; - END IF; + RAISE NOTICE 'cleanup_old_audit_logs: deleted=% max_batches=%', deleted_total, max_batches; END; $$; @@ -196,9 +197,10 @@ BEGIN EXIT WHEN updated_batch = 0; END LOOP; - IF updated_total > 0 THEN - RAISE NOTICE 'null_migrated_app_version_manifests: nulled manifest arrays on % versions', updated_total; - END IF; + RAISE NOTICE + 'null_migrated_app_version_manifests: updated=% max_batches=%', + updated_total, + max_batches; END; $$; @@ -426,6 +428,13 @@ BEGIN IF old_record_payload IS NOT NULL THEN old_record_payload := old_record_payload - 'manifest' - 'native_packages'; END IF; + + -- Dual-storage reclaim only nulls fat columns (+ auto updated_at). Skip queue fan-out. + IF TG_OP = 'UPDATE' + AND (record_payload - 'updated_at') IS NOT DISTINCT FROM (old_record_payload - 'updated_at') + THEN + RETURN NEW; + END IF; END IF; payload := pg_catalog.jsonb_build_object( @@ -539,6 +548,12 @@ SET enabled = true, updated_at = pg_catalog.now(); + +-- Bound dual-storage candidate discovery once most arrays are nulled. +CREATE INDEX IF NOT EXISTS app_versions_manifest_present_idx + ON public.app_versions USING btree (id) + WHERE manifest IS NOT NULL; + -- --------------------------------------------------------------------------- -- Allow clearing dual-storage fat columns -- after upload (null only). Unblocks on_version_update + reclaim jobs. From 0e2bbe1f4f4f126c8a10ef42878f5751a846ef68 Mon Sep 17 00:00:00 2001 From: Martin Donadieu Date: Wed, 22 Jul 2026 13:29:52 +0300 Subject: [PATCH 08/16] fix(db): count only productive batches in queue cleanup budget Empty queue probes no longer consume the global max_batches_total, so cleanup still reaches archives with old rows in one invocation. Co-authored-by: Cursor --- .../20260722082019_fix_supabase_swap_memory.sql | 15 +++++---------- 1 file changed, 5 insertions(+), 10 deletions(-) diff --git a/supabase/migrations/20260722082019_fix_supabase_swap_memory.sql b/supabase/migrations/20260722082019_fix_supabase_swap_memory.sql index a91afb826b..72345f7dbe 100644 --- a/supabase/migrations/20260722082019_fix_supabase_swap_memory.sql +++ b/supabase/migrations/20260722082019_fix_supabase_swap_memory.sql @@ -16,11 +16,10 @@ DECLARE queue_name text; cutoff timestamptz := pg_catalog.now() - INTERVAL '2 days'; batch_size integer := 10000; - -- Hard cap across ALL queues/phases for one cron/reclaim invocation. + -- Hard cap of productive batches across ALL queues/phases for one invocation. max_batches_total integer := 40; batches_used integer := 0; deleted_batch integer; - deleted_total bigint; deleted_archived_total bigint := 0; deleted_stuck_total bigint := 0; BEGIN @@ -29,10 +28,8 @@ BEGIN ) LOOP EXIT WHEN batches_used >= max_batches_total; - deleted_total := 0; LOOP EXIT WHEN batches_used >= max_batches_total; - batches_used := batches_used + 1; EXECUTE pg_catalog.format( 'DELETE FROM pgmq.a_%I @@ -48,15 +45,13 @@ BEGIN USING cutoff, batch_size; GET DIAGNOSTICS deleted_batch = ROW_COUNT; - deleted_total := deleted_total + deleted_batch; - deleted_archived_total := deleted_archived_total + deleted_batch; EXIT WHEN deleted_batch = 0; + batches_used := batches_used + 1; + deleted_archived_total := deleted_archived_total + deleted_batch; END LOOP; - deleted_total := 0; LOOP EXIT WHEN batches_used >= max_batches_total; - batches_used := batches_used + 1; EXECUTE pg_catalog.format( 'DELETE FROM pgmq.q_%I @@ -72,9 +67,9 @@ BEGIN USING batch_size; GET DIAGNOSTICS deleted_batch = ROW_COUNT; - deleted_total := deleted_total + deleted_batch; - deleted_stuck_total := deleted_stuck_total + deleted_batch; EXIT WHEN deleted_batch = 0; + batches_used := batches_used + 1; + deleted_stuck_total := deleted_stuck_total + deleted_batch; END LOOP; END LOOP; From 03ed950852c4396c95ce771278b18079a8b00998 Mon Sep 17 00:00:00 2001 From: Martin Donadieu Date: Wed, 22 Jul 2026 13:39:39 +0300 Subject: [PATCH 09/16] chore(db): sync read-replica schema for manifest present index Co-authored-by: Cursor --- read_replicate/schema_replicate.catalog.json | 7 +++++++ read_replicate/schema_replicate.sql | 7 +++++++ 2 files changed, 14 insertions(+) diff --git a/read_replicate/schema_replicate.catalog.json b/read_replicate/schema_replicate.catalog.json index 4857f5c5bb..05fb5a9845 100644 --- a/read_replicate/schema_replicate.catalog.json +++ b/read_replicate/schema_replicate.catalog.json @@ -1982,6 +1982,13 @@ "table": "app_versions", "valid": true }, + { + "constraintOwned": false, + "definition": "CREATE INDEX app_versions_manifest_present_idx ON public.app_versions USING btree (id) WHERE (manifest IS NOT NULL)", + "name": "app_versions_manifest_present_idx", + "table": "app_versions", + "valid": true + }, { "constraintOwned": true, "definition": "CREATE UNIQUE INDEX app_versions_name_app_id_key ON public.app_versions USING btree (name, app_id)", diff --git a/read_replicate/schema_replicate.sql b/read_replicate/schema_replicate.sql index 94e69f1dc1..0fc007e74b 100644 --- a/read_replicate/schema_replicate.sql +++ b/read_replicate/schema_replicate.sql @@ -607,6 +607,13 @@ ALTER TABLE ONLY public.orgs CREATE INDEX app_versions_cli_version_idx ON public.app_versions USING btree (cli_version); +-- +-- Name: app_versions_manifest_present_idx; Type: INDEX; Schema: public; Owner: - +-- + +CREATE INDEX app_versions_manifest_present_idx ON public.app_versions USING btree (id) WHERE (manifest IS NOT NULL); + + -- -- Name: app_versions_r2_path_idx; Type: INDEX; Schema: public; Owner: - -- From 1ef934d00be9ba4a01022410c26fa1436d0e000e Mon Sep 17 00:00:00 2001 From: Martin Donadieu Date: Wed, 22 Jul 2026 14:02:56 +0300 Subject: [PATCH 10/16] docs(db): clarify reclaim verify coverage and native_packages lock Document VACUUM FULL for TOAST compaction, verify all pgmq archives via list_queues, include pg_settings.unit, and align trigger comments with manifest-only reclaim. Co-authored-by: Cursor --- scripts/ops/reclaim_supabase_swap.sql | 5 +++ scripts/ops/verify_supabase_swap.sql | 44 +++++-------------- ...0260722082019_fix_supabase_swap_memory.sql | 9 ++-- 3 files changed, 21 insertions(+), 37 deletions(-) diff --git a/scripts/ops/reclaim_supabase_swap.sql b/scripts/ops/reclaim_supabase_swap.sql index 3fab978bf2..458a0c6c74 100644 --- a/scripts/ops/reclaim_supabase_swap.sql +++ b/scripts/ops/reclaim_supabase_swap.sql @@ -59,6 +59,9 @@ FROM pgmq.list_queues() SELECT public.null_migrated_app_version_manifests(); VACUUM (VERBOSE) public.app_versions; +-- Routine VACUUM does not shrink TOAST. After nulling is done, compact in the +-- maintenance window (exclusive lock): +-- VACUUM (FULL, VERBOSE) public.app_versions; -- --------------------------------------------------------------------------- -- 4) Trim audit_logs older than 30 days (bounded batches). Re-run until deleted=0. @@ -66,6 +69,8 @@ VACUUM (VERBOSE) public.app_versions; SELECT public.cleanup_old_audit_logs(); VACUUM (VERBOSE) public.audit_logs; +-- After deleted=0, compact TOAST if pg_total_relation_size must fall: +-- VACUUM (FULL, VERBOSE) public.audit_logs; -- --------------------------------------------------------------------------- -- 5) Final sizes diff --git a/scripts/ops/verify_supabase_swap.sql b/scripts/ops/verify_supabase_swap.sql index e4bafd6d42..fc686b627e 100644 --- a/scripts/ops/verify_supabase_swap.sql +++ b/scripts/ops/verify_supabase_swap.sql @@ -3,7 +3,7 @@ SELECT pg_size_pretty(pg_database_size(current_database())::bigint) AS db_size; -SELECT name, setting +SELECT name, setting, unit FROM pg_settings WHERE name IN ('shared_buffers', 'work_mem', 'max_connections'); @@ -66,42 +66,20 @@ FROM pg_indexes WHERE schemaname = 'public' AND indexname = 'app_versions_manifest_present_idx'; --- Bounded existence checks per archive queue (no full-table aggregates). -SELECT queue_name, has_rows_older_than_2d -FROM ( - SELECT 'a_on_manifest_create' AS queue_name, - EXISTS ( - SELECT 1 - FROM pgmq.a_on_manifest_create - WHERE archived_at < now() - interval '2 days' - LIMIT 1 - ) AS has_rows_older_than_2d - UNION ALL - SELECT 'a_on_version_update', - EXISTS ( - SELECT 1 - FROM pgmq.a_on_version_update - WHERE archived_at < now() - interval '2 days' - LIMIT 1 - ) - UNION ALL - SELECT 'a_webhook_dispatcher', +-- Same queue set as cleanup_queue_messages() (psql \gexec; bounded EXISTS). +SELECT format( + $fmt$SELECT %L AS queue_name, EXISTS ( SELECT 1 - FROM pgmq.a_webhook_dispatcher + FROM pgmq.a_%I WHERE archived_at < now() - interval '2 days' LIMIT 1 - ) - UNION ALL - SELECT 'a_on_channel_update', - EXISTS ( - SELECT 1 - FROM pgmq.a_on_channel_update - WHERE archived_at < now() - interval '2 days' - LIMIT 1 - ) -) AS archives -ORDER BY queue_name; + ) AS has_rows_older_than_2d;$fmt$, + queue_name, + queue_name +) +FROM pgmq.list_queues() +\gexec SELECT 'index hit rate' AS name, diff --git a/supabase/migrations/20260722082019_fix_supabase_swap_memory.sql b/supabase/migrations/20260722082019_fix_supabase_swap_memory.sql index 72345f7dbe..fe494d23fa 100644 --- a/supabase/migrations/20260722082019_fix_supabase_swap_memory.sql +++ b/supabase/migrations/20260722082019_fix_supabase_swap_memory.sql @@ -550,8 +550,8 @@ CREATE INDEX IF NOT EXISTS app_versions_manifest_present_idx WHERE manifest IS NOT NULL; -- --------------------------------------------------------------------------- --- Allow clearing dual-storage fat columns --- after upload (null only). Unblocks on_version_update + reclaim jobs. +-- Allow clearing dual-storage app_versions.manifest after upload (null only). +-- native_packages stays locked: no alternate persisted source of truth. -- --------------------------------------------------------------------------- CREATE OR REPLACE FUNCTION "public"."check_encrypted_bundle_on_insert"() RETURNS "trigger" LANGUAGE "plpgsql" SECURITY DEFINER @@ -568,8 +568,9 @@ BEGIN IF TG_OP = 'UPDATE' THEN bundle_was_ready := OLD.storage_provider IS DISTINCT FROM 'r2-direct'; - -- Nulling migrated dual-storage columns is allowed after upload completes. - -- Rewriting non-null manifest/native_packages content stays locked. + -- Nulling a fully migrated dual-storage manifest array is allowed after upload. + -- native_packages remains locked (compatibility metadata has no table copy). + -- Rewriting non-null manifest content stays locked. IF bundle_was_ready AND ( NEW.name IS DISTINCT FROM OLD.name From dd9b96a4b22bd598ba4ec7296e90dae4d1931029 Mon Sep 17 00:00:00 2001 From: Martin Donadieu Date: Wed, 22 Jul 2026 14:25:32 +0300 Subject: [PATCH 11/16] fix(db): tighten dual-storage reclaim safety and cleanup fairness Match each legacy manifest entry before nulling, skip queue fan-out only for reclaim nulling, round-robin the global queue cleanup budget, and move the candidate index to a concurrent ops step. Co-authored-by: Cursor --- read_replicate/schema_replicate.catalog.json | 7 -- read_replicate/schema_replicate.sql | 7 -- scripts/ops/reclaim_supabase_swap.sql | 33 ++++--- scripts/ops/verify_supabase_swap.sql | 65 ++++++++++---- ...0260722082019_fix_supabase_swap_memory.sql | 89 +++++++++++-------- 5 files changed, 119 insertions(+), 82 deletions(-) diff --git a/read_replicate/schema_replicate.catalog.json b/read_replicate/schema_replicate.catalog.json index 05fb5a9845..4857f5c5bb 100644 --- a/read_replicate/schema_replicate.catalog.json +++ b/read_replicate/schema_replicate.catalog.json @@ -1982,13 +1982,6 @@ "table": "app_versions", "valid": true }, - { - "constraintOwned": false, - "definition": "CREATE INDEX app_versions_manifest_present_idx ON public.app_versions USING btree (id) WHERE (manifest IS NOT NULL)", - "name": "app_versions_manifest_present_idx", - "table": "app_versions", - "valid": true - }, { "constraintOwned": true, "definition": "CREATE UNIQUE INDEX app_versions_name_app_id_key ON public.app_versions USING btree (name, app_id)", diff --git a/read_replicate/schema_replicate.sql b/read_replicate/schema_replicate.sql index 0fc007e74b..94e69f1dc1 100644 --- a/read_replicate/schema_replicate.sql +++ b/read_replicate/schema_replicate.sql @@ -607,13 +607,6 @@ ALTER TABLE ONLY public.orgs CREATE INDEX app_versions_cli_version_idx ON public.app_versions USING btree (cli_version); --- --- Name: app_versions_manifest_present_idx; Type: INDEX; Schema: public; Owner: - --- - -CREATE INDEX app_versions_manifest_present_idx ON public.app_versions USING btree (id) WHERE (manifest IS NOT NULL); - - -- -- Name: app_versions_r2_path_idx; Type: INDEX; Schema: public; Owner: - -- diff --git a/scripts/ops/reclaim_supabase_swap.sql b/scripts/ops/reclaim_supabase_swap.sql index 458a0c6c74..f3600295e5 100644 --- a/scripts/ops/reclaim_supabase_swap.sql +++ b/scripts/ops/reclaim_supabase_swap.sql @@ -1,7 +1,8 @@ -- Capgo-EU Phase A reclaim (run manually in a maintenance window). --- Prefer psql (VACUUM cannot run inside a transaction / SQL-editor DO block). +-- REQUIRED: psql (uses \gexec; VACUUM cannot run inside a transaction). +-- Prefer ~/.pgpass / PGPASSFILE instead of putting the password on the CLI. -- Example: --- PGPASSWORD=... psql "postgresql://..." -v ON_ERROR_STOP=1 -f scripts/ops/reclaim_supabase_swap.sql +-- psql "postgresql://postgres@HOST:5432/postgres?sslmode=require" -v ON_ERROR_STOP=1 -f scripts/ops/reclaim_supabase_swap.sql -- Safe order: truncate empty bloat -> batched archive deletes -> null dual manifests -> trim audit. -- Each statement commits separately. Re-run until cleanup notices report deleted/updated = 0 -- (functions always emit a notice, including zero totals). @@ -34,27 +35,29 @@ ORDER BY pg_total_relation_size(format('%I.%I', schemaname, relname)::regclass) TRUNCATE TABLE net._http_response; -- --------------------------------------------------------------------------- --- 2) Purge pgmq archives/stuck messages (global batch budget per call). --- Re-run this SELECT until the notice shows archived_deleted=0 and stuck_deleted=0. +-- 2) Purge pgmq archives/stuck messages (global round-robin batch budget). +-- Re-run this SELECT until the notice shows archived_deleted=0 and +-- stuck_deleted=0. -- --------------------------------------------------------------------------- SELECT public.cleanup_queue_messages(); --- Vacuum every pgmq archive + queue table (psql \gexec; VACUUM cannot run in DO/tx). -SELECT format('VACUUM (VERBOSE) pgmq.a_%I;', queue_name) +-- Vacuum every pgmq archive + queue table (quote full physical table name). +SELECT format('VACUUM (VERBOSE) pgmq.%I;', 'a_' || pg_catalog.lower(queue_name)) FROM pgmq.list_queues() \gexec -SELECT format('VACUUM (VERBOSE) pgmq.q_%I;', queue_name) +SELECT format('VACUUM (VERBOSE) pgmq.%I;', 'q_' || pg_catalog.lower(queue_name)) FROM pgmq.list_queues() \gexec -- Optional hard reclaim (stronger locks): --- SELECT format('VACUUM (FULL, VERBOSE) pgmq.a_%I;', queue_name) FROM pgmq.list_queues() \gexec --- SELECT format('VACUUM (FULL, VERBOSE) pgmq.q_%I;', queue_name) FROM pgmq.list_queues() \gexec +-- SELECT format('VACUUM (FULL, VERBOSE) pgmq.%I;', 'a_' || pg_catalog.lower(queue_name)) +-- FROM pgmq.list_queues() +-- \gexec -- --------------------------------------------------------------------------- -- 3) Null fully migrated app_versions.manifest arrays --- Requires every expected legacy entry to exist in public.manifest. --- Re-run until notice shows updated=0. +-- Requires every legacy entry to exist in public.manifest by +-- file_name/s3_path/file_hash. Re-run until notice shows updated=0. -- --------------------------------------------------------------------------- SELECT public.null_migrated_app_version_manifests(); @@ -63,8 +66,14 @@ VACUUM (VERBOSE) public.app_versions; -- maintenance window (exclusive lock): -- VACUUM (FULL, VERBOSE) public.app_versions; +-- Non-blocking candidate index for ongoing hourly cleanup (outside a tx): +-- CREATE INDEX CONCURRENTLY IF NOT EXISTS app_versions_manifest_present_idx +-- ON public.app_versions USING btree (id) +-- WHERE manifest IS NOT NULL; + -- --------------------------------------------------------------------------- --- 4) Trim audit_logs older than 30 days (bounded batches). Re-run until deleted=0. +-- 4) Trim audit_logs older than 30 days (bounded batches). +-- Re-run until deleted=0. -- --------------------------------------------------------------------------- SELECT public.cleanup_old_audit_logs(); diff --git a/scripts/ops/verify_supabase_swap.sql b/scripts/ops/verify_supabase_swap.sql index fc686b627e..588fa93e21 100644 --- a/scripts/ops/verify_supabase_swap.sql +++ b/scripts/ops/verify_supabase_swap.sql @@ -1,9 +1,13 @@ -- Post-deploy / post-reclaim verification for Capgo-EU swap pressure. --- Prefer psql. Avoids unbounded whole-table counts where possible. +-- REQUIRED: psql (uses \gexec). Example: +-- psql "$DATABASE_URL" -v ON_ERROR_STOP=1 -f scripts/ops/verify_supabase_swap.sql SELECT pg_size_pretty(pg_database_size(current_database())::bigint) AS db_size; -SELECT name, setting, unit +SELECT + name, + setting, + unit FROM pg_settings WHERE name IN ('shared_buffers', 'work_mem', 'max_connections'); @@ -24,7 +28,7 @@ WHERE (schemaname, relname) IN ( ) ORDER BY pg_total_relation_size(format('%I.%I', schemaname, relname)::regclass) DESC; --- Bound candidate discovery first, then evaluate eligibility inside the sample. +-- Sample first 1000 non-null manifests; a zero does not prove global completion. SELECT count(*) AS eligible_dual_storage_sample FROM ( SELECT sample.id @@ -36,20 +40,29 @@ FROM ( LIMIT 1000 ) AS sample WHERE cardinality(sample.manifest) > 0 - AND ( - SELECT count(*)::integer - FROM public.manifest AS m - WHERE m.app_version_id = sample.id - ) >= ( - CASE - WHEN COALESCE(sample.manifest_count, 0) >= cardinality(sample.manifest) - THEN COALESCE(sample.manifest_count, 0) - ELSE cardinality(sample.manifest) - END + AND NOT EXISTS ( + SELECT 1 + FROM unnest(sample.manifest) AS entry(file_name, s3_path, file_hash) + WHERE NOT EXISTS ( + SELECT 1 + FROM public.manifest AS m + WHERE m.app_version_id = sample.id + AND m.file_name = entry.file_name + AND m.s3_path = entry.s3_path + AND m.file_hash = entry.file_hash + ) ) ) AS eligible; -SELECT name, enabled, hour_interval, run_at_hour, run_at_minute, target, updated_at +SELECT + name, + enabled, + hour_interval, + run_at_hour, + run_at_minute, + target, + description, + updated_at FROM public.cron_tasks WHERE name IN ( 'cleanup_queue_messages', @@ -59,24 +72,38 @@ WHERE name IN ( ) ORDER BY name; --- process_all_cron_tasks() swallows per-task errors; cron.job_run_details only --- reflects the outer job. Prefer Postgres logs / healthchecks for task failures. +-- process_all_cron_tasks() swallows per-task errors; prefer Postgres logs / +-- healthchecks for task failures. SELECT indexname FROM pg_indexes WHERE schemaname = 'public' AND indexname = 'app_versions_manifest_present_idx'; --- Same queue set as cleanup_queue_messages() (psql \gexec; bounded EXISTS). +-- Same queue set as cleanup_queue_messages() (archives + stuck). SELECT format( $fmt$SELECT %L AS queue_name, EXISTS ( SELECT 1 - FROM pgmq.a_%I + FROM pgmq.%I WHERE archived_at < now() - interval '2 days' LIMIT 1 ) AS has_rows_older_than_2d;$fmt$, queue_name, - queue_name + 'a_' || pg_catalog.lower(queue_name) +) +FROM pgmq.list_queues() +\gexec + +SELECT format( + $fmt$SELECT %L AS queue_name, + EXISTS ( + SELECT 1 + FROM pgmq.%I + WHERE read_ct > 5 + LIMIT 1 + ) AS has_stuck_read_ct_gt_5;$fmt$, + queue_name, + 'q_' || pg_catalog.lower(queue_name) ) FROM pgmq.list_queues() \gexec diff --git a/supabase/migrations/20260722082019_fix_supabase_swap_memory.sql b/supabase/migrations/20260722082019_fix_supabase_swap_memory.sql index fe494d23fa..acf2b1e69f 100644 --- a/supabase/migrations/20260722082019_fix_supabase_swap_memory.sql +++ b/supabase/migrations/20260722082019_fix_supabase_swap_memory.sql @@ -22,13 +22,17 @@ DECLARE deleted_batch integer; deleted_archived_total bigint := 0; deleted_stuck_total bigint := 0; + did_work boolean; BEGIN - FOR queue_name IN ( - SELECT q.queue_name FROM pgmq.list_queues() q - ) LOOP + -- Round-robin: at most one archive batch and one stuck batch per queue per pass, + -- so a busy first queue cannot starve later queues within the global budget. + LOOP EXIT WHEN batches_used >= max_batches_total; + did_work := false; - LOOP + FOR queue_name IN ( + SELECT q.queue_name FROM pgmq.list_queues() q + ) LOOP EXIT WHEN batches_used >= max_batches_total; EXECUTE pg_catalog.format( @@ -45,12 +49,12 @@ BEGIN USING cutoff, batch_size; GET DIAGNOSTICS deleted_batch = ROW_COUNT; - EXIT WHEN deleted_batch = 0; - batches_used := batches_used + 1; - deleted_archived_total := deleted_archived_total + deleted_batch; - END LOOP; + IF deleted_batch > 0 THEN + batches_used := batches_used + 1; + deleted_archived_total := deleted_archived_total + deleted_batch; + did_work := true; + END IF; - LOOP EXIT WHEN batches_used >= max_batches_total; EXECUTE pg_catalog.format( @@ -67,10 +71,14 @@ BEGIN USING batch_size; GET DIAGNOSTICS deleted_batch = ROW_COUNT; - EXIT WHEN deleted_batch = 0; - batches_used := batches_used + 1; - deleted_stuck_total := deleted_stuck_total + deleted_batch; + IF deleted_batch > 0 THEN + batches_used := batches_used + 1; + deleted_stuck_total := deleted_stuck_total + deleted_batch; + did_work := true; + END IF; END LOOP; + + EXIT WHEN NOT did_work; END LOOP; RAISE NOTICE @@ -168,16 +176,17 @@ BEGIN FROM public.app_versions AS av WHERE av.manifest IS NOT NULL AND pg_catalog.cardinality(av.manifest) > 0 - AND ( - SELECT count(*)::integer - FROM public.manifest AS m - WHERE m.app_version_id = av.id - ) >= ( - CASE - WHEN COALESCE(av.manifest_count, 0) >= pg_catalog.cardinality(av.manifest) - THEN COALESCE(av.manifest_count, 0) - ELSE pg_catalog.cardinality(av.manifest) - END + AND NOT EXISTS ( + SELECT 1 + FROM pg_catalog.unnest(av.manifest) AS entry(file_name, s3_path, file_hash) + WHERE NOT EXISTS ( + SELECT 1 + FROM public.manifest AS m + WHERE m.app_version_id = av.id + AND m.file_name = entry.file_name + AND m.s3_path = entry.s3_path + AND m.file_hash = entry.file_hash + ) ) ORDER BY av.id LIMIT batch_size @@ -424,8 +433,11 @@ BEGIN old_record_payload := old_record_payload - 'manifest' - 'native_packages'; END IF; - -- Dual-storage reclaim only nulls fat columns (+ auto updated_at). Skip queue fan-out. + -- Skip queue fan-out only for dual-storage reclaim: non-null manifest -> NULL. + -- Do not skip upload-time manifest writes that still need on_version_update migration. IF TG_OP = 'UPDATE' + AND OLD.manifest IS NOT NULL + AND NEW.manifest IS NULL AND (record_payload - 'updated_at') IS NOT DISTINCT FROM (old_record_payload - 'updated_at') THEN RETURN NEW; @@ -544,10 +556,12 @@ SET updated_at = pg_catalog.now(); --- Bound dual-storage candidate discovery once most arrays are nulled. -CREATE INDEX IF NOT EXISTS app_versions_manifest_present_idx - ON public.app_versions USING btree (id) - WHERE manifest IS NOT NULL; + +UPDATE public.cron_tasks +SET + description = 'Delete audit_logs older than 30 days in bounded batches', + updated_at = pg_catalog.now() +WHERE name = 'cleanup_old_audit_logs'; -- --------------------------------------------------------------------------- -- Allow clearing dual-storage app_versions.manifest after upload (null only). @@ -586,16 +600,17 @@ BEGIN OR ( NEW.manifest IS NULL AND OLD.manifest IS NOT NULL - AND ( - SELECT count(*)::integer - FROM public.manifest AS m - WHERE m.app_version_id = OLD.id - ) < ( - CASE - WHEN COALESCE(OLD.manifest_count, 0) >= COALESCE(pg_catalog.cardinality(OLD.manifest), 0) - THEN COALESCE(OLD.manifest_count, 0) - ELSE COALESCE(pg_catalog.cardinality(OLD.manifest), 0) - END + AND EXISTS ( + SELECT 1 + FROM pg_catalog.unnest(OLD.manifest) AS entry(file_name, s3_path, file_hash) + WHERE NOT EXISTS ( + SELECT 1 + FROM public.manifest AS m + WHERE m.app_version_id = OLD.id + AND m.file_name = entry.file_name + AND m.s3_path = entry.s3_path + AND m.file_hash = entry.file_hash + ) ) ) OR NEW.native_packages IS DISTINCT FROM OLD.native_packages From 140138f40a40f29ca43e31e19ae09e3ed3c1391b Mon Sep 17 00:00:00 2001 From: Martin Donadieu Date: Wed, 22 Jul 2026 14:44:36 +0300 Subject: [PATCH 12/16] fix(db): block incomplete manifest nulling on in-progress bundles Require full public.manifest entry coverage before encryption bypass or reclaim nulling, keep upload fat-field names in audit changed_fields, and only skip queue/audit for true reclaim null transitions. Co-authored-by: Cursor --- ...0260722082019_fix_supabase_swap_memory.sql | 65 +++++++++++++++---- tests/cleanup_swap_memory.test.ts | 5 +- 2 files changed, 54 insertions(+), 16 deletions(-) diff --git a/supabase/migrations/20260722082019_fix_supabase_swap_memory.sql b/supabase/migrations/20260722082019_fix_supabase_swap_memory.sql index acf2b1e69f..fb1eebd61a 100644 --- a/supabase/migrations/20260722082019_fix_supabase_swap_memory.sql +++ b/supabase/migrations/20260722082019_fix_supabase_swap_memory.sql @@ -237,7 +237,6 @@ DECLARE v_stats_refresh_fields constant text[] := ARRAY['stats_refresh_requested_at', 'stats_updated_at', 'updated_at']; v_background_counter_fields constant text[] := ARRAY['channel_device_count', 'manifest_bundle_count', 'updated_at']; v_fat_app_version_fields constant text[] := ARRAY['manifest', 'native_packages']; - v_noise_app_version_fields constant text[] := ARRAY['manifest', 'native_packages', 'updated_at']; BEGIN SELECT auth.uid() INTO v_actor_user_id; @@ -318,6 +317,7 @@ BEGIN END IF; -- Never persist multi-MB array/json columns in audit TOAST. + -- Keep fat field names in changed_fields so upload-time edits remain visible. IF TG_TABLE_NAME = 'app_versions' THEN IF v_old_record IS NOT NULL THEN v_old_record := v_old_record - v_fat_app_version_fields; @@ -325,19 +325,16 @@ BEGIN IF v_new_record IS NOT NULL THEN v_new_record := v_new_record - v_fat_app_version_fields; END IF; - IF v_changed_fields IS NOT NULL THEN - SELECT pg_catalog.array_agg(field_name) - INTO v_changed_fields - FROM pg_catalog.unnest(v_changed_fields) AS changed_field(field_name) - WHERE changed_field.field_name <> ALL(v_fat_app_version_fields); - END IF; - -- Skip updates that only touched stripped fat columns / auto timestamps. + -- Skip audit only for dual-storage reclaim: non-null manifest -> NULL (+ updated_at). IF TG_OP = 'UPDATE' + AND OLD.manifest IS NOT NULL + AND NEW.manifest IS NULL + AND NEW.native_packages IS NOT DISTINCT FROM OLD.native_packages AND NOT EXISTS ( SELECT 1 FROM pg_catalog.unnest(COALESCE(v_changed_fields, ARRAY[]::text[])) AS changed_field(field_name) - WHERE changed_field.field_name <> ALL(v_noise_app_version_fields) + WHERE changed_field.field_name <> ALL(ARRAY['manifest', 'updated_at']::text[]) ) THEN RETURN NEW; END IF; @@ -434,10 +431,11 @@ BEGIN END IF; -- Skip queue fan-out only for dual-storage reclaim: non-null manifest -> NULL. - -- Do not skip upload-time manifest writes that still need on_version_update migration. + -- native_packages is stripped from payloads, so compare it explicitly. IF TG_OP = 'UPDATE' AND OLD.manifest IS NOT NULL AND NEW.manifest IS NULL + AND NEW.native_packages IS NOT DISTINCT FROM OLD.native_packages AND (record_payload - 'updated_at') IS NOT DISTINCT FROM (old_record_payload - 'updated_at') THEN RETURN NEW; @@ -580,6 +578,27 @@ DECLARE bundle_was_ready boolean; BEGIN IF TG_OP = 'UPDATE' THEN + -- Never drop the only copy of legacy file metadata, ready or not. + IF NEW.manifest IS NULL + AND OLD.manifest IS NOT NULL + AND EXISTS ( + SELECT 1 + FROM pg_catalog.unnest(OLD.manifest) AS entry(file_name, s3_path, file_hash) + WHERE NOT EXISTS ( + SELECT 1 + FROM public.manifest AS m + WHERE m.app_version_id = OLD.id + AND m.file_name = entry.file_name + AND m.s3_path = entry.s3_path + AND m.file_hash = entry.file_hash + ) + ) + THEN + RAISE EXCEPTION '%', + 'bundle_manifest_not_migrated: Cannot clear app_versions.manifest ' + || 'until every entry exists in public.manifest.'; + END IF; + bundle_was_ready := OLD.storage_provider IS DISTINCT FROM 'r2-direct'; -- Nulling a fully migrated dual-storage manifest array is allowed after upload. @@ -632,9 +651,9 @@ BEGIN END IF; END IF; - -- Manifest/native_packages nulling must not re-run encryption enforcement. - -- Legacy rows can predate org encryption requirements; reclaim only clears - -- dual-storage columns and must not abort on those orgs. + -- Fully migrated dual-storage nulling must not re-run encryption enforcement. + -- Incomplete nulling (still missing public.manifest rows) must not bypass checks, + -- including for in-progress r2-direct uploads. IF TG_OP = 'UPDATE' AND NEW.session_key IS NOT DISTINCT FROM OLD.session_key AND NEW.key_id IS NOT DISTINCT FROM OLD.key_id @@ -645,7 +664,25 @@ BEGIN AND NEW.external_url IS NOT DISTINCT FROM OLD.external_url AND NEW.checksum IS NOT DISTINCT FROM OLD.checksum AND NEW.native_packages IS NOT DISTINCT FROM OLD.native_packages - AND (NEW.manifest IS NULL OR NEW.manifest IS NOT DISTINCT FROM OLD.manifest) + AND ( + NEW.manifest IS NOT DISTINCT FROM OLD.manifest + OR ( + NEW.manifest IS NULL + AND OLD.manifest IS NOT NULL + AND NOT EXISTS ( + SELECT 1 + FROM pg_catalog.unnest(OLD.manifest) AS entry(file_name, s3_path, file_hash) + WHERE NOT EXISTS ( + SELECT 1 + FROM public.manifest AS m + WHERE m.app_version_id = OLD.id + AND m.file_name = entry.file_name + AND m.s3_path = entry.s3_path + AND m.file_hash = entry.file_hash + ) + ) + ) + ) THEN RETURN NEW; END IF; diff --git a/tests/cleanup_swap_memory.test.ts b/tests/cleanup_swap_memory.test.ts index 8e7043e17b..6e3d00f73a 100644 --- a/tests/cleanup_swap_memory.test.ts +++ b/tests/cleanup_swap_memory.test.ts @@ -158,8 +158,9 @@ describe('swap memory cleanup functions', () => { expect(logs[0]?.new_record?.native_packages).toBeUndefined() expect(logs[0]?.new_record?.comment).toBe('after') expect(logs[0]?.changed_fields).toContain('comment') - expect(logs[0]?.changed_fields ?? []).not.toContain('manifest') - expect(logs[0]?.changed_fields ?? []).not.toContain('native_packages') + // Fat payloads stay stripped, but field names remain for upload-time history. + expect(logs[0]?.changed_fields).toContain('manifest') + expect(logs[0]?.changed_fields).toContain('native_packages') await executeSQL(`DELETE FROM public.audit_logs WHERE record_id = $1 AND table_name = 'app_versions'`, [String(versionId)]) await executeSQL(`DELETE FROM public.app_versions WHERE id = $1`, [versionId]) From f537d99138acbd66182e2928c5e4a27abd4a98b7 Mon Sep 17 00:00:00 2001 From: Martin Donadieu Date: Wed, 22 Jul 2026 14:44:57 +0300 Subject: [PATCH 13/16] chore(ops): drop unused manifest_count from verify sample Co-authored-by: Cursor --- scripts/ops/verify_supabase_swap.sql | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/scripts/ops/verify_supabase_swap.sql b/scripts/ops/verify_supabase_swap.sql index 588fa93e21..471dd02aa4 100644 --- a/scripts/ops/verify_supabase_swap.sql +++ b/scripts/ops/verify_supabase_swap.sql @@ -33,7 +33,7 @@ SELECT count(*) AS eligible_dual_storage_sample FROM ( SELECT sample.id FROM ( - SELECT av.id, av.manifest, av.manifest_count + SELECT av.id, av.manifest FROM public.app_versions AS av WHERE av.manifest IS NOT NULL ORDER BY av.id From a3b7649c533dad246e852873c15cb4797e4b6316 Mon Sep 17 00:00:00 2001 From: Martin Donadieu Date: Wed, 22 Jul 2026 14:50:38 +0300 Subject: [PATCH 14/16] fix(db): add manifest candidate index and harden cleanup budget exits Restore the partial app_versions index before hourly nulling runs, and exit the queue loop immediately when the global batch budget is spent. Co-authored-by: Cursor --- read_replicate/schema_replicate.catalog.json | 7 +++++++ read_replicate/schema_replicate.sql | 7 +++++++ scripts/ops/reclaim_supabase_swap.sql | 3 ++- .../20260722082019_fix_supabase_swap_memory.sql | 11 ++++++++++- 4 files changed, 26 insertions(+), 2 deletions(-) diff --git a/read_replicate/schema_replicate.catalog.json b/read_replicate/schema_replicate.catalog.json index 4857f5c5bb..05fb5a9845 100644 --- a/read_replicate/schema_replicate.catalog.json +++ b/read_replicate/schema_replicate.catalog.json @@ -1982,6 +1982,13 @@ "table": "app_versions", "valid": true }, + { + "constraintOwned": false, + "definition": "CREATE INDEX app_versions_manifest_present_idx ON public.app_versions USING btree (id) WHERE (manifest IS NOT NULL)", + "name": "app_versions_manifest_present_idx", + "table": "app_versions", + "valid": true + }, { "constraintOwned": true, "definition": "CREATE UNIQUE INDEX app_versions_name_app_id_key ON public.app_versions USING btree (name, app_id)", diff --git a/read_replicate/schema_replicate.sql b/read_replicate/schema_replicate.sql index 94e69f1dc1..0fc007e74b 100644 --- a/read_replicate/schema_replicate.sql +++ b/read_replicate/schema_replicate.sql @@ -607,6 +607,13 @@ ALTER TABLE ONLY public.orgs CREATE INDEX app_versions_cli_version_idx ON public.app_versions USING btree (cli_version); +-- +-- Name: app_versions_manifest_present_idx; Type: INDEX; Schema: public; Owner: - +-- + +CREATE INDEX app_versions_manifest_present_idx ON public.app_versions USING btree (id) WHERE (manifest IS NOT NULL); + + -- -- Name: app_versions_r2_path_idx; Type: INDEX; Schema: public; Owner: - -- diff --git a/scripts/ops/reclaim_supabase_swap.sql b/scripts/ops/reclaim_supabase_swap.sql index f3600295e5..c83b4c7425 100644 --- a/scripts/ops/reclaim_supabase_swap.sql +++ b/scripts/ops/reclaim_supabase_swap.sql @@ -66,7 +66,8 @@ VACUUM (VERBOSE) public.app_versions; -- maintenance window (exclusive lock): -- VACUUM (FULL, VERBOSE) public.app_versions; --- Non-blocking candidate index for ongoing hourly cleanup (outside a tx): +-- Candidate index for hourly cleanup. Migration creates it non-concurrently; +-- if deploying via ops only, prefer CONCURRENTLY outside a transaction: -- CREATE INDEX CONCURRENTLY IF NOT EXISTS app_versions_manifest_present_idx -- ON public.app_versions USING btree (id) -- WHERE manifest IS NOT NULL; diff --git a/supabase/migrations/20260722082019_fix_supabase_swap_memory.sql b/supabase/migrations/20260722082019_fix_supabase_swap_memory.sql index fb1eebd61a..17137cc81b 100644 --- a/supabase/migrations/20260722082019_fix_supabase_swap_memory.sql +++ b/supabase/migrations/20260722082019_fix_supabase_swap_memory.sql @@ -55,7 +55,9 @@ BEGIN did_work := true; END IF; - EXIT WHEN batches_used >= max_batches_total; + IF batches_used >= max_batches_total THEN + EXIT; -- leave queue loop; do not start stuck deletes over budget + END IF; EXECUTE pg_catalog.format( 'DELETE FROM pgmq.q_%I @@ -561,6 +563,13 @@ SET updated_at = pg_catalog.now() WHERE name = 'cleanup_old_audit_logs'; + +-- Bound dual-storage candidate discovery for hourly reclaim. +-- Maintenance-window deploy: brief lock on app_versions is expected. +CREATE INDEX IF NOT EXISTS app_versions_manifest_present_idx + ON public.app_versions USING btree (id) + WHERE manifest IS NOT NULL; + -- --------------------------------------------------------------------------- -- Allow clearing dual-storage app_versions.manifest after upload (null only). -- native_packages stays locked: no alternate persisted source of truth. From 6a6fd8a50d7993922c0d569d3e411b3d715cf17e Mon Sep 17 00:00:00 2001 From: Martin Donadieu Date: Wed, 22 Jul 2026 15:06:46 +0300 Subject: [PATCH 15/16] fix(db): match dual-storage reclaim by s3_path and file_hash file_name may be normalized during table migration, so entry equality must not require the legacy array's original file_name string. Co-authored-by: Cursor --- scripts/ops/verify_supabase_swap.sql | 1 - .../migrations/20260722082019_fix_supabase_swap_memory.sql | 4 ++-- 2 files changed, 2 insertions(+), 3 deletions(-) diff --git a/scripts/ops/verify_supabase_swap.sql b/scripts/ops/verify_supabase_swap.sql index 471dd02aa4..44fcbc5d1a 100644 --- a/scripts/ops/verify_supabase_swap.sql +++ b/scripts/ops/verify_supabase_swap.sql @@ -47,7 +47,6 @@ FROM ( SELECT 1 FROM public.manifest AS m WHERE m.app_version_id = sample.id - AND m.file_name = entry.file_name AND m.s3_path = entry.s3_path AND m.file_hash = entry.file_hash ) diff --git a/supabase/migrations/20260722082019_fix_supabase_swap_memory.sql b/supabase/migrations/20260722082019_fix_supabase_swap_memory.sql index 17137cc81b..1ecc9d6220 100644 --- a/supabase/migrations/20260722082019_fix_supabase_swap_memory.sql +++ b/supabase/migrations/20260722082019_fix_supabase_swap_memory.sql @@ -185,7 +185,7 @@ BEGIN SELECT 1 FROM public.manifest AS m WHERE m.app_version_id = av.id - AND m.file_name = entry.file_name + -- Match stable identity; file_name may have been normalized at migrate time. AND m.s3_path = entry.s3_path AND m.file_hash = entry.file_hash ) @@ -685,7 +685,7 @@ BEGIN SELECT 1 FROM public.manifest AS m WHERE m.app_version_id = OLD.id - AND m.file_name = entry.file_name + -- Match stable identity; file_name may have been normalized at migrate time. AND m.s3_path = entry.s3_path AND m.file_hash = entry.file_hash ) From 4757e8aaee9d491866741b38d3bcf6126f11728c Mon Sep 17 00:00:00 2001 From: Martin Donadieu Date: Wed, 22 Jul 2026 17:22:14 +0300 Subject: [PATCH 16/16] fix(db): align trigger reclaim guards with s3_path/hash identity Drop write-blocking index creation from the migration (ops uses CREATE INDEX CONCURRENTLY) and stop requiring normalized file_name matches in check_encrypted_bundle_on_insert. Co-authored-by: Cursor --- read_replicate/schema_replicate.catalog.json | 7 --- read_replicate/schema_replicate.sql | 7 --- scripts/ops/reclaim_supabase_swap.sql | 63 +++++++++---------- scripts/ops/verify_supabase_swap.sql | 17 ++--- ...0260722082019_fix_supabase_swap_memory.sql | 25 +++----- 5 files changed, 49 insertions(+), 70 deletions(-) diff --git a/read_replicate/schema_replicate.catalog.json b/read_replicate/schema_replicate.catalog.json index 05fb5a9845..4857f5c5bb 100644 --- a/read_replicate/schema_replicate.catalog.json +++ b/read_replicate/schema_replicate.catalog.json @@ -1982,13 +1982,6 @@ "table": "app_versions", "valid": true }, - { - "constraintOwned": false, - "definition": "CREATE INDEX app_versions_manifest_present_idx ON public.app_versions USING btree (id) WHERE (manifest IS NOT NULL)", - "name": "app_versions_manifest_present_idx", - "table": "app_versions", - "valid": true - }, { "constraintOwned": true, "definition": "CREATE UNIQUE INDEX app_versions_name_app_id_key ON public.app_versions USING btree (name, app_id)", diff --git a/read_replicate/schema_replicate.sql b/read_replicate/schema_replicate.sql index 0fc007e74b..94e69f1dc1 100644 --- a/read_replicate/schema_replicate.sql +++ b/read_replicate/schema_replicate.sql @@ -607,13 +607,6 @@ ALTER TABLE ONLY public.orgs CREATE INDEX app_versions_cli_version_idx ON public.app_versions USING btree (cli_version); --- --- Name: app_versions_manifest_present_idx; Type: INDEX; Schema: public; Owner: - --- - -CREATE INDEX app_versions_manifest_present_idx ON public.app_versions USING btree (id) WHERE (manifest IS NOT NULL); - - -- -- Name: app_versions_r2_path_idx; Type: INDEX; Schema: public; Owner: - -- diff --git a/scripts/ops/reclaim_supabase_swap.sql b/scripts/ops/reclaim_supabase_swap.sql index c83b4c7425..9717a8478b 100644 --- a/scripts/ops/reclaim_supabase_swap.sql +++ b/scripts/ops/reclaim_supabase_swap.sql @@ -3,10 +3,12 @@ -- Prefer ~/.pgpass / PGPASSFILE instead of putting the password on the CLI. -- Example: -- psql "postgresql://postgres@HOST:5432/postgres?sslmode=require" -v ON_ERROR_STOP=1 -f scripts/ops/reclaim_supabase_swap.sql --- Safe order: truncate empty bloat -> batched archive deletes -> null dual manifests -> trim audit. --- Each statement commits separately. Re-run until cleanup notices report deleted/updated = 0 +-- Safe order: index -> truncate -> archives -> null manifests -> audit trim. +-- Re-run the FULL script until cleanup notices report deleted/updated = 0 -- (functions always emit a notice, including zero totals). +SET lock_timeout = '5s'; + -- --------------------------------------------------------------------------- -- 0) Baseline sizes -- --------------------------------------------------------------------------- @@ -29,57 +31,52 @@ WHERE (schemaname, relname) IN ( ) ORDER BY pg_total_relation_size(format('%I.%I', schemaname, relname)::regclass) DESC; +-- --------------------------------------------------------------------------- +-- 0b) Candidate index for hourly nulling (non-blocking; must be outside a tx) +-- --------------------------------------------------------------------------- +CREATE INDEX CONCURRENTLY IF NOT EXISTS app_versions_manifest_present_idx + ON public.app_versions USING btree (id) + WHERE manifest IS NOT NULL; + -- --------------------------------------------------------------------------- -- 1) Truncate pg_net response bloat -- --------------------------------------------------------------------------- TRUNCATE TABLE net._http_response; -- --------------------------------------------------------------------------- --- 2) Purge pgmq archives/stuck messages (global round-robin batch budget). --- Re-run this SELECT until the notice shows archived_deleted=0 and --- stuck_deleted=0. +-- 2) Purge pgmq archives/stuck messages. +-- Re-run the FULL script until archived_deleted=0 and stuck_deleted=0. -- --------------------------------------------------------------------------- SELECT public.cleanup_queue_messages(); --- Vacuum every pgmq archive + queue table (quote full physical table name). -SELECT format('VACUUM (VERBOSE) pgmq.%I;', 'a_' || pg_catalog.lower(queue_name)) -FROM pgmq.list_queues() -\gexec -SELECT format('VACUUM (VERBOSE) pgmq.%I;', 'q_' || pg_catalog.lower(queue_name)) -FROM pgmq.list_queues() -\gexec - --- Optional hard reclaim (stronger locks): --- SELECT format('VACUUM (FULL, VERBOSE) pgmq.%I;', 'a_' || pg_catalog.lower(queue_name)) --- FROM pgmq.list_queues() --- \gexec +-- Vacuum Capgo-EU evidenced bloated queues only. +VACUUM (VERBOSE) pgmq.a_on_version_update; +VACUUM (VERBOSE) pgmq.a_on_manifest_create; +VACUUM (VERBOSE) pgmq.a_webhook_dispatcher; +VACUUM (VERBOSE) pgmq.a_on_channel_update; +VACUUM (VERBOSE) pgmq.q_on_version_update; +VACUUM (VERBOSE) pgmq.q_on_manifest_create; +VACUUM (VERBOSE) pgmq.q_webhook_dispatcher; +VACUUM (VERBOSE) pgmq.q_on_channel_update; -- --------------------------------------------------------------------------- --- 3) Null fully migrated app_versions.manifest arrays --- Requires every legacy entry to exist in public.manifest by --- file_name/s3_path/file_hash. Re-run until notice shows updated=0. +-- 3) Null fully migrated app_versions.manifest arrays (s3_path + file_hash). +-- Re-run the FULL script until updated=0. -- --------------------------------------------------------------------------- SELECT public.null_migrated_app_version_manifests(); -VACUUM (VERBOSE) public.app_versions; --- Routine VACUUM does not shrink TOAST. After nulling is done, compact in the --- maintenance window (exclusive lock): +VACUUM (ANALYZE, VERBOSE) public.app_versions; +-- Optional TOAST compaction after updated=0: -- VACUUM (FULL, VERBOSE) public.app_versions; --- Candidate index for hourly cleanup. Migration creates it non-concurrently; --- if deploying via ops only, prefer CONCURRENTLY outside a transaction: --- CREATE INDEX CONCURRENTLY IF NOT EXISTS app_versions_manifest_present_idx --- ON public.app_versions USING btree (id) --- WHERE manifest IS NOT NULL; - -- --------------------------------------------------------------------------- --- 4) Trim audit_logs older than 30 days (bounded batches). --- Re-run until deleted=0. +-- 4) Trim audit_logs older than 30 days. +-- Re-run the FULL script until deleted=0. -- --------------------------------------------------------------------------- SELECT public.cleanup_old_audit_logs(); -VACUUM (VERBOSE) public.audit_logs; --- After deleted=0, compact TOAST if pg_total_relation_size must fall: +VACUUM (ANALYZE, VERBOSE) public.audit_logs; +-- Optional TOAST compaction after deleted=0: -- VACUUM (FULL, VERBOSE) public.audit_logs; -- --------------------------------------------------------------------------- diff --git a/scripts/ops/verify_supabase_swap.sql b/scripts/ops/verify_supabase_swap.sql index 44fcbc5d1a..732908b759 100644 --- a/scripts/ops/verify_supabase_swap.sql +++ b/scripts/ops/verify_supabase_swap.sql @@ -39,8 +39,7 @@ FROM ( ORDER BY av.id LIMIT 1000 ) AS sample - WHERE cardinality(sample.manifest) > 0 - AND NOT EXISTS ( + WHERE NOT EXISTS ( SELECT 1 FROM unnest(sample.manifest) AS entry(file_name, s3_path, file_hash) WHERE NOT EXISTS ( @@ -71,14 +70,18 @@ WHERE name IN ( ) ORDER BY name; --- process_all_cron_tasks() swallows per-task errors; prefer Postgres logs / --- healthchecks for task failures. SELECT indexname FROM pg_indexes WHERE schemaname = 'public' AND indexname = 'app_versions_manifest_present_idx'; --- Same queue set as cleanup_queue_messages() (archives + stuck). +SELECT EXISTS ( + SELECT 1 + FROM public.audit_logs + WHERE created_at < now() - interval '30 days' + LIMIT 1 +) AS has_audit_logs_older_than_30d; + SELECT format( $fmt$SELECT %L AS queue_name, EXISTS ( @@ -109,10 +112,10 @@ FROM pgmq.list_queues() SELECT 'index hit rate' AS name, - ROUND((sum(idx_blks_hit) / nullif(sum(idx_blks_hit + idx_blks_read), 0) * 100)::numeric, 2) AS ratio + ROUND((sum(idx_blks_hit)::numeric / nullif(sum(idx_blks_hit + idx_blks_read), 0) * 100), 2) AS ratio FROM pg_statio_user_indexes UNION ALL SELECT 'table hit rate', - ROUND((sum(heap_blks_hit) / nullif(sum(heap_blks_hit) + sum(heap_blks_read), 0) * 100)::numeric, 2) + ROUND((sum(heap_blks_hit)::numeric / nullif(sum(heap_blks_hit) + sum(heap_blks_read), 0) * 100), 2) FROM pg_statio_user_tables; diff --git a/supabase/migrations/20260722082019_fix_supabase_swap_memory.sql b/supabase/migrations/20260722082019_fix_supabase_swap_memory.sql index 1ecc9d6220..8b88de00e2 100644 --- a/supabase/migrations/20260722082019_fix_supabase_swap_memory.sql +++ b/supabase/migrations/20260722082019_fix_supabase_swap_memory.sql @@ -177,7 +177,6 @@ BEGIN SELECT av.id FROM public.app_versions AS av WHERE av.manifest IS NOT NULL - AND pg_catalog.cardinality(av.manifest) > 0 AND NOT EXISTS ( SELECT 1 FROM pg_catalog.unnest(av.manifest) AS entry(file_name, s3_path, file_hash) @@ -564,12 +563,6 @@ SET WHERE name = 'cleanup_old_audit_logs'; --- Bound dual-storage candidate discovery for hourly reclaim. --- Maintenance-window deploy: brief lock on app_versions is expected. -CREATE INDEX IF NOT EXISTS app_versions_manifest_present_idx - ON public.app_versions USING btree (id) - WHERE manifest IS NOT NULL; - -- --------------------------------------------------------------------------- -- Allow clearing dual-storage app_versions.manifest after upload (null only). -- native_packages stays locked: no alternate persisted source of truth. @@ -597,7 +590,7 @@ BEGIN SELECT 1 FROM public.manifest AS m WHERE m.app_version_id = OLD.id - AND m.file_name = entry.file_name + -- Match stable identity; file_name may have been normalized at migrate time. AND m.s3_path = entry.s3_path AND m.file_hash = entry.file_hash ) @@ -635,7 +628,7 @@ BEGIN SELECT 1 FROM public.manifest AS m WHERE m.app_version_id = OLD.id - AND m.file_name = entry.file_name + -- Match stable identity; file_name may have been normalized at migrate time. AND m.s3_path = entry.s3_path AND m.file_hash = entry.file_hash ) @@ -645,7 +638,7 @@ BEGIN ) THEN PERFORM public.pg_log('deny: BUNDLE_CONTENT_LOCKED_TRIGGER', - jsonb_build_object( + pg_catalog.jsonb_build_object( 'org_id', OLD.owner_org, 'app_id', OLD.app_id, 'version_name', OLD.name, @@ -721,11 +714,11 @@ BEGIN END IF; bundle_is_encrypted := public.is_bundle_encrypted(NEW.session_key); - bundle_key_id := NULLIF(btrim(NEW.key_id), '')::varchar(20); + bundle_key_id := NULLIF(pg_catalog.btrim(NEW.key_id), '')::varchar(20); IF NOT bundle_is_encrypted THEN PERFORM public.pg_log('deny: ORG_REQUIRES_ENCRYPTED_BUNDLES_TRIGGER', - jsonb_build_object( + pg_catalog.jsonb_build_object( 'org_id', org_id, 'app_id', NEW.app_id, 'version_name', NEW.name, @@ -740,7 +733,7 @@ BEGIN IF org_required_key IS NOT NULL AND org_required_key <> '' THEN IF bundle_key_id IS NULL THEN PERFORM public.pg_log('deny: ORG_REQUIRES_SPECIFIC_ENCRYPTION_KEY_TRIGGER', - jsonb_build_object( + pg_catalog.jsonb_build_object( 'org_id', org_id, 'app_id', NEW.app_id, 'version_name', NEW.name, @@ -757,11 +750,11 @@ BEGIN -- key_id is 20 chars and required_encryption_key may be 20 or 21 chars. IF NOT ( - bundle_key_id = LEFT(org_required_key, 20) - OR LEFT(bundle_key_id, LENGTH(org_required_key)) = org_required_key + bundle_key_id = pg_catalog.left(org_required_key, 20) + OR pg_catalog.left(bundle_key_id, pg_catalog.length(org_required_key)) = org_required_key ) THEN PERFORM public.pg_log('deny: ORG_REQUIRES_SPECIFIC_ENCRYPTION_KEY_TRIGGER', - jsonb_build_object( + pg_catalog.jsonb_build_object( 'org_id', org_id, 'app_id', NEW.app_id, 'version_name', NEW.name,