feat(release): prepare 1.1.0 candidate
This commit is contained in:
parent
5a34f66981
commit
5205dcdfa9
736 changed files with 115667 additions and 12203 deletions
|
|
@ -0,0 +1,792 @@
|
|||
BEGIN;
|
||||
|
||||
-- A dispatch reservation is a durable lease, not proof that FCM accepted the
|
||||
-- message. Existing pre-outbox rows are treated as completed so a rollout
|
||||
-- never replays historical notifications unexpectedly.
|
||||
ALTER TABLE public.push_dispatch_attempts
|
||||
ADD COLUMN status text NOT NULL DEFAULT 'pending',
|
||||
ADD COLUMN attempt_count integer NOT NULL DEFAULT 0,
|
||||
ADD COLUMN dispatch_lease_token uuid,
|
||||
ADD COLUMN lease_expires_at timestamptz,
|
||||
ADD COLUMN next_retry_at timestamptz,
|
||||
ADD COLUMN completed_at timestamptz,
|
||||
ADD COLUMN last_error_code text,
|
||||
ADD COLUMN updated_at timestamptz NOT NULL DEFAULT now();
|
||||
|
||||
UPDATE public.push_dispatch_attempts
|
||||
SET status = 'succeeded',
|
||||
completed_at = created_at,
|
||||
updated_at = now()
|
||||
WHERE status = 'pending';
|
||||
|
||||
ALTER TABLE public.push_dispatch_attempts
|
||||
ADD CONSTRAINT push_dispatch_status_check CHECK (
|
||||
status IN ('pending', 'processing', 'succeeded', 'partial', 'failed')
|
||||
),
|
||||
ADD CONSTRAINT push_dispatch_attempt_count_check CHECK (attempt_count >= 0),
|
||||
ADD CONSTRAINT push_dispatch_error_code_check CHECK (
|
||||
last_error_code IS NULL OR last_error_code ~ '^[a-z0-9_]{1,64}$'
|
||||
),
|
||||
ADD CONSTRAINT push_dispatch_lease_shape_check CHECK (
|
||||
(dispatch_lease_token IS NULL AND lease_expires_at IS NULL)
|
||||
OR (dispatch_lease_token IS NOT NULL AND lease_expires_at IS NOT NULL)
|
||||
);
|
||||
|
||||
CREATE INDEX idx_push_dispatch_due
|
||||
ON public.push_dispatch_attempts(next_retry_at, created_at)
|
||||
WHERE status IN ('pending', 'processing');
|
||||
|
||||
CREATE TABLE public.push_deliveries (
|
||||
id uuid PRIMARY KEY DEFAULT gen_random_uuid(),
|
||||
attempt_id uuid NOT NULL REFERENCES public.push_dispatch_attempts(id) ON DELETE CASCADE,
|
||||
push_token_id uuid NOT NULL,
|
||||
target_user_id uuid NOT NULL REFERENCES auth.users(id) ON DELETE CASCADE,
|
||||
provider text NOT NULL CHECK (provider IN ('fcm', 'apns', 'webpush')),
|
||||
token_fingerprint text NOT NULL CHECK (token_fingerprint ~ '^[0-9a-f]{64}$'),
|
||||
status text NOT NULL DEFAULT 'pending' CHECK (
|
||||
status IN ('pending', 'leased', 'delivered', 'stale', 'retryable', 'permanent_failed')
|
||||
),
|
||||
attempt_count integer NOT NULL DEFAULT 0 CHECK (attempt_count >= 0),
|
||||
delivery_lease_token uuid,
|
||||
lease_expires_at timestamptz,
|
||||
next_retry_at timestamptz,
|
||||
last_error_code text CHECK (
|
||||
last_error_code IS NULL OR last_error_code ~ '^[a-z0-9_]{1,64}$'
|
||||
),
|
||||
completed_at timestamptz,
|
||||
created_at timestamptz NOT NULL DEFAULT now(),
|
||||
updated_at timestamptz NOT NULL DEFAULT now(),
|
||||
UNIQUE (attempt_id, push_token_id),
|
||||
CHECK (
|
||||
(status = 'leased' AND delivery_lease_token IS NOT NULL AND lease_expires_at IS NOT NULL)
|
||||
OR (status <> 'leased' AND delivery_lease_token IS NULL AND lease_expires_at IS NULL)
|
||||
)
|
||||
);
|
||||
|
||||
CREATE INDEX idx_push_deliveries_attempt_status
|
||||
ON public.push_deliveries(attempt_id, status, next_retry_at);
|
||||
|
||||
ALTER TABLE public.push_deliveries ENABLE ROW LEVEL SECURITY;
|
||||
REVOKE ALL ON TABLE public.push_deliveries FROM PUBLIC, anon, authenticated;
|
||||
GRANT ALL ON TABLE public.push_deliveries TO service_role;
|
||||
|
||||
CREATE OR REPLACE FUNCTION public.push_dispatch_target_user(
|
||||
target_event_type text,
|
||||
target_resource_id uuid
|
||||
)
|
||||
RETURNS uuid
|
||||
LANGUAGE plpgsql
|
||||
STABLE
|
||||
SECURITY DEFINER
|
||||
SET search_path = ''
|
||||
AS $$
|
||||
DECLARE
|
||||
target_user_id uuid;
|
||||
BEGIN
|
||||
IF auth.role() <> 'service_role' THEN
|
||||
RAISE EXCEPTION 'service_role_required' USING ERRCODE = '42501';
|
||||
END IF;
|
||||
CASE target_event_type
|
||||
WHEN 'transcription.completed' THEN
|
||||
SELECT history.user_id INTO target_user_id
|
||||
FROM public.history AS history
|
||||
WHERE history.id = target_resource_id AND history.status = 'completed';
|
||||
WHEN 'team.invite.created' THEN
|
||||
SELECT account.id INTO target_user_id
|
||||
FROM public.team_invites AS invite
|
||||
JOIN auth.users AS account ON lower(account.email) = lower(invite.email)
|
||||
WHERE invite.id = target_resource_id
|
||||
AND invite.accepted_at IS NULL
|
||||
AND invite.expires_at > now()
|
||||
ORDER BY account.created_at, account.id
|
||||
LIMIT 1;
|
||||
WHEN 'billing.status.changed' THEN
|
||||
SELECT subscription.user_id INTO target_user_id
|
||||
FROM public.subscriptions AS subscription
|
||||
WHERE subscription.id = target_resource_id;
|
||||
ELSE
|
||||
RAISE EXCEPTION 'invalid_push_event' USING ERRCODE = '22023';
|
||||
END CASE;
|
||||
RETURN target_user_id;
|
||||
END;
|
||||
$$;
|
||||
|
||||
CREATE OR REPLACE FUNCTION public.reserve_push_dispatch(
|
||||
push_event_type text,
|
||||
push_resource_id uuid
|
||||
)
|
||||
RETURNS jsonb
|
||||
LANGUAGE plpgsql
|
||||
SECURITY DEFINER
|
||||
SET search_path = ''
|
||||
AS $$
|
||||
DECLARE
|
||||
current_user_id uuid := auth.uid();
|
||||
normalized_event_type text := lower(btrim(push_event_type));
|
||||
existing_attempt public.push_dispatch_attempts;
|
||||
created_attempt public.push_dispatch_attempts;
|
||||
new_lease_token uuid := gen_random_uuid();
|
||||
BEGIN
|
||||
IF current_user_id IS NULL THEN
|
||||
RAISE EXCEPTION 'authentication_required' USING ERRCODE = '42501';
|
||||
END IF;
|
||||
IF push_resource_id IS NULL OR normalized_event_type NOT IN (
|
||||
'transcription.completed', 'team.invite.created', 'billing.status.changed'
|
||||
) THEN
|
||||
RAISE EXCEPTION 'invalid_push_event' USING ERRCODE = '22023';
|
||||
END IF;
|
||||
|
||||
-- RPC callers cannot poison another user's immutable event reservation.
|
||||
IF normalized_event_type = 'transcription.completed' AND NOT EXISTS (
|
||||
SELECT 1 FROM public.history
|
||||
WHERE id = push_resource_id AND user_id = current_user_id AND status = 'completed'
|
||||
) THEN
|
||||
RAISE EXCEPTION 'push_resource_unavailable' USING ERRCODE = 'P0002';
|
||||
ELSIF normalized_event_type = 'team.invite.created' AND NOT EXISTS (
|
||||
SELECT 1 FROM public.team_invites
|
||||
WHERE id = push_resource_id AND invited_by = current_user_id
|
||||
AND accepted_at IS NULL AND expires_at > now()
|
||||
) THEN
|
||||
RAISE EXCEPTION 'push_resource_unavailable' USING ERRCODE = 'P0002';
|
||||
ELSIF normalized_event_type = 'billing.status.changed' AND NOT EXISTS (
|
||||
SELECT 1 FROM public.subscriptions
|
||||
WHERE id = push_resource_id AND user_id = current_user_id
|
||||
) THEN
|
||||
RAISE EXCEPTION 'push_resource_unavailable' USING ERRCODE = 'P0002';
|
||||
END IF;
|
||||
|
||||
PERFORM pg_advisory_xact_lock(hashtextextended(current_user_id::text, 73048));
|
||||
PERFORM pg_advisory_xact_lock(
|
||||
hashtextextended(normalized_event_type || ':' || push_resource_id::text, 73051)
|
||||
);
|
||||
|
||||
SELECT attempt.* INTO existing_attempt
|
||||
FROM public.push_dispatch_attempts AS attempt
|
||||
WHERE attempt.event_type = normalized_event_type
|
||||
AND attempt.resource_id = push_resource_id
|
||||
AND (
|
||||
normalized_event_type IN ('transcription.completed', 'team.invite.created')
|
||||
OR attempt.status IN ('pending', 'processing')
|
||||
OR attempt.created_at > now() - interval '30 seconds'
|
||||
)
|
||||
ORDER BY
|
||||
CASE WHEN attempt.status IN ('pending', 'processing') THEN 0 ELSE 1 END,
|
||||
attempt.created_at DESC
|
||||
LIMIT 1
|
||||
FOR UPDATE;
|
||||
|
||||
IF existing_attempt.id IS NOT NULL THEN
|
||||
IF existing_attempt.status IN ('succeeded', 'partial', 'failed') THEN
|
||||
RETURN jsonb_build_object(
|
||||
'reserved', false, 'duplicate', true,
|
||||
'attempt_id', existing_attempt.id, 'status', existing_attempt.status
|
||||
);
|
||||
END IF;
|
||||
IF existing_attempt.lease_expires_at IS NOT NULL
|
||||
AND existing_attempt.lease_expires_at > now() THEN
|
||||
RETURN jsonb_build_object(
|
||||
'reserved', false, 'duplicate', true,
|
||||
'attempt_id', existing_attempt.id, 'status', 'processing'
|
||||
);
|
||||
END IF;
|
||||
IF existing_attempt.next_retry_at IS NOT NULL
|
||||
AND existing_attempt.next_retry_at > now() THEN
|
||||
RETURN jsonb_build_object(
|
||||
'reserved', false, 'duplicate', true,
|
||||
'attempt_id', existing_attempt.id, 'status', 'pending',
|
||||
'next_retry_at', existing_attempt.next_retry_at
|
||||
);
|
||||
END IF;
|
||||
UPDATE public.push_dispatch_attempts
|
||||
SET status = 'processing',
|
||||
attempt_count = attempt_count + 1,
|
||||
dispatch_lease_token = new_lease_token,
|
||||
lease_expires_at = now() + interval '2 minutes',
|
||||
next_retry_at = NULL,
|
||||
completed_at = NULL,
|
||||
last_error_code = NULL,
|
||||
updated_at = now()
|
||||
WHERE id = existing_attempt.id
|
||||
RETURNING * INTO existing_attempt;
|
||||
RETURN jsonb_build_object(
|
||||
'reserved', true, 'duplicate', false,
|
||||
'attempt_id', existing_attempt.id, 'status', existing_attempt.status,
|
||||
'dispatch_lease_token', new_lease_token
|
||||
);
|
||||
END IF;
|
||||
|
||||
IF (
|
||||
SELECT count(*) FROM public.push_dispatch_attempts AS attempt
|
||||
WHERE attempt.actor_kind = 'user' AND attempt.caller_id = current_user_id
|
||||
AND attempt.created_at > now() - interval '1 minute'
|
||||
) >= 10 OR (
|
||||
SELECT count(*) FROM public.push_dispatch_attempts AS attempt
|
||||
WHERE attempt.actor_kind = 'user' AND attempt.caller_id = current_user_id
|
||||
AND attempt.created_at > now() - interval '24 hours'
|
||||
) >= 100 THEN
|
||||
RAISE EXCEPTION 'push_rate_limited' USING ERRCODE = '54000';
|
||||
END IF;
|
||||
|
||||
INSERT INTO public.push_dispatch_attempts (
|
||||
caller_id, actor_kind, event_type, resource_id, status, attempt_count,
|
||||
dispatch_lease_token, lease_expires_at
|
||||
) VALUES (
|
||||
current_user_id, 'user', normalized_event_type, push_resource_id,
|
||||
'processing', 1, new_lease_token, now() + interval '2 minutes'
|
||||
) RETURNING * INTO created_attempt;
|
||||
RETURN jsonb_build_object(
|
||||
'reserved', true, 'duplicate', false,
|
||||
'attempt_id', created_attempt.id, 'status', created_attempt.status,
|
||||
'dispatch_lease_token', new_lease_token
|
||||
);
|
||||
END;
|
||||
$$;
|
||||
|
||||
CREATE OR REPLACE FUNCTION public.reserve_system_push_dispatch(
|
||||
push_event_type text,
|
||||
push_resource_id uuid
|
||||
)
|
||||
RETURNS jsonb
|
||||
LANGUAGE plpgsql
|
||||
SECURITY DEFINER
|
||||
SET search_path = ''
|
||||
AS $$
|
||||
DECLARE
|
||||
normalized_event_type text := lower(btrim(push_event_type));
|
||||
existing_attempt public.push_dispatch_attempts;
|
||||
created_attempt public.push_dispatch_attempts;
|
||||
new_lease_token uuid := gen_random_uuid();
|
||||
target_user_id uuid;
|
||||
BEGIN
|
||||
IF auth.role() <> 'service_role' THEN
|
||||
RAISE EXCEPTION 'service_role_required' USING ERRCODE = '42501';
|
||||
END IF;
|
||||
IF push_resource_id IS NULL OR normalized_event_type NOT IN (
|
||||
'transcription.completed', 'team.invite.created', 'billing.status.changed'
|
||||
) THEN
|
||||
RAISE EXCEPTION 'invalid_push_event' USING ERRCODE = '22023';
|
||||
END IF;
|
||||
target_user_id := public.push_dispatch_target_user(normalized_event_type, push_resource_id);
|
||||
IF target_user_id IS NULL THEN
|
||||
RAISE EXCEPTION 'push_resource_unavailable' USING ERRCODE = 'P0002';
|
||||
END IF;
|
||||
|
||||
PERFORM pg_advisory_xact_lock(hashtextextended('system-push-global', 73050));
|
||||
PERFORM pg_advisory_xact_lock(
|
||||
hashtextextended(normalized_event_type || ':' || push_resource_id::text, 73051)
|
||||
);
|
||||
|
||||
SELECT attempt.* INTO existing_attempt
|
||||
FROM public.push_dispatch_attempts AS attempt
|
||||
WHERE attempt.event_type = normalized_event_type
|
||||
AND attempt.resource_id = push_resource_id
|
||||
AND (
|
||||
normalized_event_type IN ('transcription.completed', 'team.invite.created')
|
||||
OR attempt.status IN ('pending', 'processing')
|
||||
OR attempt.created_at > now() - interval '30 seconds'
|
||||
)
|
||||
ORDER BY
|
||||
CASE WHEN attempt.status IN ('pending', 'processing') THEN 0 ELSE 1 END,
|
||||
attempt.created_at DESC
|
||||
LIMIT 1
|
||||
FOR UPDATE;
|
||||
|
||||
IF existing_attempt.id IS NOT NULL THEN
|
||||
IF existing_attempt.status IN ('succeeded', 'partial', 'failed') THEN
|
||||
RETURN jsonb_build_object(
|
||||
'reserved', false, 'duplicate', true,
|
||||
'attempt_id', existing_attempt.id, 'status', existing_attempt.status
|
||||
);
|
||||
END IF;
|
||||
IF existing_attempt.lease_expires_at IS NOT NULL
|
||||
AND existing_attempt.lease_expires_at > now() THEN
|
||||
RETURN jsonb_build_object(
|
||||
'reserved', false, 'duplicate', true,
|
||||
'attempt_id', existing_attempt.id, 'status', 'processing'
|
||||
);
|
||||
END IF;
|
||||
IF existing_attempt.next_retry_at IS NOT NULL
|
||||
AND existing_attempt.next_retry_at > now() THEN
|
||||
RETURN jsonb_build_object(
|
||||
'reserved', false, 'duplicate', true,
|
||||
'attempt_id', existing_attempt.id, 'status', 'pending',
|
||||
'next_retry_at', existing_attempt.next_retry_at
|
||||
);
|
||||
END IF;
|
||||
UPDATE public.push_dispatch_attempts
|
||||
SET status = 'processing',
|
||||
attempt_count = attempt_count + 1,
|
||||
dispatch_lease_token = new_lease_token,
|
||||
lease_expires_at = now() + interval '2 minutes',
|
||||
next_retry_at = NULL,
|
||||
completed_at = NULL,
|
||||
last_error_code = NULL,
|
||||
updated_at = now()
|
||||
WHERE id = existing_attempt.id
|
||||
RETURNING * INTO existing_attempt;
|
||||
RETURN jsonb_build_object(
|
||||
'reserved', true, 'duplicate', false,
|
||||
'attempt_id', existing_attempt.id, 'status', existing_attempt.status,
|
||||
'dispatch_lease_token', new_lease_token
|
||||
);
|
||||
END IF;
|
||||
|
||||
IF (
|
||||
SELECT count(*) FROM public.push_dispatch_attempts AS attempt
|
||||
WHERE attempt.actor_kind = 'system'
|
||||
AND attempt.created_at > now() - interval '1 minute'
|
||||
) >= 500 OR (
|
||||
SELECT count(*) FROM public.push_dispatch_attempts AS attempt
|
||||
WHERE attempt.actor_kind = 'system'
|
||||
AND attempt.created_at > now() - interval '24 hours'
|
||||
) >= 10000 THEN
|
||||
RAISE EXCEPTION 'push_rate_limited' USING ERRCODE = '54000';
|
||||
END IF;
|
||||
|
||||
INSERT INTO public.push_dispatch_attempts (
|
||||
caller_id, actor_kind, event_type, resource_id, status, attempt_count,
|
||||
dispatch_lease_token, lease_expires_at
|
||||
) VALUES (
|
||||
NULL, 'system', normalized_event_type, push_resource_id,
|
||||
'processing', 1, new_lease_token, now() + interval '2 minutes'
|
||||
) RETURNING * INTO created_attempt;
|
||||
RETURN jsonb_build_object(
|
||||
'reserved', true, 'duplicate', false,
|
||||
'attempt_id', created_attempt.id, 'status', created_attempt.status,
|
||||
'dispatch_lease_token', new_lease_token
|
||||
);
|
||||
END;
|
||||
$$;
|
||||
|
||||
CREATE OR REPLACE FUNCTION public.lease_push_deliveries(
|
||||
target_attempt_id uuid,
|
||||
target_dispatch_lease_token uuid
|
||||
)
|
||||
RETURNS TABLE (
|
||||
delivery_id uuid,
|
||||
delivery_lease_token uuid,
|
||||
push_token_id uuid,
|
||||
provider text,
|
||||
registration_id text,
|
||||
last_registered_at timestamptz
|
||||
)
|
||||
LANGUAGE plpgsql
|
||||
SECURITY DEFINER
|
||||
SET search_path = ''
|
||||
AS $$
|
||||
DECLARE
|
||||
dispatch public.push_dispatch_attempts;
|
||||
resolved_target_user_id uuid;
|
||||
BEGIN
|
||||
IF auth.role() <> 'service_role' THEN
|
||||
RAISE EXCEPTION 'service_role_required' USING ERRCODE = '42501';
|
||||
END IF;
|
||||
SELECT * INTO dispatch FROM public.push_dispatch_attempts
|
||||
WHERE id = target_attempt_id FOR UPDATE;
|
||||
IF dispatch.id IS NULL THEN
|
||||
RAISE EXCEPTION 'push_dispatch_not_found' USING ERRCODE = 'P0002';
|
||||
END IF;
|
||||
IF dispatch.status IN ('succeeded', 'partial', 'failed') THEN
|
||||
RETURN;
|
||||
END IF;
|
||||
IF target_dispatch_lease_token IS NULL
|
||||
OR dispatch.dispatch_lease_token IS DISTINCT FROM target_dispatch_lease_token
|
||||
OR dispatch.lease_expires_at IS NULL
|
||||
OR dispatch.lease_expires_at <= now() THEN
|
||||
RAISE EXCEPTION 'push_dispatch_lease_conflict' USING ERRCODE = '40001';
|
||||
END IF;
|
||||
|
||||
resolved_target_user_id := public.push_dispatch_target_user(dispatch.event_type, dispatch.resource_id);
|
||||
IF resolved_target_user_id IS NULL THEN
|
||||
UPDATE public.push_dispatch_attempts
|
||||
SET status = 'succeeded', completed_at = now(), last_error_code = 'no_targets',
|
||||
dispatch_lease_token = NULL, lease_expires_at = NULL,
|
||||
next_retry_at = NULL, updated_at = now()
|
||||
WHERE id = dispatch.id;
|
||||
RETURN;
|
||||
END IF;
|
||||
|
||||
-- Refresh the delivery snapshot without duplicating raw registration ids.
|
||||
INSERT INTO public.push_deliveries (
|
||||
attempt_id, push_token_id, target_user_id, provider, token_fingerprint
|
||||
)
|
||||
SELECT
|
||||
dispatch.id,
|
||||
token.id,
|
||||
resolved_target_user_id,
|
||||
token.provider,
|
||||
encode(extensions.digest(token.token, 'sha256'), 'hex')
|
||||
FROM public.push_tokens AS token
|
||||
JOIN public.devices AS device
|
||||
ON device.id = token.device_id
|
||||
AND device.user_id = token.user_id
|
||||
AND device.revoked_at IS NULL
|
||||
WHERE token.user_id = resolved_target_user_id
|
||||
AND token.last_registered_at >= now() - interval '35 days'
|
||||
ON CONFLICT ON CONSTRAINT push_deliveries_attempt_id_push_token_id_key DO NOTHING;
|
||||
|
||||
-- Missing, revoked, or old registrations are terminal neutral outcomes.
|
||||
UPDATE public.push_deliveries AS delivery
|
||||
SET status = 'stale', completed_at = now(), last_error_code = 'registration_stale',
|
||||
delivery_lease_token = NULL, lease_expires_at = NULL,
|
||||
next_retry_at = NULL, updated_at = now()
|
||||
WHERE delivery.attempt_id = dispatch.id
|
||||
AND delivery.status IN ('pending', 'leased', 'retryable')
|
||||
AND NOT EXISTS (
|
||||
SELECT 1
|
||||
FROM public.push_tokens AS token
|
||||
JOIN public.devices AS device
|
||||
ON device.id = token.device_id
|
||||
AND device.user_id = token.user_id
|
||||
AND device.revoked_at IS NULL
|
||||
WHERE token.id = delivery.push_token_id
|
||||
AND token.user_id = resolved_target_user_id
|
||||
AND token.last_registered_at >= now() - interval '35 days'
|
||||
AND encode(extensions.digest(token.token, 'sha256'), 'hex') = delivery.token_fingerprint
|
||||
);
|
||||
|
||||
DELETE FROM public.push_tokens AS token
|
||||
WHERE token.user_id = resolved_target_user_id
|
||||
AND token.last_registered_at < now() - interval '35 days';
|
||||
|
||||
UPDATE public.push_deliveries
|
||||
SET status = CASE WHEN attempt_count >= 5 THEN 'permanent_failed' ELSE 'retryable' END,
|
||||
next_retry_at = CASE
|
||||
WHEN attempt_count >= 5 THEN NULL
|
||||
ELSE now() + make_interval(secs => least(3600, (30 * power(2, greatest(0, attempt_count - 1)))::integer))
|
||||
END,
|
||||
completed_at = CASE WHEN attempt_count >= 5 THEN now() ELSE NULL END,
|
||||
last_error_code = CASE WHEN attempt_count >= 5 THEN 'delivery_attempts_exhausted' ELSE 'delivery_lease_expired' END,
|
||||
delivery_lease_token = NULL,
|
||||
lease_expires_at = NULL,
|
||||
updated_at = now()
|
||||
WHERE attempt_id = dispatch.id AND status = 'leased' AND lease_expires_at <= now();
|
||||
|
||||
UPDATE public.push_deliveries
|
||||
SET status = 'permanent_failed', completed_at = now(), next_retry_at = NULL,
|
||||
last_error_code = 'delivery_attempts_exhausted', updated_at = now()
|
||||
WHERE attempt_id = dispatch.id AND status = 'retryable' AND attempt_count >= 5;
|
||||
|
||||
IF NOT EXISTS (
|
||||
SELECT 1 FROM public.push_deliveries WHERE attempt_id = dispatch.id
|
||||
) THEN
|
||||
UPDATE public.push_dispatch_attempts
|
||||
SET status = 'succeeded', completed_at = now(), last_error_code = 'no_targets',
|
||||
dispatch_lease_token = NULL, lease_expires_at = NULL,
|
||||
next_retry_at = NULL, updated_at = now()
|
||||
WHERE id = dispatch.id;
|
||||
RETURN;
|
||||
END IF;
|
||||
|
||||
RETURN QUERY
|
||||
WITH candidates AS (
|
||||
SELECT delivery.id
|
||||
FROM public.push_deliveries AS delivery
|
||||
WHERE delivery.attempt_id = dispatch.id
|
||||
AND (
|
||||
delivery.status = 'pending'
|
||||
OR (delivery.status = 'retryable' AND coalesce(delivery.next_retry_at, now()) <= now())
|
||||
)
|
||||
AND delivery.attempt_count < 5
|
||||
ORDER BY delivery.created_at, delivery.id
|
||||
FOR UPDATE SKIP LOCKED
|
||||
LIMIT 100
|
||||
), leased AS (
|
||||
UPDATE public.push_deliveries AS delivery
|
||||
SET status = 'leased',
|
||||
attempt_count = delivery.attempt_count + 1,
|
||||
delivery_lease_token = gen_random_uuid(),
|
||||
lease_expires_at = now() + interval '60 seconds',
|
||||
next_retry_at = NULL,
|
||||
completed_at = NULL,
|
||||
last_error_code = NULL,
|
||||
updated_at = now()
|
||||
FROM candidates
|
||||
WHERE delivery.id = candidates.id
|
||||
RETURNING delivery.*
|
||||
)
|
||||
SELECT
|
||||
leased.id,
|
||||
leased.delivery_lease_token,
|
||||
leased.push_token_id,
|
||||
leased.provider,
|
||||
token.token,
|
||||
token.last_registered_at
|
||||
FROM leased
|
||||
JOIN public.push_tokens AS token
|
||||
ON token.id = leased.push_token_id
|
||||
AND token.user_id = leased.target_user_id
|
||||
AND encode(extensions.digest(token.token, 'sha256'), 'hex') = leased.token_fingerprint
|
||||
ORDER BY leased.created_at, leased.id;
|
||||
END;
|
||||
$$;
|
||||
|
||||
CREATE OR REPLACE FUNCTION public.finalize_push_delivery(
|
||||
target_delivery_id uuid,
|
||||
target_delivery_lease_token uuid,
|
||||
delivery_outcome text,
|
||||
delivery_error_code text DEFAULT NULL
|
||||
)
|
||||
RETURNS jsonb
|
||||
LANGUAGE plpgsql
|
||||
SECURITY DEFINER
|
||||
SET search_path = ''
|
||||
AS $$
|
||||
DECLARE
|
||||
delivery public.push_deliveries;
|
||||
normalized_outcome text := lower(btrim(delivery_outcome));
|
||||
normalized_error text := nullif(lower(btrim(delivery_error_code)), '');
|
||||
next_status text;
|
||||
retry_at timestamptz;
|
||||
BEGIN
|
||||
IF auth.role() <> 'service_role' THEN
|
||||
RAISE EXCEPTION 'service_role_required' USING ERRCODE = '42501';
|
||||
END IF;
|
||||
IF normalized_outcome NOT IN (
|
||||
'delivered', 'stale', 'retryable_failure', 'permanent_failure'
|
||||
) OR (normalized_error IS NOT NULL AND normalized_error !~ '^[a-z0-9_]{1,64}$') THEN
|
||||
RAISE EXCEPTION 'invalid_delivery_outcome' USING ERRCODE = '22023';
|
||||
END IF;
|
||||
IF normalized_outcome IN ('retryable_failure', 'permanent_failure')
|
||||
AND normalized_error IS NULL THEN
|
||||
RAISE EXCEPTION 'delivery_error_code_required' USING ERRCODE = '22023';
|
||||
END IF;
|
||||
|
||||
SELECT * INTO delivery FROM public.push_deliveries
|
||||
WHERE id = target_delivery_id FOR UPDATE;
|
||||
IF delivery.id IS NULL THEN
|
||||
RAISE EXCEPTION 'push_delivery_not_found' USING ERRCODE = 'P0002';
|
||||
END IF;
|
||||
IF delivery.status <> 'leased'
|
||||
OR target_delivery_lease_token IS NULL
|
||||
OR delivery.delivery_lease_token IS DISTINCT FROM target_delivery_lease_token
|
||||
OR delivery.lease_expires_at IS NULL
|
||||
OR delivery.lease_expires_at <= now() THEN
|
||||
RAISE EXCEPTION 'push_delivery_lease_conflict' USING ERRCODE = '40001';
|
||||
END IF;
|
||||
|
||||
IF normalized_outcome = 'delivered' THEN
|
||||
next_status := 'delivered';
|
||||
ELSIF normalized_outcome = 'stale' THEN
|
||||
next_status := 'stale';
|
||||
DELETE FROM public.push_tokens AS token
|
||||
WHERE token.id = delivery.push_token_id
|
||||
AND encode(extensions.digest(token.token, 'sha256'), 'hex') = delivery.token_fingerprint;
|
||||
ELSIF normalized_outcome = 'retryable_failure' AND delivery.attempt_count < 5 THEN
|
||||
next_status := 'retryable';
|
||||
retry_at := now() + make_interval(
|
||||
secs => least(3600, (30 * power(2, greatest(0, delivery.attempt_count - 1)))::integer)
|
||||
);
|
||||
ELSE
|
||||
next_status := 'permanent_failed';
|
||||
END IF;
|
||||
|
||||
UPDATE public.push_deliveries
|
||||
SET status = next_status,
|
||||
next_retry_at = retry_at,
|
||||
last_error_code = CASE
|
||||
WHEN next_status = 'delivered' THEN NULL
|
||||
WHEN next_status = 'stale' THEN coalesce(normalized_error, 'registration_stale')
|
||||
ELSE normalized_error
|
||||
END,
|
||||
completed_at = CASE
|
||||
WHEN next_status IN ('delivered', 'stale', 'permanent_failed') THEN now()
|
||||
ELSE NULL
|
||||
END,
|
||||
delivery_lease_token = NULL,
|
||||
lease_expires_at = NULL,
|
||||
updated_at = now()
|
||||
WHERE id = delivery.id;
|
||||
|
||||
RETURN jsonb_build_object(
|
||||
'delivery_id', delivery.id,
|
||||
'status', next_status,
|
||||
'next_retry_at', retry_at
|
||||
);
|
||||
END;
|
||||
$$;
|
||||
|
||||
CREATE OR REPLACE FUNCTION public.finalize_push_dispatch(
|
||||
target_attempt_id uuid,
|
||||
target_dispatch_lease_token uuid
|
||||
)
|
||||
RETURNS jsonb
|
||||
LANGUAGE plpgsql
|
||||
SECURITY DEFINER
|
||||
SET search_path = ''
|
||||
AS $$
|
||||
DECLARE
|
||||
dispatch public.push_dispatch_attempts;
|
||||
delivered_count integer;
|
||||
stale_count integer;
|
||||
retryable_count integer;
|
||||
permanent_count integer;
|
||||
leased_count integer;
|
||||
pending_count integer;
|
||||
next_retry timestamptz;
|
||||
final_status text;
|
||||
is_complete boolean;
|
||||
BEGIN
|
||||
IF auth.role() <> 'service_role' THEN
|
||||
RAISE EXCEPTION 'service_role_required' USING ERRCODE = '42501';
|
||||
END IF;
|
||||
SELECT * INTO dispatch FROM public.push_dispatch_attempts
|
||||
WHERE id = target_attempt_id FOR UPDATE;
|
||||
IF dispatch.id IS NULL THEN
|
||||
RAISE EXCEPTION 'push_dispatch_not_found' USING ERRCODE = 'P0002';
|
||||
END IF;
|
||||
IF dispatch.status IN ('succeeded', 'partial', 'failed') THEN
|
||||
RETURN jsonb_build_object(
|
||||
'attempt_id', dispatch.id, 'status', dispatch.status, 'complete', true
|
||||
);
|
||||
END IF;
|
||||
IF target_dispatch_lease_token IS NULL
|
||||
OR dispatch.dispatch_lease_token IS DISTINCT FROM target_dispatch_lease_token
|
||||
OR dispatch.lease_expires_at IS NULL
|
||||
OR dispatch.lease_expires_at <= now() THEN
|
||||
RAISE EXCEPTION 'push_dispatch_lease_conflict' USING ERRCODE = '40001';
|
||||
END IF;
|
||||
|
||||
UPDATE public.push_deliveries
|
||||
SET status = CASE WHEN attempt_count >= 5 THEN 'permanent_failed' ELSE 'retryable' END,
|
||||
next_retry_at = CASE
|
||||
WHEN attempt_count >= 5 THEN NULL
|
||||
ELSE now() + make_interval(secs => least(3600, (30 * power(2, greatest(0, attempt_count - 1)))::integer))
|
||||
END,
|
||||
completed_at = CASE WHEN attempt_count >= 5 THEN now() ELSE NULL END,
|
||||
last_error_code = CASE WHEN attempt_count >= 5 THEN 'delivery_attempts_exhausted' ELSE 'delivery_lease_expired' END,
|
||||
delivery_lease_token = NULL, lease_expires_at = NULL, updated_at = now()
|
||||
WHERE attempt_id = dispatch.id AND status = 'leased' AND lease_expires_at <= now();
|
||||
|
||||
SELECT
|
||||
count(*) FILTER (WHERE status = 'delivered'),
|
||||
count(*) FILTER (WHERE status = 'stale'),
|
||||
count(*) FILTER (WHERE status = 'retryable'),
|
||||
count(*) FILTER (WHERE status = 'permanent_failed'),
|
||||
count(*) FILTER (WHERE status = 'leased'),
|
||||
count(*) FILTER (WHERE status = 'pending'),
|
||||
min(CASE
|
||||
WHEN status = 'pending' THEN now()
|
||||
WHEN status = 'retryable' THEN coalesce(next_retry_at, now())
|
||||
WHEN status = 'leased' THEN lease_expires_at
|
||||
END)
|
||||
INTO delivered_count, stale_count, retryable_count, permanent_count,
|
||||
leased_count, pending_count, next_retry
|
||||
FROM public.push_deliveries
|
||||
WHERE attempt_id = dispatch.id;
|
||||
|
||||
IF retryable_count + leased_count + pending_count > 0 THEN
|
||||
final_status := CASE WHEN leased_count > 0 THEN 'processing' ELSE 'pending' END;
|
||||
is_complete := false;
|
||||
UPDATE public.push_dispatch_attempts
|
||||
SET status = final_status,
|
||||
dispatch_lease_token = NULL,
|
||||
lease_expires_at = NULL,
|
||||
next_retry_at = coalesce(next_retry, now()),
|
||||
completed_at = NULL,
|
||||
last_error_code = CASE WHEN retryable_count > 0 THEN 'delivery_retry_pending' ELSE NULL END,
|
||||
updated_at = now()
|
||||
WHERE id = dispatch.id;
|
||||
ELSE
|
||||
is_complete := true;
|
||||
final_status := CASE
|
||||
WHEN delivered_count > 0 AND permanent_count = 0 THEN 'succeeded'
|
||||
WHEN delivered_count > 0 AND permanent_count > 0 THEN 'partial'
|
||||
WHEN delivered_count = 0 AND permanent_count > 0 THEN 'failed'
|
||||
ELSE 'succeeded'
|
||||
END;
|
||||
UPDATE public.push_dispatch_attempts
|
||||
SET status = final_status,
|
||||
dispatch_lease_token = NULL,
|
||||
lease_expires_at = NULL,
|
||||
next_retry_at = NULL,
|
||||
completed_at = now(),
|
||||
last_error_code = CASE
|
||||
WHEN final_status = 'partial' THEN 'delivery_partial_failure'
|
||||
WHEN final_status = 'failed' THEN 'delivery_failed'
|
||||
WHEN delivered_count = 0 THEN 'no_targets'
|
||||
ELSE NULL
|
||||
END,
|
||||
updated_at = now()
|
||||
WHERE id = dispatch.id;
|
||||
END IF;
|
||||
|
||||
RETURN jsonb_build_object(
|
||||
'attempt_id', dispatch.id,
|
||||
'status', final_status,
|
||||
'complete', is_complete,
|
||||
'delivered', delivered_count,
|
||||
'stale', stale_count,
|
||||
'retryable_failed', retryable_count,
|
||||
'permanent_failed', permanent_count,
|
||||
'next_retry_at', next_retry
|
||||
);
|
||||
END;
|
||||
$$;
|
||||
|
||||
CREATE OR REPLACE FUNCTION public.claim_due_push_dispatches(batch_limit integer DEFAULT 50)
|
||||
RETURNS TABLE (
|
||||
attempt_id uuid,
|
||||
event_type text,
|
||||
resource_id uuid,
|
||||
dispatch_lease_token uuid
|
||||
)
|
||||
LANGUAGE plpgsql
|
||||
SECURITY DEFINER
|
||||
SET search_path = ''
|
||||
AS $$
|
||||
BEGIN
|
||||
IF auth.role() <> 'service_role' THEN
|
||||
RAISE EXCEPTION 'service_role_required' USING ERRCODE = '42501';
|
||||
END IF;
|
||||
IF batch_limit IS NULL OR batch_limit NOT BETWEEN 1 AND 100 THEN
|
||||
RAISE EXCEPTION 'invalid_batch_limit' USING ERRCODE = '22023';
|
||||
END IF;
|
||||
|
||||
RETURN QUERY
|
||||
WITH due AS (
|
||||
SELECT attempt.id
|
||||
FROM public.push_dispatch_attempts AS attempt
|
||||
WHERE (
|
||||
attempt.status = 'pending' AND coalesce(attempt.next_retry_at, attempt.created_at) <= now()
|
||||
) OR (
|
||||
attempt.status = 'processing' AND attempt.lease_expires_at <= now()
|
||||
)
|
||||
ORDER BY coalesce(attempt.next_retry_at, attempt.created_at), attempt.id
|
||||
FOR UPDATE SKIP LOCKED
|
||||
LIMIT batch_limit
|
||||
), claimed AS (
|
||||
UPDATE public.push_dispatch_attempts AS attempt
|
||||
SET status = 'processing',
|
||||
attempt_count = attempt.attempt_count + 1,
|
||||
dispatch_lease_token = gen_random_uuid(),
|
||||
lease_expires_at = now() + interval '2 minutes',
|
||||
next_retry_at = NULL,
|
||||
updated_at = now()
|
||||
FROM due
|
||||
WHERE attempt.id = due.id
|
||||
RETURNING attempt.*
|
||||
)
|
||||
SELECT claimed.id, claimed.event_type, claimed.resource_id, claimed.dispatch_lease_token
|
||||
FROM claimed
|
||||
ORDER BY claimed.created_at, claimed.id;
|
||||
END;
|
||||
$$;
|
||||
|
||||
REVOKE ALL ON FUNCTION public.push_dispatch_target_user(text, uuid)
|
||||
FROM PUBLIC, anon, authenticated;
|
||||
REVOKE ALL ON FUNCTION public.lease_push_deliveries(uuid, uuid)
|
||||
FROM PUBLIC, anon, authenticated;
|
||||
REVOKE ALL ON FUNCTION public.finalize_push_delivery(uuid, uuid, text, text)
|
||||
FROM PUBLIC, anon, authenticated;
|
||||
REVOKE ALL ON FUNCTION public.finalize_push_dispatch(uuid, uuid)
|
||||
FROM PUBLIC, anon, authenticated;
|
||||
REVOKE ALL ON FUNCTION public.claim_due_push_dispatches(integer)
|
||||
FROM PUBLIC, anon, authenticated;
|
||||
GRANT EXECUTE ON FUNCTION public.push_dispatch_target_user(text, uuid) TO service_role;
|
||||
GRANT EXECUTE ON FUNCTION public.lease_push_deliveries(uuid, uuid) TO service_role;
|
||||
GRANT EXECUTE ON FUNCTION public.finalize_push_delivery(uuid, uuid, text, text) TO service_role;
|
||||
GRANT EXECUTE ON FUNCTION public.finalize_push_dispatch(uuid, uuid) TO service_role;
|
||||
GRANT EXECUTE ON FUNCTION public.claim_due_push_dispatches(integer) TO service_role;
|
||||
|
||||
COMMIT;
|
||||
Loading…
Add table
Add a link
Reference in a new issue