166 lines
4.7 KiB
PL/PgSQL
166 lines
4.7 KiB
PL/PgSQL
BEGIN;
|
|
|
|
-- Business state and notification enqueue must commit together. Edge calls
|
|
-- made after a commit can be lost if the worker crashes between the two.
|
|
CREATE OR REPLACE FUNCTION public.enqueue_immutable_push_event(
|
|
queued_event_type text,
|
|
queued_resource_id uuid
|
|
)
|
|
RETURNS void
|
|
LANGUAGE plpgsql
|
|
SECURITY DEFINER
|
|
SET search_path = ''
|
|
AS $$
|
|
BEGIN
|
|
IF queued_event_type NOT IN ('transcription.completed', 'team.invite.created')
|
|
OR queued_resource_id IS NULL THEN
|
|
RAISE EXCEPTION 'invalid_push_event' USING ERRCODE = '22023';
|
|
END IF;
|
|
INSERT INTO public.push_dispatch_attempts (
|
|
caller_id, actor_kind, event_type, resource_id, status,
|
|
attempt_count, next_retry_at
|
|
) VALUES (
|
|
NULL, 'system', queued_event_type, queued_resource_id, 'pending',
|
|
0, now()
|
|
)
|
|
ON CONFLICT DO NOTHING;
|
|
END;
|
|
$$;
|
|
|
|
CREATE OR REPLACE FUNCTION public.enqueue_billing_push_event(
|
|
queued_subscription_id uuid
|
|
)
|
|
RETURNS void
|
|
LANGUAGE plpgsql
|
|
SECURITY DEFINER
|
|
SET search_path = ''
|
|
AS $$
|
|
BEGIN
|
|
IF queued_subscription_id IS NULL THEN
|
|
RAISE EXCEPTION 'invalid_push_event' USING ERRCODE = '22023';
|
|
END IF;
|
|
PERFORM pg_advisory_xact_lock(
|
|
hashtextextended('billing.status.changed:' || queued_subscription_id::text, 73052)
|
|
);
|
|
IF EXISTS (
|
|
SELECT 1 FROM public.push_dispatch_attempts
|
|
WHERE event_type = 'billing.status.changed'
|
|
AND resource_id = queued_subscription_id
|
|
AND created_at > now() - interval '30 seconds'
|
|
) THEN
|
|
RETURN;
|
|
END IF;
|
|
INSERT INTO public.push_dispatch_attempts (
|
|
caller_id, actor_kind, event_type, resource_id, status,
|
|
attempt_count, next_retry_at
|
|
) VALUES (
|
|
NULL, 'system', 'billing.status.changed', queued_subscription_id,
|
|
'pending', 0, now()
|
|
);
|
|
END;
|
|
$$;
|
|
|
|
CREATE OR REPLACE FUNCTION public.enqueue_processing_job_push()
|
|
RETURNS trigger
|
|
LANGUAGE plpgsql
|
|
SECURITY DEFINER
|
|
SET search_path = ''
|
|
AS $$
|
|
BEGIN
|
|
IF NEW.kind = 'transcription'
|
|
AND NEW.status = 'succeeded'
|
|
AND NEW.history_id IS NOT NULL
|
|
AND OLD.status IS DISTINCT FROM 'succeeded'
|
|
AND EXISTS (
|
|
SELECT 1 FROM public.history
|
|
WHERE id = NEW.history_id AND user_id = NEW.user_id AND status = 'completed'
|
|
) THEN
|
|
PERFORM public.enqueue_immutable_push_event(
|
|
'transcription.completed', NEW.history_id
|
|
);
|
|
END IF;
|
|
RETURN NEW;
|
|
END;
|
|
$$;
|
|
|
|
CREATE OR REPLACE FUNCTION public.enqueue_history_after_worker_push()
|
|
RETURNS trigger
|
|
LANGUAGE plpgsql
|
|
SECURITY DEFINER
|
|
SET search_path = ''
|
|
AS $$
|
|
BEGIN
|
|
IF NEW.status = 'completed'
|
|
AND OLD.status IS DISTINCT FROM 'completed'
|
|
AND EXISTS (
|
|
SELECT 1 FROM public.processing_jobs
|
|
WHERE history_id = NEW.id
|
|
AND user_id = NEW.user_id
|
|
AND kind = 'transcription'
|
|
AND status = 'succeeded'
|
|
) THEN
|
|
PERFORM public.enqueue_immutable_push_event(
|
|
'transcription.completed', NEW.id
|
|
);
|
|
END IF;
|
|
RETURN NEW;
|
|
END;
|
|
$$;
|
|
|
|
CREATE OR REPLACE FUNCTION public.enqueue_team_invite_push()
|
|
RETURNS trigger
|
|
LANGUAGE plpgsql
|
|
SECURITY DEFINER
|
|
SET search_path = ''
|
|
AS $$
|
|
BEGIN
|
|
PERFORM public.enqueue_immutable_push_event('team.invite.created', NEW.id);
|
|
RETURN NEW;
|
|
END;
|
|
$$;
|
|
|
|
CREATE OR REPLACE FUNCTION public.enqueue_subscription_push()
|
|
RETURNS trigger
|
|
LANGUAGE plpgsql
|
|
SECURITY DEFINER
|
|
SET search_path = ''
|
|
AS $$
|
|
BEGIN
|
|
IF ROW(
|
|
OLD.tier, OLD.status, OLD.provider, OLD.current_period_end,
|
|
OLD.cancel_at, OLD.auto_renewing
|
|
) IS DISTINCT FROM ROW(
|
|
NEW.tier, NEW.status, NEW.provider, NEW.current_period_end,
|
|
NEW.cancel_at, NEW.auto_renewing
|
|
) THEN
|
|
PERFORM public.enqueue_billing_push_event(NEW.id);
|
|
END IF;
|
|
RETURN NEW;
|
|
END;
|
|
$$;
|
|
|
|
CREATE TRIGGER enqueue_processing_job_push_after_success
|
|
AFTER UPDATE OF status ON public.processing_jobs
|
|
FOR EACH ROW EXECUTE FUNCTION public.enqueue_processing_job_push();
|
|
|
|
CREATE TRIGGER enqueue_history_push_after_worker_success
|
|
AFTER UPDATE OF status ON public.history
|
|
FOR EACH ROW EXECUTE FUNCTION public.enqueue_history_after_worker_push();
|
|
|
|
CREATE TRIGGER enqueue_team_invite_push_after_insert
|
|
AFTER INSERT ON public.team_invites
|
|
FOR EACH ROW EXECUTE FUNCTION public.enqueue_team_invite_push();
|
|
|
|
CREATE TRIGGER enqueue_subscription_push_after_change
|
|
AFTER UPDATE OF tier, status, provider, current_period_end, cancel_at, auto_renewing
|
|
ON public.subscriptions
|
|
FOR EACH ROW EXECUTE FUNCTION public.enqueue_subscription_push();
|
|
|
|
REVOKE ALL ON FUNCTION public.enqueue_immutable_push_event(text, uuid)
|
|
FROM PUBLIC, anon, authenticated;
|
|
REVOKE ALL ON FUNCTION public.enqueue_billing_push_event(uuid)
|
|
FROM PUBLIC, anon, authenticated;
|
|
GRANT EXECUTE ON FUNCTION public.enqueue_immutable_push_event(text, uuid) TO service_role;
|
|
GRANT EXECUTE ON FUNCTION public.enqueue_billing_push_event(uuid) TO service_role;
|
|
|
|
COMMIT;
|