-- Repoint references before enforcing one row per provider subscription. WITH ranked AS ( SELECT id, FIRST_VALUE(id) OVER ( PARTITION BY provider, provider_subscription_id ORDER BY updated_at DESC, created_at DESC, id DESC ) AS keep_id, ROW_NUMBER() OVER ( PARTITION BY provider, provider_subscription_id ORDER BY updated_at DESC, created_at DESC, id DESC ) AS row_number FROM subscriptions WHERE provider_subscription_id IS NOT NULL ) UPDATE usage_periods AS usage SET subscription_id = ranked.keep_id FROM ranked WHERE ranked.row_number > 1 AND usage.subscription_id = ranked.id; WITH ranked AS ( SELECT id, FIRST_VALUE(id) OVER ( PARTITION BY provider, provider_subscription_id ORDER BY updated_at DESC, created_at DESC, id DESC ) AS keep_id, ROW_NUMBER() OVER ( PARTITION BY provider, provider_subscription_id ORDER BY updated_at DESC, created_at DESC, id DESC ) AS row_number FROM subscriptions WHERE provider_subscription_id IS NOT NULL ) UPDATE invoices AS invoice SET subscription_id = ranked.keep_id FROM ranked WHERE ranked.row_number > 1 AND invoice.subscription_id = ranked.id; WITH ranked AS ( SELECT id, ROW_NUMBER() OVER ( PARTITION BY provider, provider_subscription_id ORDER BY updated_at DESC, created_at DESC, id DESC ) AS row_number FROM subscriptions WHERE provider_subscription_id IS NOT NULL ) DELETE FROM subscriptions AS subscription USING ranked WHERE ranked.row_number > 1 AND subscription.id = ranked.id; CREATE UNIQUE INDEX IF NOT EXISTS idx_subscriptions_provider_object_unique ON subscriptions(provider, provider_subscription_id); CREATE TABLE IF NOT EXISTS provider_object_event_watermarks ( provider VARCHAR(20) NOT NULL, object_type VARCHAR(50) NOT NULL, provider_object_id VARCHAR(200) NOT NULL, last_event_created BIGINT NOT NULL, last_event_rank SMALLINT NOT NULL DEFAULT 0, last_event_id VARCHAR(200) NOT NULL, is_deleted BOOLEAN NOT NULL DEFAULT false, updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), PRIMARY KEY (provider, object_type, provider_object_id), CONSTRAINT provider_object_event_watermarks_rank_check CHECK (last_event_rank BETWEEN 0 AND 100) ); CREATE INDEX IF NOT EXISTS idx_provider_object_event_watermarks_updated ON provider_object_event_watermarks(updated_at); -- Seed rollout watermarks so a delayed pre-deployment event cannot revive an -- already-canceled Stripe subscription before the first new event arrives. INSERT INTO provider_object_event_watermarks ( provider, object_type, provider_object_id, last_event_created, last_event_rank, last_event_id, is_deleted, updated_at ) SELECT 'stripe', 'subscription', provider_subscription_id, EXTRACT(EPOCH FROM updated_at)::bigint, CASE WHEN status = 'canceled' THEN 2 ELSE 1 END, 'migration:' || id::text, status = 'canceled', updated_at FROM subscriptions WHERE provider = 'stripe' AND provider_subscription_id IS NOT NULL ON CONFLICT (provider, object_type, provider_object_id) DO NOTHING;