| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346 |
- 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()
|