20260818_549_metering_showback_integrity.py 11 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394
  1. """Close WP12 aggregation, correction and allocation integrity gaps."""
  2. from alembic import op
  3. revision = "20260818_549"
  4. down_revision = "20260818_548"
  5. branch_labels = None
  6. depends_on = None
  7. def upgrade() -> None:
  8. op.get_bind().exec_driver_sql(r'''
  9. ALTER FUNCTION public.metering_showback_control_write(text,text,jsonb) RENAME TO metering_showback_control_write_legacy;
  10. 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 $guard$
  11. DECLARE n integer; distinct_n integer; valid_n integer;
  12. BEGIN
  13. IF p_action='allocation' THEN
  14. IF jsonb_typeof(p_payload)<>'object' OR jsonb_typeof(p_payload->'allocations')<>'array' THEN RAISE EXCEPTION 'metering_allocation_closed'; END IF;
  15. 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) INTO n,distinct_n,valid_n FROM jsonb_array_elements(p_payload->'allocations');
  16. IF n<>distinct_n THEN RAISE EXCEPTION 'metering_allocation_target_duplicate'; END IF;
  17. IF n<>valid_n THEN RAISE EXCEPTION 'metering_allocation_weight_invalid'; END IF;
  18. END IF;
  19. RETURN public.metering_showback_control_write_legacy(p_action,p_principal,p_payload);
  20. END; $guard$;
  21. ALTER FUNCTION public.metering_showback_control_write(text,text,jsonb) OWNER TO dataops_tenant_foundation_owner;
  22. REVOKE ALL ON FUNCTION public.metering_showback_control_write_legacy(text,text,jsonb),public.metering_showback_control_write(text,text,jsonb) FROM PUBLIC,dataops_app,dataops_app_runtime,dataops_bi_ai_catalog_control;
  23. GRANT EXECUTE ON FUNCTION public.metering_showback_control_write(text,text,jsonb) TO dataops_bi_ai_catalog_control;
  24. CREATE FUNCTION public.metering_showback_correction_guard() RETURNS trigger LANGUAGE plpgsql SECURITY DEFINER SET search_path=pg_catalog,public AS $correction$
  25. DECLARE original record;
  26. BEGIN
  27. IF NEW.evidence_reference !~ '^local-fixture://wp12/v1(/[a-z][a-z0-9-]{0,62})?$' OR NEW.evidence_reference ~ '(sql|script|select|insert|update|delete|drop|exec|curl|wget|token|secret|credential|password)' THEN RAISE EXCEPTION 'metering_evidence_reference_invalid'; END IF;
  28. IF NEW.correction_of IS NULL THEN RETURN NEW; END IF;
  29. SELECT * INTO original FROM public.metering_events WHERE event_uid=NEW.correction_of FOR KEY SHARE;
  30. IF NOT FOUND OR original.tenant_ref<>NEW.tenant_ref OR original.domain_ref<>NEW.domain_ref OR original.correction_of IS NOT NULL OR original.event_kind<>NEW.event_kind OR original.unit<>NEW.unit OR original.window_start<>NEW.window_start OR original.window_end<>NEW.window_end OR original.department_ref<>NEW.department_ref OR original.project_ref<>NEW.project_ref OR original.cost_center_ref<>NEW.cost_center_ref THEN RAISE EXCEPTION 'metering_correction_scope_invalid'; END IF;
  31. RETURN NEW;
  32. END; $correction$;
  33. CREATE TRIGGER metering_showback_correction_guard BEFORE INSERT ON public.metering_events FOR EACH ROW EXECUTE FUNCTION public.metering_showback_correction_guard();
  34. CREATE FUNCTION public.metering_showback_rule_window_guard() RETURNS trigger LANGUAGE plpgsql SECURITY DEFINER SET search_path=pg_catalog,public AS $window$
  35. BEGIN
  36. IF EXISTS(SELECT 1 FROM public.metering_allocation_rules r WHERE r.tenant_ref=NEW.tenant_ref AND r.domain_ref=NEW.domain_ref AND r.department_ref=NEW.department_ref AND r.project_ref=NEW.project_ref AND r.cost_center_ref=NEW.cost_center_ref AND (r.rule_uid,r.rule_version)<>(NEW.rule_uid,NEW.rule_version) AND tstzrange(r.effective_start,r.effective_end,'[)') && tstzrange(NEW.effective_start,NEW.effective_end,'[)')) THEN RAISE EXCEPTION 'metering_allocation_window_overlap'; END IF;
  37. RETURN NEW;
  38. END; $window$;
  39. CREATE TRIGGER metering_showback_rule_window_guard BEFORE INSERT OR UPDATE ON public.metering_allocation_rules FOR EACH ROW EXECUTE FUNCTION public.metering_showback_rule_window_guard();
  40. CREATE FUNCTION public.metering_showback_rollup(p_payload jsonb) RETURNS jsonb LANGUAGE plpgsql SECURITY DEFINER SET search_path=pg_catalog,public AS $rollup$
  41. DECLARE c record; source_micros bigint; allocated_micros bigint; details jsonb;
  42. BEGIN
  43. IF NOT pg_has_role(session_user,'dataops_app_runtime','MEMBER') OR jsonb_typeof(p_payload)<>'object' OR p_payload-ARRAY['request_claim','window']<>'{}'::jsonb OR p_payload->>'request_claim' !~ '^[0-9a-f-]{36}$' OR p_payload->>'window' !~ '^[0-9]{4}-(0[1-9]|1[0-2])$' THEN RAISE EXCEPTION 'metering_read_denied'; END IF;
  44. SELECT * INTO c FROM public.metering_runtime_claims WHERE claim_uid=(p_payload->>'request_claim')::uuid AND action_name='read' AND consumed_at IS NULL AND expires_at>clock_timestamp() FOR UPDATE;
  45. IF NOT FOUND THEN RAISE EXCEPTION 'metering_claim_denied'; END IF;
  46. WITH scoped AS (SELECT e.*,r.rule_uid,r.rule_version,r.rule_digest FROM public.metering_events e LEFT JOIN LATERAL (SELECT x.rule_uid,x.rule_version,x.rule_digest FROM public.metering_allocation_rules x WHERE x.tenant_ref=e.tenant_ref AND x.domain_ref=e.domain_ref AND x.department_ref=e.department_ref AND x.project_ref=e.project_ref AND x.cost_center_ref=e.cost_center_ref AND x.effective_start<=e.window_start AND x.effective_end>e.window_start ORDER BY x.rule_version DESC,x.rule_uid DESC LIMIT 1) r ON true WHERE e.tenant_ref=c.tenant_ref AND e.domain_ref=c.domain_ref AND to_char(e.window_start AT TIME ZONE 'UTC','YYYY-MM')=p_payload->>'window'), rules AS (SELECT rule_uid,rule_version,rule_digest,department_ref,project_ref,cost_center_ref,sum(quantity_micros)::bigint AS source FROM scoped WHERE rule_uid IS NOT NULL GROUP BY rule_uid,rule_version,rule_digest,department_ref,project_ref,cost_center_ref), weights AS (SELECT rs.*,a.target_ref,(rs.source*a.weight_micros/1000000)::bigint AS base FROM rules rs JOIN public.metering_allocations a ON a.rule_uid=rs.rule_uid AND a.rule_version=rs.rule_version), expanded AS (SELECT w.*,w.base+CASE WHEN w.target_ref=min(w.target_ref) OVER (PARTITION BY w.rule_uid,w.rule_version) THEN w.source-sum(w.base) OVER (PARTITION BY w.rule_uid,w.rule_version) ELSE 0 END AS amount FROM weights w), per_rule AS (SELECT rule_uid,rule_version,rule_digest,department_ref,project_ref,cost_center_ref,max(source) AS source,sum(amount)::bigint AS allocated,jsonb_agg(jsonb_build_object('target',target_ref,'quantity_micros',amount) ORDER BY target_ref) AS allocations FROM expanded GROUP BY rule_uid,rule_version,rule_digest,department_ref,project_ref,cost_center_ref) SELECT COALESCE((SELECT sum(quantity_micros) FROM scoped),0),COALESCE(sum(allocated),0),COALESCE(jsonb_agg(jsonb_build_object('rule_uid',rule_uid,'rule_version',rule_version,'rule_digest',rule_digest,'mapping',jsonb_build_object('department',department_ref,'business_domain',c.domain_ref,'project',project_ref,'cost_center',cost_center_ref),'source_micros',source,'allocated_micros',allocated,'variance_micros',source-allocated,'allocations',allocations) ORDER BY rule_uid,rule_version),'[]'::jsonb) INTO source_micros,allocated_micros,details FROM per_rule;
  47. UPDATE public.metering_runtime_claims SET consumed_at=clock_timestamp() WHERE claim_uid=c.claim_uid;
  48. INSERT INTO public.metering_audit_events(tenant_ref,domain_ref,event_type,actor_ref,payload_digest) VALUES(c.tenant_ref,c.domain_ref,'showback.read',c.principal_ref,encode(sha256(convert_to('rollup|'||p_payload::text,'utf8')),'hex'));
  49. RETURN jsonb_build_object('window',p_payload->>'window','source_micros',source_micros,'allocated_micros',allocated_micros,'unallocated_micros',source_micros-allocated_micros,'variance_micros',source_micros-allocated_micros,'allocation_rules',details,'mode','ENGINEERING_EVIDENCE_ONLY','chargeback_enabled',false);
  50. END; $rollup$;
  51. ALTER FUNCTION public.metering_showback_rollup(jsonb) OWNER TO dataops_tenant_foundation_owner;
  52. REVOKE ALL ON FUNCTION public.metering_showback_rollup(jsonb),public.metering_showback_correction_guard(),public.metering_showback_rule_window_guard() FROM PUBLIC,dataops_app,dataops_bi_ai_catalog_control;
  53. GRANT EXECUTE ON FUNCTION public.metering_showback_rollup(jsonb) TO dataops_app_runtime;
  54. ALTER FUNCTION public.metering_showback_runtime_read(text,jsonb) RENAME TO metering_showback_runtime_read_legacy;
  55. CREATE FUNCTION public.metering_showback_runtime_read(p_action text,p_payload jsonb) RETURNS jsonb LANGUAGE plpgsql SECURITY DEFINER SET search_path=pg_catalog,public AS $read_guard$
  56. BEGIN
  57. IF p_action='chargeback' THEN RAISE EXCEPTION 'chargeback_disabled'; END IF;
  58. IF p_action IN ('showback','reconciliation') THEN RAISE EXCEPTION 'metering_rollup_required'; END IF;
  59. IF p_action<>'audit' THEN RAISE EXCEPTION 'metering_read_denied'; END IF;
  60. RETURN public.metering_showback_runtime_read_legacy(p_action,p_payload);
  61. END; $read_guard$;
  62. ALTER FUNCTION public.metering_showback_runtime_read(text,jsonb) OWNER TO dataops_tenant_foundation_owner;
  63. REVOKE ALL ON FUNCTION public.metering_showback_runtime_read_legacy(text,jsonb),public.metering_showback_runtime_read(text,jsonb) FROM PUBLIC,dataops_app,dataops_bi_ai_catalog_control;
  64. GRANT EXECUTE ON FUNCTION public.metering_showback_runtime_read(text,jsonb) TO dataops_app_runtime;
  65. ''')
  66. def downgrade() -> None:
  67. op.get_bind().exec_driver_sql("""
  68. DO $$ BEGIN
  69. IF to_regprocedure('public.metering_showback_runtime_read_legacy(text,jsonb)') IS NOT NULL THEN
  70. DROP FUNCTION public.metering_showback_runtime_read(text,jsonb);
  71. ALTER FUNCTION public.metering_showback_runtime_read_legacy(text,jsonb) RENAME TO metering_showback_runtime_read;
  72. GRANT EXECUTE ON FUNCTION public.metering_showback_runtime_read(text,jsonb) TO dataops_app_runtime;
  73. END IF;
  74. IF to_regprocedure('public.metering_showback_control_write_legacy(text,text,jsonb)') IS NOT NULL THEN
  75. DROP FUNCTION public.metering_showback_control_write(text,text,jsonb);
  76. ALTER FUNCTION public.metering_showback_control_write_legacy(text,text,jsonb) RENAME TO metering_showback_control_write;
  77. GRANT EXECUTE ON FUNCTION public.metering_showback_control_write(text,text,jsonb) TO dataops_bi_ai_catalog_control;
  78. END IF;
  79. END $$;
  80. DROP FUNCTION public.metering_showback_rollup(jsonb);
  81. DROP TRIGGER metering_showback_rule_window_guard ON public.metering_allocation_rules;
  82. DROP FUNCTION public.metering_showback_rule_window_guard();
  83. DROP TRIGGER metering_showback_correction_guard ON public.metering_events;
  84. DROP FUNCTION public.metering_showback_correction_guard();
  85. """)