test_dataflow_create_saga.py 11 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347
  1. from __future__ import annotations
  2. import copy
  3. import pytest
  4. from app.core.common.identifiers import new_governance_uid
  5. from app.core.data_flow.create_reconciliation import DataFlowCreateReconciler
  6. from app.core.data_flow.dataflows import DataFlowService
  7. class GraphResult:
  8. def __init__(self, record=None):
  9. self.record = record
  10. def single(self):
  11. return self.record
  12. class GraphSession:
  13. def __init__(self, *, existing=None, conflict=False):
  14. self.existing = existing
  15. self.conflict = conflict
  16. self.calls = []
  17. def run(self, query, parameters=None, **kwargs):
  18. values = parameters or kwargs
  19. self.calls.append((query, values))
  20. if query.startswith("CREATE CONSTRAINT"):
  21. return GraphResult()
  22. if "WHERE n.uid IS NULL OR n.uid <> $uid" in query:
  23. return GraphResult({"uid": "other"}) if self.conflict else GraphResult()
  24. if query.startswith("MERGE"):
  25. node = copy.deepcopy(self.existing or values["properties"])
  26. self.existing = node
  27. return GraphResult({"n": node, "node_id": 73})
  28. raise AssertionError(query)
  29. def __enter__(self):
  30. return self
  31. def __exit__(self, *_args):
  32. return None
  33. class GraphDriver:
  34. def __init__(self, session):
  35. self.graph_session = session
  36. def session(self):
  37. return self.graph_session
  38. def close(self):
  39. return None
  40. def governed_node():
  41. return {
  42. "uid": new_governance_uid(),
  43. "name_zh": "客户治理生产线",
  44. "name_en": "customer_line",
  45. "script_type": "governed",
  46. "script_requirement": '{"dataflow_spec":{"schema_version":"2.0"}}',
  47. "script_path": "",
  48. }
  49. def test_governed_graph_create_installs_constraint_and_reconciles_same_uid(
  50. monkeypatch,
  51. ):
  52. node = governed_node()
  53. graph = GraphSession()
  54. monkeypatch.setattr(
  55. "app.core.data_flow.dataflows.connect_graph",
  56. lambda: GraphDriver(graph),
  57. )
  58. first_id, first = DataFlowService._merge_governed_dataflow(node)
  59. second_id, second = DataFlowService._merge_governed_dataflow(node)
  60. assert first_id == second_id == 73
  61. assert first == second
  62. assert any(
  63. call[0]
  64. == "CREATE CONSTRAINT data_flow_uid IF NOT EXISTS "
  65. "FOR (n:DataFlow) REQUIRE n.uid IS UNIQUE"
  66. for call in graph.calls
  67. )
  68. assert sum(call[0].startswith("MERGE") for call in graph.calls) == 2
  69. def test_governed_graph_reconcile_fails_closed_on_uid_or_name_conflict(
  70. monkeypatch,
  71. ):
  72. node = governed_node()
  73. graph = GraphSession(existing={**node, "name_zh": "被篡改的名称"})
  74. monkeypatch.setattr(
  75. "app.core.data_flow.dataflows.connect_graph",
  76. lambda: GraphDriver(graph),
  77. )
  78. with pytest.raises(ValueError, match="dataflow_uid_conflict"):
  79. DataFlowService._merge_governed_dataflow(node)
  80. name_conflict = GraphSession(conflict=True)
  81. monkeypatch.setattr(
  82. "app.core.data_flow.dataflows.connect_graph",
  83. lambda: GraphDriver(name_conflict),
  84. )
  85. with pytest.raises(ValueError, match="dataflow_uid_conflict"):
  86. DataFlowService._merge_governed_dataflow(node)
  87. def test_completed_saga_replays_result_without_another_neo4j_write(monkeypatch):
  88. expected = {"id": 73, **governed_node()}
  89. class Repository:
  90. def load_published_assets(self, _flow):
  91. return {}, {}
  92. def begin_dataflow_create(
  93. self,
  94. _receipt,
  95. *,
  96. actor_uid,
  97. create_request,
  98. create_intent,
  99. ):
  100. assert actor_uid
  101. assert create_request["intent"] == create_intent
  102. return {
  103. "status": "completed",
  104. "dataflow_uid": expected["uid"],
  105. "result": expected,
  106. "request_digest": "a" * 64,
  107. }
  108. flow = {
  109. "schema_version": "2.0",
  110. "dataflow_uid": expected["uid"],
  111. "name": "客户治理生产线",
  112. "input_schema_refs": ["bd:customer:v1"],
  113. "output_schema_ref": "bd:customer_clean:v1",
  114. "components": [],
  115. }
  116. # Use the project's validator fixture shape through the repository-level
  117. # contract; the replay must happen before any graph write.
  118. from tests.core.data_rules.test_contracts import valid_dataflow_spec
  119. flow = valid_dataflow_spec()
  120. flow["dataflow_uid"] = expected["uid"]
  121. envelope = {
  122. "dataflow_spec": flow,
  123. "dataset_edges": {
  124. "source_table": flow["input_schema_refs"],
  125. "target_table": flow["output_schema_ref"],
  126. },
  127. "migration_metadata": {
  128. "status": "migrated",
  129. "legacy_fields_present": False,
  130. "preserved_for_read_only": True,
  131. "governed_semantics": "dataflow_spec",
  132. },
  133. }
  134. monkeypatch.setattr(
  135. DataFlowService,
  136. "_merge_governed_dataflow",
  137. lambda _node: pytest.fail("Neo4j was called during completed replay"),
  138. )
  139. result = DataFlowService.create_dataflow(
  140. {
  141. "name_zh": "客户治理生产线",
  142. "describe": "响应丢失重试",
  143. "script_type": "governed",
  144. "script_requirement": envelope,
  145. "draft_reservation": {
  146. "reservation_id": new_governance_uid(),
  147. "dataflow_uid": expected["uid"],
  148. "nonce": "response-loss",
  149. },
  150. },
  151. repository=Repository(),
  152. actor_uid=new_governance_uid(),
  153. )
  154. assert result == expected
  155. def test_graph_failure_is_persisted_but_finalize_failure_leaves_reconcilable_lease(
  156. monkeypatch,
  157. ):
  158. from tests.core.data_rules.test_contracts import valid_dataflow_spec
  159. from tests.test_legacy_governance_cutover import PublishedAssetRepository
  160. flow = valid_dataflow_spec()
  161. actor = new_governance_uid()
  162. def payload(repository):
  163. receipt = repository.reserve_dataflow_draft(actor_uid=actor)
  164. receipt["dataflow_uid"] = flow["dataflow_uid"]
  165. return {
  166. "name_zh": "故障注入生产线",
  167. "describe": "故障注入",
  168. "script_type": "governed",
  169. "script_requirement": {
  170. "dataflow_spec": flow,
  171. "dataset_edges": {
  172. "source_table": flow["input_schema_refs"],
  173. "target_table": flow["output_schema_ref"],
  174. },
  175. "migration_metadata": {
  176. "status": "migrated",
  177. "legacy_fields_present": False,
  178. "preserved_for_read_only": True,
  179. "governed_semantics": "dataflow_spec",
  180. },
  181. },
  182. "draft_reservation": {
  183. key: receipt[key]
  184. for key in ("reservation_id", "dataflow_uid", "nonce")
  185. },
  186. }
  187. graph_failure = PublishedAssetRepository()
  188. failures = []
  189. graph_failure.commit_dataflow_create_failure = (
  190. lambda **kwargs: failures.append(kwargs)
  191. )
  192. monkeypatch.setattr(
  193. DataFlowService,
  194. "_merge_governed_dataflow",
  195. lambda _node: (_ for _ in ()).throw(RuntimeError("neo4j unavailable")),
  196. )
  197. with pytest.raises(RuntimeError, match="neo4j unavailable"):
  198. DataFlowService.create_dataflow(
  199. payload(graph_failure), repository=graph_failure, actor_uid=actor
  200. )
  201. assert failures[0]["error_code"] == "neo4j_create_failed"
  202. finalize_failure = PublishedAssetRepository()
  203. monkeypatch.setattr(
  204. DataFlowService,
  205. "_merge_governed_dataflow",
  206. lambda node: (73, {**node, "id": 73}),
  207. )
  208. finalize_failure.complete_dataflow_create = (
  209. lambda **_kwargs: (_ for _ in ()).throw(
  210. RuntimeError("postgres finalize failed")
  211. )
  212. )
  213. with pytest.raises(RuntimeError, match="postgres finalize failed"):
  214. DataFlowService.create_dataflow(
  215. payload(finalize_failure),
  216. repository=finalize_failure,
  217. actor_uid=actor,
  218. )
  219. def test_reconciler_dry_run_is_read_only_and_digest_failure_is_isolated(
  220. monkeypatch,
  221. ):
  222. good_uid = new_governance_uid()
  223. bad_uid = new_governance_uid()
  224. class Session:
  225. def __init__(self):
  226. self.commits = 0
  227. self.rollbacks = 0
  228. def commit(self):
  229. self.commits += 1
  230. def rollback(self):
  231. self.rollbacks += 1
  232. class Repository:
  233. def __init__(self):
  234. self.session = Session()
  235. self.claim_calls = 0
  236. self.failures = []
  237. self.completions = []
  238. def list_reconcilable_dataflow_creates(self, *, limit):
  239. assert limit == 2
  240. return [{"reservation_id": "preview", "state": "failed"}]
  241. def claim_reconcilable_dataflow_creates(
  242. self, *, limit, lease_seconds
  243. ):
  244. assert limit == 2
  245. assert lease_seconds == 60
  246. self.claim_calls += 1
  247. return [
  248. {
  249. "reservation_id": new_governance_uid(),
  250. "dataflow_uid": bad_uid,
  251. "lease_token": new_governance_uid(),
  252. "request_digest": "a" * 64,
  253. "create_intent": {},
  254. "attempt": 2,
  255. "integrity_valid": False,
  256. },
  257. {
  258. "reservation_id": new_governance_uid(),
  259. "dataflow_uid": good_uid,
  260. "lease_token": new_governance_uid(),
  261. "request_digest": "b" * 64,
  262. "create_intent": {"dataflow_uid": good_uid},
  263. "attempt": 2,
  264. "integrity_valid": True,
  265. },
  266. ]
  267. def commit_dataflow_create_failure(self, **kwargs):
  268. self.failures.append(kwargs)
  269. self.session.commit()
  270. def complete_dataflow_create(self, **kwargs):
  271. self.completions.append(kwargs)
  272. return {"id": 73, "uid": kwargs["dataflow_uid"]}
  273. repository = Repository()
  274. reconciler = DataFlowCreateReconciler(repository)
  275. preview = reconciler.run(dry_run=True, limit=2)
  276. assert preview["candidate_count"] == 1
  277. assert repository.claim_calls == 0
  278. assert repository.session.commits == 0
  279. monkeypatch.setattr(
  280. DataFlowService,
  281. "validate_governed_create_intent",
  282. lambda value, *, repository: {
  283. "node": {"uid": value["dataflow_uid"]},
  284. "tags": [],
  285. },
  286. )
  287. monkeypatch.setattr(
  288. DataFlowService,
  289. "_merge_governed_dataflow",
  290. lambda node: (73, {"id": 73, **node}),
  291. )
  292. report = reconciler.run(dry_run=False, limit=2)
  293. assert report["reconciled_count"] == 1
  294. assert report["failed_count"] == 1
  295. assert repository.failures[0]["error_code"] == (
  296. "create_intent_integrity_failed"
  297. )
  298. assert repository.completions[0]["dataflow_uid"] == good_uid