from __future__ import annotations import os import uuid from datetime import UTC, datetime, timedelta import pytest from sqlalchemy import create_engine, text from sqlalchemy.orm import Session from app.core.data_service.product_governance import ProductGovernanceService from app.core.data_service.product_governance_repository import ( SqlAlchemyProductGovernanceRepository, WorkCenterProductApprovalGateway, ) from app.core.governance.work_center import UnifiedWorkCenterService from app.core.governance.work_center_repository import SqlAlchemyWorkCenterRepository pytestmark = pytest.mark.integration def _uid(): return str(uuid.uuid4()) def _seed_certificate_evidence(session, owner_uid, domain_uid): now = datetime.now(UTC) ids = {name: _uid() for name in ( "source", "plan", "metadata_run", "asset", "lineage", "template", "template_version", "quality_run", "sla_freshness", "sla_quality", "rule", "rule_version", "dataflow", "workflow_version", "workflow_run", )} session.execute(text(""" INSERT INTO public.ingestion_sources ( uid, source_type, name, config, permission_scope, status, created_by ) VALUES ( CAST(:source AS uuid), 'database', 'WP08 second-domain source', '{}'::jsonb, '{}'::jsonb, 'active', :owner ) """), {**ids, "owner": owner_uid}) session.execute(text(""" INSERT INTO public.active_metadata_plans ( uid, source_uid, name, source_kind, schedule_type, discovery_mode, scope, cursor_state, owner_uid, enabled, current_version, created_by ) VALUES ( CAST(:plan AS uuid), CAST(:source AS uuid), 'WP08 metadata plan', 'database', 'manual', 'snapshot', '{}'::jsonb, '{}'::jsonb, CAST(:owner AS uuid), TRUE, 1, CAST(:owner AS uuid) ) """), {**ids, "owner": owner_uid}) session.execute(text(""" INSERT INTO public.active_metadata_runs ( uid, plan_uid, batch_key, status, attempt_count, cursor_before, cursor_after, snapshot_hash, statistics, actor_uid, started_at, finished_at ) VALUES ( CAST(:metadata_run AS uuid), CAST(:plan AS uuid), 'wp08-batch', 'completed', 1, '{}'::jsonb, '{}'::jsonb, :hash, '{}'::jsonb, CAST(:owner AS uuid), :started_at, :finished_at ) """), {**ids, "owner": owner_uid, "hash": "a" * 64, "started_at": now - timedelta(minutes=5), "finished_at": now}) 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(:asset AS uuid), CAST(:source AS uuid), :asset_key, 'equipment_ops', 'device_health', 'table', 'active', 1, :hash, CAST(:snapshot AS jsonb), '{}'::jsonb, CAST(:metadata_run AS uuid) ) """), {**ids, "asset_key": f"{ids['source']}:equipment_ops.device_health", "hash": "b" * 64, "snapshot": f'{{"business_domain_uid":"{domain_uid}"}}'}) session.execute(text(""" INSERT INTO public.active_metadata_lineage ( uid, run_uid, parse_status, source_asset, source_field, target_asset, target_field, relation_type, evidence ) VALUES ( CAST(:lineage AS uuid), CAST(:metadata_run AS uuid), 'resolved', 'iot.device_events', 'device_id', 'equipment_ops.device_health', 'device_id', 'derived_from', '{}'::jsonb ) """), ids) session.execute(text(""" INSERT INTO public.quality_templates ( uid, code, name, owner_uid, status, current_version, created_by ) VALUES ( CAST(:template AS uuid), :code, 'WP08 quality template', CAST(:owner AS uuid), 'published', 1, CAST(:owner AS uuid) ) """), {**ids, "owner": owner_uid, "code": f"WP08_{ids['template'][:8].upper()}"}) session.execute(text(""" INSERT INTO public.quality_template_versions ( uid, template_uid, version, status, definition, content_hash, created_by, published_by, published_at ) VALUES ( CAST(:template_version AS uuid), CAST(:template AS uuid), 1, 'published', '{}'::jsonb, :hash, CAST(:owner AS uuid), CAST(:owner AS uuid), :now ) """), {**ids, "owner": owner_uid, "hash": "c" * 64, "now": now}) session.execute(text(""" UPDATE public.quality_templates SET active_version_uid = CAST(:template_version AS uuid) WHERE uid = CAST(:template AS uuid) """), ids) session.execute(text(""" INSERT INTO public.quality_profile_runs ( uid, template_uid, template_version_uid, template_hash, asset_uid, source_uid, business_domain_uid, batch_key, status, row_count, score, source_observed_at, comparison, profile, field_bindings, finding_count, created_by ) VALUES ( CAST(:quality_run AS uuid), CAST(:template AS uuid), CAST(:template_version AS uuid), :hash, CAST(:asset AS uuid), CAST(:source AS uuid), :domain, 'wp08-product-batch', 'success', 2, 98.50, :now, '{}'::jsonb, '{}'::jsonb, '{}'::jsonb, 0, CAST(:owner AS uuid) ) """), {**ids, "owner": owner_uid, "domain": domain_uid, "hash": "c" * 64, "now": now}) for uid_key, sla_type, actual, threshold in ( ("sla_freshness", "freshness", 2, 24), ("sla_quality", "quality_score", 98.5, 90), ): session.execute(text(""" INSERT INTO public.quality_sla_events ( uid, run_uid, asset_uid, sla_type, status, severity, actual, threshold, owner_uid, escalation_level, evidence ) VALUES ( CAST(:uid AS uuid), CAST(:quality_run AS uuid), CAST(:asset AS uuid), :sla_type, 'met', 'info', :actual, :threshold, CAST(:owner AS uuid), 0, '{}'::jsonb ) """), {**ids, "uid": ids[uid_key], "sla_type": sla_type, "actual": actual, "threshold": threshold, "owner": owner_uid}) session.execute(text(""" INSERT INTO public.data_rules ( id, rule_uid, name, category, owner_uid, status ) VALUES ( CAST(:rule AS uuid), CAST(:rule AS uuid), 'WP08 product rule', 'quality', CAST(:owner AS uuid), 'active' ) """), {**ids, "owner": owner_uid}) session.execute(text(""" INSERT INTO public.data_rule_versions ( id, rule_uid, version_no, source_text, source_language, rule_spec, spec_hash, generated_kind, status, created_by, published_at ) VALUES ( CAST(:rule_version AS uuid), CAST(:rule AS uuid), 1, '设备编码不能为空', 'zh-CN', '{}'::jsonb, :hash, 'rulespec', 'published', CAST(:owner AS uuid), :now ) """), {**ids, "owner": owner_uid, "hash": "d" * 64, "now": now}) session.execute(text(""" INSERT INTO public.dataflow_workflow_versions ( id, dataflow_uid, environment, version_no, n8n_workflow_id, n8n_workflow_name, definition_hash, definition_snapshot, status, engine_definition_id, created_by, activated_by, activated_at ) VALUES ( CAST(:workflow_version AS uuid), CAST(:dataflow AS uuid), 'production', 1, :workflow_id, 'WP08 product workflow', :hash, '{}'::jsonb, 'active', :workflow_id, CAST(:owner AS uuid), CAST(:owner AS uuid), :now ) """), {**ids, "owner": owner_uid, "workflow_id": f"wp08-{ids['dataflow'][:8]}", "hash": "e" * 64, "now": now}) session.execute(text(""" INSERT INTO public.workflow_runs ( id, workflow_version_id, engine_type, engine_execution_id, trigger_type, status, correlation_id, started_at, finished_at ) VALUES ( CAST(:workflow_run AS uuid), CAST(:workflow_version AS uuid), 'n8n', :execution_id, 'manual', 'success', CAST(:workflow_run AS uuid), :started_at, :finished_at ) """), {**ids, "execution_id": f"wp08-{ids['workflow_run']}", "started_at": now - timedelta(minutes=2), "finished_at": now}) session.flush() return ids def test_second_domain_product_completes_governed_lifecycle_in_postgres(): database_url = os.environ.get("TEST_DATABASE_URL") if not database_url: pytest.skip("TEST_DATABASE_URL is required") engine = create_engine(database_url) connection = engine.connect() transaction = connection.begin() session = Session(bind=connection) owner_uid, reviewer_uid, domain_uid = _uid(), _uid(), _uid() try: for uid, role in ((owner_uid, "owner"), (reviewer_uid, "reviewer")): session.execute(text(""" INSERT INTO public.users ( id, username, display_name, password_hash, status ) VALUES ( CAST(:uid AS uuid), :username, :display_name, 'wp08-integration-only', 'active' ) """), {"uid": uid, "username": f"wp08-{role}-{uid[:8]}", "display_name": f"WP08 {role}"}) legacy_product_id = session.execute(text(""" INSERT INTO public.data_products ( product_name, product_name_en, description, target_table, target_schema, status, created_by ) VALUES ( '设备健康数据产品', 'Device Health Product', 'P2-WP08 second-domain integration product', :table_name, 'equipment_ops', 'active', 'wp08' ) RETURNING id """), {"table_name": f"device_health_{domain_uid[:8]}"}).scalar_one() evidence = _seed_certificate_evidence(session, owner_uid, domain_uid) center = UnifiedWorkCenterService( SqlAlchemyWorkCenterRepository(session), commit=session.flush, rollback=session.rollback, ) workflow = center.create_workflow({ "code": f"WP08_{domain_uid[:8].upper()}", "name": "WP08 data-product approval", "subject_types": ["data_product"], "routes": [], "default_route": { "approval_mode": "any", "reviewer_uids": [reviewer_uid], "min_approvals": 1, "due_hours": 24, "timeout_action": "close", "notification_channels": ["in_app"], }, }, actor_uid=owner_uid) workflow = center.publish_workflow( workflow["uid"], expected_version=1, actor_uid=owner_uid ) service = ProductGovernanceService( SqlAlchemyProductGovernanceRepository(session), approval_gateway=WorkCenterProductApprovalGateway(session), commit=session.flush, rollback=session.rollback, ) product = service.register_product({ "legacy_product_id": legacy_product_id, "product_code": f"DEVICE_HEALTH_{domain_uid[:8].upper()}", "name": "设备健康数据产品", "product_type": "data_product", "owner_uid": owner_uid, "business_domain_uid": domain_uid, "description": "第二业务域设备运行与维护产品", "quality_target": 90, "sla": {"availability_target": 99, "freshness_hours": 24, "support_tier": "business_hours"}, }, actor_uid=owner_uid) application = service.create_application({ "title": "申请设备健康 API 数据", "source_type": "api", "source_ref": {"endpoint": "/device-health"}, "business_domain_uid": domain_uid, "purpose": "设备预测性维护", "requested_fields": ["device_id", "health_score"], }, actor_uid=reviewer_uid) application = service.submit_application( application["uid"], {"workflow_uid": workflow["uid"]}, expected_version=1, actor_uid=reviewer_uid, ) task = center.review_task( application["approval_task_uid"], {"decision": "approve", "reason": "第二业务域用途与责任清晰"}, expected_version=1, actor_uid=reviewer_uid, ) application = service.reconcile_application( application["uid"], expected_version=2, actor_uid=owner_uid ) application = service.fulfill_application( application["uid"], product["uid"], expected_version=3, actor_uid=owner_uid, ) definition = { "schema": {"fields": [ {"name": "device_id", "type": "string", "nullable": False}, {"name": "health_score", "type": "number", "nullable": False}, ]}, "delivery": {"mode": "product", "format": "table"}, "quality_terms": {"minimum_score": 90}, "sla_terms": {"freshness_hours": 24, "availability_target": 99}, "usage_terms": {"purpose": "设备预测性维护", "retention_days": 365}, "compatibility_mode": "backward", "change_reason": "建立第二业务域首版合同", } contract = service.create_contract(product["uid"], definition, actor_uid=owner_uid) contract = service.publish_contract( contract["uid"], expected_version=1, actor_uid=owner_uid ) certificate = service.generate_certificate(product["uid"], { "approval_task_uid": task["uid"], "evidence": { "quality_run_uid": evidence["quality_run"], "lineage_uids": [evidence["lineage"]], "rule_version_uids": [evidence["rule_version"]], "workflow_run_uids": [evidence["workflow_run"]], }, }, actor_uid=owner_uid) product = service.transition_product( product["uid"], {"action": "activate", "reason": "合同与合格证门禁已满足"}, expected_version=1, actor_uid=owner_uid, ) feedback = service.create_feedback(product["uid"], { "category": "usability", "rating": 4, "summary": "增加健康分解释", "details": "需要说明评分组成", }, actor_uid=reviewer_uid) feedback = service.transition_feedback(feedback["uid"], { "action": "triage", "assignee_uid": owner_uid, "note": "纳入改进", }, expected_version=1, actor_uid=owner_uid) feedback = service.transition_feedback(feedback["uid"], { "action": "start", "note": "开始补充说明", }, expected_version=2, actor_uid=owner_uid) feedback = service.transition_feedback(feedback["uid"], { "action": "resolve", "note": "产品说明已补充评分口径", "evidence_refs": ["doc://device-health/score-definition"], }, expected_version=3, actor_uid=owner_uid) feedback = service.transition_feedback(feedback["uid"], { "action": "close", "note": "用户确认改进有效", }, expected_version=4, actor_uid=reviewer_uid) detail = service.product_detail(product["uid"]) assert application["status"] == "fulfilled" assert contract["status"] == "active" assert certificate["status"] == "qualified" assert product["status"] == "active" assert feedback["status"] == "closed" assert detail["business_domain_uid"] == domain_uid assert detail["certificates"][0]["evidence_snapshot"]["quality"]["uid"] == evidence["quality_run"] assert detail["feedback"][0]["evidence_refs"] == ["doc://device-health/score-definition"] outbox_types = { row[0] for row in session.execute( text(""" SELECT event_type FROM public.outbox_events WHERE aggregate_type = 'governed_data_product' AND aggregate_id = :product_uid """), {"product_uid": product["uid"]}, ) } assert { "data_product.contract.published", "data_product.certificate.issued", "data_product.feedback.resolved", } <= outbox_types finally: session.close() transaction.rollback() connection.close() engine.dispose()