from __future__ import annotations import json import os import pytest from sqlalchemy import text pytestmark = pytest.mark.integration def test_semantic_governance_publishes_maps_and_rolls_back(monkeypatch): database_url = os.environ.get("TEST_DATABASE_URL") if not database_url: pytest.skip("TEST_DATABASE_URL is required") monkeypatch.setenv("DATABASE_URL", database_url) from app import create_app, db from app.core.common.identifiers import new_governance_uid from app.core.data_research.semantic_governance import ( SemanticGovernanceService, ) from app.core.data_research.semantic_repository import ( SqlAlchemySemanticGovernanceRepository, ) from app.core.data_rules.repository import DataRuleRepository from app.core.events.outbox import claim_outbox, enqueue_outbox app = create_app() app.config.update(TESTING=True) actor_uid = new_governance_uid() owner_uid = new_governance_uid() reviewer_uid = new_governance_uid() source_uid = new_governance_uid() plan_uid = new_governance_uid() run_uid = new_governance_uid() physical_asset_uid = new_governance_uid() element_uid = new_governance_uid() element_version_uid = new_governance_uid() standard_uid = new_governance_uid() standard_id = new_governance_uid() standard_version_uid = new_governance_uid() semantic_uids: list[str] = [] try: with app.app_context(): for user_uid, label in ( (actor_uid, "editor"), (owner_uid, "owner"), (reviewer_uid, "reviewer"), ): db.session.execute( text( """ INSERT INTO public.users ( id, username, display_name, password_hash, status ) VALUES ( CAST(:uid AS uuid), :username, :username, 'p2-wp03-integration-hash', 'active' ) """ ), { "uid": user_uid, "username": f"wp03-{label}-{user_uid[:8]}", }, ) db.session.execute( text( """ INSERT INTO public.ingestion_sources ( uid, source_type, name, config, permission_scope, status, created_by ) VALUES ( CAST(:uid AS uuid), 'database', :name, '{}'::jsonb, '{}'::jsonb, 'active', :created_by ) """ ), { "uid": source_uid, "name": f"WP03 source {source_uid[:8]}", "created_by": actor_uid, }, ) db.session.execute( text( """ INSERT INTO public.active_metadata_plans ( uid, source_uid, name, source_kind, schedule_type, schedule_expression, discovery_mode, scope, owner_uid, enabled, cursor_state, created_by ) VALUES ( CAST(:uid AS uuid), CAST(:source_uid AS uuid), :name, 'database', 'manual', NULL, 'snapshot', '{}'::jsonb, CAST(:owner_uid AS uuid), TRUE, '{}'::jsonb, CAST(:actor_uid AS uuid) ) """ ), { "uid": plan_uid, "source_uid": source_uid, "name": "WP03 semantic mapping source", "owner_uid": owner_uid, "actor_uid": actor_uid, }, ) db.session.execute( text( """ INSERT INTO public.active_metadata_runs ( uid, plan_uid, batch_key, status, attempt_count, cursor_before, cursor_after, statistics, actor_uid, started_at, finished_at ) VALUES ( CAST(:uid AS uuid), CAST(:plan_uid AS uuid), :batch_key, 'completed', 1, '{}'::jsonb, '{}'::jsonb, '{}'::jsonb, CAST(:actor_uid AS uuid), CURRENT_TIMESTAMP, CURRENT_TIMESTAMP ) """ ), { "uid": run_uid, "plan_uid": plan_uid, "batch_key": f"wp03-{run_uid}", "actor_uid": actor_uid, }, ) snapshot = { "name": "customers", "namespace": "public", "asset_type": "table", "fields": [{"name": "customer_level", "data_type": "varchar"}], } db.session.execute( text( """ INSERT INTO public.active_metadata_assets ( uid, source_uid, asset_key, namespace, name, asset_type, lifecycle_status, current_version, content_hash, snapshot, health, last_run_uid ) VALUES ( CAST(:uid AS uuid), CAST(:source_uid AS uuid), :asset_key, 'public', 'customers', 'table', 'active', 1, :content_hash, CAST(:snapshot AS jsonb), '{}'::jsonb, CAST(:run_uid AS uuid) ) """ ), { "uid": physical_asset_uid, "source_uid": source_uid, "asset_key": f"{source_uid}:public.customers", "content_hash": "a" * 64, "snapshot": json.dumps(snapshot), "run_uid": run_uid, }, ) db.session.execute( text( """ INSERT INTO public.data_elements ( uid, code, current_version, status, business_domain_uids, created_by ) VALUES ( CAST(:uid AS uuid), :code, 1, 'published', CAST(:domains AS jsonb), :created_by ); INSERT INTO public.data_element_versions ( uid, data_element_uid, version, status, snapshot, evidence_uids, created_by ) VALUES ( CAST(:version_uid AS uuid), CAST(:uid AS uuid), 1, 'published', CAST(:snapshot AS jsonb), '[]'::jsonb, :created_by ) """ ), { "uid": element_uid, "version_uid": element_version_uid, "code": f"CUSTOMER_LEVEL_{element_uid[:8]}", "domains": json.dumps([source_uid]), "created_by": actor_uid, "snapshot": json.dumps({"name": "客户等级"}), }, ) db.session.execute( text( """ INSERT INTO public.data_standards ( id, standard_uid, name, owner_uid, status ) VALUES ( CAST(:id AS uuid), CAST(:standard_uid AS uuid), '客户主数据标准', CAST(:owner_uid AS uuid), 'active' ); INSERT INTO public.data_standard_versions ( id, standard_uid, version_no, source_text, standard_spec, spec_hash, scope, status, created_by, published_at ) VALUES ( CAST(:version_uid AS uuid), CAST(:standard_uid AS uuid), 1, '客户等级必须使用统一代码', '{}'::jsonb, :spec_hash, '{}'::jsonb, 'validated', CAST(:actor_uid AS uuid), NULL ) """ ), { "id": standard_id, "standard_uid": standard_uid, "version_uid": standard_version_uid, "owner_uid": owner_uid, "actor_uid": actor_uid, "spec_hash": "b" * 64, }, ) db.session.commit() standard = DataRuleRepository(db.session).publish_standard_version( version_id=standard_version_uid, published_by=actor_uid, ) db.session.commit() assert standard["status"] == "published" repository = SqlAlchemySemanticGovernanceRepository(db.session) service = SemanticGovernanceService( repository, outbox_enqueue=lambda **event: enqueue_outbox( db.session, **event ), ) code_set = service.create_draft( "code_set", { "code": f"CUSTOMER_LEVEL_{source_uid[:8]}", "name": "客户等级代码集", "owner_uid": owner_uid, "business_domain_uid": source_uid, "definition": "客户等级统一参考数据", "values": [ {"code": "VIP", "name": "重要客户"}, { "code": "VIP_A", "name": "A级重要客户", "parent_code": "VIP", }, ], "standard_version_uids": [standard_version_uid], }, actor_uid=actor_uid, ) semantic_uids.append(code_set["uid"]) service.submit(code_set["uid"], expected_version=1, actor_uid=actor_uid) service.review( code_set["uid"], expected_version=1, decision="approve", reason="参考数据责任人确认", actor_uid=reviewer_uid, ) service.publish( code_set["uid"], expected_version=1, actor_uid=reviewer_uid ) term = service.create_draft( "business_term", { "code": f"TERM_CUSTOMER_LEVEL_{source_uid[:8]}", "name": "客户等级", "owner_uid": owner_uid, "business_domain_uid": source_uid, "definition": "客户运营分层等级", "aliases": ["客群等级", "客户分层"], "related_asset_uids": [physical_asset_uid], "related_data_element_uids": [element_uid], "standard_version_uids": [standard_version_uid], "code_references": [ {"code_set_uid": code_set["uid"], "code": "VIP"} ], }, actor_uid=actor_uid, ) semantic_uids.append(term["uid"]) service.submit(term["uid"], expected_version=1, actor_uid=actor_uid) service.review( term["uid"], expected_version=1, decision="approve", reason="术语定义审核通过", actor_uid=reviewer_uid, ) service.publish(term["uid"], expected_version=1, actor_uid=reviewer_uid) metric_payload = { "code": f"METRIC_VIP_RATE_{source_uid[:8]}", "name": "重要客户占比", "owner_uid": owner_uid, "business_domain_uid": source_uid, "definition": "重要客户数占全部客户数的比例", "formula": "VIP客户数 / 全部客户数", "unit": "%", "dimensions": [{"code": "region", "name": "区域"}], "related_asset_uids": [physical_asset_uid], "standard_version_uids": [standard_version_uid], } metric = service.create_draft( "metric", metric_payload, actor_uid=actor_uid ) semantic_uids.append(metric["uid"]) for semantic_uid in (metric["uid"],): service.submit( semantic_uid, expected_version=1, actor_uid=actor_uid ) service.review( semantic_uid, expected_version=1, decision="approve", reason="指标口径审核通过", actor_uid=reviewer_uid, ) service.publish( semantic_uid, expected_version=1, actor_uid=reviewer_uid, ) mapping = service.map_physical_field( { "asset_uid": physical_asset_uid, "field_name": "customer_level", "data_element_uid": element_uid, "owner_uid": owner_uid, "evidence": {"source": "catalog"}, }, actor_uid=actor_uid, ) revised = service.revise( metric["uid"], {**metric_payload, "formula": "VIP客户数 / 有效客户数"}, expected_version=1, reason="排除无效客户", actor_uid=actor_uid, ) service.submit( metric["uid"], expected_version=2, actor_uid=actor_uid, ) service.review( metric["uid"], expected_version=2, decision="approve", reason="口径修订审核通过", actor_uid=reviewer_uid, ) service.publish( metric["uid"], expected_version=2, actor_uid=reviewer_uid, ) rolled_back = service.rollback( metric["uid"], target_version=1, expected_version=2, reason="恢复已验收口径", actor_uid=reviewer_uid, ) db.session.commit() assert mapping["data_element_version"] == 1 assert revised["current_version"] == 2 assert rolled_back["current_version"] == 3 assert rolled_back["definition"]["formula"] == ( "VIP客户数 / 全部客户数" ) assert [ item["status"] for item in service.list_versions(metric["uid"]) ] == ["superseded", "superseded", "published"] assert len(service.list_audits(metric["uid"])) == 9 event_count = db.session.execute( text( """ SELECT count(*) FROM public.outbox_events WHERE aggregate_type = 'semantic_asset' AND aggregate_id = :aggregate_id AND event_type = 'semantic_asset.version_published' """ ), {"aggregate_id": metric["uid"]}, ).scalar_one() assert int(event_count) == 3 standard_event_count = db.session.execute( text( """ SELECT count(*) FROM public.outbox_events WHERE aggregate_type = 'data_standard' AND aggregate_id = :aggregate_id AND event_type = 'data_standard.version_published' """ ), {"aggregate_id": standard_uid}, ).scalar_one() assert int(standard_event_count) == 1 claimed = claim_outbox( db.session, limit=10, event_types=["data_standard.version_published"], ) assert len(claimed) == 1 assert claimed[0].event_type == "data_standard.version_published" db.session.rollback() finally: with app.app_context(): db.session.rollback() if semantic_uids: params = {"uids": semantic_uids} db.session.execute( text( "DELETE FROM public.outbox_events " "WHERE aggregate_type = 'semantic_asset' " "AND aggregate_id = ANY(:uids)" ), params, ) for table in ( "semantic_asset_reviews", "semantic_asset_links", "semantic_asset_versions", ): db.session.execute( text( f"DELETE FROM public.{table} " "WHERE asset_uid = ANY(CAST(:uids AS uuid[]))" ), params, ) db.session.execute( text( "DELETE FROM public.semantic_publication_audits " "WHERE target_uid = ANY(CAST(:uids AS uuid[]))" ), params, ) db.session.execute( text( "DELETE FROM public.semantic_assets " "WHERE uid = ANY(CAST(:uids AS uuid[]))" ), params, ) db.session.execute( text( "DELETE FROM public.semantic_publication_audits " "WHERE target_type = 'field_mapping' " "AND target_uid IN (" "SELECT uid FROM public.data_element_field_mappings " "WHERE asset_uid = CAST(:asset_uid AS uuid))" ), {"asset_uid": physical_asset_uid}, ) db.session.execute( text( "DELETE FROM public.data_element_field_mappings " "WHERE asset_uid = CAST(:asset_uid AS uuid)" ), {"asset_uid": physical_asset_uid}, ) db.session.execute( text( "DELETE FROM public.outbox_events " "WHERE aggregate_type = 'data_standard' " "AND aggregate_id = :aggregate_id" ), {"aggregate_id": standard_uid}, ) db.session.execute( text( "DELETE FROM public.semantic_publication_audits " "WHERE target_type = 'data_standard' " "AND target_uid = CAST(:target_uid AS uuid)" ), {"target_uid": standard_uid}, ) db.session.execute( text( "DELETE FROM public.data_standard_versions " "WHERE id = CAST(:uid AS uuid)" ), {"uid": standard_version_uid}, ) db.session.execute( text( "DELETE FROM public.data_standards " "WHERE id = CAST(:uid AS uuid)" ), {"uid": standard_id}, ) db.session.execute( text( "DELETE FROM public.data_element_versions " "WHERE data_element_uid = CAST(:uid AS uuid)" ), {"uid": element_uid}, ) db.session.execute( text( "DELETE FROM public.data_elements " "WHERE uid = CAST(:uid AS uuid)" ), {"uid": element_uid}, ) db.session.execute( text( "DELETE FROM public.active_metadata_assets " "WHERE uid = CAST(:uid AS uuid)" ), {"uid": physical_asset_uid}, ) db.session.execute( text( "DELETE FROM public.active_metadata_runs " "WHERE uid = CAST(:uid AS uuid)" ), {"uid": run_uid}, ) db.session.execute( text( "DELETE FROM public.active_metadata_plans " "WHERE uid = CAST(:uid AS uuid)" ), {"uid": plan_uid}, ) db.session.execute( text( "DELETE FROM public.ingestion_sources " "WHERE uid = CAST(:uid AS uuid)" ), {"uid": source_uid}, ) db.session.execute( text( "DELETE FROM public.users " "WHERE id = ANY(CAST(:uids AS uuid[]))" ), {"uids": [actor_uid, owner_uid, reviewer_uid]}, ) db.session.commit()