"""Make WP12 allocation publication atomic and concurrency-safe.""" from alembic import op revision = "20260818_551" down_revision = "20260818_550" branch_labels = None depends_on = None def upgrade() -> None: op.get_bind().exec_driver_sql(r''' DO $preflight$ BEGIN IF EXISTS( SELECT 1 FROM public.metering_allocation_rules r LEFT JOIN public.metering_allocations a ON (a.rule_uid,a.rule_version)=(r.rule_uid,r.rule_version) GROUP BY r.rule_uid,r.rule_version HAVING COALESCE(sum(a.weight_micros),0)<>1000000 ) THEN RAISE EXCEPTION 'upgrade refused: WP12 allocation integrity preflight failed'; END IF; END; $preflight$; CREATE EXTENSION IF NOT EXISTS btree_gist; ALTER TABLE public.metering_allocation_rules ADD CONSTRAINT metering_allocation_rules_scope_window_excl EXCLUDE USING gist (tenant_ref WITH =,domain_ref WITH =,department_ref WITH =,project_ref WITH =,cost_center_ref WITH =,tstzrange(effective_start,effective_end,'[)') WITH &&); CREATE FUNCTION public.metering_showback_allocation_total_guard() RETURNS trigger LANGUAGE plpgsql SECURITY DEFINER SET search_path=pg_catalog,public AS $total$ DECLARE uid text; version integer; total bigint; BEGIN IF TG_OP='DELETE' THEN uid:=OLD.rule_uid; version:=OLD.rule_version; ELSE uid:=NEW.rule_uid; version:=NEW.rule_version; END IF; SELECT COALESCE(sum(weight_micros),0) INTO total FROM public.metering_allocations WHERE rule_uid=uid AND rule_version=version; IF total<>1000000 THEN RAISE EXCEPTION 'metering_allocation_weight_invalid'; END IF; RETURN NULL; END; $total$; CREATE CONSTRAINT TRIGGER metering_allocation_rule_total_guard AFTER INSERT OR UPDATE ON public.metering_allocation_rules DEFERRABLE INITIALLY DEFERRED FOR EACH ROW EXECUTE FUNCTION public.metering_showback_allocation_total_guard(); CREATE CONSTRAINT TRIGGER metering_allocation_row_total_guard AFTER INSERT OR UPDATE OR DELETE ON public.metering_allocations DEFERRABLE INITIALLY DEFERRED FOR EACH ROW EXECUTE FUNCTION public.metering_showback_allocation_total_guard(); ALTER FUNCTION public.metering_showback_control_write(text,text,jsonb) RENAME TO metering_showback_control_write_550; CREATE FUNCTION public.metering_showback_control_write(p_action text,p_principal text,p_payload jsonb) RETURNS jsonb LANGUAGE plpgsql SECURITY DEFINER SET search_path=pg_catalog,public AS $control$ DECLARE t text; d text; r text; n integer; distinct_n integer; valid_n integer; total integer; existing_digest text; payload_digest text; BEGIN IF NOT pg_has_role(session_user,'dataops_bi_ai_catalog_issuer','MEMBER') OR p_action NOT IN ('allocation','budget') OR jsonb_typeof(p_payload)<>'object' OR p_principal !~ '^[A-Za-z0-9_.:-]{1,120}$' THEN RAISE EXCEPTION 'metering_control_denied'; END IF; SELECT tenant_ref,domain_ref,role_name INTO t,d,r FROM public.metering_scope_grants WHERE principal_ref=p_principal AND active AND revoked_at IS NULL ORDER BY CASE role_name WHEN 'admin' THEN 1 WHEN 'operator' THEN 2 ELSE 3 END LIMIT 1; IF NOT FOUND OR r NOT IN ('operator','admin') THEN RAISE EXCEPTION 'metering_control_denied'; END IF; IF p_action='budget' THEN RETURN public.metering_showback_control_write_legacy(p_action,p_principal,p_payload); END IF; IF p_payload-ARRAY['rule_uid','rule_version','effective_start','effective_end','mapping','allocations']<>'{}'::jsonb OR jsonb_typeof(p_payload->'mapping')<>'object' OR jsonb_typeof(p_payload->'allocations')<>'array' OR jsonb_array_length(p_payload->'allocations') NOT BETWEEN 1 AND 32 OR p_payload->>'rule_uid' !~ '^[A-Za-z0-9][A-Za-z0-9._:-]{0,119}$' OR p_payload->>'rule_version' !~ '^[1-9][0-9]{0,6}$' OR (p_payload->'mapping')-ARRAY['department','business_domain','project','cost_center']<>'{}'::jsonb OR p_payload->'mapping'->>'business_domain'<>d OR p_payload->'mapping'->>'department' !~ '^[a-z][a-z0-9-]{0,62}$' OR p_payload->'mapping'->>'project' !~ '^[a-z][a-z0-9-]{0,62}$' OR p_payload->'mapping'->>'cost_center' !~ '^[a-z][a-z0-9-]{0,62}$' OR (p_payload->>'effective_start')::timestamptz >= (p_payload->>'effective_end')::timestamptz THEN RAISE EXCEPTION 'metering_allocation_closed'; END IF; SELECT count(*),count(DISTINCT value->>'target'),count(*) FILTER (WHERE jsonb_typeof(value)='object' AND value-ARRAY['target','weight_micros']='{}'::jsonb AND value->>'target' ~ '^[a-z][a-z0-9-]{0,62}$' AND value->>'weight_micros' ~ '^[1-9][0-9]{0,6}$' AND (value->>'weight_micros')::integer<=1000000),COALESCE(sum(CASE WHEN jsonb_typeof(value)='object' AND value-ARRAY['target','weight_micros']='{}'::jsonb AND value->>'weight_micros' ~ '^[1-9][0-9]{0,6}$' THEN (value->>'weight_micros')::integer ELSE 0 END),0) INTO n,distinct_n,valid_n,total FROM jsonb_array_elements(p_payload->'allocations'); IF n<>distinct_n THEN RAISE EXCEPTION 'metering_allocation_target_duplicate'; END IF; IF n<>valid_n OR total<>1000000 THEN RAISE EXCEPTION 'metering_allocation_weight_invalid'; END IF; payload_digest:=encode(sha256(convert_to(p_payload::text,'utf8')),'hex'); PERFORM pg_advisory_xact_lock(hashtextextended('metering-rule|'||(p_payload->>'rule_uid')||'|'||(p_payload->>'rule_version'),0)); PERFORM pg_advisory_xact_lock(hashtextextended('metering-scope|'||t||'|'||d||'|'||(p_payload->'mapping'->>'department')||'|'||(p_payload->'mapping'->>'project')||'|'||(p_payload->'mapping'->>'cost_center'),0)); SELECT rule_digest INTO existing_digest FROM public.metering_allocation_rules WHERE rule_uid=p_payload->>'rule_uid' AND rule_version=(p_payload->>'rule_version')::integer FOR UPDATE; IF FOUND THEN IF existing_digest<>payload_digest THEN RAISE EXCEPTION 'metering_allocation_conflict'; END IF; RETURN jsonb_build_object('rule_uid',p_payload->>'rule_uid','rule_version',(p_payload->>'rule_version')::integer,'persisted_before_ack',true,'replay',true); END IF; INSERT INTO public.metering_allocation_rules(rule_uid,rule_version,tenant_ref,domain_ref,department_ref,project_ref,cost_center_ref,effective_start,effective_end,rule_digest) VALUES(p_payload->>'rule_uid',(p_payload->>'rule_version')::integer,t,d,p_payload->'mapping'->>'department',p_payload->'mapping'->>'project',p_payload->'mapping'->>'cost_center',(p_payload->>'effective_start')::timestamptz,(p_payload->>'effective_end')::timestamptz,payload_digest); INSERT INTO public.metering_allocations(rule_uid,rule_version,target_ref,weight_micros) SELECT p_payload->>'rule_uid',(p_payload->>'rule_version')::integer,value->>'target',(value->>'weight_micros')::integer FROM jsonb_array_elements(p_payload->'allocations'); RETURN jsonb_build_object('rule_uid',p_payload->>'rule_uid','rule_version',(p_payload->>'rule_version')::integer,'persisted_before_ack',true,'replay',false); END; $control$; ALTER FUNCTION public.metering_showback_control_write(text,text,jsonb) OWNER TO dataops_tenant_foundation_owner; REVOKE ALL ON FUNCTION public.metering_showback_control_write_550(text,text,jsonb),public.metering_showback_control_write(text,text,jsonb),public.metering_showback_allocation_total_guard() FROM PUBLIC,dataops_app,dataops_app_runtime,dataops_bi_ai_catalog_control; GRANT EXECUTE ON FUNCTION public.metering_showback_control_write(text,text,jsonb) TO dataops_bi_ai_catalog_control; ''') def downgrade() -> None: bind = op.get_bind() if bind.exec_driver_sql("SELECT EXISTS(SELECT 1 FROM public.metering_allocation_rules LIMIT 1) OR EXISTS(SELECT 1 FROM public.metering_allocations LIMIT 1)").scalar(): raise RuntimeError("downgrade refused: WP12 allocation facts are nonempty") bind.exec_driver_sql(""" DROP FUNCTION public.metering_showback_control_write(text,text,jsonb); ALTER FUNCTION public.metering_showback_control_write_550(text,text,jsonb) RENAME TO metering_showback_control_write; GRANT EXECUTE ON FUNCTION public.metering_showback_control_write(text,text,jsonb) TO dataops_bi_ai_catalog_control; DROP TRIGGER metering_allocation_rule_total_guard ON public.metering_allocation_rules; DROP TRIGGER metering_allocation_row_total_guard ON public.metering_allocations; DROP FUNCTION public.metering_showback_allocation_total_guard(); ALTER TABLE public.metering_allocation_rules DROP CONSTRAINT metering_allocation_rules_scope_window_excl; """)