diff --git a/.github/workflows/tests.yml b/.github/workflows/tests.yml index 247c80d9e2..addd8113d6 100644 --- a/.github/workflows/tests.yml +++ b/.github/workflows/tests.yml @@ -132,6 +132,8 @@ jobs: FILTER_CLI: ${{ steps.filter.outputs.cli }} FILTER_NOTIFICATIONS: ${{ steps.filter.outputs.notifications }} FILTER_SHARED: ${{ steps.filter.outputs.shared }} + GH_TOKEN: ${{ github.token }} + BRANCH_NAME: ${{ github.ref_name }} run: | if [ "$EVENT_NAME" = "workflow_call" ]; then echo "capgo=$INPUT_RUN_CAPGO" >> "$GITHUB_OUTPUT" @@ -140,6 +142,35 @@ jobs: exit 0 fi + # Push + pull_request both run this workflow on feature branches. When an open + # PR already covers the branch, the push suite is redundant and its flakes + # (Docker port binds, transient 502/503) stick on the PR check list. + if [ "$EVENT_NAME" = "push" ]; then + set +e + pr_json="$(gh pr list \ + --repo "$GITHUB_REPOSITORY" \ + --state open \ + --json number,headRefName,headRepository 2>/tmp/gh-pr-list.err)" + gh_status=$? + set -e + if [ "$gh_status" -ne 0 ]; then + echo "Warning: gh pr list failed; not skipping push tests" + cat /tmp/gh-pr-list.err || true + else + pr_count="$(printf '%s' "$pr_json" | jq \ + --arg branch "$BRANCH_NAME" \ + --arg repo "$GITHUB_REPOSITORY" \ + '[.[] | select(.headRefName == $branch and .headRepository.nameWithOwner == $repo)] | length')" + if [ "${pr_count:-0}" -gt 0 ]; then + echo "Skipping push test suite; open PR already covers branch $BRANCH_NAME in $GITHUB_REPOSITORY" + echo "capgo=false" >> "$GITHUB_OUTPUT" + echo "cli=false" >> "$GITHUB_OUTPUT" + echo "notifications=false" >> "$GITHUB_OUTPUT" + exit 0 + fi + fi + fi + run_capgo=false run_cli=false run_notifications=false @@ -431,6 +462,7 @@ jobs: warm_get '/organization?orgId=00000000-0000-0000-0000-000000000000' warm_post /organization warm_post /organization/members + warm_post /apikey - name: Run backend integration tests env: VITEST_SHARD: ${{ matrix.shard }} @@ -1213,7 +1245,8 @@ jobs: if: needs.changes.outputs.run_cli == 'true' name: CLI POSIX paths (${{ matrix.os }}) runs-on: ${{ matrix.os }} - timeout-minutes: 5 + # windows-2022 often needs >5m for bun install + CLI zip under shared runners. + timeout-minutes: 10 strategy: matrix: os: [windows-2025, windows-2022] diff --git a/playwright/e2e/compatibility-events.spec.ts b/playwright/e2e/compatibility-events.spec.ts index f781ece260..8c30663650 100644 --- a/playwright/e2e/compatibility-events.spec.ts +++ b/playwright/e2e/compatibility-events.spec.ts @@ -6,6 +6,7 @@ test.use({ screenshot: 'off', trace: 'off', video: 'off' }) // The seeded demo app owned by the `test@capgo.app` user (see supabase/seed.sql). // Sibling specs implicitly rely on this same login, so we reuse its demo app id. const APP_ID = 'com.demo.app' +const TEST_USER_ID = '6aa76066-55ef-4238-ade6-0b32334a4097' // A single unresolved, incompatible event used across the history + accept flows. // `id` is the PostgREST primary key the accept RPC is called with; the bundle ids @@ -107,6 +108,10 @@ async function mockStoreReleaseValidationStatus(page: Page) { test.describe('Compatibility events', () => { test.beforeEach(async ({ page }) => { await page.login('test@capgo.app', 'testtest') + // Keep the support-usernames prompt from covering store-release alert actions. + await page.evaluate((userId) => { + localStorage.setItem(`capgo.supportUsernames.dismissed.${userId}`, '1') + }, TEST_USER_ID) }) test('shows the store release validation alert before opening the modal', async ({ page }) => { diff --git a/supabase/functions/_backend/triggers/canceled_org_retention_alerts.ts b/supabase/functions/_backend/triggers/canceled_org_retention_alerts.ts new file mode 100644 index 0000000000..cefc9fd8a1 --- /dev/null +++ b/supabase/functions/_backend/triggers/canceled_org_retention_alerts.ts @@ -0,0 +1,130 @@ +import type { MiddlewareKeyVariables } from '../utils/hono.ts' +import { Hono } from 'hono/tiny' +import { BRES, middlewareAPISecret, parseBody, simpleError } from '../utils/hono.ts' +import { cloudlog } from '../utils/logging.ts' +import { sendEventToTracking } from '../utils/tracking.ts' + +type RetentionAlertType = 'bundles_deletion_warning' | 'app_deletion_warning' + +interface CanceledOrgRetentionAlertPayload { + org_id: string + org_name?: string + management_email?: string + alert_type: RetentionAlertType + access_end?: string + days_until_deletion?: number + app_ids?: string[] +} + +const ALERT_CONFIG = { + bundles_deletion_warning: { + bentoEvent: 'org:bundles_will_be_deleted', + trackingEvent: 'Bundles will be deleted', + icon: '📦', + }, + app_deletion_warning: { + bentoEvent: 'org:apps_will_be_deleted', + trackingEvent: 'Apps will be deleted', + icon: '🗑️', + }, +} as const satisfies Record + +const ORG_ID_UUID_RE = /^[0-9a-f]{8}-[0-9a-f]{4}-[1-5][0-9a-f]{3}-[89ab][0-9a-f]{3}-[0-9a-f]{12}$/i + +function isRetentionAlertType(value: unknown): value is RetentionAlertType { + return value === 'bundles_deletion_warning' || value === 'app_deletion_warning' +} + +/** Cycle key matching SQL to_char(access_end AT TIME ZONE 'UTC', 'YYYY-MM-DD"T"HH24:MI:SS"Z"'). */ +function accessEndCycleKey(accessEnd: unknown): string { + if (typeof accessEnd !== 'string' || !accessEnd.trim()) + return 'unknown' + const parsed = new Date(accessEnd) + if (Number.isNaN(parsed.getTime())) + return 'invalid' + return parsed.toISOString().replace(/\.\d{3}Z$/, 'Z') +} + +export const app = new Hono() + +app.post('/', middlewareAPISecret, async (c) => { + const payload = await parseBody(c) + if (!payload || typeof payload !== 'object') { + throw simpleError('invalid_payload', 'Missing retention alert payload', { + hasPayload: false, + }) + } + + const orgId = typeof payload.org_id === 'string' ? payload.org_id.trim() : '' + const alertType = payload.alert_type + + if (!orgId || !ORG_ID_UUID_RE.test(orgId) || !isRetentionAlertType(alertType)) { + throw simpleError( + 'invalid_payload', + 'Missing or invalid org_id/alert_type in retention alert payload', + { + hasOrgId: orgId.length > 0, + orgIdValid: orgId.length > 0 && ORG_ID_UUID_RE.test(orgId), + alertType: typeof alertType === 'string' ? alertType : typeof alertType, + }, + ) + } + + const config = ALERT_CONFIG[alertType] + + const parsedDays = Number(payload.days_until_deletion ?? 5) + const daysUntilDeletion = Number.isFinite(parsedDays) ? parsedDays : 5 + const appIds = Array.isArray(payload.app_ids) ? payload.app_ids : [] + const accessEndKey = accessEndCycleKey(payload.access_end) + const uniqId = `retention:${alertType}:${accessEndKey}` + + const metadata = { + org_id: orgId, + org_name: payload.org_name ?? '', + management_email: payload.management_email ?? '', + alert_type: alertType, + access_end: typeof payload.access_end === 'string' ? payload.access_end : '', + days_until_deletion: daysUntilDeletion, + app_ids: appIds, + app_count: appIds.length, + } + + cloudlog({ + requestId: c.get('requestId'), + message: 'canceled org retention alert', + eventName: config.bentoEvent, + orgId, + alertType, + daysUntilDeletion, + appCount: appIds.length, + }) + + await sendEventToTracking(c, { + bento: { + once: true, + data: metadata, + event: config.bentoEvent, + preferenceKey: 'usage_limit', + uniqId, + audience: 'billing', + }, + channel: 'usage', + event: config.trackingEvent, + icon: config.icon, + user_id: orgId, + groups: { organization: orgId }, + notify: false, + sentToBento: true, + tags: { + alert_type: alertType, + days_until_deletion: String(daysUntilDeletion), + app_count: String(appIds.length), + }, + }, { background: false, strict: true }) + + return c.json(BRES) +}) diff --git a/supabase/functions/_backend/utils/org_email_notifications.ts b/supabase/functions/_backend/utils/org_email_notifications.ts index ff6ceacde7..23a5f8f375 100644 --- a/supabase/functions/_backend/utils/org_email_notifications.ts +++ b/supabase/functions/_backend/utils/org_email_notifications.ts @@ -683,6 +683,8 @@ export async function sendNotifToOrgMembersOnce( preferenceKey, orgId, }) + // Claim anyway so once-dedup (SQL + retries) does not re-queue forever. + await claimNotifOrgOnce(c, eventName, orgId, uniqId, writeClient) return false } diff --git a/supabase/functions/_backend/utils/tracking.ts b/supabase/functions/_backend/utils/tracking.ts index 2487fc5b1a..955e9be48a 100644 --- a/supabase/functions/_backend/utils/tracking.ts +++ b/supabase/functions/_backend/utils/tracking.ts @@ -107,7 +107,7 @@ async function executeTracking(c: Context, payload: SendEventToTrackingPayload, await Promise.all(tasks) } -async function executeBentoTracking(c: Context, payload: SendEventToTrackingPayload) { +async function executeBentoTracking(c: Context, payload: SendEventToTrackingPayload, strict = false) { if (!payload.sentToBento) return @@ -119,6 +119,8 @@ async function executeBentoTracking(c: Context, payload: SendEventToTrackingPayl event: payload.event, user_id: payload.user_id, }) + if (strict) + throw new Error('sendEventToTracking missing Bento payload') return } @@ -133,6 +135,8 @@ async function executeBentoTracking(c: Context, payload: SendEventToTrackingPayl event: payload.event, user_id: payload.user_id, }) + if (strict) + throw new Error('sendEventToTracking missing org id for Bento notification') return } @@ -142,6 +146,8 @@ async function executeBentoTracking(c: Context, payload: SendEventToTrackingPayl if (bento.once) { // Permanent per-(event, org, uniqId) claim: per-entity alerts (e.g. an // incompatible bundle version) must not re-email org admins on retries. + // Discard the boolean: false is often benign (already claimed, no + // recipients, Bento unset). Under strict, only thrown errors fail closed. await sendNotifToOrgMembersOnce( c, bento.event, @@ -152,35 +158,34 @@ async function executeBentoTracking(c: Context, payload: SendEventToTrackingPayl getDrizzleClient(pgClient), bento.audience, ) + return } - else { - await sendNotifToOrgMembers( - c, - bento.event, - bento.preferenceKey, - bento.data, - orgId, - bento.uniqId, - bento.cron ?? '* * * * *', - getDrizzleClient(pgClient), - bento.audience, - ) - } + await sendNotifToOrgMembers( + c, + bento.event, + bento.preferenceKey, + bento.data, + orgId, + bento.uniqId, + bento.cron ?? '* * * * *', + getDrizzleClient(pgClient), + bento.audience, + ) } finally { await pgClient.end() } - }) + }, strict) } export async function sendEventToTracking(c: Context, payload: SendEventToTrackingPayload, options: SendEventToTrackingOptions = {}) { const trackingTask = executeTracking(c, payload, options) if (options.background === false) { await trackingTask - await executeBentoTracking(c, payload) + await executeBentoTracking(c, payload, options.strict === true) return } await backgroundTask(c, trackingTask) - await backgroundTask(c, executeBentoTracking(c, payload)) + await backgroundTask(c, executeBentoTracking(c, payload, options.strict === true)) } diff --git a/supabase/functions/triggers/index.ts b/supabase/functions/triggers/index.ts index e4bf159f9f..83b7d48728 100644 --- a/supabase/functions/triggers/index.ts +++ b/supabase/functions/triggers/index.ts @@ -1,3 +1,4 @@ +import { app as canceled_org_retention_alerts } from '../_backend/triggers/canceled_org_retention_alerts.ts' import { app as credit_usage_alerts } from '../_backend/triggers/credit_usage_alerts.ts' import { app as credit_usage_posthog } from '../_backend/triggers/credit_usage_posthog.ts' import { app as cron_clean_orphan_images } from '../_backend/triggers/cron_clean_orphan_images.ts' @@ -77,6 +78,7 @@ appGlobal.route('/cron_clear_versions', cron_clear_versions) appGlobal.route('/cron_clean_orphan_images', cron_clean_orphan_images) appGlobal.route('/cron_reconcile_build_status', cron_reconcile_build_status) appGlobal.route('/cron_rollout_auto_pause', cron_rollout_auto_pause) +appGlobal.route('/canceled_org_retention_alerts', canceled_org_retention_alerts) appGlobal.route('/credit_usage_alerts', credit_usage_alerts) appGlobal.route('/credit_usage_posthog', credit_usage_posthog) appGlobal.route('/on_organization_delete', on_organization_delete) diff --git a/supabase/migrations/20260807182308_canceled_org_app_retention_lifecycle.sql b/supabase/migrations/20260807182308_canceled_org_app_retention_lifecycle.sql new file mode 100644 index 0000000000..a5e717dab5 --- /dev/null +++ b/supabase/migrations/20260807182308_canceled_org_app_retention_lifecycle.sql @@ -0,0 +1,409 @@ +-- Canceled-org retention lifecycle: +-- 85 days: warn that bundles will be deleted in 5 days (Bento / tracking) +-- 90 days: soft-delete bundles (existing) + warn that apps will be deleted in 5 days +-- 95 days: delete apps and archive app_id + creator email into old_apps +-- +-- Alerts are queued to pgmq and delivered by the canceled_org_retention_alerts +-- edge trigger. Dedup uses public.notifications once-claims so daily cron +-- catch-up does not re-queue forever. + +-- Permanent archive of apps removed by the 95-day unpaid retention path. +CREATE TABLE IF NOT EXISTS public.old_apps ( + id bigint GENERATED BY DEFAULT AS IDENTITY PRIMARY KEY, + app_id character varying NOT NULL, + email text NOT NULL, + owner_org uuid NOT NULL, + user_id uuid, + app_created_at timestamp with time zone, + deleted_at timestamp with time zone DEFAULT now() NOT NULL, + created_at timestamp with time zone DEFAULT now() NOT NULL, + CONSTRAINT old_apps_app_id_owner_org_key UNIQUE (app_id, owner_org) +); + +ALTER TABLE public.old_apps OWNER TO postgres; + +COMMENT ON TABLE public.old_apps IS + 'Archive of apps deleted after the unpaid/canceled retention window (95 days). Keeps app_id and creator email for support memory; not purged by delete_old_deleted_apps.'; + +COMMENT ON COLUMN public.old_apps.email IS + 'Email of the app creator (public.users via apps.user_id), falling back to org management_email.'; + +COMMENT ON COLUMN public.old_apps.app_created_at IS + 'Original apps.created_at at archival time.'; + +CREATE INDEX IF NOT EXISTS old_apps_email_idx + ON public.old_apps (email); + +CREATE INDEX IF NOT EXISTS old_apps_deleted_at_idx + ON public.old_apps (deleted_at); + +ALTER TABLE public.old_apps ENABLE ROW LEVEL SECURITY; + +DROP POLICY IF EXISTS "Deny all access on old_apps" ON public.old_apps; +CREATE POLICY "Deny all access on old_apps" + ON public.old_apps + FOR ALL + TO anon, authenticated + USING (false) + WITH CHECK (false); + +GRANT ALL ON TABLE public.old_apps TO service_role; +REVOKE ALL ON TABLE public.old_apps FROM PUBLIC; +REVOKE ALL ON TABLE public.old_apps FROM anon; +REVOKE ALL ON TABLE public.old_apps FROM authenticated; +GRANT ALL ON SEQUENCE public.old_apps_id_seq TO service_role; +REVOKE ALL ON SEQUENCE public.old_apps_id_seq FROM PUBLIC; +REVOKE ALL ON SEQUENCE public.old_apps_id_seq FROM anon; +REVOKE ALL ON SEQUENCE public.old_apps_id_seq FROM authenticated; + +-- Parameterized eligibility: canceled/deleted orgs past N days of end-of-access. +CREATE OR REPLACE FUNCTION public.canceled_org_ids_past_grace(p_days integer) +RETURNS SETOF uuid +LANGUAGE sql +STABLE +SECURITY DEFINER +SET search_path = '' +AS $$ + SELECT o.id + FROM public.stripe_info AS si + JOIN public.orgs AS o ON o.customer_id = si.customer_id + WHERE si.status IN ('canceled', 'deleted') + AND GREATEST(si.canceled_at, si.subscription_anchor_end, si.trial_at) + <= pg_catalog.now() - make_interval(days => GREATEST(0, COALESCE(p_days, 0))); +$$; + +ALTER FUNCTION public.canceled_org_ids_past_grace(integer) OWNER TO postgres; +REVOKE ALL ON FUNCTION public.canceled_org_ids_past_grace(integer) FROM PUBLIC; +REVOKE ALL ON FUNCTION public.canceled_org_ids_past_grace(integer) FROM anon; +REVOKE ALL ON FUNCTION public.canceled_org_ids_past_grace(integer) FROM authenticated; +GRANT ALL ON FUNCTION public.canceled_org_ids_past_grace(integer) TO service_role; + +COMMENT ON FUNCTION public.canceled_org_ids_past_grace(integer) IS + 'Org ids whose stripe_info is canceled/deleted and GREATEST(canceled_at, subscription_anchor_end, trial_at) is older than p_days.'; + +-- Keep the existing 90-day helper as a thin wrapper for call sites / tests. +CREATE OR REPLACE FUNCTION public.long_canceled_org_ids() +RETURNS SETOF uuid +LANGUAGE sql +STABLE +SECURITY DEFINER +SET search_path = '' +AS $$ + SELECT public.canceled_org_ids_past_grace(90); +$$; + +ALTER FUNCTION public.long_canceled_org_ids() OWNER TO postgres; +REVOKE ALL ON FUNCTION public.long_canceled_org_ids() FROM PUBLIC; +REVOKE ALL ON FUNCTION public.long_canceled_org_ids() FROM anon; +REVOKE ALL ON FUNCTION public.long_canceled_org_ids() FROM authenticated; +GRANT ALL ON FUNCTION public.long_canceled_org_ids() TO service_role; + +COMMENT ON FUNCTION public.long_canceled_org_ids() IS + 'Org ids whose stripe_info is canceled/deleted and GREATEST(canceled_at, subscription_anchor_end, trial_at) is older than 90 days.'; + +-- Queue retention warning events (85d bundles / 90d apps). +-- Window is [p_min_days, p_min_days+5) so warnings never fire in the same cron +-- pass as the next destructive step (90d soft-delete / 95d app delete). +-- Dedup: skip if already claimed in notifications OR still +-- pending in the queue. +CREATE OR REPLACE FUNCTION public.queue_canceled_org_retention_alerts( + p_alert_type text, + p_min_days integer, + p_batch_size integer DEFAULT 500 +) +RETURNS bigint +LANGUAGE plpgsql +SECURITY DEFINER +SET search_path = '' +AS $$ +DECLARE + v_batch_size integer := GREATEST(1, COALESCE(p_batch_size, 500)); + v_min_days integer := GREATEST(0, COALESCE(p_min_days, 0)); + v_max_days integer := v_min_days + 5; + v_event text; + v_queued bigint := 0; + org_record RECORD; + v_days_until integer; +BEGIN + IF p_alert_type = 'bundles_deletion_warning' THEN + v_event := 'org:bundles_will_be_deleted'; + ELSIF p_alert_type = 'app_deletion_warning' THEN + v_event := 'org:apps_will_be_deleted'; + ELSE + RAISE EXCEPTION 'unsupported retention alert type: %', p_alert_type; + END IF; + + FOR org_record IN + SELECT + o.id AS org_id, + o.name AS org_name, + o.management_email, + GREATEST(si.canceled_at, si.subscription_anchor_end, si.trial_at) AS access_end, + COALESCE( + ( + SELECT jsonb_agg(a.app_id ORDER BY a.app_id) + FROM public.apps AS a + WHERE a.owner_org = o.id + ), + '[]'::jsonb + ) AS app_ids + FROM public.stripe_info AS si + JOIN public.orgs AS o ON o.customer_id = si.customer_id + WHERE si.status IN ('canceled', 'deleted') + AND GREATEST(si.canceled_at, si.subscription_anchor_end, si.trial_at) + <= pg_catalog.now() - make_interval(days => v_min_days) + AND GREATEST(si.canceled_at, si.subscription_anchor_end, si.trial_at) + > pg_catalog.now() - make_interval(days => v_max_days) + AND ( + CASE + WHEN p_alert_type = 'bundles_deletion_warning' THEN EXISTS ( + SELECT 1 + FROM public.app_versions AS av + WHERE av.owner_org = o.id + AND av.deleted = false + AND av.name NOT IN ('builtin', 'unknown') + ) + ELSE EXISTS ( + SELECT 1 + FROM public.apps AS a + WHERE a.owner_org = o.id + ) + END + ) + AND NOT EXISTS ( + SELECT 1 + FROM public.notifications AS n + WHERE n.owner_org = o.id + AND n.event = v_event + AND n.uniq_id = ( + 'retention:' + || p_alert_type + || ':' + || to_char( + GREATEST(si.canceled_at, si.subscription_anchor_end, si.trial_at) AT TIME ZONE 'UTC', + 'YYYY-MM-DD"T"HH24:MI:SS"Z"' + ) + ) + ) + AND NOT EXISTS ( + SELECT 1 + FROM pgmq.q_canceled_org_retention_alerts AS q + WHERE q.message -> 'payload' ->> 'org_id' = o.id::text + AND q.message -> 'payload' ->> 'alert_type' = p_alert_type + ) + ORDER BY o.id + LIMIT v_batch_size + LOOP + v_days_until := GREATEST( + 0, + CEIL( + EXTRACT( + EPOCH FROM ( + org_record.access_end + + make_interval(days => v_max_days) + - pg_catalog.now() + ) + ) / 86400.0 + )::integer + ); + + PERFORM pgmq.send( + 'canceled_org_retention_alerts', + jsonb_build_object( + 'function_name', 'canceled_org_retention_alerts', + 'function_type', 'cloudflare', + 'payload', jsonb_build_object( + 'org_id', org_record.org_id, + 'org_name', org_record.org_name, + 'management_email', org_record.management_email, + 'alert_type', p_alert_type, + 'access_end', org_record.access_end, + 'days_until_deletion', v_days_until, + 'app_ids', org_record.app_ids + ) + ) + ); + v_queued := v_queued + 1; + END LOOP; + + IF v_queued > 0 THEN + RAISE NOTICE + 'queue_canceled_org_retention_alerts: type=% queued=% window=[%,%)', + p_alert_type, + v_queued, + v_min_days, + v_max_days; + END IF; + + RETURN v_queued; +END; +$$; + +ALTER FUNCTION public.queue_canceled_org_retention_alerts(text, integer, integer) OWNER TO postgres; +REVOKE ALL ON FUNCTION public.queue_canceled_org_retention_alerts(text, integer, integer) FROM PUBLIC; +REVOKE ALL ON FUNCTION public.queue_canceled_org_retention_alerts(text, integer, integer) FROM anon; +REVOKE ALL ON FUNCTION public.queue_canceled_org_retention_alerts(text, integer, integer) FROM authenticated; +GRANT ALL ON FUNCTION public.queue_canceled_org_retention_alerts(text, integer, integer) TO service_role; + +COMMENT ON FUNCTION public.queue_canceled_org_retention_alerts(text, integer, integer) IS + 'Queues once-per-cancel-cycle Bento/tracking warnings for canceled orgs ' + 'in [p_min_days, p_min_days+5). Bundles require a deletable app_versions row; ' + 'apps require an apps row. Dedup uniq_id uses access_end UTC timestamp.'; + +-- Hard-delete apps for orgs past the 95-day unpaid window; archive to old_apps first. +CREATE OR REPLACE FUNCTION public.delete_apps_for_long_canceled_orgs( + p_batch_size integer DEFAULT 500 +) +RETURNS bigint +LANGUAGE plpgsql +SECURITY DEFINER +SET search_path = '' +AS $$ +DECLARE + v_batch_size integer := GREATEST(1, COALESCE(p_batch_size, 500)); + deleted_count bigint := 0; +BEGIN + WITH candidates AS ( + SELECT + a.app_id, + a.owner_org, + a.user_id, + a.created_at AS app_created_at, + COALESCE(u.email::text, o.management_email) AS email + FROM public.apps AS a + JOIN public.orgs AS o ON o.id = a.owner_org + LEFT JOIN public.users AS u ON u.id = a.user_id + WHERE a.owner_org IN (SELECT public.canceled_org_ids_past_grace(95)) + ORDER BY a.app_id + LIMIT v_batch_size + ), + archived AS ( + INSERT INTO public.old_apps (app_id, email, owner_org, user_id, app_created_at) + SELECT + c.app_id, + c.email, + c.owner_org, + c.user_id, + c.app_created_at + FROM candidates AS c + ON CONFLICT (app_id, owner_org) DO UPDATE + SET + email = EXCLUDED.email, + user_id = EXCLUDED.user_id, + app_created_at = EXCLUDED.app_created_at, + deleted_at = pg_catalog.now() + RETURNING old_apps.app_id + ), + deleted AS ( + DELETE FROM public.apps AS a + USING candidates AS c + WHERE a.app_id = c.app_id + AND a.owner_org = c.owner_org + RETURNING a.app_id + ) + SELECT COUNT(*) INTO deleted_count FROM deleted; + + IF deleted_count > 0 THEN + RAISE NOTICE + 'delete_apps_for_long_canceled_orgs: deleted_apps=%', + deleted_count; + END IF; + + RETURN deleted_count; +END; +$$; + +ALTER FUNCTION public.delete_apps_for_long_canceled_orgs(integer) OWNER TO postgres; +REVOKE ALL ON FUNCTION public.delete_apps_for_long_canceled_orgs(integer) FROM PUBLIC; +REVOKE ALL ON FUNCTION public.delete_apps_for_long_canceled_orgs(integer) FROM anon; +REVOKE ALL ON FUNCTION public.delete_apps_for_long_canceled_orgs(integer) FROM authenticated; +GRANT ALL ON FUNCTION public.delete_apps_for_long_canceled_orgs(integer) TO service_role; + +COMMENT ON FUNCTION public.delete_apps_for_long_canceled_orgs(integer) IS + 'Deletes apps for orgs past the 95-day canceled grace window after archiving app_id + creator email into old_apps. Bounded by p_batch_size.'; + +-- Daily retention orchestrator: warnings + version soft-delete + app delete + audit purge. +CREATE OR REPLACE FUNCTION public.cleanup_long_canceled_org_data() +RETURNS void +LANGUAGE plpgsql +SECURITY DEFINER +SET search_path = '' +AS $$ +BEGIN + -- 85 days: warn that bundles will be deleted in 5 days. + PERFORM public.queue_canceled_org_retention_alerts('bundles_deletion_warning', 85); + -- 90 days: soft-delete bundles for long-canceled orgs. + PERFORM public.soft_delete_versions_for_long_canceled_orgs(); + -- 90 days: warn that apps will be deleted in 5 days. + PERFORM public.queue_canceled_org_retention_alerts('app_deletion_warning', 90); + -- 95 days: delete apps and keep app_id + creator email in old_apps. + PERFORM public.delete_apps_for_long_canceled_orgs(); + -- Purge audit logs for orgs past the 90-day grace window. + PERFORM public.cleanup_audit_logs_for_long_canceled_orgs(); +END; +$$; + +ALTER FUNCTION public.cleanup_long_canceled_org_data() OWNER TO postgres; +REVOKE ALL ON FUNCTION public.cleanup_long_canceled_org_data() FROM PUBLIC; +REVOKE ALL ON FUNCTION public.cleanup_long_canceled_org_data() FROM anon; +REVOKE ALL ON FUNCTION public.cleanup_long_canceled_org_data() FROM authenticated; +GRANT ALL ON FUNCTION public.cleanup_long_canceled_org_data() TO service_role; + +COMMENT ON FUNCTION public.cleanup_long_canceled_org_data() IS + 'Daily canceled-org retention: 85d bundle warning, 90d soft-delete versions + app warning, 95d delete apps into old_apps, then purge audit_logs. Steps are batch/runtime bounded.'; + +-- Alert delivery queue + high-frequency drain registration. +DO $$ +BEGIN + IF NOT EXISTS ( + SELECT 1 + FROM pgmq.list_queues() + WHERE queue_name = 'canceled_org_retention_alerts' + ) THEN + PERFORM pgmq.create('canceled_org_retention_alerts'); + END IF; +END; +$$; + +DO $$ +DECLARE + high_frequency_task_type public.cron_task_type; + high_frequency_target jsonb; +BEGIN + SELECT cron.task_type, cron.target::jsonb + INTO high_frequency_task_type, high_frequency_target + FROM public.cron_tasks AS cron + WHERE cron.name = 'high_frequency_queues' + FOR UPDATE; + + IF NOT FOUND THEN + RAISE EXCEPTION 'Required cron task high_frequency_queues is missing'; + END IF; + + IF high_frequency_task_type + IS DISTINCT FROM 'function_queue'::public.cron_task_type THEN + RAISE EXCEPTION 'Cron task high_frequency_queues must use task type function_queue'; + END IF; + + IF pg_catalog.jsonb_typeof(high_frequency_target) + IS DISTINCT FROM 'array' THEN + RAISE EXCEPTION 'Cron task high_frequency_queues target must be a JSON array'; + END IF; + + IF NOT (high_frequency_target ? 'canceled_org_retention_alerts') THEN + UPDATE public.cron_tasks + SET + target = ( + high_frequency_target || '["canceled_org_retention_alerts"]'::jsonb + )::text, + updated_at = pg_catalog.now() + WHERE name = 'high_frequency_queues'; + END IF; +END; +$$; + +UPDATE public.cron_tasks +SET + description = 'Canceled-org retention: 85d bundle warn, 90d soft-delete + app warn, 95d delete apps into old_apps, purge audit logs', + updated_at = now() +WHERE name = 'canceled_org_version_cleanup'; diff --git a/supabase/tests/64_test_canceled_org_version_cleanup.sql b/supabase/tests/64_test_canceled_org_version_cleanup.sql index ecb3e81878..b1e867296e 100644 --- a/supabase/tests/64_test_canceled_org_version_cleanup.sql +++ b/supabase/tests/64_test_canceled_org_version_cleanup.sql @@ -1,6 +1,55 @@ BEGIN; -SELECT plan(25); +SELECT plan(39); + +-- pgmq schema is not granted to service_role; use postgres-owned helpers. +CREATE OR REPLACE FUNCTION pg_temp.delete_canceled_org_retention_alerts( + p_org_ids text[] +) +RETURNS void +LANGUAGE plpgsql +SECURITY DEFINER +SET search_path = '' +AS $$ +BEGIN + IF pg_catalog.to_regclass('pgmq.q_canceled_org_retention_alerts') IS NULL THEN + RETURN; + END IF; + + EXECUTE + 'DELETE FROM pgmq.q_canceled_org_retention_alerts + WHERE message -> ''payload'' ->> ''org_id'' = ANY($1)' + USING p_org_ids; +END; +$$; + +CREATE OR REPLACE FUNCTION pg_temp.count_canceled_org_retention_alerts( + p_alert_type text, + p_org_ids uuid[] +) +RETURNS integer +LANGUAGE plpgsql +SECURITY DEFINER +SET search_path = '' +AS $$ +DECLARE + v_count integer := 0; +BEGIN + IF pg_catalog.to_regclass('pgmq.q_canceled_org_retention_alerts') IS NULL THEN + RETURN 0; + END IF; + + EXECUTE + 'SELECT count(*)::integer + FROM pgmq.q_canceled_org_retention_alerts + WHERE ($1 IS NULL OR message -> ''payload'' ->> ''alert_type'' = $1) + AND (message -> ''payload'' ->> ''org_id'')::uuid = ANY($2)' + INTO v_count + USING p_alert_type, p_org_ids; + + RETURN COALESCE(v_count, 0); +END; +$$; SELECT tests.authenticate_as_service_role(); @@ -9,6 +58,11 @@ SELECT ok( 'long_canceled_org_ids exists' ); +SELECT ok( + to_regprocedure('public.canceled_org_ids_past_grace(integer)') IS NOT NULL, + 'canceled_org_ids_past_grace exists' +); + SELECT ok( to_regprocedure('public.soft_delete_versions_for_long_canceled_orgs(integer)') IS NOT NULL, 'soft_delete_versions_for_long_canceled_orgs exists' @@ -21,11 +75,26 @@ SELECT ok( 'cleanup_audit_logs_for_long_canceled_orgs exists' ); +SELECT ok( + to_regprocedure('public.queue_canceled_org_retention_alerts(text, integer, integer)') IS NOT NULL, + 'queue_canceled_org_retention_alerts exists' +); + +SELECT ok( + to_regprocedure('public.delete_apps_for_long_canceled_orgs(integer)') IS NOT NULL, + 'delete_apps_for_long_canceled_orgs exists' +); + SELECT ok( to_regprocedure('public.cleanup_long_canceled_org_data()') IS NOT NULL, 'cleanup_long_canceled_org_data exists' ); +SELECT ok( + to_regclass('public.old_apps') IS NOT NULL, + 'old_apps table exists' +); + SELECT is( has_function_privilege( 'anon', @@ -71,6 +140,9 @@ SELECT ok( ); -- Dedicated fixtures (unique customer/app ids for parallel safety) +-- long = 92d (past 90, before 95) for version soft-delete without app delete +-- warn85 = 87d for bundle-deletion warning queue +-- ultra = 100d (past 95) for app delete + old_apps archive CREATE TEMP TABLE canceled_cleanup_ctx AS SELECT 'a0c1e2f3-1111-4aaa-8bbb-000000000001'::uuid AS long_canceled_org, @@ -78,16 +150,22 @@ SELECT 'a0c1e2f3-1111-4aaa-8bbb-000000000003'::uuid AS trial_org, 'a0c1e2f3-1111-4aaa-8bbb-000000000004'::uuid AS paying_org, 'a0c1e2f3-1111-4aaa-8bbb-000000000005'::uuid AS early_cancel_org, + 'a0c1e2f3-1111-4aaa-8bbb-000000000006'::uuid AS warn85_org, + 'a0c1e2f3-1111-4aaa-8bbb-000000000007'::uuid AS ultra_canceled_org, 'cus_canceled_cleanup_long'::varchar AS long_customer, 'cus_canceled_cleanup_recent'::varchar AS recent_customer, 'cus_canceled_cleanup_trial'::varchar AS trial_customer, 'cus_canceled_cleanup_paying'::varchar AS paying_customer, 'cus_canceled_cleanup_early'::varchar AS early_customer, + 'cus_canceled_cleanup_warn85'::varchar AS warn85_customer, + 'cus_canceled_cleanup_ultra'::varchar AS ultra_customer, 'com.test.canceled.cleanup.long'::varchar AS long_app, 'com.test.canceled.cleanup.recent'::varchar AS recent_app, 'com.test.canceled.cleanup.trial'::varchar AS trial_app, 'com.test.canceled.cleanup.paying'::varchar AS paying_app, 'com.test.canceled.cleanup.early'::varchar AS early_app, + 'com.test.canceled.cleanup.warn85'::varchar AS warn85_app, + 'com.test.canceled.cleanup.ultra'::varchar AS ultra_app, '6aa76066-55ef-4238-ade6-0b32334a4097'::uuid AS user_id, 'prod_LQIregjtNduh4q'::varchar AS product_id; @@ -108,8 +186,8 @@ SELECT now() - interval '200 days', false, now() - interval '200 days', - now() - interval '100 days', - now() - interval '100 days' + now() - interval '92 days', + now() - interval '92 days' FROM canceled_cleanup_ctx UNION ALL SELECT @@ -155,6 +233,30 @@ SELECT now() - interval '40 days', now() - interval '10 days', now() - interval '100 days' +FROM canceled_cleanup_ctx +UNION ALL +-- 87 days: past 85 bundle warning, before 90 soft-delete. +SELECT + warn85_customer, + 'canceled'::public.stripe_status, + product_id, + now() - interval '200 days', + false, + now() - interval '200 days', + now() - interval '87 days', + now() - interval '87 days' +FROM canceled_cleanup_ctx +UNION ALL +-- 100 days: past 95 app-delete window. +SELECT + ultra_customer, + 'canceled'::public.stripe_status, + product_id, + now() - interval '200 days', + false, + now() - interval '200 days', + now() - interval '100 days', + now() - interval '100 days' FROM canceled_cleanup_ctx; INSERT INTO public.orgs (id, created_by, name, management_email, customer_id) @@ -171,6 +273,12 @@ SELECT paying_org, user_id, 'Paying Cleanup Org', 'canceled-paying@test.local', FROM canceled_cleanup_ctx UNION ALL SELECT early_cancel_org, user_id, 'Early Cancel Cleanup Org', 'canceled-early@test.local', early_customer +FROM canceled_cleanup_ctx +UNION ALL +SELECT warn85_org, user_id, 'Warn85 Canceled Cleanup Org', 'canceled-warn85@test.local', warn85_customer +FROM canceled_cleanup_ctx +UNION ALL +SELECT ultra_canceled_org, user_id, 'Ultra Canceled Cleanup Org', 'canceled-ultra@test.local', ultra_customer FROM canceled_cleanup_ctx; INSERT INTO public.apps (app_id, icon_url, owner_org, name, user_id) @@ -182,7 +290,11 @@ SELECT trial_app, '', trial_org, 'Trial App', user_id FROM canceled_cleanup_ctx UNION ALL SELECT paying_app, '', paying_org, 'Paying App', user_id FROM canceled_cleanup_ctx UNION ALL -SELECT early_app, '', early_cancel_org, 'Early Cancel App', user_id FROM canceled_cleanup_ctx; +SELECT early_app, '', early_cancel_org, 'Early Cancel App', user_id FROM canceled_cleanup_ctx +UNION ALL +SELECT warn85_app, '', warn85_org, 'Warn85 App', user_id FROM canceled_cleanup_ctx +UNION ALL +SELECT ultra_app, '', ultra_canceled_org, 'Ultra Canceled App', user_id FROM canceled_cleanup_ctx; INSERT INTO public.app_versions (id, app_id, name, storage_provider, owner_org, user_id, deleted) SELECT 970101, long_app, '1.0.0', 'r2', long_canceled_org, user_id, false FROM canceled_cleanup_ctx @@ -197,7 +309,11 @@ SELECT 970301, trial_app, '1.0.0', 'r2', trial_org, user_id, false FROM canceled UNION ALL SELECT 970401, paying_app, '1.0.0', 'r2', paying_org, user_id, false FROM canceled_cleanup_ctx UNION ALL -SELECT 970501, early_app, '1.0.0', 'r2', early_cancel_org, user_id, false FROM canceled_cleanup_ctx; +SELECT 970501, early_app, '1.0.0', 'r2', early_cancel_org, user_id, false FROM canceled_cleanup_ctx +UNION ALL +SELECT 970601, warn85_app, '1.0.0', 'r2', warn85_org, user_id, false FROM canceled_cleanup_ctx +UNION ALL +SELECT 970701, ultra_app, '1.0.0', 'r2', ultra_canceled_org, user_id, false FROM canceled_cleanup_ctx; SELECT set_config('capgo.seed_channel_targets', 'true', true); @@ -349,7 +465,14 @@ SELECT is( 'process_free_trial_expired leaves paying orgs unchanged' ); --- Full canceled-org cleanup (versions + audit logs) +-- Full canceled-org cleanup (warnings + versions + app delete + audit logs) +SELECT pg_temp.delete_canceled_org_retention_alerts(ARRAY[ + (SELECT long_canceled_org::text FROM canceled_cleanup_ctx), + (SELECT warn85_org::text FROM canceled_cleanup_ctx), + (SELECT ultra_canceled_org::text FROM canceled_cleanup_ctx), + (SELECT recent_canceled_org::text FROM canceled_cleanup_ctx) +]); + SELECT public.cleanup_long_canceled_org_data(); SELECT is( @@ -416,6 +539,146 @@ SELECT is( 'early cancel still inside period-end grace keeps versions' ); +SELECT is( + (SELECT deleted FROM public.app_versions WHERE id = 970601), + false, + '87-day org versions are kept until the 90-day soft-delete window' +); + +SELECT ok( + EXISTS ( + SELECT 1 + FROM public.apps + WHERE app_id = (SELECT long_app FROM canceled_cleanup_ctx) + ), + '92-day org app is kept until the 95-day app-delete window' +); + +SELECT ok( + NOT EXISTS ( + SELECT 1 + FROM public.apps + WHERE app_id = (SELECT ultra_app FROM canceled_cleanup_ctx) + ), + '100-day org app is deleted after 95 days' +); + +SELECT ok( + EXISTS ( + SELECT 1 + FROM public.old_apps + WHERE app_id = (SELECT ultra_app FROM canceled_cleanup_ctx) + AND owner_org = (SELECT ultra_canceled_org FROM canceled_cleanup_ctx) + AND email = 'test@capgo.app' + ), + 'deleted 95-day app is archived into old_apps with creator email' +); + +SELECT is( + pg_temp.count_canceled_org_retention_alerts( + 'bundles_deletion_warning', + ARRAY[ + (SELECT warn85_org FROM canceled_cleanup_ctx), + (SELECT long_canceled_org FROM canceled_cleanup_ctx), + (SELECT ultra_canceled_org FROM canceled_cleanup_ctx) + ] + ), + 1, + 'queues one 85-day bundle warning for warn85 only' +); + +SELECT is( + pg_temp.count_canceled_org_retention_alerts( + 'app_deletion_warning', + ARRAY[ + (SELECT long_canceled_org FROM canceled_cleanup_ctx), + (SELECT ultra_canceled_org FROM canceled_cleanup_ctx) + ] + ), + 1, + 'queues one 90-day app warning for long only' +); + +SELECT is( + pg_temp.count_canceled_org_retention_alerts( + NULL, + ARRAY[(SELECT recent_canceled_org FROM canceled_cleanup_ctx)] + ), + 0, + 'recently canceled orgs do not get retention deletion warnings' +); + +SELECT public.cleanup_long_canceled_org_data(); + +SELECT is( + pg_temp.count_canceled_org_retention_alerts( + NULL, + ARRAY[ + (SELECT warn85_org FROM canceled_cleanup_ctx), + (SELECT long_canceled_org FROM canceled_cleanup_ctx), + (SELECT ultra_canceled_org FROM canceled_cleanup_ctx), + (SELECT recent_canceled_org FROM canceled_cleanup_ctx) + ] + ), + 2, + 'second cleanup does not re-queue pending retention warnings' +); + +-- Drain pending messages, claim via notifications, assert claim-based dedup. +SELECT pg_temp.delete_canceled_org_retention_alerts(ARRAY[ + (SELECT warn85_org::text FROM canceled_cleanup_ctx), + (SELECT long_canceled_org::text FROM canceled_cleanup_ctx) +]); + +INSERT INTO public.notifications (owner_org, event, uniq_id) +SELECT + warn85_org, + 'org:bundles_will_be_deleted', + 'retention:bundles_deletion_warning:' + || to_char( + (SELECT GREATEST(si.canceled_at, si.subscription_anchor_end, si.trial_at) + FROM public.stripe_info AS si + WHERE si.customer_id = warn85_customer) AT TIME ZONE 'UTC', + 'YYYY-MM-DD"T"HH24:MI:SS"Z"' + ) +FROM canceled_cleanup_ctx; + +INSERT INTO public.notifications (owner_org, event, uniq_id) +SELECT + long_canceled_org, + 'org:apps_will_be_deleted', + 'retention:app_deletion_warning:' + || to_char( + (SELECT GREATEST(si.canceled_at, si.subscription_anchor_end, si.trial_at) + FROM public.stripe_info AS si + WHERE si.customer_id = long_customer) AT TIME ZONE 'UTC', + 'YYYY-MM-DD"T"HH24:MI:SS"Z"' + ) +FROM canceled_cleanup_ctx; + +SELECT public.cleanup_long_canceled_org_data(); + +SELECT is( + pg_temp.count_canceled_org_retention_alerts( + NULL, + ARRAY[ + (SELECT warn85_org FROM canceled_cleanup_ctx), + (SELECT long_canceled_org FROM canceled_cleanup_ctx) + ] + ), + 0, + 'notifications claim prevents re-queue after pending messages are drained' +); + +SELECT ok( + ( + SELECT cron.target::jsonb ? 'canceled_org_retention_alerts' + FROM public.cron_tasks AS cron + WHERE cron.name = 'high_frequency_queues' + ), + 'high_frequency_queues drains canceled_org_retention_alerts' +); + SELECT is( ( SELECT count(*)::int diff --git a/supabase/tests/66_test_on_user_org_access_queue.sql b/supabase/tests/66_test_on_user_org_access_queue.sql index cf5a45baf2..b894e383a4 100644 --- a/supabase/tests/66_test_on_user_org_access_queue.sql +++ b/supabase/tests/66_test_on_user_org_access_queue.sql @@ -124,9 +124,10 @@ SELECT is( "webhook_dispatcher", "webhook_delivery", "credit_usage_posthog", - "on_user_org_access" + "on_user_org_access", + "canceled_org_retention_alerts" ]'::jsonb, - 'high-frequency queues retain order and append on_user_org_access' + 'high-frequency queues retain order and append on_user_org_access then canceled_org_retention_alerts' ); SELECT is( diff --git a/tests/apikeys-expiration.test.ts b/tests/apikeys-expiration.test.ts index df46ed8749..f4c49809ea 100644 --- a/tests/apikeys-expiration.test.ts +++ b/tests/apikeys-expiration.test.ts @@ -3,7 +3,7 @@ import { randomUUID } from 'node:crypto' import { env } from 'node:process' import { createClient } from '@supabase/supabase-js' import { afterAll, beforeAll, describe, expect, it } from 'vitest' -import { appApiKeyBindings, BASE_URL, createDirectApiKeyWithBindings, executeSQL, fetchTestRequest, getAuthHeadersForCredentials, getSupabaseClient, normalizeLocalhostUrl, orgApiKeyBindings, resetAndSeedAppData, resetAppData, TEST_EMAIL, USER_EMAIL_APIKEY_EXPIRATION, USER_ID_APIKEY_EXPIRATION, USER_PASSWORD } from './test-utils.ts' +import { appApiKeyBindings, BASE_URL, createDirectApiKeyWithBindings, executeSQL, fetchTestRequest, getAuthHeadersForCredentials, getSupabaseClient, normalizeLocalhostUrl, orgApiKeyBindings, resetAndSeedAppData, resetAppData, TEST_EMAIL, USER_EMAIL_APIKEY_EXPIRATION, USER_ID_APIKEY_EXPIRATION, USER_PASSWORD, warmEdgeEndpoint } from './test-utils.ts' const id = randomUUID() const BASE_ORG_ID = randomUUID() @@ -271,6 +271,13 @@ beforeAll(async () => { throw policyAppError ?? new Error('Missing policy app') policyAppUuid = policyApp.id + + // Load the apikey isolate before concurrent POSTs from this file. + await warmEdgeEndpoint('/apikey', { + method: 'POST', + headers: authHeaders, + body: JSON.stringify({ name: keyName('warm-apikey-isolate') }), + }) }) afterAll(async () => { diff --git a/tests/canceled-org-retention-alerts.unit.test.ts b/tests/canceled-org-retention-alerts.unit.test.ts new file mode 100644 index 0000000000..777b788f04 --- /dev/null +++ b/tests/canceled-org-retention-alerts.unit.test.ts @@ -0,0 +1,205 @@ +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' + +const { + cloudlogMock, + sendEventToTrackingMock, +} = vi.hoisted(() => ({ + cloudlogMock: vi.fn(), + sendEventToTrackingMock: vi.fn(async () => undefined), +})) + +vi.mock('../supabase/functions/_backend/utils/logging.ts', () => ({ + cloudlog: cloudlogMock, +})) + +vi.mock('../supabase/functions/_backend/utils/tracking.ts', () => ({ + sendEventToTracking: sendEventToTrackingMock, +})) + +vi.mock('../supabase/functions/_backend/utils/hono.ts', async () => { + const actual = await vi.importActual('../supabase/functions/_backend/utils/hono.ts') + return { + ...actual, + middlewareAPISecret: async (_c: unknown, next: () => Promise) => await next(), + } +}) + +const ORG_BUNDLES = 'a0c1e2f3-1111-4aaa-8bbb-000000000101' +const ORG_APPS = 'a0c1e2f3-1111-4aaa-8bbb-000000000102' +const ORG_BAD = 'a0c1e2f3-1111-4aaa-8bbb-000000000103' +const ORG_PROTO = 'a0c1e2f3-1111-4aaa-8bbb-000000000104' +const ORG_FALLBACK = 'a0c1e2f3-1111-4aaa-8bbb-000000000105' + +async function postRetentionAlert(body: unknown) { + const { app } = await import('../supabase/functions/_backend/triggers/canceled_org_retention_alerts.ts') + return app.request('http://localhost/', { + method: 'POST', + headers: { 'content-type': 'application/json' }, + body: JSON.stringify(body), + }) +} + +describe('canceled_org_retention_alerts', () => { + beforeEach(() => { + vi.clearAllMocks() + sendEventToTrackingMock.mockResolvedValue(undefined) + }) + + afterEach(() => { + vi.resetModules() + }) + + it('sends a once Bento/tracking event for bundle deletion warnings', async () => { + const response = await postRetentionAlert({ + org_id: ORG_BUNDLES, + org_name: 'Retention Org', + management_email: 'billing@example.com', + alert_type: 'bundles_deletion_warning', + access_end: '2026-05-01T00:00:00.000Z', + days_until_deletion: 5, + app_ids: ['com.example.app'], + }) + + expect(response.status).toBe(200) + expect(await response.json()).toEqual({ status: 'ok' }) + expect(sendEventToTrackingMock).toHaveBeenCalledWith( + expect.anything(), + expect.objectContaining({ + channel: 'usage', + event: 'Bundles will be deleted', + sentToBento: true, + user_id: ORG_BUNDLES, + groups: { organization: ORG_BUNDLES }, + bento: expect.objectContaining({ + once: true, + event: 'org:bundles_will_be_deleted', + preferenceKey: 'usage_limit', + audience: 'billing', + uniqId: 'retention:bundles_deletion_warning:2026-05-01T00:00:00Z', + data: expect.objectContaining({ + org_id: ORG_BUNDLES, + app_ids: ['com.example.app'], + days_until_deletion: 5, + }), + }), + }), + { background: false, strict: true }, + ) + }) + + it('sends a once Bento/tracking event for app deletion warnings', async () => { + const response = await postRetentionAlert({ + org_id: ORG_APPS, + org_name: 'Retention Org', + alert_type: 'app_deletion_warning', + access_end: '2026-04-20T12:30:00.000Z', + app_ids: ['com.example.one', 'com.example.two'], + }) + + expect(response.status).toBe(200) + expect(sendEventToTrackingMock).toHaveBeenCalledWith( + expect.anything(), + expect.objectContaining({ + event: 'Apps will be deleted', + bento: expect.objectContaining({ + once: true, + event: 'org:apps_will_be_deleted', + uniqId: 'retention:app_deletion_warning:2026-04-20T12:30:00Z', + data: expect.objectContaining({ + app_count: 2, + }), + }), + }), + { background: false, strict: true }, + ) + }) + + it('rejects unsupported alert types', async () => { + const response = await postRetentionAlert({ + org_id: ORG_BAD, + alert_type: 'not_a_real_alert', + }) + + expect(response.status).toBeGreaterThanOrEqual(400) + expect(sendEventToTrackingMock).not.toHaveBeenCalled() + }) + + it.each(['constructor', 'toString', 'valueOf'])('rejects prototype key %s as alert type', async (alertType) => { + const response = await postRetentionAlert({ + org_id: ORG_PROTO, + alert_type: alertType, + }) + + expect(response.status).toBeGreaterThanOrEqual(400) + expect(sendEventToTrackingMock).not.toHaveBeenCalled() + }) + + it('rejects missing org_id', async () => { + const response = await postRetentionAlert({ + alert_type: 'bundles_deletion_warning', + }) + + expect(response.status).toBeGreaterThanOrEqual(400) + expect(sendEventToTrackingMock).not.toHaveBeenCalled() + }) + + it('rejects null JSON body without 500', async () => { + const response = await postRetentionAlert(null) + + expect(response.status).toBeGreaterThanOrEqual(400) + expect(response.status).toBeLessThan(500) + expect(sendEventToTrackingMock).not.toHaveBeenCalled() + }) + + it('rejects non-uuid org_id', async () => { + const response = await postRetentionAlert({ + org_id: 'not-a-uuid', + alert_type: 'bundles_deletion_warning', + }) + + expect(response.status).toBeGreaterThanOrEqual(400) + expect(sendEventToTrackingMock).not.toHaveBeenCalled() + }) + + it('falls back access_end uniqId and NaN days_until_deletion', async () => { + const response = await postRetentionAlert({ + org_id: ORG_FALLBACK, + alert_type: 'app_deletion_warning', + access_end: 'not-a-date', + days_until_deletion: 'nope', + }) + + expect(response.status).toBe(200) + expect(sendEventToTrackingMock).toHaveBeenCalledWith( + expect.anything(), + expect.objectContaining({ + bento: expect.objectContaining({ + uniqId: 'retention:app_deletion_warning:invalid', + data: expect.objectContaining({ + days_until_deletion: 5, + }), + }), + }), + { background: false, strict: true }, + ) + }) + + it('falls back to unknown uniqId for non-string access_end', async () => { + const response = await postRetentionAlert({ + org_id: ORG_FALLBACK, + alert_type: 'bundles_deletion_warning', + access_end: { nested: true }, + }) + + expect(response.status).toBe(200) + expect(sendEventToTrackingMock).toHaveBeenCalledWith( + expect.anything(), + expect.objectContaining({ + bento: expect.objectContaining({ + uniqId: 'retention:bundles_deletion_warning:unknown', + }), + }), + { background: false, strict: true }, + ) + }) +}) diff --git a/tests/security-definer-execute-hardening.test.ts b/tests/security-definer-execute-hardening.test.ts index 846143eea0..57fd760ec3 100644 --- a/tests/security-definer-execute-hardening.test.ts +++ b/tests/security-definer-execute-hardening.test.ts @@ -29,10 +29,13 @@ const SERVICE_ONLY_PROCS = [ 'public.cleanup_onboarding_app_data_on_complete()', 'public.delete_old_deleted_versions()', 'public.long_canceled_org_ids()', + 'public.canceled_org_ids_past_grace(integer)', 'public.process_free_trial_expired()', 'public.soft_delete_versions_for_long_canceled_orgs(integer)', 'public.cleanup_audit_logs_for_long_canceled_orgs(integer, integer, integer)', 'public.cleanup_long_canceled_org_data()', + 'public.queue_canceled_org_retention_alerts(text, integer, integer)', + 'public.delete_apps_for_long_canceled_orgs(integer)', 'public.enqueue_credit_usage_posthog_event()', 'public.generate_org_user_stripe_info_on_org_create()', 'public.get_apikey()', diff --git a/tests/test-utils.ts b/tests/test-utils.ts index f938f76087..24d8595c22 100644 --- a/tests/test-utils.ts +++ b/tests/test-utils.ts @@ -487,19 +487,37 @@ export const headersInternal = { 'apisecret': API_SECRET, } +/** Kong proxy body when the Deno isolate dies mid-request under shard load. */ +const KONG_UPSTREAM_INVALID_RESPONSE = 'An invalid response was received from the upstream server' + /** - * Send one request. A transient failure is test evidence, not a reason to rerun it. + * Send one request. Application 4xx/5xx are test evidence and are not retried. + * Only Kong's upstream-invalid 502/503 (isolate crash/reload) is retried — same + * signal the CI warm step already treats as non-ready. */ export async function fetchTestRequest( url: string, options?: RequestInit, ): Promise { - const response = await fetch(url, options) - if (response.status === 502 || response.status === 503) { + const maxAttempts = 3 + let lastResponse: Response | undefined + for (let attempt = 1; attempt <= maxAttempts; attempt++) { + const response = await fetch(url, options) + lastResponse = response + if (response.status !== 502 && response.status !== 503) + return response + const body = await response.clone().text().catch(() => '') - console.error(`[fetchTestRequest] gateway status=${response.status} url=${url} body=${body.slice(0, 800)}`) + console.error(`[fetchTestRequest] gateway status=${response.status} attempt=${attempt}/${maxAttempts} url=${url} body=${body.slice(0, 800)}`) + + const isKongUpstreamDeath = body.includes(KONG_UPSTREAM_INVALID_RESPONSE) + if (!isKongUpstreamDeath || attempt === maxAttempts) + return response + + await new Promise(resolve => setTimeout(resolve, 250 * attempt)) } - return response + + return lastResponse! } /** diff --git a/tests/webhooks.test.ts b/tests/webhooks.test.ts index 70ccf7e3a3..ba2b8fb8b6 100644 --- a/tests/webhooks.test.ts +++ b/tests/webhooks.test.ts @@ -228,7 +228,7 @@ describeBackend('[GET] /webhooks', () => { describeBackend('[POST] /webhooks', () => { it('create webhook', async () => { - const response = await fetch(webhookEndpoint(), { + const response = await fetchTestRequest(webhookEndpoint(), { method: 'POST', headers: webhookHeaders, body: JSON.stringify({ @@ -253,7 +253,7 @@ describeBackend('[POST] /webhooks', () => { }) it('create webhook with Standard Webhooks delivery version', async () => { - const response = await fetch(webhookEndpoint(), { + const response = await fetchTestRequest(webhookEndpoint(), { method: 'POST', headers: webhookHeaders, body: JSON.stringify({ @@ -395,7 +395,7 @@ describeBackend('[POST] /webhooks', () => { }) it('create webhook with missing required fields', async () => { - const response = await fetch(webhookEndpoint(), { + const response = await fetchTestRequest(webhookEndpoint(), { method: 'POST', headers: webhookHeaders, body: JSON.stringify({ @@ -410,7 +410,7 @@ describeBackend('[POST] /webhooks', () => { }) it('create webhook with invalid URL', async () => { - const response = await fetch(webhookEndpoint(), { + const response = await fetchTestRequest(webhookEndpoint(), { method: 'POST', headers: webhookHeaders, body: JSON.stringify({ @@ -426,7 +426,7 @@ describeBackend('[POST] /webhooks', () => { }) it('create webhook with HTTP URL (non-HTTPS)', async () => { - const response = await fetch(webhookEndpoint(), { + const response = await fetchTestRequest(webhookEndpoint(), { method: 'POST', headers: webhookHeaders, body: JSON.stringify({ @@ -443,7 +443,7 @@ describeBackend('[POST] /webhooks', () => { }) it('create webhook with invalid events', async () => { - const response = await fetch(webhookEndpoint(), { + const response = await fetchTestRequest(webhookEndpoint(), { method: 'POST', headers: webhookHeaders, body: JSON.stringify({ @@ -459,7 +459,7 @@ describeBackend('[POST] /webhooks', () => { }) it('create webhook with invalid delivery version', async () => { - const response = await fetch(webhookEndpoint(), { + const response = await fetchTestRequest(webhookEndpoint(), { method: 'POST', headers: webhookHeaders, body: JSON.stringify({ @@ -476,7 +476,7 @@ describeBackend('[POST] /webhooks', () => { }) it('create webhook with empty events array', async () => { - const response = await fetch(webhookEndpoint(), { + const response = await fetchTestRequest(webhookEndpoint(), { method: 'POST', headers: webhookHeaders, body: JSON.stringify({ @@ -493,7 +493,7 @@ describeBackend('[POST] /webhooks', () => { it('create webhook with invalid orgId', async () => { const invalidOrgId = randomUUID() - const response = await fetch(webhookEndpoint(), { + const response = await fetchTestRequest(webhookEndpoint(), { method: 'POST', headers: webhookHeaders, body: JSON.stringify({ @@ -507,7 +507,7 @@ describeBackend('[POST] /webhooks', () => { }) it('create webhook rejects localhost URL', async () => { - const response = await fetch(webhookEndpoint(), { + const response = await fetchTestRequest(webhookEndpoint(), { method: 'POST', headers: webhookHeaders, body: JSON.stringify({ @@ -557,7 +557,7 @@ describeBackend('[PUT] /webhooks', () => { throw new Error('Webhook was not created in previous test') const newName = `Updated Webhook ${globalId}` - const response = await fetch(webhookEndpoint(), { + const response = await fetchTestRequest(webhookEndpoint(), { method: 'PUT', headers: webhookHeaders, body: JSON.stringify({ @@ -578,7 +578,7 @@ describeBackend('[PUT] /webhooks', () => { throw new Error('Webhook was not created in previous test') const newUrl = 'https://updated.example.com/webhook' - const response = await fetch(webhookEndpoint(), { + const response = await fetchTestRequest(webhookEndpoint(), { method: 'PUT', headers: webhookHeaders, body: JSON.stringify({ @@ -596,7 +596,7 @@ describeBackend('[PUT] /webhooks', () => { if (!createdWebhookId) throw new Error('Webhook was not created in previous test') - const response = await fetch(webhookEndpoint(), { + const response = await fetchTestRequest(webhookEndpoint(), { method: 'PUT', headers: webhookHeaders, body: JSON.stringify({ @@ -615,7 +615,7 @@ describeBackend('[PUT] /webhooks', () => { if (!createdWebhookId) throw new Error('Webhook was not created in previous test') - const response = await fetch(webhookEndpoint(), { + const response = await fetchTestRequest(webhookEndpoint(), { method: 'PUT', headers: webhookHeaders, body: JSON.stringify({ @@ -629,7 +629,7 @@ describeBackend('[PUT] /webhooks', () => { expect(data.webhook.enabled).toBe(false) // Re-enable for subsequent tests - await fetch(webhookEndpoint(), { + await fetchTestRequest(webhookEndpoint(), { method: 'PUT', headers: webhookHeaders, body: JSON.stringify({ @@ -644,7 +644,7 @@ describeBackend('[PUT] /webhooks', () => { if (!createdWebhookId) throw new Error('Webhook was not created in previous test') - const standardResponse = await fetch(webhookEndpoint(), { + const standardResponse = await fetchTestRequest(webhookEndpoint(), { method: 'PUT', headers: webhookHeaders, body: JSON.stringify({ @@ -657,7 +657,7 @@ describeBackend('[PUT] /webhooks', () => { const standardData = await standardResponse.json() as { webhook: { delivery_version: string } } expect(standardData.webhook.delivery_version).toBe('standard') - const legacyResponse = await fetch(webhookEndpoint(), { + const legacyResponse = await fetchTestRequest(webhookEndpoint(), { method: 'PUT', headers: webhookHeaders, body: JSON.stringify({ @@ -675,7 +675,7 @@ describeBackend('[PUT] /webhooks', () => { if (!createdWebhookId) throw new Error('Webhook was not created in previous test') - const response = await fetch(webhookEndpoint(), { + const response = await fetchTestRequest(webhookEndpoint(), { method: 'PUT', headers: webhookHeaders, body: JSON.stringify({ @@ -690,7 +690,7 @@ describeBackend('[PUT] /webhooks', () => { it('update webhook with invalid webhookId', async () => { const invalidWebhookId = randomUUID() - const response = await fetch(webhookEndpoint(), { + const response = await fetchTestRequest(webhookEndpoint(), { method: 'PUT', headers: webhookHeaders, body: JSON.stringify({ @@ -708,7 +708,7 @@ describeBackend('[PUT] /webhooks', () => { if (!createdWebhookId) throw new Error('Webhook was not created in previous test') - const response = await fetch(webhookEndpoint(), { + const response = await fetchTestRequest(webhookEndpoint(), { method: 'PUT', headers: webhookHeaders, body: JSON.stringify({ @@ -726,7 +726,7 @@ describeBackend('[PUT] /webhooks', () => { if (!createdWebhookId) throw new Error('Webhook was not created in previous test') - const response = await fetch(webhookEndpoint(), { + const response = await fetchTestRequest(webhookEndpoint(), { method: 'PUT', headers: webhookHeaders, body: JSON.stringify({ @@ -747,7 +747,7 @@ describeBackend('[POST] /webhooks/test', () => { if (!createdWebhookId) throw new Error('Webhook was not created in previous test') - const response = await fetch(webhookEndpoint('/test'), { + const response = await fetchTestRequest(webhookEndpoint('/test'), { method: 'POST', headers: webhookHeaders, body: JSON.stringify({ @@ -782,7 +782,7 @@ describeBackend('[POST] /webhooks/test', () => { it('test webhook with invalid webhookId', async () => { const invalidWebhookId = randomUUID() - const response = await fetch(webhookEndpoint('/test'), { + const response = await fetchTestRequest(webhookEndpoint('/test'), { method: 'POST', headers: webhookHeaders, body: JSON.stringify({ @@ -806,7 +806,7 @@ describeBackend('[POST] /webhooks/test', () => { expect(updateError).toBeNull() try { - const response = await fetch(webhookEndpoint('/test'), { + const response = await fetchTestRequest(webhookEndpoint('/test'), { method: 'POST', headers: webhookHeaders, body: JSON.stringify({ @@ -828,7 +828,7 @@ describeBackend('[POST] /webhooks/test', () => { }) it('test webhook with missing body', async () => { - const response = await fetch(webhookEndpoint('/test'), { + const response = await fetchTestRequest(webhookEndpoint('/test'), { method: 'POST', headers: webhookHeaders, body: JSON.stringify({}), @@ -843,7 +843,7 @@ describeBackend('[POST] /webhooks/test', () => { if (!createdWebhookId || !appScopedKey) throw new Error('Webhook test prerequisites were not created') - const response = await fetch(webhookEndpoint('/test'), { + const response = await fetchTestRequest(webhookEndpoint('/test'), { method: 'POST', headers: { 'Content-Type': 'application/json', @@ -865,7 +865,7 @@ describeBackend('[POST] /webhooks/test', () => { if (!createdWebhookId || !appScopedKey || !orgScopedSubkeyId) throw new Error('Webhook subkey test prerequisites were not created') - const response = await fetch(webhookEndpoint('/test'), { + const response = await fetchTestRequest(webhookEndpoint('/test'), { method: 'POST', headers: { 'Content-Type': 'application/json', @@ -975,7 +975,7 @@ describeBackend('[GET] /webhooks/deliveries', () => { describeBackend('[POST] /webhooks/deliveries/retry', () => { it('retry delivery with invalid deliveryId', async () => { const invalidDeliveryId = randomUUID() - const response = await fetch(webhookEndpoint('/deliveries/retry'), { + const response = await fetchTestRequest(webhookEndpoint('/deliveries/retry'), { method: 'POST', headers: webhookHeaders, body: JSON.stringify({ @@ -992,7 +992,7 @@ describeBackend('[POST] /webhooks/deliveries/retry', () => { if (!lastDeliveryId || !appScopedKey) throw new Error('Delivery retry prerequisites were not created') - const response = await fetch(webhookEndpoint('/deliveries/retry'), { + const response = await fetchTestRequest(webhookEndpoint('/deliveries/retry'), { method: 'POST', headers: { 'Content-Type': 'application/json', @@ -1014,7 +1014,7 @@ describeBackend('[POST] /webhooks/deliveries/retry', () => { if (!lastDeliveryId || !appScopedKey || !orgScopedSubkeyId) throw new Error('Delivery retry subkey prerequisites were not created') - const response = await fetch(webhookEndpoint('/deliveries/retry'), { + const response = await fetchTestRequest(webhookEndpoint('/deliveries/retry'), { method: 'POST', headers: { 'Content-Type': 'application/json', @@ -1065,7 +1065,7 @@ describeBackend('[POST] /webhooks/deliveries/retry', () => { expect(insertError).toBeNull() try { - const response = await fetch(webhookEndpoint('/deliveries/retry'), { + const response = await fetchTestRequest(webhookEndpoint('/deliveries/retry'), { method: 'POST', headers: webhookHeaders, body: JSON.stringify({ @@ -1126,7 +1126,7 @@ describeBackend('[POST] /webhooks/deliveries/retry', () => { expect(updateError).toBeNull() try { - const response = await fetch(webhookEndpoint('/deliveries/retry'), { + const response = await fetchTestRequest(webhookEndpoint('/deliveries/retry'), { method: 'POST', headers: webhookHeaders, body: JSON.stringify({ @@ -1152,7 +1152,7 @@ describeBackend('[POST] /webhooks/deliveries/retry', () => { }) it('retry delivery with missing body', async () => { - const response = await fetch(webhookEndpoint('/deliveries/retry'), { + const response = await fetchTestRequest(webhookEndpoint('/deliveries/retry'), { method: 'POST', headers: webhookHeaders, body: JSON.stringify({}), @@ -1167,7 +1167,7 @@ describeBackend('[POST] /webhooks/deliveries/retry', () => { describeBackend('[DELETE] /webhooks', () => { it('delete webhook with invalid webhookId', async () => { const invalidWebhookId = randomUUID() - const response = await fetch(webhookEndpoint(`?orgId=${WEBHOOK_TEST_ORG_ID}&webhookId=${invalidWebhookId}`), { + const response = await fetchTestRequest(webhookEndpoint(`?orgId=${WEBHOOK_TEST_ORG_ID}&webhookId=${invalidWebhookId}`), { method: 'DELETE', headers: webhookHeaders, }) @@ -1177,7 +1177,7 @@ describeBackend('[DELETE] /webhooks', () => { }) it('delete webhook with missing body', async () => { - const response = await fetch(webhookEndpoint(), { + const response = await fetchTestRequest(webhookEndpoint(), { method: 'DELETE', headers: webhookHeaders, }) @@ -1191,7 +1191,7 @@ describeBackend('[DELETE] /webhooks', () => { if (!createdWebhookId) throw new Error('Webhook was not created in previous test') - const response = await fetch(webhookEndpoint(`?orgId=${WEBHOOK_TEST_ORG_ID}&webhookId=${createdWebhookId}`), { + const response = await fetchTestRequest(webhookEndpoint(`?orgId=${WEBHOOK_TEST_ORG_ID}&webhookId=${createdWebhookId}`), { method: 'DELETE', headers: webhookHeaders, }) @@ -1201,7 +1201,7 @@ describeBackend('[DELETE] /webhooks', () => { expect(data.webhookId).toBe(createdWebhookId) // Verify deletion - const getResponse = await fetch(webhookEndpoint(`?orgId=${WEBHOOK_TEST_ORG_ID}&webhookId=${createdWebhookId}`), { + const getResponse = await fetchTestRequest(webhookEndpoint(`?orgId=${WEBHOOK_TEST_ORG_ID}&webhookId=${createdWebhookId}`), { headers: webhookHeaders, }) expect(getResponse.status).toBe(400) @@ -1211,7 +1211,7 @@ describeBackend('[DELETE] /webhooks', () => { it('delete already deleted webhook', async () => { // Create a new webhook to delete - const createResponse = await fetch(webhookEndpoint(), { + const createResponse = await fetchTestRequest(webhookEndpoint(), { method: 'POST', headers: webhookHeaders, body: JSON.stringify({ @@ -1225,13 +1225,13 @@ describeBackend('[DELETE] /webhooks', () => { const webhookId = createData.webhook.id // First deletion - await fetch(webhookEndpoint(`?orgId=${WEBHOOK_TEST_ORG_ID}&webhookId=${webhookId}`), { + await fetchTestRequest(webhookEndpoint(`?orgId=${WEBHOOK_TEST_ORG_ID}&webhookId=${webhookId}`), { method: 'DELETE', headers: webhookHeaders, }) // Second deletion attempt - const response = await fetch(webhookEndpoint(`?orgId=${WEBHOOK_TEST_ORG_ID}&webhookId=${webhookId}`), { + const response = await fetchTestRequest(webhookEndpoint(`?orgId=${WEBHOOK_TEST_ORG_ID}&webhookId=${webhookId}`), { method: 'DELETE', headers: webhookHeaders, })