test_product_governance_postgres.py 16 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346
  1. from __future__ import annotations
  2. import os
  3. import uuid
  4. from datetime import UTC, datetime, timedelta
  5. import pytest
  6. from sqlalchemy import create_engine, text
  7. from sqlalchemy.orm import Session
  8. from app.core.data_service.product_governance import ProductGovernanceService
  9. from app.core.data_service.product_governance_repository import (
  10. SqlAlchemyProductGovernanceRepository,
  11. WorkCenterProductApprovalGateway,
  12. )
  13. from app.core.governance.work_center import UnifiedWorkCenterService
  14. from app.core.governance.work_center_repository import SqlAlchemyWorkCenterRepository
  15. pytestmark = pytest.mark.integration
  16. def _uid():
  17. return str(uuid.uuid4())
  18. def _seed_certificate_evidence(session, owner_uid, domain_uid):
  19. now = datetime.now(UTC)
  20. ids = {name: _uid() for name in (
  21. "source", "plan", "metadata_run", "asset", "lineage", "template",
  22. "template_version", "quality_run", "sla_freshness", "sla_quality",
  23. "rule", "rule_version", "dataflow", "workflow_version", "workflow_run",
  24. )}
  25. session.execute(text("""
  26. INSERT INTO public.ingestion_sources (
  27. uid, source_type, name, config, permission_scope, status, created_by
  28. ) VALUES (
  29. CAST(:source AS uuid), 'database', 'WP08 second-domain source',
  30. '{}'::jsonb, '{}'::jsonb, 'active', :owner
  31. )
  32. """), {**ids, "owner": owner_uid})
  33. session.execute(text("""
  34. INSERT INTO public.active_metadata_plans (
  35. uid, source_uid, name, source_kind, schedule_type, discovery_mode,
  36. scope, cursor_state, owner_uid, enabled, current_version, created_by
  37. ) VALUES (
  38. CAST(:plan AS uuid), CAST(:source AS uuid), 'WP08 metadata plan',
  39. 'database', 'manual', 'snapshot', '{}'::jsonb, '{}'::jsonb,
  40. CAST(:owner AS uuid), TRUE, 1, CAST(:owner AS uuid)
  41. )
  42. """), {**ids, "owner": owner_uid})
  43. session.execute(text("""
  44. INSERT INTO public.active_metadata_runs (
  45. uid, plan_uid, batch_key, status, attempt_count, cursor_before,
  46. cursor_after, snapshot_hash, statistics, actor_uid, started_at, finished_at
  47. ) VALUES (
  48. CAST(:metadata_run AS uuid), CAST(:plan AS uuid), 'wp08-batch',
  49. 'completed', 1, '{}'::jsonb, '{}'::jsonb, :hash,
  50. '{}'::jsonb, CAST(:owner AS uuid), :started_at, :finished_at
  51. )
  52. """), {**ids, "owner": owner_uid, "hash": "a" * 64, "started_at": now - timedelta(minutes=5), "finished_at": now})
  53. session.execute(text("""
  54. INSERT INTO public.active_metadata_assets (
  55. uid, source_uid, asset_key, namespace, name, asset_type,
  56. lifecycle_status, current_version, content_hash, snapshot,
  57. health, last_run_uid
  58. ) VALUES (
  59. CAST(:asset AS uuid), CAST(:source AS uuid), :asset_key,
  60. 'equipment_ops', 'device_health', 'table', 'active', 1, :hash,
  61. CAST(:snapshot AS jsonb), '{}'::jsonb, CAST(:metadata_run AS uuid)
  62. )
  63. """), {**ids, "asset_key": f"{ids['source']}:equipment_ops.device_health", "hash": "b" * 64, "snapshot": f'{{"business_domain_uid":"{domain_uid}"}}'})
  64. session.execute(text("""
  65. INSERT INTO public.active_metadata_lineage (
  66. uid, run_uid, parse_status, source_asset, source_field,
  67. target_asset, target_field, relation_type, evidence
  68. ) VALUES (
  69. CAST(:lineage AS uuid), CAST(:metadata_run AS uuid), 'resolved',
  70. 'iot.device_events', 'device_id', 'equipment_ops.device_health',
  71. 'device_id', 'derived_from', '{}'::jsonb
  72. )
  73. """), ids)
  74. session.execute(text("""
  75. INSERT INTO public.quality_templates (
  76. uid, code, name, owner_uid, status, current_version, created_by
  77. ) VALUES (
  78. CAST(:template AS uuid), :code, 'WP08 quality template',
  79. CAST(:owner AS uuid), 'published', 1, CAST(:owner AS uuid)
  80. )
  81. """), {**ids, "owner": owner_uid, "code": f"WP08_{ids['template'][:8].upper()}"})
  82. session.execute(text("""
  83. INSERT INTO public.quality_template_versions (
  84. uid, template_uid, version, status, definition, content_hash,
  85. created_by, published_by, published_at
  86. ) VALUES (
  87. CAST(:template_version AS uuid), CAST(:template AS uuid), 1,
  88. 'published', '{}'::jsonb, :hash, CAST(:owner AS uuid),
  89. CAST(:owner AS uuid), :now
  90. )
  91. """), {**ids, "owner": owner_uid, "hash": "c" * 64, "now": now})
  92. session.execute(text("""
  93. UPDATE public.quality_templates
  94. SET active_version_uid = CAST(:template_version AS uuid)
  95. WHERE uid = CAST(:template AS uuid)
  96. """), ids)
  97. session.execute(text("""
  98. INSERT INTO public.quality_profile_runs (
  99. uid, template_uid, template_version_uid, template_hash, asset_uid,
  100. source_uid, business_domain_uid, batch_key, status, row_count,
  101. score, source_observed_at, comparison, profile, field_bindings,
  102. finding_count, created_by
  103. ) VALUES (
  104. CAST(:quality_run AS uuid), CAST(:template AS uuid),
  105. CAST(:template_version AS uuid), :hash, CAST(:asset AS uuid),
  106. CAST(:source AS uuid), :domain, 'wp08-product-batch', 'success', 2,
  107. 98.50, :now, '{}'::jsonb, '{}'::jsonb, '{}'::jsonb, 0,
  108. CAST(:owner AS uuid)
  109. )
  110. """), {**ids, "owner": owner_uid, "domain": domain_uid, "hash": "c" * 64, "now": now})
  111. for uid_key, sla_type, actual, threshold in (
  112. ("sla_freshness", "freshness", 2, 24),
  113. ("sla_quality", "quality_score", 98.5, 90),
  114. ):
  115. session.execute(text("""
  116. INSERT INTO public.quality_sla_events (
  117. uid, run_uid, asset_uid, sla_type, status, severity, actual,
  118. threshold, owner_uid, escalation_level, evidence
  119. ) VALUES (
  120. CAST(:uid AS uuid), CAST(:quality_run AS uuid), CAST(:asset AS uuid),
  121. :sla_type, 'met', 'info', :actual, :threshold,
  122. CAST(:owner AS uuid), 0, '{}'::jsonb
  123. )
  124. """), {**ids, "uid": ids[uid_key], "sla_type": sla_type, "actual": actual, "threshold": threshold, "owner": owner_uid})
  125. session.execute(text("""
  126. INSERT INTO public.data_rules (
  127. id, rule_uid, name, category, owner_uid, status
  128. ) VALUES (
  129. CAST(:rule AS uuid), CAST(:rule AS uuid), 'WP08 product rule',
  130. 'quality', CAST(:owner AS uuid), 'active'
  131. )
  132. """), {**ids, "owner": owner_uid})
  133. session.execute(text("""
  134. INSERT INTO public.data_rule_versions (
  135. id, rule_uid, version_no, source_text, source_language, rule_spec,
  136. spec_hash, generated_kind, status, created_by, published_at
  137. ) VALUES (
  138. CAST(:rule_version AS uuid), CAST(:rule AS uuid), 1,
  139. '设备编码不能为空', 'zh-CN', '{}'::jsonb, :hash, 'rulespec',
  140. 'published', CAST(:owner AS uuid), :now
  141. )
  142. """), {**ids, "owner": owner_uid, "hash": "d" * 64, "now": now})
  143. session.execute(text("""
  144. INSERT INTO public.dataflow_workflow_versions (
  145. id, dataflow_uid, environment, version_no, n8n_workflow_id,
  146. n8n_workflow_name, definition_hash, definition_snapshot, status,
  147. engine_definition_id, created_by, activated_by, activated_at
  148. ) VALUES (
  149. CAST(:workflow_version AS uuid), CAST(:dataflow AS uuid), 'production',
  150. 1, :workflow_id, 'WP08 product workflow', :hash, '{}'::jsonb,
  151. 'active', :workflow_id, CAST(:owner AS uuid),
  152. CAST(:owner AS uuid), :now
  153. )
  154. """), {**ids, "owner": owner_uid, "workflow_id": f"wp08-{ids['dataflow'][:8]}", "hash": "e" * 64, "now": now})
  155. session.execute(text("""
  156. INSERT INTO public.workflow_runs (
  157. id, workflow_version_id, engine_type, engine_execution_id,
  158. trigger_type, status, correlation_id, started_at, finished_at
  159. ) VALUES (
  160. CAST(:workflow_run AS uuid), CAST(:workflow_version AS uuid), 'n8n',
  161. :execution_id, 'manual', 'success', CAST(:workflow_run AS uuid),
  162. :started_at, :finished_at
  163. )
  164. """), {**ids, "execution_id": f"wp08-{ids['workflow_run']}", "started_at": now - timedelta(minutes=2), "finished_at": now})
  165. session.flush()
  166. return ids
  167. def test_second_domain_product_completes_governed_lifecycle_in_postgres():
  168. database_url = os.environ.get("TEST_DATABASE_URL")
  169. if not database_url:
  170. pytest.skip("TEST_DATABASE_URL is required")
  171. engine = create_engine(database_url)
  172. connection = engine.connect()
  173. transaction = connection.begin()
  174. session = Session(bind=connection)
  175. owner_uid, reviewer_uid, domain_uid = _uid(), _uid(), _uid()
  176. try:
  177. for uid, role in ((owner_uid, "owner"), (reviewer_uid, "reviewer")):
  178. session.execute(text("""
  179. INSERT INTO public.users (
  180. id, username, display_name, password_hash, status
  181. ) VALUES (
  182. CAST(:uid AS uuid), :username, :display_name,
  183. 'wp08-integration-only', 'active'
  184. )
  185. """), {"uid": uid, "username": f"wp08-{role}-{uid[:8]}", "display_name": f"WP08 {role}"})
  186. legacy_product_id = session.execute(text("""
  187. INSERT INTO public.data_products (
  188. product_name, product_name_en, description,
  189. target_table, target_schema, status, created_by
  190. ) VALUES (
  191. '设备健康数据产品', 'Device Health Product',
  192. 'P2-WP08 second-domain integration product',
  193. :table_name, 'equipment_ops', 'active', 'wp08'
  194. ) RETURNING id
  195. """), {"table_name": f"device_health_{domain_uid[:8]}"}).scalar_one()
  196. evidence = _seed_certificate_evidence(session, owner_uid, domain_uid)
  197. center = UnifiedWorkCenterService(
  198. SqlAlchemyWorkCenterRepository(session),
  199. commit=session.flush,
  200. rollback=session.rollback,
  201. )
  202. workflow = center.create_workflow({
  203. "code": f"WP08_{domain_uid[:8].upper()}",
  204. "name": "WP08 data-product approval",
  205. "subject_types": ["data_product"],
  206. "routes": [],
  207. "default_route": {
  208. "approval_mode": "any", "reviewer_uids": [reviewer_uid],
  209. "min_approvals": 1, "due_hours": 24,
  210. "timeout_action": "close", "notification_channels": ["in_app"],
  211. },
  212. }, actor_uid=owner_uid)
  213. workflow = center.publish_workflow(
  214. workflow["uid"], expected_version=1, actor_uid=owner_uid
  215. )
  216. service = ProductGovernanceService(
  217. SqlAlchemyProductGovernanceRepository(session),
  218. approval_gateway=WorkCenterProductApprovalGateway(session),
  219. commit=session.flush,
  220. rollback=session.rollback,
  221. )
  222. product = service.register_product({
  223. "legacy_product_id": legacy_product_id,
  224. "product_code": f"DEVICE_HEALTH_{domain_uid[:8].upper()}",
  225. "name": "设备健康数据产品",
  226. "product_type": "data_product",
  227. "owner_uid": owner_uid,
  228. "business_domain_uid": domain_uid,
  229. "description": "第二业务域设备运行与维护产品",
  230. "quality_target": 90,
  231. "sla": {"availability_target": 99, "freshness_hours": 24, "support_tier": "business_hours"},
  232. }, actor_uid=owner_uid)
  233. application = service.create_application({
  234. "title": "申请设备健康 API 数据",
  235. "source_type": "api",
  236. "source_ref": {"endpoint": "/device-health"},
  237. "business_domain_uid": domain_uid,
  238. "purpose": "设备预测性维护",
  239. "requested_fields": ["device_id", "health_score"],
  240. }, actor_uid=reviewer_uid)
  241. application = service.submit_application(
  242. application["uid"], {"workflow_uid": workflow["uid"]},
  243. expected_version=1, actor_uid=reviewer_uid,
  244. )
  245. task = center.review_task(
  246. application["approval_task_uid"],
  247. {"decision": "approve", "reason": "第二业务域用途与责任清晰"},
  248. expected_version=1, actor_uid=reviewer_uid,
  249. )
  250. application = service.reconcile_application(
  251. application["uid"], expected_version=2, actor_uid=owner_uid
  252. )
  253. application = service.fulfill_application(
  254. application["uid"], product["uid"],
  255. expected_version=3, actor_uid=owner_uid,
  256. )
  257. definition = {
  258. "schema": {"fields": [
  259. {"name": "device_id", "type": "string", "nullable": False},
  260. {"name": "health_score", "type": "number", "nullable": False},
  261. ]},
  262. "delivery": {"mode": "product", "format": "table"},
  263. "quality_terms": {"minimum_score": 90},
  264. "sla_terms": {"freshness_hours": 24, "availability_target": 99},
  265. "usage_terms": {"purpose": "设备预测性维护", "retention_days": 365},
  266. "compatibility_mode": "backward",
  267. "change_reason": "建立第二业务域首版合同",
  268. }
  269. contract = service.create_contract(product["uid"], definition, actor_uid=owner_uid)
  270. contract = service.publish_contract(
  271. contract["uid"], expected_version=1, actor_uid=owner_uid
  272. )
  273. certificate = service.generate_certificate(product["uid"], {
  274. "approval_task_uid": task["uid"],
  275. "evidence": {
  276. "quality_run_uid": evidence["quality_run"],
  277. "lineage_uids": [evidence["lineage"]],
  278. "rule_version_uids": [evidence["rule_version"]],
  279. "workflow_run_uids": [evidence["workflow_run"]],
  280. },
  281. }, actor_uid=owner_uid)
  282. product = service.transition_product(
  283. product["uid"], {"action": "activate", "reason": "合同与合格证门禁已满足"},
  284. expected_version=1, actor_uid=owner_uid,
  285. )
  286. feedback = service.create_feedback(product["uid"], {
  287. "category": "usability", "rating": 4,
  288. "summary": "增加健康分解释", "details": "需要说明评分组成",
  289. }, actor_uid=reviewer_uid)
  290. feedback = service.transition_feedback(feedback["uid"], {
  291. "action": "triage", "assignee_uid": owner_uid, "note": "纳入改进",
  292. }, expected_version=1, actor_uid=owner_uid)
  293. feedback = service.transition_feedback(feedback["uid"], {
  294. "action": "start", "note": "开始补充说明",
  295. }, expected_version=2, actor_uid=owner_uid)
  296. feedback = service.transition_feedback(feedback["uid"], {
  297. "action": "resolve", "note": "产品说明已补充评分口径",
  298. "evidence_refs": ["doc://device-health/score-definition"],
  299. }, expected_version=3, actor_uid=owner_uid)
  300. feedback = service.transition_feedback(feedback["uid"], {
  301. "action": "close", "note": "用户确认改进有效",
  302. }, expected_version=4, actor_uid=reviewer_uid)
  303. detail = service.product_detail(product["uid"])
  304. assert application["status"] == "fulfilled"
  305. assert contract["status"] == "active"
  306. assert certificate["status"] == "qualified"
  307. assert product["status"] == "active"
  308. assert feedback["status"] == "closed"
  309. assert detail["business_domain_uid"] == domain_uid
  310. assert detail["certificates"][0]["evidence_snapshot"]["quality"]["uid"] == evidence["quality_run"]
  311. assert detail["feedback"][0]["evidence_refs"] == ["doc://device-health/score-definition"]
  312. outbox_types = {
  313. row[0]
  314. for row in session.execute(
  315. text("""
  316. SELECT event_type FROM public.outbox_events
  317. WHERE aggregate_type = 'governed_data_product'
  318. AND aggregate_id = :product_uid
  319. """),
  320. {"product_uid": product["uid"]},
  321. )
  322. }
  323. assert {
  324. "data_product.contract.published",
  325. "data_product.certificate.issued",
  326. "data_product.feedback.resolved",
  327. } <= outbox_types
  328. finally:
  329. session.close()
  330. transaction.rollback()
  331. connection.close()
  332. engine.dispose()