Browse Source

fix: fence dataflow create reconciliation

马小龙 4 weeks ago
parent
commit
21c4f89883

+ 46 - 0
.superpowers/sdd/task-8-report.md

@@ -325,6 +325,52 @@ Fourth-review acceptance:
 - full repository suite:
   `673 passed, 33 skipped, 59 subtests passed`.
 
+## Fenced reconciliation and strict-intent closeout
+
+The fifth review removed the remaining concurrency and poison-record
+failure modes from orphan reconciliation:
+
+- apply mode is bounded by the requested total but claims only one row at a
+  time, commits that claim, processes it and only then claims the next row;
+  a slow first item therefore cannot leave a batch of later items waiting on
+  leases that expire before their processing begins;
+- after the trusted intent is semantically validated and immediately before
+  the Neo4j side effect, the reconciler conditionally renews the lease using
+  the reservation ID and current lease token. Lost ownership prevents the
+  side effect and is reported as `lease_recovery_required`;
+- reconciliation parsing is isolated per claimed row. Scalar, malformed-shape
+  or oversized request/intent JSON becomes terminal
+  `create_intent_integrity_failed` without aborting the next valid item;
+- semantic intent validation failures become terminal
+  `create_intent_validation_failed`. Both terminal codes are excluded from
+  later list and claim queries, while infrastructure failures remain
+  retryable;
+- failure persistence uses `UPDATE ... RETURNING`; a stale worker that no
+  longer owns the lease cannot report a false successful failure transition;
+- governed tag relationships now use one atomic Neo4j
+  `MERGE (DataFlow)-[:LABEL]->(DataLabel)` path shared by normal creation and
+  reconciliation. Governed tag failures are strict: PostgreSQL is not
+  finalized, the Saga remains retryable, and a later retry completes without
+  duplicate relationships. Legacy tag handling retains its previous
+  best-effort behavior.
+
+Fifth-review acceptance:
+
+- focused Saga/repository/cutover contracts: `40 passed`;
+- real PostgreSQL/API/Neo4j integration: `3 passed`;
+- full repository run completed with
+  `680 passed, 33 skipped, 59 subtests passed` and one unrelated existing
+  publication-receipt assertion mismatch: the tampered receipt was correctly
+  rejected as non-canonical encoding while that randomized test expected the
+  word `signature`;
+- selected Ruff `F`, `I`, and `B` checks passed;
+- the production frontend build completed with zero errors and the same
+  existing 20 console warnings;
+- the backend image was rebuilt from the reviewed source; backend and frontend
+  containers are healthy, Alembic reports `20260724_230 (head)`, application
+  health reports database and Neo4j healthy with code 200, and frontend HTTP
+  returns 200.
+
 ## Residual scope
 
 - Production-line cross-station compatibility remains ultimately authoritative

+ 36 - 7
app/core/data_flow/create_reconciliation.py

@@ -32,12 +32,16 @@ class DataFlowCreateReconciler:
                 "items": candidates,
             }
 
-        claims = self.repository.claim_reconcilable_dataflow_creates(
-            limit=limit, lease_seconds=lease_seconds
-        )
-        self.repository.session.commit()
         items = []
