test_dataflow_create_saga.py 23 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732
  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_governed_tag_failure_never_finalizes_completed(monkeypatch):
  220. from tests.core.data_rules.test_contracts import valid_dataflow_spec
  221. from tests.test_legacy_governance_cutover import PublishedAssetRepository
  222. flow = valid_dataflow_spec()
  223. actor = new_governance_uid()
  224. repository = PublishedAssetRepository()
  225. receipt = repository.reserve_dataflow_draft(actor_uid=actor)
  226. receipt["dataflow_uid"] = flow["dataflow_uid"]
  227. failures = []
  228. repository.commit_dataflow_create_failure = (
  229. lambda **kwargs: failures.append(kwargs)
  230. )
  231. monkeypatch.setattr(
  232. DataFlowService,
  233. "_merge_governed_dataflow",
  234. lambda node: (73, {"id": 73, **node}),
  235. )
  236. def fail_tags(_dataflow_id, _tags, *, strict):
  237. assert strict is True
  238. raise RuntimeError("neo4j tag merge failed")
  239. monkeypatch.setattr(
  240. DataFlowService, "_handle_tag_relationships", fail_tags
  241. )
  242. with pytest.raises(RuntimeError, match="neo4j tag merge failed"):
  243. DataFlowService.create_dataflow(
  244. {
  245. "name_zh": "标签严格生产线",
  246. "describe": "严格标签合并验收",
  247. "script_type": "governed",
  248. "script_requirement": {
  249. "dataflow_spec": flow,
  250. "dataset_edges": {
  251. "source_table": flow["input_schema_refs"],
  252. "target_table": flow["output_schema_ref"],
  253. },
  254. "migration_metadata": {
  255. "status": "migrated",
  256. "legacy_fields_present": False,
  257. "preserved_for_read_only": True,
  258. "governed_semantics": "dataflow_spec",
  259. },
  260. },
  261. "tag": [{"id": 81}],
  262. "draft_reservation": {
  263. key: receipt[key]
  264. for key in ("reservation_id", "dataflow_uid", "nonce")
  265. },
  266. },
  267. repository=repository,
  268. actor_uid=actor,
  269. )
  270. assert repository.completed == {}
  271. assert failures[0]["error_code"] == "neo4j_tag_merge_failed"
  272. def test_reconciler_dry_run_is_read_only_and_digest_failure_is_isolated(
  273. monkeypatch,
  274. ):
  275. good_uid = new_governance_uid()
  276. bad_uid = new_governance_uid()
  277. class Session:
  278. def __init__(self):
  279. self.commits = 0
  280. self.rollbacks = 0
  281. def commit(self):
  282. self.commits += 1
  283. def rollback(self):
  284. self.rollbacks += 1
  285. class Repository:
  286. def __init__(self):
  287. self.session = Session()
  288. self.claim_calls = 0
  289. self.failures = []
  290. self.completions = []
  291. self.renewals = []
  292. self.claims = [
  293. {
  294. "reservation_id": new_governance_uid(),
  295. "dataflow_uid": bad_uid,
  296. "lease_token": new_governance_uid(),
  297. "request_digest": "a" * 64,
  298. "create_intent": {},
  299. "attempt": 2,
  300. "integrity_valid": False,
  301. },
  302. {
  303. "reservation_id": new_governance_uid(),
  304. "dataflow_uid": good_uid,
  305. "lease_token": new_governance_uid(),
  306. "request_digest": "b" * 64,
  307. "create_intent": {"dataflow_uid": good_uid},
  308. "attempt": 2,
  309. "integrity_valid": True,
  310. },
  311. ]
  312. def list_reconcilable_dataflow_creates(self, *, limit):
  313. assert limit == 2
  314. return [{"reservation_id": "preview", "state": "failed"}]
  315. def claim_reconcilable_dataflow_creates(
  316. self, *, limit, lease_seconds
  317. ):
  318. assert limit == 1
  319. assert lease_seconds == 60
  320. self.claim_calls += 1
  321. return [self.claims.pop(0)] if self.claims else []
  322. def commit_dataflow_create_claim(self):
  323. self.session.commit()
  324. def renew_dataflow_create_lease(self, **kwargs):
  325. self.renewals.append(kwargs)
  326. def commit_dataflow_create_failure(self, **kwargs):
  327. self.failures.append(kwargs)
  328. self.session.commit()
  329. def complete_dataflow_create(self, **kwargs):
  330. self.completions.append(kwargs)
  331. return {"id": 73, "uid": kwargs["dataflow_uid"]}
  332. repository = Repository()
  333. reconciler = DataFlowCreateReconciler(repository)
  334. preview = reconciler.run(dry_run=True, limit=2)
  335. assert preview["candidate_count"] == 1
  336. assert repository.claim_calls == 0
  337. assert repository.session.commits == 0
  338. monkeypatch.setattr(
  339. DataFlowService,
  340. "validate_governed_create_intent",
  341. lambda value, *, repository: {
  342. "node": {"uid": value["dataflow_uid"]},
  343. "tags": [],
  344. },
  345. )
  346. monkeypatch.setattr(
  347. DataFlowService,
  348. "_merge_governed_dataflow",
  349. lambda node: (73, {"id": 73, **node}),
  350. )
  351. report = reconciler.run(dry_run=False, limit=2)
  352. assert report["reconciled_count"] == 1
  353. assert report["failed_count"] == 1
  354. assert repository.failures[0]["error_code"] == (
  355. "create_intent_integrity_failed"
  356. )
  357. assert repository.completions[0]["dataflow_uid"] == good_uid
  358. assert len(repository.renewals) == 1
  359. def test_reconciler_validation_is_terminal_and_lost_lease_is_reported(
  360. monkeypatch,
  361. ):
  362. flow_uid = new_governance_uid()
  363. class Session:
  364. def commit(self):
  365. return None
  366. def rollback(self):
  367. return None
  368. class Repository:
  369. session = Session()
  370. def __init__(self):
  371. self.claimed = False
  372. self.failure_codes = []
  373. def claim_reconcilable_dataflow_creates(self, **_kwargs):
  374. if self.claimed:
  375. return []
  376. self.claimed = True
  377. return [
  378. {
  379. "reservation_id": new_governance_uid(),
  380. "dataflow_uid": flow_uid,
  381. "lease_token": new_governance_uid(),
  382. "request_digest": "c" * 64,
  383. "create_intent": {"dataflow_uid": flow_uid},
  384. "attempt": 1,
  385. "integrity_valid": True,
  386. }
  387. ]
  388. def commit_dataflow_create_claim(self):
  389. return None
  390. def commit_dataflow_create_failure(self, **kwargs):
  391. self.failure_codes.append(kwargs["error_code"])
  392. repository = Repository()
  393. monkeypatch.setattr(
  394. DataFlowService,
  395. "validate_governed_create_intent",
  396. lambda *_args, **_kwargs: (_ for _ in ()).throw(
  397. ValueError("published asset is unavailable")
  398. ),
  399. )
  400. report = DataFlowCreateReconciler(repository).run(
  401. dry_run=False, limit=2
  402. )
  403. assert report["items"][0] == {
  404. "reservation_id": report["items"][0]["reservation_id"],
  405. "dataflow_uid": flow_uid,
  406. "attempt": 1,
  407. "request_digest": "c" * 64,
  408. "status": "failed",
  409. "error_code": "create_intent_validation_failed",
  410. }
  411. assert repository.failure_codes == ["create_intent_validation_failed"]
  412. repository = Repository()
  413. monkeypatch.setattr(
  414. DataFlowService,
  415. "validate_governed_create_intent",
  416. lambda value, **_kwargs: {"node": value, "tags": []},
  417. )
  418. repository.renew_dataflow_create_lease = lambda **_kwargs: (_ for _ in ()).throw(
  419. RuntimeError("dataflow create lease was lost")
  420. )
  421. repository.commit_dataflow_create_failure = lambda **_kwargs: (
  422. _ for _ in ()
  423. ).throw(RuntimeError("dataflow create lease was lost"))
  424. monkeypatch.setattr(
  425. DataFlowService,
  426. "_merge_governed_dataflow",
  427. lambda _node: pytest.fail("side effect ran after lease loss"),
  428. )
  429. report = DataFlowCreateReconciler(repository).run(
  430. dry_run=False, limit=1
  431. )
  432. assert report["items"][0]["status"] == "lease_recovery_required"
  433. assert report["items"][0]["error_code"] == (
  434. "failure_state_persistence_failed"
  435. )
  436. def test_tag_relationship_uses_one_atomic_merge(monkeypatch):
  437. calls = []
  438. class Result:
  439. def single(self):
  440. return {"merged": 1}
  441. class Session:
  442. def run(self, query, **values):
  443. calls.append((query, values))
  444. return Result()
  445. def __enter__(self):
  446. return self
  447. def __exit__(self, *_args):
  448. return None
  449. class Driver:
  450. def session(self):
  451. return Session()
  452. def close(self):
  453. return None
  454. monkeypatch.setattr(
  455. "app.core.data_flow.dataflows.connect_graph", lambda: Driver()
  456. )
  457. DataFlowService._handle_single_tag_relationship(73, 81)
  458. DataFlowService._handle_single_tag_relationship(73, 81)
  459. assert len(calls) == 2
  460. assert all("MERGE (a)-[:LABEL]->(b)" in query for query, _ in calls)
  461. assert all("CREATE (a)-[:LABEL]->(b)" not in query for query, _ in calls)
  462. def test_reconciler_tag_failure_does_not_complete_and_retry_succeeds(
  463. monkeypatch,
  464. ):
  465. flow_uid = new_governance_uid()
  466. reservation_id = new_governance_uid()
  467. digest = "d" * 64
  468. class Session:
  469. def commit(self):
  470. return None
  471. def rollback(self):
  472. return None
  473. class Repository:
  474. session = Session()
  475. def __init__(self):
  476. self.attempt = 0
  477. self.available = True
  478. self.failures = []
  479. self.completions = []
  480. def claim_reconcilable_dataflow_creates(self, **_kwargs):
  481. if not self.available:
  482. return []
  483. self.available = False
  484. self.attempt += 1
  485. return [
  486. {
  487. "reservation_id": reservation_id,
  488. "dataflow_uid": flow_uid,
  489. "lease_token": new_governance_uid(),
  490. "request_digest": digest,
  491. "create_intent": {"dataflow_uid": flow_uid},
  492. "attempt": self.attempt,
  493. "integrity_valid": True,
  494. }
  495. ]
  496. def commit_dataflow_create_claim(self):
  497. return None
  498. def renew_dataflow_create_lease(self, **_kwargs):
  499. return None
  500. def commit_dataflow_create_failure(self, **kwargs):
  501. self.failures.append(kwargs)
  502. self.available = True
  503. def complete_dataflow_create(self, **kwargs):
  504. self.completions.append(kwargs)
  505. return {"id": 73, "uid": flow_uid}
  506. repository = Repository()
  507. monkeypatch.setattr(
  508. DataFlowService,
  509. "validate_governed_create_intent",
  510. lambda _value, **_kwargs: {
  511. "node": {"uid": flow_uid},
  512. "tags": [{"id": 81}],
  513. },
  514. )
  515. monkeypatch.setattr(
  516. DataFlowService,
  517. "_merge_governed_dataflow",
  518. lambda node: (73, {"id": 73, **node}),
  519. )
  520. attempts = {"count": 0}
  521. def merge_tags(_node_id, _tags, *, strict):
  522. assert strict is True
  523. attempts["count"] += 1
  524. if attempts["count"] == 1:
  525. raise RuntimeError("neo4j tag write failed")
  526. monkeypatch.setattr(
  527. DataFlowService, "_handle_tag_relationships", merge_tags
  528. )
  529. reconciler = DataFlowCreateReconciler(repository)
  530. failed = reconciler.run(dry_run=False, limit=1)
  531. assert failed["items"][0]["status"] == "failed"
  532. assert repository.completions == []
  533. assert repository.failures[0]["error_code"] == "reconciliation_failed"
  534. completed = reconciler.run(dry_run=False, limit=1)
  535. assert completed["items"][0]["status"] == "completed"
  536. assert len(repository.completions) == 1
  537. assert attempts["count"] == 2
  538. def test_slow_reconciler_does_not_preclaim_next_item_from_second_worker(
  539. monkeypatch,
  540. ):
  541. first_uid = new_governance_uid()
  542. second_uid = new_governance_uid()
  543. shared = {
  544. first_uid: {"state": "failed", "attempt": 0},
  545. second_uid: {"state": "failed", "attempt": 0},
  546. }
  547. processed = []
  548. class Session:
  549. def commit(self):
  550. return None
  551. def rollback(self):
  552. return None
  553. class Repository:
  554. session = Session()
  555. def __init__(self, worker):
  556. self.worker = worker
  557. def claim_reconcilable_dataflow_creates(
  558. self, *, limit, lease_seconds
  559. ):
  560. assert limit == 1
  561. assert lease_seconds == 10
  562. available = [
  563. uid
  564. for uid, record in shared.items()
  565. if record["state"] == "failed"
  566. ]
  567. if not available:
  568. return []
  569. uid = available[0]
  570. shared[uid]["state"] = f"creating:{self.worker}"
  571. shared[uid]["attempt"] += 1
  572. return [
  573. {
  574. "reservation_id": new_governance_uid(),
  575. "dataflow_uid": uid,
  576. "lease_token": new_governance_uid(),
  577. "request_digest": ("a" if uid == first_uid else "b") * 64,
  578. "create_intent": {"dataflow_uid": uid},
  579. "attempt": shared[uid]["attempt"],
  580. "integrity_valid": True,
  581. }
  582. ]
  583. def commit_dataflow_create_claim(self):
  584. return None
  585. def renew_dataflow_create_lease(self, **_kwargs):
  586. return None
  587. def complete_dataflow_create(self, *, dataflow_uid, **_kwargs):
  588. assert shared[dataflow_uid]["state"] == f"creating:{self.worker}"
  589. shared[dataflow_uid]["state"] = "completed"
  590. processed.append((self.worker, dataflow_uid))
  591. return {"id": 73, "uid": dataflow_uid}
  592. def commit_dataflow_create_failure(self, **_kwargs):
  593. pytest.fail("unexpected reconciliation failure")
  594. monkeypatch.setattr(
  595. DataFlowService,
  596. "validate_governed_create_intent",
  597. lambda value, **_kwargs: {
  598. "node": {"uid": value["dataflow_uid"]},
  599. "tags": [],
  600. },
  601. )
  602. worker_two = DataFlowCreateReconciler(Repository("worker-2"))
  603. entered_slow_call = {"value": False}
  604. def slow_merge(node):
  605. if node["uid"] == first_uid and not entered_slow_call["value"]:
  606. entered_slow_call["value"] = True
  607. # This models worker 1 crossing its original lease boundary while
  608. # processing item A. Item B was never preclaimed, so worker 2 can
  609. # safely process B without ever sharing ownership of A.
  610. second_report = worker_two.run(
  611. dry_run=False, limit=1, lease_seconds=10
  612. )
  613. assert second_report["items"][0]["dataflow_uid"] == second_uid
  614. return 73, {"id": 73, **node}
  615. monkeypatch.setattr(
  616. DataFlowService, "_merge_governed_dataflow", slow_merge
  617. )
  618. first_report = DataFlowCreateReconciler(Repository("worker-1")).run(
  619. dry_run=False, limit=2, lease_seconds=10
  620. )
  621. assert first_report["reconciled_count"] == 1
  622. assert processed == [
  623. ("worker-2", second_uid),
  624. ("worker-1", first_uid),
  625. ]
  626. assert all(record["attempt"] == 1 for record in shared.values())