test_dataflow_create_saga.py 24 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766
  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. def consume(self):
  13. return None
  14. class GraphSession:
  15. def __init__(
  16. self, *, existing=None, conflict=False, constraint_error=None
  17. ):
  18. self.existing = existing
  19. self.conflict = conflict
  20. self.constraint_error = constraint_error
  21. self.calls = []
  22. def run(self, query, parameters=None, **kwargs):
  23. values = parameters or kwargs
  24. self.calls.append((query, values))
  25. if query.startswith("CREATE CONSTRAINT"):
  26. if (
  27. self.constraint_error is not None
  28. and "data_flow_name_zh" in query
  29. ):
  30. raise self.constraint_error
  31. return GraphResult()
  32. if "WHERE n.uid IS NULL OR n.uid <> $uid" in query:
  33. return GraphResult({"uid": "other"}) if self.conflict else GraphResult()
  34. if query.startswith("MERGE"):
  35. node = copy.deepcopy(self.existing or values["properties"])
  36. self.existing = node
  37. return GraphResult({"n": node, "node_id": 73})
  38. raise AssertionError(query)
  39. def __enter__(self):
  40. return self
  41. def __exit__(self, *_args):
  42. return None
  43. class GraphDriver:
  44. def __init__(self, session):
  45. self.graph_session = session
  46. def session(self):
  47. return self.graph_session
  48. def close(self):
  49. return None
  50. def governed_node():
  51. return {
  52. "uid": new_governance_uid(),
  53. "name_zh": "客户治理生产线",
  54. "name_en": "customer_line",
  55. "script_type": "governed",
  56. "script_requirement": '{"dataflow_spec":{"schema_version":"2.0"}}',
  57. "script_path": "",
  58. }
  59. def test_governed_graph_create_installs_constraint_and_reconciles_same_uid(
  60. monkeypatch,
  61. ):
  62. node = governed_node()
  63. graph = GraphSession()
  64. monkeypatch.setattr(
  65. "app.core.data_flow.dataflows.connect_graph",
  66. lambda: GraphDriver(graph),
  67. )
  68. first_id, first = DataFlowService._merge_governed_dataflow(node)
  69. second_id, second = DataFlowService._merge_governed_dataflow(node)
  70. assert first_id == second_id == 73
  71. assert first == second
  72. assert any(
  73. call[0]
  74. == "CREATE CONSTRAINT data_flow_uid IF NOT EXISTS "
  75. "FOR (n:DataFlow) REQUIRE n.uid IS UNIQUE"
  76. for call in graph.calls
  77. )
  78. assert any(
  79. call[0]
  80. == "CREATE CONSTRAINT data_flow_name_zh IF NOT EXISTS "
  81. "FOR (n:DataFlow) REQUIRE n.name_zh IS UNIQUE"
  82. for call in graph.calls
  83. )
  84. assert sum(call[0].startswith("MERGE") for call in graph.calls) == 2
  85. def test_governed_graph_reconcile_fails_closed_on_uid_or_name_conflict(
  86. monkeypatch,
  87. ):
  88. node = governed_node()
  89. graph = GraphSession(existing={**node, "name_zh": "被篡改的名称"})
  90. monkeypatch.setattr(
  91. "app.core.data_flow.dataflows.connect_graph",
  92. lambda: GraphDriver(graph),
  93. )
  94. with pytest.raises(ValueError, match="dataflow_uid_conflict"):
  95. DataFlowService._merge_governed_dataflow(node)
  96. name_conflict = GraphSession(conflict=True)
  97. monkeypatch.setattr(
  98. "app.core.data_flow.dataflows.connect_graph",
  99. lambda: GraphDriver(name_conflict),
  100. )
  101. with pytest.raises(ValueError, match="dataflow_uid_conflict"):
  102. DataFlowService._merge_governed_dataflow(node)
  103. def test_governed_graph_fails_closed_when_name_constraint_finds_duplicates(
  104. monkeypatch,
  105. ):
  106. class ConstraintFailure(RuntimeError):
  107. code = "Neo.ClientError.Schema.ConstraintValidationFailed"
  108. graph = GraphSession(constraint_error=ConstraintFailure("duplicates"))
  109. monkeypatch.setattr(
  110. "app.core.data_flow.dataflows.connect_graph",
  111. lambda: GraphDriver(graph),
  112. )
  113. with pytest.raises(ValueError, match="dataflow_uid_conflict"):
  114. DataFlowService._merge_governed_dataflow(governed_node())
  115. assert not any(call[0].startswith("MERGE") for call in graph.calls)
  116. def test_completed_saga_replays_result_without_another_neo4j_write(monkeypatch):
  117. expected = {"id": 73, **governed_node()}
  118. class Repository:
  119. def load_published_assets(self, _flow):
  120. return {}, {}
  121. def begin_dataflow_create(
  122. self,
  123. _receipt,
  124. *,
  125. actor_uid,
  126. create_request,
  127. create_intent,
  128. ):
  129. assert actor_uid
  130. assert create_request["intent"] == create_intent
  131. return {
  132. "status": "completed",
  133. "dataflow_uid": expected["uid"],
  134. "result": expected,
  135. "request_digest": "a" * 64,
  136. }
  137. flow = {
  138. "schema_version": "2.0",
  139. "dataflow_uid": expected["uid"],
  140. "name": "客户治理生产线",
  141. "input_schema_refs": ["bd:customer:v1"],
  142. "output_schema_ref": "bd:customer_clean:v1",
  143. "components": [],
  144. }
  145. # Use the project's validator fixture shape through the repository-level
  146. # contract; the replay must happen before any graph write.
  147. from tests.core.data_rules.test_contracts import valid_dataflow_spec
  148. flow = valid_dataflow_spec()
  149. flow["dataflow_uid"] = expected["uid"]
  150. envelope = {
  151. "dataflow_spec": flow,
  152. "dataset_edges": {
  153. "source_table": flow["input_schema_refs"],
  154. "target_table": flow["output_schema_ref"],
  155. },
  156. "migration_metadata": {
  157. "status": "migrated",
  158. "legacy_fields_present": False,
  159. "preserved_for_read_only": True,
  160. "governed_semantics": "dataflow_spec",
  161. },
  162. }
  163. monkeypatch.setattr(
  164. DataFlowService,
  165. "_merge_governed_dataflow",
  166. lambda _node: pytest.fail("Neo4j was called during completed replay"),
  167. )
  168. result = DataFlowService.create_dataflow(
  169. {
  170. "name_zh": "客户治理生产线",
  171. "describe": "响应丢失重试",
  172. "script_type": "governed",
  173. "script_requirement": envelope,
  174. "draft_reservation": {
  175. "reservation_id": new_governance_uid(),
  176. "dataflow_uid": expected["uid"],
  177. "nonce": "response-loss",
  178. },
  179. },
  180. repository=Repository(),
  181. actor_uid=new_governance_uid(),
  182. )
  183. assert result == expected
  184. def test_graph_failure_is_persisted_but_finalize_failure_leaves_reconcilable_lease(
  185. monkeypatch,
  186. ):
  187. from tests.core.data_rules.test_contracts import valid_dataflow_spec
  188. from tests.test_legacy_governance_cutover import PublishedAssetRepository
  189. flow = valid_dataflow_spec()
  190. actor = new_governance_uid()
  191. def payload(repository):
  192. receipt = repository.reserve_dataflow_draft(actor_uid=actor)
  193. receipt["dataflow_uid"] = flow["dataflow_uid"]
  194. return {
  195. "name_zh": "故障注入生产线",
  196. "describe": "故障注入",
  197. "script_type": "governed",
  198. "script_requirement": {
  199. "dataflow_spec": flow,
  200. "dataset_edges": {
  201. "source_table": flow["input_schema_refs"],
  202. "target_table": flow["output_schema_ref"],
  203. },
  204. "migration_metadata": {
  205. "status": "migrated",
  206. "legacy_fields_present": False,
  207. "preserved_for_read_only": True,
  208. "governed_semantics": "dataflow_spec",
  209. },
  210. },
  211. "draft_reservation": {
  212. key: receipt[key]
  213. for key in ("reservation_id", "dataflow_uid", "nonce")
  214. },
  215. }
  216. graph_failure = PublishedAssetRepository()
  217. failures = []
  218. graph_failure.commit_dataflow_create_failure = (
  219. lambda **kwargs: failures.append(kwargs)
  220. )
  221. monkeypatch.setattr(
  222. DataFlowService,
  223. "_merge_governed_dataflow",
  224. lambda _node: (_ for _ in ()).throw(RuntimeError("neo4j unavailable")),
  225. )
  226. with pytest.raises(RuntimeError, match="neo4j unavailable"):
  227. DataFlowService.create_dataflow(
  228. payload(graph_failure), repository=graph_failure, actor_uid=actor
  229. )
  230. assert failures[0]["error_code"] == "neo4j_create_failed"
  231. finalize_failure = PublishedAssetRepository()
  232. monkeypatch.setattr(
  233. DataFlowService,
  234. "_merge_governed_dataflow",
  235. lambda node: (73, {**node, "id": 73}),
  236. )
  237. finalize_failure.complete_dataflow_create = (
  238. lambda **_kwargs: (_ for _ in ()).throw(
  239. RuntimeError("postgres finalize failed")
  240. )
  241. )
  242. with pytest.raises(RuntimeError, match="postgres finalize failed"):
  243. DataFlowService.create_dataflow(
  244. payload(finalize_failure),
  245. repository=finalize_failure,
  246. actor_uid=actor,
  247. )
  248. def test_governed_tag_failure_never_finalizes_completed(monkeypatch):
  249. from tests.core.data_rules.test_contracts import valid_dataflow_spec
  250. from tests.test_legacy_governance_cutover import PublishedAssetRepository
  251. flow = valid_dataflow_spec()
  252. actor = new_governance_uid()
  253. repository = PublishedAssetRepository()
  254. receipt = repository.reserve_dataflow_draft(actor_uid=actor)
  255. receipt["dataflow_uid"] = flow["dataflow_uid"]
  256. failures = []
  257. repository.commit_dataflow_create_failure = (
  258. lambda **kwargs: failures.append(kwargs)
  259. )
  260. monkeypatch.setattr(
  261. DataFlowService,
  262. "_merge_governed_dataflow",
  263. lambda node: (73, {"id": 73, **node}),
  264. )
  265. def fail_tags(_dataflow_id, _tags, *, strict):
  266. assert strict is True
  267. raise RuntimeError("neo4j tag merge failed")
  268. monkeypatch.setattr(
  269. DataFlowService, "_handle_tag_relationships", fail_tags
  270. )
  271. with pytest.raises(RuntimeError, match="neo4j tag merge failed"):
  272. DataFlowService.create_dataflow(
  273. {
  274. "name_zh": "标签严格生产线",
  275. "describe": "严格标签合并验收",
  276. "script_type": "governed",
  277. "script_requirement": {
  278. "dataflow_spec": flow,
  279. "dataset_edges": {
  280. "source_table": flow["input_schema_refs"],
  281. "target_table": flow["output_schema_ref"],
  282. },
  283. "migration_metadata": {
  284. "status": "migrated",
  285. "legacy_fields_present": False,
  286. "preserved_for_read_only": True,
  287. "governed_semantics": "dataflow_spec",
  288. },
  289. },
  290. "tag": [{"id": 81}],
  291. "draft_reservation": {
  292. key: receipt[key]
  293. for key in ("reservation_id", "dataflow_uid", "nonce")
  294. },
  295. },
  296. repository=repository,
  297. actor_uid=actor,
  298. )
  299. assert repository.completed == {}
  300. assert failures[0]["error_code"] == "neo4j_tag_merge_failed"
  301. def test_reconciler_dry_run_is_read_only_and_digest_failure_is_isolated(
  302. monkeypatch,
  303. ):
  304. good_uid = new_governance_uid()
  305. bad_uid = new_governance_uid()
  306. class Session:
  307. def __init__(self):
  308. self.commits = 0
  309. self.rollbacks = 0
  310. def commit(self):
  311. self.commits += 1
  312. def rollback(self):
  313. self.rollbacks += 1
  314. class Repository:
  315. def __init__(self):
  316. self.session = Session()
  317. self.claim_calls = 0
  318. self.failures = []
  319. self.completions = []
  320. self.renewals = []
  321. self.claims = [
  322. {
  323. "reservation_id": new_governance_uid(),
  324. "dataflow_uid": bad_uid,
  325. "lease_token": new_governance_uid(),
  326. "request_digest": "a" * 64,
  327. "create_intent": {},
  328. "attempt": 2,
  329. "integrity_valid": False,
  330. },
  331. {
  332. "reservation_id": new_governance_uid(),
  333. "dataflow_uid": good_uid,
  334. "lease_token": new_governance_uid(),
  335. "request_digest": "b" * 64,
  336. "create_intent": {"dataflow_uid": good_uid},
  337. "attempt": 2,
  338. "integrity_valid": True,
  339. },
  340. ]
  341. def list_reconcilable_dataflow_creates(self, *, limit):
  342. assert limit == 2
  343. return [{"reservation_id": "preview", "state": "failed"}]
  344. def claim_reconcilable_dataflow_creates(
  345. self, *, limit, lease_seconds
  346. ):
  347. assert limit == 1
  348. assert lease_seconds == 60
  349. self.claim_calls += 1
  350. return [self.claims.pop(0)] if self.claims else []
  351. def commit_dataflow_create_claim(self):
  352. self.session.commit()
  353. def renew_dataflow_create_lease(self, **kwargs):
  354. self.renewals.append(kwargs)
  355. def commit_dataflow_create_failure(self, **kwargs):
  356. self.failures.append(kwargs)
  357. self.session.commit()
  358. def complete_dataflow_create(self, **kwargs):
  359. self.completions.append(kwargs)
  360. return {"id": 73, "uid": kwargs["dataflow_uid"]}
  361. repository = Repository()
  362. reconciler = DataFlowCreateReconciler(repository)
  363. preview = reconciler.run(dry_run=True, limit=2)
  364. assert preview["candidate_count"] == 1
  365. assert repository.claim_calls == 0
  366. assert repository.session.commits == 0
  367. monkeypatch.setattr(
  368. DataFlowService,
  369. "validate_governed_create_intent",
  370. lambda value, *, repository: {
  371. "node": {"uid": value["dataflow_uid"]},
  372. "tags": [],
  373. },
  374. )
  375. monkeypatch.setattr(
  376. DataFlowService,
  377. "_merge_governed_dataflow",
  378. lambda node: (73, {"id": 73, **node}),
  379. )
  380. report = reconciler.run(dry_run=False, limit=2)
  381. assert report["reconciled_count"] == 1
  382. assert report["failed_count"] == 1
  383. assert repository.failures[0]["error_code"] == (
  384. "create_intent_integrity_failed"
  385. )
  386. assert repository.completions[0]["dataflow_uid"] == good_uid
  387. assert len(repository.renewals) == 1
  388. def test_reconciler_validation_is_terminal_and_lost_lease_is_reported(
  389. monkeypatch,
  390. ):
  391. flow_uid = new_governance_uid()
  392. class Session:
  393. def commit(self):
  394. return None
  395. def rollback(self):
  396. return None
  397. class Repository:
  398. session = Session()
  399. def __init__(self):
  400. self.claimed = False
  401. self.failure_codes = []
  402. def claim_reconcilable_dataflow_creates(self, **_kwargs):
  403. if self.claimed:
  404. return []
  405. self.claimed = True
  406. return [
  407. {
  408. "reservation_id": new_governance_uid(),
  409. "dataflow_uid": flow_uid,
  410. "lease_token": new_governance_uid(),
  411. "request_digest": "c" * 64,
  412. "create_intent": {"dataflow_uid": flow_uid},
  413. "attempt": 1,
  414. "integrity_valid": True,
  415. }
  416. ]
  417. def commit_dataflow_create_claim(self):
  418. return None
  419. def commit_dataflow_create_failure(self, **kwargs):
  420. self.failure_codes.append(kwargs["error_code"])
  421. repository = Repository()
  422. monkeypatch.setattr(
  423. DataFlowService,
  424. "validate_governed_create_intent",
  425. lambda *_args, **_kwargs: (_ for _ in ()).throw(
  426. ValueError("published asset is unavailable")
  427. ),
  428. )
  429. report = DataFlowCreateReconciler(repository).run(
  430. dry_run=False, limit=2
  431. )
  432. assert report["items"][0] == {
  433. "reservation_id": report["items"][0]["reservation_id"],
  434. "dataflow_uid": flow_uid,
  435. "attempt": 1,
  436. "request_digest": "c" * 64,
  437. "status": "failed",
  438. "error_code": "create_intent_validation_failed",
  439. }
  440. assert repository.failure_codes == ["create_intent_validation_failed"]
  441. repository = Repository()
  442. monkeypatch.setattr(
  443. DataFlowService,
  444. "validate_governed_create_intent",
  445. lambda value, **_kwargs: {"node": value, "tags": []},
  446. )
  447. repository.renew_dataflow_create_lease = lambda **_kwargs: (_ for _ in ()).throw(
  448. RuntimeError("dataflow create lease was lost")
  449. )
  450. repository.commit_dataflow_create_failure = lambda **_kwargs: (
  451. _ for _ in ()
  452. ).throw(RuntimeError("dataflow create lease was lost"))
  453. monkeypatch.setattr(
  454. DataFlowService,
  455. "_merge_governed_dataflow",
  456. lambda _node: pytest.fail("side effect ran after lease loss"),
  457. )
  458. report = DataFlowCreateReconciler(repository).run(
  459. dry_run=False, limit=1
  460. )
  461. assert report["items"][0]["status"] == "lease_recovery_required"
  462. assert report["items"][0]["error_code"] == (
  463. "failure_state_persistence_failed"
  464. )
  465. def test_tag_relationship_uses_one_atomic_merge(monkeypatch):
  466. calls = []
  467. class Result:
  468. def single(self):
  469. return {"merged": 1}
  470. class Session:
  471. def run(self, query, **values):
  472. calls.append((query, values))
  473. return Result()
  474. def __enter__(self):
  475. return self
  476. def __exit__(self, *_args):
  477. return None
  478. class Driver:
  479. def session(self):
  480. return Session()
  481. def close(self):
  482. return None
  483. monkeypatch.setattr(
  484. "app.core.data_flow.dataflows.connect_graph", lambda: Driver()
  485. )
  486. DataFlowService._handle_single_tag_relationship(73, 81)
  487. DataFlowService._handle_single_tag_relationship(73, 81)
  488. assert len(calls) == 2
  489. assert all("MERGE (a)-[:LABEL]->(b)" in query for query, _ in calls)
  490. assert all("CREATE (a)-[:LABEL]->(b)" not in query for query, _ in calls)
  491. def test_reconciler_tag_failure_does_not_complete_and_retry_succeeds(
  492. monkeypatch,
  493. ):
  494. flow_uid = new_governance_uid()
  495. reservation_id = new_governance_uid()
  496. digest = "d" * 64
  497. class Session:
  498. def commit(self):
  499. return None
  500. def rollback(self):
  501. return None
  502. class Repository:
  503. session = Session()
  504. def __init__(self):
  505. self.attempt = 0
  506. self.available = True
  507. self.failures = []
  508. self.completions = []
  509. def claim_reconcilable_dataflow_creates(self, **_kwargs):
  510. if not self.available:
  511. return []
  512. self.available = False
  513. self.attempt += 1
  514. return [
  515. {
  516. "reservation_id": reservation_id,
  517. "dataflow_uid": flow_uid,
  518. "lease_token": new_governance_uid(),
  519. "request_digest": digest,
  520. "create_intent": {"dataflow_uid": flow_uid},
  521. "attempt": self.attempt,
  522. "integrity_valid": True,
  523. }
  524. ]
  525. def commit_dataflow_create_claim(self):
  526. return None
  527. def renew_dataflow_create_lease(self, **_kwargs):
  528. return None
  529. def commit_dataflow_create_failure(self, **kwargs):
  530. self.failures.append(kwargs)
  531. self.available = True
  532. def complete_dataflow_create(self, **kwargs):
  533. self.completions.append(kwargs)
  534. return {"id": 73, "uid": flow_uid}
  535. repository = Repository()
  536. monkeypatch.setattr(
  537. DataFlowService,
  538. "validate_governed_create_intent",
  539. lambda _value, **_kwargs: {
  540. "node": {"uid": flow_uid},
  541. "tags": [{"id": 81}],
  542. },
  543. )
  544. monkeypatch.setattr(
  545. DataFlowService,
  546. "_merge_governed_dataflow",
  547. lambda node: (73, {"id": 73, **node}),
  548. )
  549. attempts = {"count": 0}
  550. def merge_tags(_node_id, _tags, *, strict):
  551. assert strict is True
  552. attempts["count"] += 1
  553. if attempts["count"] == 1:
  554. raise RuntimeError("neo4j tag write failed")
  555. monkeypatch.setattr(
  556. DataFlowService, "_handle_tag_relationships", merge_tags
  557. )
  558. reconciler = DataFlowCreateReconciler(repository)
  559. failed = reconciler.run(dry_run=False, limit=1)
  560. assert failed["items"][0]["status"] == "failed"
  561. assert repository.completions == []
  562. assert repository.failures[0]["error_code"] == "reconciliation_failed"
  563. completed = reconciler.run(dry_run=False, limit=1)
  564. assert completed["items"][0]["status"] == "completed"
  565. assert len(repository.completions) == 1
  566. assert attempts["count"] == 2
  567. def test_slow_reconciler_does_not_preclaim_next_item_from_second_worker(
  568. monkeypatch,
  569. ):
  570. first_uid = new_governance_uid()
  571. second_uid = new_governance_uid()
  572. shared = {
  573. first_uid: {"state": "failed", "attempt": 0},
  574. second_uid: {"state": "failed", "attempt": 0},
  575. }
  576. processed = []
  577. class Session:
  578. def commit(self):
  579. return None
  580. def rollback(self):
  581. return None
  582. class Repository:
  583. session = Session()
  584. def __init__(self, worker):
  585. self.worker = worker
  586. def claim_reconcilable_dataflow_creates(
  587. self, *, limit, lease_seconds
  588. ):
  589. assert limit == 1
  590. assert lease_seconds == 10
  591. available = [
  592. uid
  593. for uid, record in shared.items()
  594. if record["state"] == "failed"
  595. ]
  596. if not available:
  597. return []
  598. uid = available[0]
  599. shared[uid]["state"] = f"creating:{self.worker}"
  600. shared[uid]["attempt"] += 1
  601. return [
  602. {
  603. "reservation_id": new_governance_uid(),
  604. "dataflow_uid": uid,
  605. "lease_token": new_governance_uid(),
  606. "request_digest": ("a" if uid == first_uid else "b") * 64,
  607. "create_intent": {"dataflow_uid": uid},
  608. "attempt": shared[uid]["attempt"],
  609. "integrity_valid": True,
  610. }
  611. ]
  612. def commit_dataflow_create_claim(self):
  613. return None
  614. def renew_dataflow_create_lease(self, **_kwargs):
  615. return None
  616. def complete_dataflow_create(self, *, dataflow_uid, **_kwargs):
  617. assert shared[dataflow_uid]["state"] == f"creating:{self.worker}"
  618. shared[dataflow_uid]["state"] = "completed"
  619. processed.append((self.worker, dataflow_uid))
  620. return {"id": 73, "uid": dataflow_uid}
  621. def commit_dataflow_create_failure(self, **_kwargs):
  622. pytest.fail("unexpected reconciliation failure")
  623. monkeypatch.setattr(
  624. DataFlowService,
  625. "validate_governed_create_intent",
  626. lambda value, **_kwargs: {
  627. "node": {"uid": value["dataflow_uid"]},
  628. "tags": [],
  629. },
  630. )
  631. worker_two = DataFlowCreateReconciler(Repository("worker-2"))
  632. entered_slow_call = {"value": False}
  633. def slow_merge(node):
  634. if node["uid"] == first_uid and not entered_slow_call["value"]:
  635. entered_slow_call["value"] = True
  636. # This models worker 1 crossing its original lease boundary while
  637. # processing item A. Item B was never preclaimed, so worker 2 can
  638. # safely process B without ever sharing ownership of A.
  639. second_report = worker_two.run(
  640. dry_run=False, limit=1, lease_seconds=10
  641. )
  642. assert second_report["items"][0]["dataflow_uid"] == second_uid
  643. return 73, {"id": 73, **node}
  644. monkeypatch.setattr(
  645. DataFlowService, "_merge_governed_dataflow", slow_merge
  646. )
  647. first_report = DataFlowCreateReconciler(Repository("worker-1")).run(
  648. dry_run=False, limit=2, lease_seconds=10
  649. )
  650. assert first_report["reconciled_count"] == 1
  651. assert processed == [
  652. ("worker-2", second_uid),
  653. ("worker-1", first_uid),
  654. ]
  655. assert all(record["attempt"] == 1 for record in shared.values())