-        for claim in claims:
+        for _ in range(limit):
+            claims = self.repository.claim_reconcilable_dataflow_creates(
+                limit=1, lease_seconds=lease_seconds
+            )
+            if not claims:
+                self.repository.session.rollback()
+                break
+            claim = claims[0]
+            self.repository.commit_dataflow_create_claim()
             base = {
                 "reservation_id": claim["reservation_id"],
                 "dataflow_uid": claim["dataflow_uid"],
@@ -65,12 +69,37 @@ class DataFlowCreateReconciler:
                 intent = DataFlowService.validate_governed_create_intent(
                     claim["create_intent"], repository=self.repository
                 )
+            except (TypeError, ValueError):
+                self.repository.session.rollback()
+                try:
+                    self.repository.commit_dataflow_create_failure(
+                        reservation_id=claim["reservation_id"],
+                        lease_token=claim["lease_token"],
+                        error_code="create_intent_validation_failed",
+                    )
+                    status = "failed"
+                    error_code = "create_intent_validation_failed"
+                except Exception:
+                    self.repository.session.rollback()
+                    status = "lease_recovery_required"
+                    error_code = "failure_state_persistence_failed"
+                items.append(
+                    {**base, "status": status, "error_code": error_code}
+                )
+                continue
+            try:
+                self.repository.renew_dataflow_create_lease(
+                    reservation_id=claim["reservation_id"],
+                    lease_token=claim["lease_token"],
+                    lease_seconds=lease_seconds,
+                )
+                self.repository.commit_dataflow_create_claim()
                 node_id, result = DataFlowService._merge_governed_dataflow(
                     intent["node"]
                 )
                 if intent["tags"]:
                     DataFlowService._handle_tag_relationships(
-                        node_id, intent["tags"]
+                        node_id, intent["tags"], strict=True
                     )
                 result = self.repository.complete_dataflow_create(
                     reservation_id=claim["reservation_id"],
@@ -108,7 +137,7 @@ class DataFlowCreateReconciler:
         completed = sum(item["status"] == "completed" for item in items)
         return {
             "dry_run": False,
-            "candidate_count": len(claims),
+            "candidate_count": len(items),
             "reconciled_count": completed,
             "failed_count": len(items) - completed,
             "items": items,

+ 41 - 23
app/core/data_flow/dataflows.py

@@ -606,8 +606,22 @@ class DataFlowService:
             tag_list = tag_list if governed else data.get("tag", [])
             if tag_list:
                 try:
-                    DataFlowService._handle_tag_relationships(dataflow_id, tag_list)
+                    DataFlowService._handle_tag_relationships(
+                        dataflow_id, tag_list, strict=governed
+                    )
                 except Exception as e:
+                    if governed:
+                        try:
+                            repository.commit_dataflow_create_failure(
+                                reservation_id=receipt["reservation_id"],
+                                lease_token=saga_claim["lease_token"],
+                                error_code="neo4j_tag_merge_failed",
+                            )
+                        except Exception:
+                            logger.exception(
+                                "failed to persist governed DataFlow tag failure"
+                            )
+                        raise
                     logger.warning(f"处理标签关系时出错: {str(e)}")
 
             # Governed definitions are control-plane assets. Data Factory owns
@@ -1286,7 +1300,7 @@ class DataFlowService:
                 logger.warning(f"创建子节点关系失败 {child_id}: {str(e)}")
 
     @staticmethod
-    def _handle_tag_relationships(dataflow_id, tag_list):
+    def _handle_tag_relationships(dataflow_id, tag_list, *, strict=False):
         """
         处理多标签关系
 
@@ -1307,33 +1321,37 @@ class DataFlowService:
                     tag_id = int(tag_item)
 
             if tag_id:
-                DataFlowService._handle_single_tag_relationship(dataflow_id, tag_id)
+                DataFlowService._handle_single_tag_relationship(
+                    dataflow_id, tag_id, strict=strict
+                )
+            elif strict:
+                raise ValueError("dataflow create intent tag is invalid")
 
     @staticmethod
-    def _handle_single_tag_relationship(dataflow_id, tag_id):
+    def _handle_single_tag_relationship(dataflow_id, tag_id, *, strict=False):
         """处理单个标签关系"""
+        driver = None
         try:
-            # 查找标签节点
-            query = "MATCH (n:DataLabel) WHERE id(n) = $tag_id RETURN n"
-            with connect_graph().session() as session:
-                result = session.run(query, tag_id=tag_id).data()
-
-                # 创建关系 - 使用ID调用relationship_exists
-                if (
-                    result
-                    and dataflow_id
-                    and not relationship_exists(dataflow_id, "LABEL", tag_id)
-                ):
-                    session.run(
-                        "MATCH (a), (b) WHERE id(a) = $dataflow_id "
-                        "AND id(b) = $tag_id "
-                        "CREATE (a)-[:LABEL]->(b)",
-                        dataflow_id=dataflow_id,
-                        tag_id=tag_id,
-                    )
-                    logger.info(f"创建标签关系: {dataflow_id} -> {tag_id}")
+            driver = connect_graph()
+            with driver.session() as session:
+                merged = session.run(
+                    "MATCH (a:DataFlow), (b:DataLabel) "
+                    "WHERE id(a) = $dataflow_id AND id(b) = $tag_id "
+                    "MERGE (a)-[:LABEL]->(b) "
+                    "RETURN count(*) AS merged",
+                    dataflow_id=dataflow_id,
+                    tag_id=tag_id,
+                ).single()
+                if merged is None or merged["merged"] != 1:
+                    raise ValueError("dataflow tag target does not exist")
+                logger.info(f"合并标签关系: {dataflow_id} -> {tag_id}")
         except Exception as e:
+            if strict:
+                raise
             logger.warning(f"创建标签关系失败 {tag_id}: {str(e)}")
+        finally:
+            if driver is not None:
+                driver.close()
 
     @staticmethod
     def update_dataflow_script_path(

+ 93 - 32
app/core/data_rules/repository.py

@@ -45,6 +45,10 @@ PLAN_BACKENDS = {
     "generated_python",
     "external_adapter",
 }
+TERMINAL_DATAFLOW_CREATE_ERRORS = {
+    "create_intent_integrity_failed",
+    "create_intent_validation_failed",
+}
 
 
 def _uid(value: Any, label: str) -> str:
@@ -404,28 +408,74 @@ class DataRuleRepository:
     def fail_dataflow_create(
         self, *, reservation_id: str, lease_token: str, error_code: str
     ) -> None:
-        self.session.execute(
-            text(
-                "UPDATE public.dataflow_draft_reservations "
-                "SET state = 'failed', error_code = :error_code, "
-                "failed_at = CURRENT_TIMESTAMP, lease_token = NULL, "
-                "lease_expires_at = NULL, updated_at = CURRENT_TIMESTAMP "
-                "WHERE id = CAST(:id AS uuid) "
-                "AND state = 'creating' "
-                "AND lease_token = CAST(:lease_token AS uuid) "
-                "/* fail_dataflow_create */"
-            ),
-            {
-                "id": _uid(reservation_id, "reservation_id"),
-                "lease_token": _uid(lease_token, "lease_token"),
-                "error_code": _text(error_code, "error_code", 100),
-            },
+        failed = (
+            self.session.execute(
+                text(
+                    "UPDATE public.dataflow_draft_reservations "
+                    "SET state = 'failed', error_code = :error_code, "
+                    "failed_at = CURRENT_TIMESTAMP, lease_token = NULL, "
+                    "lease_expires_at = NULL, updated_at = CURRENT_TIMESTAMP "
+                    "WHERE id = CAST(:id AS uuid) "
+                    "AND state = 'creating' "
+                    "AND lease_token = CAST(:lease_token AS uuid) "
+                    "RETURNING id /* fail_dataflow_create */"
+                ),
+                {
+                    "id": _uid(reservation_id, "reservation_id"),
+                    "lease_token": _uid(lease_token, "lease_token"),
+                    "error_code": _text(error_code, "error_code", 100),
+                },
+            )
+            .mappings()
+            .one_or_none()
         )
+        if failed is None:
+            raise RuntimeError("dataflow create lease was lost")
 
     def commit_dataflow_create_claim(self) -> None:
         """Durably publish the lease before the Neo4j side effect starts."""
         self.session.commit()
 
+    def renew_dataflow_create_lease(
+        self,
+        *,
+        reservation_id: str,
+        lease_token: str,
+        lease_seconds: int = 60,
+    ) -> None:
+        """Fence a side effect with a fresh lease owned by this worker."""
+        if (
+            isinstance(lease_seconds, bool)
+            or not isinstance(lease_seconds, int)
+            or lease_seconds < 10
+            or lease_seconds > 300
+        ):
+            raise ValueError("dataflow create lease is invalid")
+        renewed = (
+            self.session.execute(
+                text(
+                    "UPDATE public.dataflow_draft_reservations "
+                    "SET lease_expires_at = CURRENT_TIMESTAMP + "
+                    "(:lease_seconds * INTERVAL '1 second'), "
+                    "updated_at = CURRENT_TIMESTAMP "
+                    "WHERE id = CAST(:id AS uuid) "
+                    "AND state = 'creating' "
+                    "AND lease_token = CAST(:lease_token AS uuid) "
+                    "AND lease_expires_at > CURRENT_TIMESTAMP "
+                    "RETURNING id /* renew_dataflow_create_lease */"
+                ),
+                {
+                    "id": _uid(reservation_id, "reservation_id"),
+                    "lease_token": _uid(lease_token, "lease_token"),
+                    "lease_seconds": lease_seconds,
+                },
+            )
+            .mappings()
+            .one_or_none()
+        )
+        if renewed is None:
+            raise RuntimeError("dataflow create lease was lost")
+
     def commit_dataflow_create_failure(
         self, *, reservation_id: str, lease_token: str, error_code: str
     ) -> None:
@@ -456,8 +506,9 @@ class DataRuleRepository:
                     "WHERE request_digest IS NOT NULL "
                     "AND create_request IS NOT NULL "
                     "AND create_intent IS NOT NULL "
-                    "AND error_code IS DISTINCT FROM "
-                    "'create_intent_integrity_failed' "
+                    "AND (error_code IS NULL OR error_code NOT IN "
+                    "('create_intent_integrity_failed', "
+                    "'create_intent_validation_failed')) "
                     "AND (state = 'failed' OR (state = 'creating' "
                     "AND lease_expires_at <= CURRENT_TIMESTAMP)) "
                     "ORDER BY updated_at, id LIMIT :limit "
@@ -520,8 +571,9 @@ class DataRuleRepository:
                     "WHERE request_digest IS NOT NULL "
                     "AND create_request IS NOT NULL "
                     "AND create_intent IS NOT NULL "
-                    "AND error_code IS DISTINCT FROM "
-                    "'create_intent_integrity_failed' "
+                    "AND (error_code IS NULL OR error_code NOT IN "
+                    "('create_intent_integrity_failed', "
+                    "'create_intent_validation_failed')) "
                     "AND (state = 'failed' OR (state = 'creating' "
                     "AND lease_expires_at <= CURRENT_TIMESTAMP)) "
                     "ORDER BY updated_at, id "
@@ -551,19 +603,27 @@ class DataRuleRepository:
         )
         claimed = []
         for row in rows:
-            create_request = _bounded_object(
-                row["create_request"], "dataflow create request"
-            )
-            create_intent = _bounded_object(
-                row["create_intent"], "dataflow create intent"
-            )
             request_digest = str(row["request_digest"])
-            integrity_valid = not (
-                _canonical_hash(create_request) != request_digest
-                or create_request.get("intent") != create_intent
-                or create_intent.get("dataflow_uid")
-                != str(row["dataflow_uid"])
-            )
+            create_intent = None
+            integrity_error = None
+            try:
+                create_request = _bounded_object(
+                    row["create_request"], "dataflow create request"
+                )
+                create_intent = _bounded_object(
+                    row["create_intent"], "dataflow create intent"
+                )
+                integrity_valid = not (
+                    _canonical_hash(create_request) != request_digest
+                    or create_request.get("intent") != create_intent
+                    or create_intent.get("dataflow_uid")
+                    != str(row["dataflow_uid"])
+                )
+                if not integrity_valid:
+                    integrity_error = "dataflow create intent digest mismatch"
+            except Exception as exc:
+                integrity_valid = False
+                integrity_error = str(exc)
             claimed.append(
                 {
                     "reservation_id": str(row["reservation_id"]),
@@ -573,6 +633,7 @@ class DataRuleRepository:
                     "create_intent": create_intent,
                     "attempt": int(row["attempt_count"]),
                     "integrity_valid": integrity_valid,
+                    "integrity_error": integrity_error,
                 }
             )
         return claimed

+ 88 - 0
tests/core/data_rules/test_data_rule_repository.py

@@ -297,6 +297,94 @@ def test_dataflow_create_saga_claim_replay_and_integrity_are_closed():
     assert "nonce_hash = :nonce_hash" in sql
 
 
+def test_reconciliation_claim_isolates_poison_json_and_keeps_valid_row():
+    from app.core.data_rules.repository import (
+        DataRuleRepository,
+        _canonical_hash,
+    )
+
+    poison_uid = new_governance_uid()
+    valid_uid = new_governance_uid()
+    valid_intent = {
+        "dataflow_uid": valid_uid,
+        "node": {"uid": valid_uid},
+        "tags": [],
+    }
+    valid_request = {"payload": {}, "intent": valid_intent}
+
+    class ClaimSession:
+        def execute(self, statement, _params=None):
+            assert "claim_reconcilable_dataflow_creates" in str(statement)
+            return FakeResult(
+                rows=[
+                    {
+                        "reservation_id": new_governance_uid(),
+                        "dataflow_uid": poison_uid,
+                        "create_request": "poison-scalar",
+                        "create_intent": {},
+                        "request_digest": "a" * 64,
+                        "attempt_count": 2,
+                    },
+                    {
+                        "reservation_id": new_governance_uid(),
+                        "dataflow_uid": valid_uid,
+                        "create_request": valid_request,
+                        "create_intent": valid_intent,
+                        "request_digest": _canonical_hash(valid_request),
+                        "attempt_count": 3,
+                    },
+                ]
+            )
+
+    claims = DataRuleRepository(
+        ClaimSession()
+    ).claim_reconcilable_dataflow_creates(limit=2)
+    assert [claim["integrity_valid"] for claim in claims] == [False, True]
+    assert claims[0]["create_intent"] is None
+    assert claims[0]["integrity_error"]
+    assert claims[1]["create_intent"] == valid_intent
+
+
+def test_failure_and_lease_renewal_require_current_ownership():
+    from app.core.data_rules.repository import DataRuleRepository
+
+    class LostLeaseSession:
+        def execute(self, statement, _params=None):
+            sql = str(statement)
+            assert (
+                "fail_dataflow_create" in sql
+                or "renew_dataflow_create_lease" in sql
+            )
+            return FakeResult()
+
+    repository = DataRuleRepository(LostLeaseSession())
+    reservation_id = new_governance_uid()
+    lease_token = new_governance_uid()
+    with pytest.raises(RuntimeError, match="lease was lost"):
+        repository.renew_dataflow_create_lease(
+            reservation_id=reservation_id,
+            lease_token=lease_token,
+        )
+    with pytest.raises(RuntimeError, match="lease was lost"):
+        repository.fail_dataflow_create(
+            reservation_id=reservation_id,
+            lease_token=lease_token,
+            error_code="reconciliation_failed",
+        )
+
+
+def test_terminal_reconciliation_failures_are_excluded_from_queries():
+    from app.core.data_rules.repository import DataRuleRepository
+
+    session = FakeSession()
+    repository = DataRuleRepository(session)
+    assert repository.list_reconcilable_dataflow_creates(limit=1) == []
+    repository.claim_reconcilable_dataflow_creates(limit=1)
+    sql = _sql(session)
+    assert sql.count("create_intent_integrity_failed") == 2
+    assert sql.count("create_intent_validation_failed") == 2
+
+
 def test_create_rule_version_rejects_legacy_v1_rule_specs():
     from app.core.data_rules.repository import DataRuleRepository
 

+ 30 - 9
tests/integration/test_dataflow_create_reconciliation.py

@@ -40,7 +40,15 @@ def test_real_postgres_neo4j_orphan_is_reconciled_from_stored_intent(
     actor = new_governance_uid()
     reservation_id = None
     dataflow_uid = None
+    tag_uid = new_governance_uid()
+    tag_node_id = None
     try:
+        with graph.session() as graph_session:
+            tag_node_id = graph_session.run(
+                "CREATE (t:DataLabel {acceptance_uid: $uid}) "
+                "RETURN id(t) AS node_id",
+                {"uid": tag_uid},
+            ).single()["node_id"]
         with app.app_context(), Session(pg) as session:
             session.execute(
                 text(
@@ -75,7 +83,7 @@ def test_real_postgres_neo4j_orphan_is_reconciled_from_stored_intent(
             intent = {
                 "dataflow_uid": dataflow_uid,
                 "node": node,
-                "tags": [],
+                "tags": [{"id": tag_node_id}],
             }
             create_request = {
                 "payload": {"name_zh": node["name_zh"]},
@@ -143,18 +151,31 @@ def test_real_postgres_neo4j_orphan_is_reconciled_from_stored_intent(
             assert state["state"] == "completed"
             assert state["request_digest"]
             assert state["result_digest"]
+            DataFlowService._handle_tag_relationships(
+                report["items"][0]["dataflow_node_id"],
+                [{"id": tag_node_id}],
+                strict=True,
+            )
         with graph.session() as graph_session:
-            count = graph_session.run(
-                "MATCH (n:DataFlow {uid: $uid}) RETURN count(n) AS count",
-                {"uid": dataflow_uid},
-            ).single()["count"]
-        assert count == 1
+            record = graph_session.run(
+                "MATCH (n:DataFlow {uid: $uid}) "
+                "OPTIONAL MATCH (n)-[r:LABEL]->"
+                "(:DataLabel {acceptance_uid: $tag_uid}) "
+                "RETURN count(DISTINCT n) AS node_count, "
+                "count(r) AS relationship_count",
+                {"uid": dataflow_uid, "tag_uid": tag_uid},
+            ).single()
+        assert record["node_count"] == 1
+        assert record["relationship_count"] == 1
     finally:
-        if dataflow_uid:
+        if dataflow_uid or tag_uid:
             with graph.session() as graph_session:
                 graph_session.run(
-                    "MATCH (n:DataFlow {uid: $uid}) DETACH DELETE n",
-                    {"uid": dataflow_uid},
+                    "MATCH (n) WHERE "
+                    "(n:DataFlow AND n.uid = $uid) OR "
+                    "(n:DataLabel AND n.acceptance_uid = $tag_uid) "
+                    "DETACH DELETE n",
+                    {"uid": dataflow_uid, "tag_uid": tag_uid},
                 )
         graph.close()
         with Session(pg) as session:

+ 397 - 12
tests/test_dataflow_create_saga.py

@@ -255,6 +255,64 @@ def test_graph_failure_is_persisted_but_finalize_failure_leaves_reconcilable_lea
         )
 
 
+def test_governed_tag_failure_never_finalizes_completed(monkeypatch):
+    from tests.core.data_rules.test_contracts import valid_dataflow_spec
+    from tests.test_legacy_governance_cutover import PublishedAssetRepository
+
+    flow = valid_dataflow_spec()
+    actor = new_governance_uid()
+    repository = PublishedAssetRepository()
+    receipt = repository.reserve_dataflow_draft(actor_uid=actor)
+    receipt["dataflow_uid"] = flow["dataflow_uid"]
+    failures = []
+    repository.commit_dataflow_create_failure = (
+        lambda **kwargs: failures.append(kwargs)
+    )
+    monkeypatch.setattr(
+        DataFlowService,
+        "_merge_governed_dataflow",
+        lambda node: (73, {"id": 73, **node}),
+    )
+
+    def fail_tags(_dataflow_id, _tags, *, strict):
+        assert strict is True
+        raise RuntimeError("neo4j tag merge failed")
+
+    monkeypatch.setattr(
+        DataFlowService, "_handle_tag_relationships", fail_tags
+    )
+    with pytest.raises(RuntimeError, match="neo4j tag merge failed"):
+        DataFlowService.create_dataflow(
+            {
+                "name_zh": "标签严格生产线",
+                "describe": "严格标签合并验收",
+                "script_type": "governed",
+                "script_requirement": {
+                    "dataflow_spec": flow,
+                    "dataset_edges": {
+                        "source_table": flow["input_schema_refs"],
+                        "target_table": flow["output_schema_ref"],
+                    },
+                    "migration_metadata": {
+                        "status": "migrated",
+                        "legacy_fields_present": False,
+                        "preserved_for_read_only": True,
+                        "governed_semantics": "dataflow_spec",
+                    },
+                },
+                "tag": [{"id": 81}],
+                "draft_reservation": {
+                    key: receipt[key]
+                    for key in ("reservation_id", "dataflow_uid", "nonce")
+                },
+            },
+            repository=repository,
+            actor_uid=actor,
+        )
+    assert repository.completed == {}
+    assert failures[0]["error_code"] == "neo4j_tag_merge_failed"
+
+
 def test_reconciler_dry_run_is_read_only_and_digest_failure_is_isolated(
     monkeypatch,
 ):
@@ -278,18 +336,8 @@ def test_reconciler_dry_run_is_read_only_and_digest_failure_is_isolated(
             self.claim_calls = 0
             self.failures = []
             self.completions = []
-
-        def list_reconcilable_dataflow_creates(self, *, limit):
-            assert limit == 2
-            return [{"reservation_id": "preview", "state": "failed"}]
-
-        def claim_reconcilable_dataflow_creates(
-            self, *, limit, lease_seconds
-        ):
-            assert limit == 2
-            assert lease_seconds == 60
-            self.claim_calls += 1
-            return [
+            self.renewals = []
+            self.claims = [
                 {
                     "reservation_id": new_governance_uid(),
                     "dataflow_uid": bad_uid,
@@ -310,6 +358,24 @@ def test_reconciler_dry_run_is_read_only_and_digest_failure_is_isolated(
                 },
             ]
 
+        def list_reconcilable_dataflow_creates(self, *, limit):
+            assert limit == 2
+            return [{"reservation_id": "preview", "state": "failed"}]
+
+        def claim_reconcilable_dataflow_creates(
+            self, *, limit, lease_seconds
+        ):
+            assert limit == 1
+            assert lease_seconds == 60
+            self.claim_calls += 1
+            return [self.claims.pop(0)] if self.claims else []
+
+        def commit_dataflow_create_claim(self):
+            self.session.commit()
+
+        def renew_dataflow_create_lease(self, **kwargs):
+            self.renewals.append(kwargs)
+
         def commit_dataflow_create_failure(self, **kwargs):
             self.failures.append(kwargs)
             self.session.commit()
@@ -345,3 +411,322 @@ def test_reconciler_dry_run_is_read_only_and_digest_failure_is_isolated(
         "create_intent_integrity_failed"
     )
     assert repository.completions[0]["dataflow_uid"] == good_uid
+    assert len(repository.renewals) == 1
+
+
+def test_reconciler_validation_is_terminal_and_lost_lease_is_reported(
+    monkeypatch,
+):
+    flow_uid = new_governance_uid()
+
+    class Session:
+        def commit(self):
+            return None
+
+        def rollback(self):
+            return None
+
+    class Repository:
+        session = Session()
+
+        def __init__(self):
+            self.claimed = False
+            self.failure_codes = []
+
+        def claim_reconcilable_dataflow_creates(self, **_kwargs):
+            if self.claimed:
+                return []
+            self.claimed = True
+            return [
+                {
+                    "reservation_id": new_governance_uid(),
+                    "dataflow_uid": flow_uid,
+                    "lease_token": new_governance_uid(),
+                    "request_digest": "c" * 64,
+                    "create_intent": {"dataflow_uid": flow_uid},
+                    "attempt": 1,
+                    "integrity_valid": True,
+                }
+            ]
+
+        def commit_dataflow_create_claim(self):
+            return None
+
+        def commit_dataflow_create_failure(self, **kwargs):
+            self.failure_codes.append(kwargs["error_code"])
+
+    repository = Repository()
+    monkeypatch.setattr(
+        DataFlowService,
+        "validate_governed_create_intent",
+        lambda *_args, **_kwargs: (_ for _ in ()).throw(
+            ValueError("published asset is unavailable")
+        ),
+    )
+    report = DataFlowCreateReconciler(repository).run(
+        dry_run=False, limit=2
+    )
+    assert report["items"][0] == {
+        "reservation_id": report["items"][0]["reservation_id"],
+        "dataflow_uid": flow_uid,
+        "attempt": 1,
+        "request_digest": "c" * 64,
+        "status": "failed",
+        "error_code": "create_intent_validation_failed",
+    }
+    assert repository.failure_codes == ["create_intent_validation_failed"]
+
+    repository = Repository()
+    monkeypatch.setattr(
+        DataFlowService,
+        "validate_governed_create_intent",
+        lambda value, **_kwargs: {"node": value, "tags": []},
+    )
+    repository.renew_dataflow_create_lease = lambda **_kwargs: (_ for _ in ()).throw(
+        RuntimeError("dataflow create lease was lost")
+    )
+    repository.commit_dataflow_create_failure = lambda **_kwargs: (
+        _ for _ in ()
+    ).throw(RuntimeError("dataflow create lease was lost"))
+    monkeypatch.setattr(
+        DataFlowService,
+        "_merge_governed_dataflow",
+        lambda _node: pytest.fail("side effect ran after lease loss"),
+    )
+    report = DataFlowCreateReconciler(repository).run(
+        dry_run=False, limit=1
+    )
+    assert report["items"][0]["status"] == "lease_recovery_required"
+    assert report["items"][0]["error_code"] == (
+        "failure_state_persistence_failed"
+    )
+
+
+def test_tag_relationship_uses_one_atomic_merge(monkeypatch):
+    calls = []
+
+    class Result:
+        def single(self):
+            return {"merged": 1}
+
+    class Session:
+        def run(self, query, **values):
+            calls.append((query, values))
+            return Result()
+
+        def __enter__(self):
+            return self
+
+        def __exit__(self, *_args):
+            return None
+
+    class Driver:
+        def session(self):
+            return Session()
+
+        def close(self):
+            return None
+
+    monkeypatch.setattr(
+        "app.core.data_flow.dataflows.connect_graph", lambda: Driver()
+    )
+    DataFlowService._handle_single_tag_relationship(73, 81)
+    DataFlowService._handle_single_tag_relationship(73, 81)
+
+    assert len(calls) == 2
+    assert all("MERGE (a)-[:LABEL]->(b)" in query for query, _ in calls)
+    assert all("CREATE (a)-[:LABEL]->(b)" not in query for query, _ in calls)
+
+
+def test_reconciler_tag_failure_does_not_complete_and_retry_succeeds(
+    monkeypatch,
+):
+    flow_uid = new_governance_uid()
+    reservation_id = new_governance_uid()
+    digest = "d" * 64
+
+    class Session:
+        def commit(self):
+            return None
+
+        def rollback(self):
+            return None
+
+    class Repository:
+        session = Session()
+
+        def __init__(self):
+            self.attempt = 0
+            self.available = True
+            self.failures = []
+            self.completions = []
+
+        def claim_reconcilable_dataflow_creates(self, **_kwargs):
+            if not self.available:
+                return []
+            self.available = False
+            self.attempt += 1
+            return [
+                {
+                    "reservation_id": reservation_id,
+                    "dataflow_uid": flow_uid,
+                    "lease_token": new_governance_uid(),
+                    "request_digest": digest,
+                    "create_intent": {"dataflow_uid": flow_uid},
+                    "attempt": self.attempt,
+                    "integrity_valid": True,
+                }
+            ]
+
+        def commit_dataflow_create_claim(self):
+            return None
+
+        def renew_dataflow_create_lease(self, **_kwargs):
+            return None
+
+        def commit_dataflow_create_failure(self, **kwargs):
+            self.failures.append(kwargs)
+            self.available = True
+
+        def complete_dataflow_create(self, **kwargs):
+            self.completions.append(kwargs)
+            return {"id": 73, "uid": flow_uid}
+
+    repository = Repository()
+    monkeypatch.setattr(
+        DataFlowService,
+        "validate_governed_create_intent",
+        lambda _value, **_kwargs: {
+            "node": {"uid": flow_uid},
+            "tags": [{"id": 81}],
+        },
+    )
+    monkeypatch.setattr(
+        DataFlowService,
+        "_merge_governed_dataflow",
+        lambda node: (73, {"id": 73, **node}),
+    )
+    attempts = {"count": 0}
+
+    def merge_tags(_node_id, _tags, *, strict):
+        assert strict is True
+        attempts["count"] += 1
+        if attempts["count"] == 1:
+            raise RuntimeError("neo4j tag write failed")
+
+    monkeypatch.setattr(
+        DataFlowService, "_handle_tag_relationships", merge_tags
+    )
+    reconciler = DataFlowCreateReconciler(repository)
+    failed = reconciler.run(dry_run=False, limit=1)
+    assert failed["items"][0]["status"] == "failed"
+    assert repository.completions == []
+    assert repository.failures[0]["error_code"] == "reconciliation_failed"
+
+    completed = reconciler.run(dry_run=False, limit=1)
+    assert completed["items"][0]["status"] == "completed"
+    assert len(repository.completions) == 1
+    assert attempts["count"] == 2
+
+
+def test_slow_reconciler_does_not_preclaim_next_item_from_second_worker(
+    monkeypatch,
+):
+    first_uid = new_governance_uid()
+    second_uid = new_governance_uid()
+    shared = {
+        first_uid: {"state": "failed", "attempt": 0},
+        second_uid: {"state": "failed", "attempt": 0},
+    }
+    processed = []
+
+    class Session:
+        def commit(self):
+            return None
+
+        def rollback(self):
+            return None
+
+    class Repository:
+        session = Session()
+
+        def __init__(self, worker):
+            self.worker = worker
+
+        def claim_reconcilable_dataflow_creates(
+            self, *, limit, lease_seconds
+        ):
+            assert limit == 1
+            assert lease_seconds == 10
+            available = [
+                uid
+                for uid, record in shared.items()
+                if record["state"] == "failed"
+            ]
+            if not available:
+                return []
+            uid = available[0]
+            shared[uid]["state"] = f"creating:{self.worker}"
+            shared[uid]["attempt"] += 1
+            return [
+                {
+                    "reservation_id": new_governance_uid(),
+                    "dataflow_uid": uid,
+                    "lease_token": new_governance_uid(),
+                    "request_digest": ("a" if uid == first_uid else "b") * 64,
+                    "create_intent": {"dataflow_uid": uid},
+                    "attempt": shared[uid]["attempt"],
+                    "integrity_valid": True,
+                }
+            ]
+
+        def commit_dataflow_create_claim(self):
+            return None
+
+        def renew_dataflow_create_lease(self, **_kwargs):
+            return None
+
+        def complete_dataflow_create(self, *, dataflow_uid, **_kwargs):
+            assert shared[dataflow_uid]["state"] == f"creating:{self.worker}"
+            shared[dataflow_uid]["state"] = "completed"
+            processed.append((self.worker, dataflow_uid))
+            return {"id": 73, "uid": dataflow_uid}
+
+        def commit_dataflow_create_failure(self, **_kwargs):
+            pytest.fail("unexpected reconciliation failure")
+
+    monkeypatch.setattr(
+        DataFlowService,
+        "validate_governed_create_intent",
+        lambda value, **_kwargs: {
+            "node": {"uid": value["dataflow_uid"]},
+            "tags": [],
+        },
+    )
+    worker_two = DataFlowCreateReconciler(Repository("worker-2"))
+    entered_slow_call = {"value": False}
+
+    def slow_merge(node):
+        if node["uid"] == first_uid and not entered_slow_call["value"]:
+            entered_slow_call["value"] = True
+            # This models worker 1 crossing its original lease boundary while
+            # processing item A. Item B was never preclaimed, so worker 2 can
+            # safely process B without ever sharing ownership of A.
+            second_report = worker_two.run(
+                dry_run=False, limit=1, lease_seconds=10
+            )
+            assert second_report["items"][0]["dataflow_uid"] == second_uid
+        return 73, {"id": 73, **node}
+
+    monkeypatch.setattr(
+        DataFlowService, "_merge_governed_dataflow", slow_merge
+    )
+    first_report = DataFlowCreateReconciler(Repository("worker-1")).run(
+        dry_run=False, limit=2, lease_seconds=10
+    )
+    assert first_report["reconciled_count"] == 1
+    assert processed == [
+        ("worker-2", second_uid),
+        ("worker-1", first_uid),
+    ]
+    assert all(record["attempt"] == 1 for record in shared.values())