Просмотр исходного кода

fix: bind and reconcile governed dataflow creates

马小龙 4 недель назад
Родитель
Сommit
99411f370a

+ 53 - 2
.superpowers/sdd/task-8-report.md

@@ -120,7 +120,7 @@ Final:
 - Task 8 API/repository/frontend, Saga and legacy-cutover contracts:
 - Task 8 API/repository/frontend, Saga and legacy-cutover contracts:
   `64 passed`;
   `64 passed`;
 - full repository suite:
 - full repository suite:
-  `672 passed, 31 skipped, 59 subtests passed`;
+  `673 passed, 33 skipped, 59 subtests passed`;
 - `git diff --check`: passed;
 - `git diff --check`: passed;
 - frontend production build: completed with `0 errors`;
 - frontend production build: completed with `0 errors`;
 - build retained 20 pre-existing `no-console` warnings plus existing CSS
 - build retained 20 pre-existing `no-console` warnings plus existing CSS
@@ -271,9 +271,60 @@ Additional acceptance evidence:
   PostgreSQL finalize failure after Neo4j success, completed-response replay,
   PostgreSQL finalize failure after Neo4j success, completed-response replay,
   conflict closure and lease recovery;
   conflict closure and lease recovery;
 - real PostgreSQL and Neo4j integration: `2 passed`;
 - real PostgreSQL and Neo4j integration: `2 passed`;
-- rebuilt backend and frontend images report migration 220, healthy database
+- rebuilt backend and frontend images report migration 230, healthy database
   and Neo4j checks, application health code 200, and frontend HTTP 200.
   and Neo4j checks, application health code 200, and frontend HTTP 200.
 
 
+## Durable snapshot and orphan-reconciliation closeout
+
+The fourth review closed the remaining transaction and recovery boundaries:
+
+- catalog endpoints now treat `SchemaResolver` as an explicit read-through
+  persistence boundary: successful context-aware GETs commit newly created
+  snapshots before returning their IDs, and every error path rolls back;
+- real PostgreSQL/API acceptance issued two separate catalog requests,
+  observed one stable snapshot ID, and queried that ID from an independent
+  database session;
+- forward-only migration `20260724_230` adds the canonical create request,
+  server-owned create intent and SHA-256 request digest without modifying
+  migration 220;
+- first claim atomically freezes request, intent and digest; every later
+  claim, completed replay and finalize must match that digest;
+- changing the name or any other canonical payload/intent value while
+  reusing a receipt fails with `dataflow_create_request_conflict` instead of
+  silently returning an older result;
+- governed creation uses the supplied English name or a deterministic
+  UID-derived fallback and never calls non-deterministic translation before
+  computing the request digest.
+
+`DataFlowCreateReconciler` and
+`python -m app.commands.reconcile_dataflow_creates` recover orphaned graph
+nodes without accepting a nonce or inline business payload:
+
+- dry-run is the default and reports bounded safe IDs, state, digest and
+  attempt information without changing state or attempt count;
+- apply mode claims expired `creating` or retryable `failed` rows with
+  `FOR UPDATE SKIP LOCKED` and an expiring lease;
+- canonical request digest, duplicated intent and DataFlow UID are verified
+  before any Neo4j write;
+- the existing UID-only graph `MERGE` either finds the orphan or creates the
+  missing node, verifies immutable properties, and then finalizes PostgreSQL;
+- corrupt intent rows receive a terminal integrity error and do not block
+  valid rows; other failures are isolated per reservation and remain
+  retryable;
+- two real concurrent reconcilers produced exactly one claim;
+- real PostgreSQL/Neo4j acceptance created a graph node, deliberately omitted
+  PostgreSQL finalize, discarded the receipt, expired the lease, and then
+  completed the Saga from stored intent with exactly one graph node;
+- the reconciliation CLI dry-run completed against the local database with
+  an empty, non-mutating report.
+
+Fourth-review acceptance:
+
+- focused repository/API/Saga/cutover contracts: `65 passed`;
+- real PostgreSQL/API/Neo4j integration: `4 passed`;
+- full repository suite:
+  `673 passed, 33 skipped, 59 subtests passed`.
+
 ## Residual scope
 ## Residual scope
 
 
 - Production-line cross-station compatibility remains ultimately authoritative
 - Production-line cross-station compatibility remains ultimately authoritative

+ 19 - 13
app/api/data_rules/routes.py

@@ -517,20 +517,23 @@ def catalog_asset(asset_type: str, version_id: str):
         }:
         }:
             raise ValueError("catalog asset query contains unsupported fields")
             raise ValueError("catalog asset query contains unsupported fields")
         inputs, output = _catalog_schema_context()
         inputs, output = _catalog_schema_context()
-        return jsonify(
-            success(
-                _repository().get_published_asset(
-                    asset_type=asset_type,
-                    version_id=version_id,
-                    input_schema_refs=inputs,
-                    output_schema_ref=output,
-                    schema_resolver=_schema_resolver() if inputs is not None else None,
-                )
-            )
+        result = _repository().get_published_asset(
+            asset_type=asset_type,
+            version_id=version_id,
+            input_schema_refs=inputs,
+            output_schema_ref=output,
+            schema_resolver=_schema_resolver() if inputs is not None else None,
         )
         )
+        if inputs is not None:
+            # SchemaResolver is a read-through snapshot boundary. Persist the
+            # IDs returned to this response before making them observable.
+            db.session.commit()
+        return jsonify(success(result))
     except (TypeError, ValueError, json.JSONDecodeError):
     except (TypeError, ValueError, json.JSONDecodeError):
+        db.session.rollback()
         return jsonify(failed("已发布资产不存在或上下文无效", code=404)), 404
         return jsonify(failed("已发布资产不存在或上下文无效", code=404)), 404
     except Exception:
     except Exception:
+        db.session.rollback()
         current_app.logger.exception("load exact catalog asset failed")
         current_app.logger.exception("load exact catalog asset failed")
         return jsonify(failed("规则目录暂时不可用", code=503)), 503
         return jsonify(failed("规则目录暂时不可用", code=503)), 503
 
 
@@ -580,12 +583,15 @@ def published_rule_catalog():
                     "schema_resolver": _schema_resolver(),
                     "schema_resolver": _schema_resolver(),
                 }
                 }
             )
             )
-        return jsonify(
-            success(_repository().search_published_assets(**catalog_args))
-        )
+        result = _repository().search_published_assets(**catalog_args)
+        if inputs is not None:
+            db.session.commit()
+        return jsonify(success(result))
     except (TypeError, ValueError):
     except (TypeError, ValueError):
+        db.session.rollback()
         return _bad_request("规则目录查询无效")
         return _bad_request("规则目录查询无效")
     except Exception:
     except Exception:
+        db.session.rollback()
         current_app.logger.exception("load rule catalog failed")
         current_app.logger.exception("load rule catalog failed")
         return jsonify(failed("规则目录暂时不可用", code=503)), 503
         return jsonify(failed("规则目录暂时不可用", code=503)), 503
 
 

+ 37 - 0
app/commands/reconcile_dataflow_creates.py

@@ -0,0 +1,37 @@
+"""Inspect or repair orphaned governed DataFlow create Sagas."""
+
+from __future__ import annotations
+
+import argparse
+import json
+
+from app import create_app, db
+from app.core.data_flow.create_reconciliation import DataFlowCreateReconciler
+from app.core.data_rules.repository import DataRuleRepository
+
+
+def main() -> None:
+    parser = argparse.ArgumentParser(
+        description="Reconcile database-owned governed DataFlow create intents"
+    )
+    mode = parser.add_mutually_exclusive_group()
+    mode.add_argument("--apply", action="store_true")
+    mode.add_argument("--dry-run", action="store_true")
+    parser.add_argument("--limit", type=int, default=100)
+    parser.add_argument("--lease-seconds", type=int, default=60)
+    args = parser.parse_args()
+
+    app = create_app()
+    with app.app_context():
+        report = DataFlowCreateReconciler(
+            DataRuleRepository(db.session)
+        ).run(
+            dry_run=not args.apply,
+            limit=args.limit,
+            lease_seconds=args.lease_seconds,
+        )
+        print(json.dumps(report, ensure_ascii=False, sort_keys=True))
+
+
+if __name__ == "__main__":
+    main()

+ 115 - 0
app/core/data_flow/create_reconciliation.py

@@ -0,0 +1,115 @@
+"""Recover governed DataFlow creates from database-owned canonical intents."""
+
+from __future__ import annotations
+
+from typing import Any
+
+from app.core.data_flow.dataflows import DataFlowService
+
+
+class DataFlowCreateReconciler:
+    """Lease, isolate, and reconcile orphaned PostgreSQL/Neo4j creates."""
+
+    def __init__(self, repository):
+        self.repository = repository
+
+    def run(
+        self,
+        *,
+        dry_run: bool = True,
+        limit: int = 100,
+        lease_seconds: int = 60,
+    ) -> dict[str, Any]:
+        if dry_run:
+            candidates = self.repository.list_reconcilable_dataflow_creates(
+                limit=limit
+            )
+            return {
+                "dry_run": True,
+                "candidate_count": len(candidates),
+                "reconciled_count": 0,
+                "failed_count": 0,
+                "items": candidates,
+            }
+
+        claims = self.repository.claim_reconcilable_dataflow_creates(
+            limit=limit, lease_seconds=lease_seconds
+        )
+        self.repository.session.commit()
+        items = []
+        for claim in claims:
+            base = {
+                "reservation_id": claim["reservation_id"],
+                "dataflow_uid": claim["dataflow_uid"],
+                "attempt": claim["attempt"],
+                "request_digest": claim["request_digest"],
+            }
+            if not claim["integrity_valid"]:
+                try:
+                    self.repository.commit_dataflow_create_failure(
+                        reservation_id=claim["reservation_id"],
+                        lease_token=claim["lease_token"],
+                        error_code="create_intent_integrity_failed",
+                    )
+                    status = "failed"
+                    error_code = "create_intent_integrity_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:
+                intent = DataFlowService.validate_governed_create_intent(
+                    claim["create_intent"], repository=self.repository
+                )
+                node_id, result = DataFlowService._merge_governed_dataflow(
+                    intent["node"]
+                )
+                if intent["tags"]:
+                    DataFlowService._handle_tag_relationships(
+                        node_id, intent["tags"]
+                    )
+                result = self.repository.complete_dataflow_create(
+                    reservation_id=claim["reservation_id"],
+                    dataflow_uid=claim["dataflow_uid"],
+                    lease_token=claim["lease_token"],
+                    request_digest=claim["request_digest"],
+                    dataflow_node_id=node_id,
+                    result=result,
+                )
+                self.repository.session.commit()
+                items.append(
+                    {
+                        **base,
+                        "status": "completed",
+                        "dataflow_node_id": result.get("id"),
+                    }
+                )
+            except Exception:
+                self.repository.session.rollback()
+                try:
+                    self.repository.commit_dataflow_create_failure(
+                        reservation_id=claim["reservation_id"],
+                        lease_token=claim["lease_token"],
+                        error_code="reconciliation_failed",
+                    )
+                    status = "failed"
+                    error_code = "reconciliation_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}
+                )
+        completed = sum(item["status"] == "completed" for item in items)
+        return {
+            "dry_run": False,
+            "candidate_count": len(claims),
+            "reconciled_count": completed,
+            "failed_count": len(items) - completed,
+            "items": items,
+        }

+ 135 - 29
app/core/data_flow/dataflows.py

@@ -56,6 +56,21 @@ class DataFlowService:
         "task_list",
         "task_list",
         "workflow",
         "workflow",
     }
     }
+    _GOVERNED_NODE_KEYS = {
+        "uid",
+        "name_zh",
+        "name_en",
+        "category",
+        "organization",
+        "leader",
+        "frequency",
+        "describe",
+        "status",
+        "update_mode",
+        "script_type",
+        "script_requirement",
+        "script_path",
+    }
 
 
     @staticmethod
     @staticmethod
     def _decode_script_requirement(value: Any) -> Any:
     def _decode_script_requirement(value: Any) -> Any:
@@ -142,6 +157,9 @@ class DataFlowService:
     @staticmethod
     @staticmethod
     def _merge_governed_dataflow(node_data: dict[str, Any]) -> tuple[int, dict]:
     def _merge_governed_dataflow(node_data: dict[str, Any]) -> tuple[int, dict]:
         """Idempotently create by stable UID and reject immutable conflicts."""
         """Idempotently create by stable UID and reject immutable conflicts."""
+        node_data = copy.deepcopy(node_data)
+        node_data.setdefault("created_at", get_formatted_time())
+        node_data.setdefault("updated_at", node_data["created_at"])
         immutable = {
         immutable = {
             key: node_data[key]
             key: node_data[key]
             for key in ("uid", "name_zh", "script_type", "script_requirement")
             for key in ("uid", "name_zh", "script_type", "script_requirement")
@@ -184,6 +202,55 @@ class DataFlowService:
         finally:
         finally:
             driver.close()
             driver.close()
 
 
+    @classmethod
+    def validate_governed_create_intent(
+        cls, value: Any, *, repository
+    ) -> dict[str, Any]:
+        if not isinstance(value, dict) or set(value) != {
+            "dataflow_uid",
+            "node",
+            "tags",
+        }:
+            raise ValueError("dataflow create intent is invalid")
+        node = value.get("node")
+        tags = value.get("tags")
+        if not isinstance(node, dict) or set(node) != cls._GOVERNED_NODE_KEYS:
+            raise ValueError("dataflow create intent node is invalid")
+        if not isinstance(tags, list) or len(tags) > 100:
+            raise ValueError("dataflow create intent tags are invalid")
+        flow_uid = ensure_governance_uid(
+            {"uid": str(value.get("dataflow_uid"))}
+        )
+        if node.get("uid") != flow_uid:
+            raise ValueError("dataflow create intent UID mismatch")
+        if node.get("script_type") != "governed" or node.get(
+            "script_path"
+        ) not in {"", None}:
+            raise ValueError("dataflow create intent execution mode is invalid")
+        requirement = cls.validate_governed_requirement(
+            node.get("script_requirement"), repository=repository
+        )
+        if requirement["dataflow_spec"]["dataflow_uid"] != flow_uid:
+            raise ValueError("dataflow create intent requirement UID mismatch")
+        normalized_node = copy.deepcopy(node)
+        normalized_node["script_requirement"] = json.dumps(
+            requirement,
+            ensure_ascii=False,
+            sort_keys=True,
+            separators=(",", ":"),
+        )
+        for key in cls._GOVERNED_NODE_KEYS - {
+            "script_path",
+            "script_requirement",
+        }:
+            if not isinstance(normalized_node.get(key), str):
+                raise ValueError("dataflow create intent node field is invalid")
+        return {
+            "dataflow_uid": flow_uid,
+            "node": normalized_node,
+            "tags": copy.deepcopy(tags),
+        }
+
     @staticmethod
     @staticmethod
     def get_dataflows(
     def get_dataflows(
         page: int = 1,
         page: int = 1,
@@ -380,18 +447,6 @@ class DataFlowService:
 
 
             dataflow_name = data["name_zh"]
             dataflow_name = data["name_zh"]
 
 
-            # 使用LLM翻译名称生成英文名
-            try:
-                result_list = translate_and_parse(dataflow_name)
-                name_en = (
-                    result_list[0]
-                    if result_list
-                    else dataflow_name.lower().replace(" ", "_")
-                )
-            except Exception as e:
-                logger.warning(f"翻译失败,使用默认英文名: {str(e)}")
-                name_en = dataflow_name.lower().replace(" ", "_")
-
             script_requirement = data.get("script_requirement")
             script_requirement = data.get("script_requirement")
             decoded_requirement = DataFlowService._decode_script_requirement(
             decoded_requirement = DataFlowService._decode_script_requirement(
                 script_requirement
                 script_requirement
@@ -399,6 +454,35 @@ class DataFlowService:
             governed = DataFlowService._signals_governed(
             governed = DataFlowService._signals_governed(
                 data, decoded_requirement
                 data, decoded_requirement
             )
             )
+            name_en = data.get("name_en")
+            if governed and (
+                not isinstance(name_en, str) or not name_en.strip()
+            ):
+                flow_value = (
+                    decoded_requirement.get("dataflow_spec")
+                    if isinstance(decoded_requirement, dict)
+                    else None
+                )
+                flow_uid = (
+                    flow_value.get("dataflow_uid")
+                    if isinstance(flow_value, dict)
+                    else None
+                )
+                name_en = f"dataflow_{str(flow_uid).replace('-', '')}"
+            elif not isinstance(name_en, str) or not name_en.strip():
+                # 使用LLM翻译名称生成英文名
+                try:
+                    result_list = translate_and_parse(dataflow_name)
+                    name_en = (
+                        result_list[0]
+                        if result_list
+                        else dataflow_name.lower().replace(" ", "_")
+                    )
+                except Exception as e:
+                    logger.warning(f"翻译失败,使用默认英文名: {str(e)}")
+                    name_en = dataflow_name.lower().replace(" ", "_")
+            name_en = name_en.strip()
+
             saga_claim = None
             saga_claim = None
             receipt = None
             receipt = None
             if governed:
             if governed:
@@ -417,20 +501,6 @@ class DataFlowService:
                         "governed DataFlow requires an authenticated actor"
                         "governed DataFlow requires an authenticated actor"
                     )
                     )
                 receipt = data.get("draft_reservation")
                 receipt = data.get("draft_reservation")
-                saga_claim = repository.begin_dataflow_create(
-                    receipt, actor_uid=actor_uid
-                )
-                if saga_claim["status"] == "completed":
-                    return copy.deepcopy(saga_claim["result"])
-                reserved_uid = saga_claim["dataflow_uid"]
-                if (
-                    reserved_uid
-                    != decoded_requirement["dataflow_spec"]["dataflow_uid"]
-                ):
-                    raise ValueError(
-                        "draft reservation does not match dataflow_uid"
-                    )
-                repository.commit_dataflow_create_claim()
                 script_requirement = decoded_requirement
                 script_requirement = decoded_requirement
 
 
             # 处理 script_requirement,将其转换为 JSON 字符串
             # 处理 script_requirement,将其转换为 JSON 字符串
@@ -463,9 +533,10 @@ class DataFlowService:
                 "script_type": data.get("script_type", "python"),
                 "script_type": data.get("script_type", "python"),
                 "script_requirement": script_requirement_str,
                 "script_requirement": script_requirement_str,
                 "script_path": "",  # 脚本路径,任务完成后更新
                 "script_path": "",  # 脚本路径,任务完成后更新
-                "created_at": get_formatted_time(),
-                "updated_at": get_formatted_time(),
             }
             }
+            if not governed:
+                node_data["created_at"] = get_formatted_time()
+                node_data["updated_at"] = get_formatted_time()
             if governed:
             if governed:
                 node_data["uid"] = decoded_requirement["dataflow_spec"][
                 node_data["uid"] = decoded_requirement["dataflow_spec"][
                     "dataflow_uid"
                     "dataflow_uid"
@@ -475,6 +546,40 @@ class DataFlowService:
 
 
             # Governed nodes use the reservation UID as the only create key.
             # Governed nodes use the reservation UID as the only create key.
             if governed:
             if governed:
+                tag_list = copy.deepcopy(data.get("tag", []))
+                if not isinstance(tag_list, list):
+                    tag_list = [tag_list] if tag_list else []
+                create_intent = {
+                    "dataflow_uid": node_data["uid"],
+                    "node": copy.deepcopy(node_data),
+                    "tags": tag_list,
+                }
+                request_payload = copy.deepcopy(data)
+                request_payload.pop("draft_reservation", None)
+                request_payload["script_requirement"] = decoded_requirement
+                create_request = {
+                    "payload": request_payload,
+                    "intent": create_intent,
+                }
+                saga_claim = repository.begin_dataflow_create(
+                    receipt,
+                    actor_uid=actor_uid,
+                    create_request=create_request,
+                    create_intent=create_intent,
+                )
+                if saga_claim["status"] == "completed":
+                    return copy.deepcopy(saga_claim["result"])
+                reserved_uid = saga_claim["dataflow_uid"]
+                if reserved_uid != node_data["uid"]:
+                    raise ValueError(
+                        "draft reservation does not match dataflow_uid"
+                    )
+                stored_intent = DataFlowService.validate_governed_create_intent(
+                    saga_claim["create_intent"], repository=repository
+                )
+                node_data = stored_intent["node"]
+                tag_list = stored_intent["tags"]
+                repository.commit_dataflow_create_claim()
                 try:
                 try:
                     dataflow_id, result = (
                     dataflow_id, result = (
                         DataFlowService._merge_governed_dataflow(node_data)
                         DataFlowService._merge_governed_dataflow(node_data)
@@ -498,7 +603,7 @@ class DataFlowService:
                 dataflow_id = create_or_get_node("DataFlow", **node_data)
                 dataflow_id = create_or_get_node("DataFlow", **node_data)
 
 
             # 处理标签关系(支持多标签数组)
             # 处理标签关系(支持多标签数组)
-            tag_list = data.get("tag", [])
+            tag_list = tag_list if governed else data.get("tag", [])
             if tag_list:
             if tag_list:
                 try:
                 try:
                     DataFlowService._handle_tag_relationships(dataflow_id, tag_list)
                     DataFlowService._handle_tag_relationships(dataflow_id, tag_list)
@@ -534,6 +639,7 @@ class DataFlowService:
                     reservation_id=receipt["reservation_id"],
                     reservation_id=receipt["reservation_id"],
                     dataflow_uid=node_data["uid"],
                     dataflow_uid=node_data["uid"],
                     lease_token=saga_claim["lease_token"],
                     lease_token=saga_claim["lease_token"],
+                    request_digest=saga_claim["request_digest"],
                     dataflow_node_id=dataflow_id,
                     dataflow_node_id=dataflow_id,
                     result=result,
                     result=result,
                 )
                 )

+ 198 - 4
app/core/data_rules/repository.py

@@ -110,6 +110,15 @@ def _canonical_hash(value: Any) -> str:
     return hashlib.sha256(canonical.encode("utf-8")).hexdigest()
     return hashlib.sha256(canonical.encode("utf-8")).hexdigest()
 
 
 
 
+def _bounded_object(
+    value: Any, label: str, *, maximum_bytes: int = 131_072
+) -> dict[str, Any]:
+    normalized = copy.deepcopy(_object(value, label))
+    if len(_json(normalized).encode("utf-8")) > maximum_bytes:
+        raise ValueError(f"{label} is too large")
+    return normalized
+
+
 class DataRuleRepository:
 class DataRuleRepository:
     """Persist versions and enforce all lifecycle transitions server-side."""
     """Persist versions and enforce all lifecycle transitions server-side."""
 
 
@@ -195,6 +204,8 @@ class DataRuleRepository:
         receipt: dict[str, Any],
         receipt: dict[str, Any],
         *,
         *,
         actor_uid: str,
         actor_uid: str,
+        create_request: dict[str, Any],
+        create_intent: dict[str, Any],
         lease_seconds: int = 60,
         lease_seconds: int = 60,
     ) -> dict[str, Any]:
     ) -> dict[str, Any]:
         """Atomically claim a create attempt or replay its committed result."""
         """Atomically claim a create attempt or replay its committed result."""
@@ -206,6 +217,15 @@ class DataRuleRepository:
         ):
         ):
             raise ValueError("dataflow create lease is invalid")
             raise ValueError("dataflow create lease is invalid")
         values = self._draft_receipt(receipt, actor_uid=actor_uid)
         values = self._draft_receipt(receipt, actor_uid=actor_uid)
+        request_value = _bounded_object(
+            create_request, "dataflow create request"
+        )
+        intent_value = _bounded_object(create_intent, "dataflow create intent")
+        if request_value.get("intent") != intent_value:
+            raise ValueError("dataflow create request intent is invalid")
+        if intent_value.get("dataflow_uid") != values["dataflow_uid"]:
+            raise ValueError("dataflow create intent UID does not match receipt")
+        request_digest = _canonical_hash(request_value)
         lease_token = new_governance_uid()
         lease_token = new_governance_uid()
         claimed = (
         claimed = (
             self.session.execute(
             self.session.execute(
@@ -215,6 +235,11 @@ class DataRuleRepository:
                     "lease_expires_at = CURRENT_TIMESTAMP + "
                     "lease_expires_at = CURRENT_TIMESTAMP + "
                     "(:lease_seconds * INTERVAL '1 second'), "
                     "(:lease_seconds * INTERVAL '1 second'), "
                     "attempt_count = attempt_count + 1, "
                     "attempt_count = attempt_count + 1, "
+                    "create_request = COALESCE("
+                    "create_request, CAST(:create_request AS jsonb)), "
+                    "create_intent = COALESCE("
+                    "create_intent, CAST(:create_intent AS jsonb)), "
+                    "request_digest = COALESCE(request_digest, :request_digest), "
                     "create_started_at = COALESCE("
                     "create_started_at = COALESCE("
                     "create_started_at, CURRENT_TIMESTAMP), "
                     "create_started_at, CURRENT_TIMESTAMP), "
                     "error_code = NULL, updated_at = CURRENT_TIMESTAMP "
                     "error_code = NULL, updated_at = CURRENT_TIMESTAMP "
@@ -222,17 +247,23 @@ class DataRuleRepository:
                     "AND dataflow_uid = CAST(:dataflow_uid AS uuid) "
                     "AND dataflow_uid = CAST(:dataflow_uid AS uuid) "
                     "AND actor_uid = CAST(:actor_uid AS uuid) "
                     "AND actor_uid = CAST(:actor_uid AS uuid) "
                     "AND nonce_hash = :nonce_hash "
                     "AND nonce_hash = :nonce_hash "
+                    "AND (request_digest IS NULL "
+                    "OR request_digest = :request_digest) "
                     "AND consumed_at IS NULL "
                     "AND consumed_at IS NULL "
                     "AND (expires_at > CURRENT_TIMESTAMP OR attempt_count > 0) "
                     "AND (expires_at > CURRENT_TIMESTAMP OR attempt_count > 0) "
                     "AND (state IN ('reserved', 'failed') OR "
                     "AND (state IN ('reserved', 'failed') OR "
                     "(state = 'creating' AND lease_expires_at <= CURRENT_TIMESTAMP)) "
                     "(state = 'creating' AND lease_expires_at <= CURRENT_TIMESTAMP)) "
-                    "RETURNING dataflow_uid::text, attempt_count "
+                    "RETURNING dataflow_uid::text, attempt_count, "
+                    "create_intent, request_digest "
                     "/* begin_dataflow_create */"
                     "/* begin_dataflow_create */"
                 ),
                 ),
                 {
                 {
                     **values,
                     **values,
                     "lease_token": lease_token,
                     "lease_token": lease_token,
                     "lease_seconds": lease_seconds,
                     "lease_seconds": lease_seconds,
+                    "create_request": _json(request_value),
+                    "create_intent": _json(intent_value),
+                    "request_digest": request_digest,
                 },
                 },
             )
             )
             .mappings()
             .mappings()
@@ -244,13 +275,18 @@ class DataRuleRepository:
                 "dataflow_uid": str(claimed["dataflow_uid"]),
                 "dataflow_uid": str(claimed["dataflow_uid"]),
                 "lease_token": lease_token,
                 "lease_token": lease_token,
                 "attempt": int(claimed["attempt_count"]),
                 "attempt": int(claimed["attempt_count"]),
+                "create_intent": _object(
+                    claimed["create_intent"], "dataflow create intent"
+                ),
+                "request_digest": str(claimed["request_digest"]),
             }
             }
 
 
         row = (
         row = (
             self.session.execute(
             self.session.execute(
                 text(
                 text(
                     "SELECT state, dataflow_uid::text, result, "
                     "SELECT state, dataflow_uid::text, result, "
-                    "result_digest, lease_expires_at, expires_at "
+                    "result_digest, lease_expires_at, expires_at, "
+                    "request_digest, create_intent "
                     "FROM public.dataflow_draft_reservations "
                     "FROM public.dataflow_draft_reservations "
                     "WHERE id = CAST(:id AS uuid) "
                     "WHERE id = CAST(:id AS uuid) "
                     "AND dataflow_uid = CAST(:dataflow_uid AS uuid) "
                     "AND dataflow_uid = CAST(:dataflow_uid AS uuid) "
@@ -263,7 +299,14 @@ class DataRuleRepository:
             .mappings()
             .mappings()
             .one_or_none()
             .one_or_none()
         )
         )
+        if row is not None and row["request_digest"] not in {
+            None,
+            request_digest,
+        }:
+            raise ValueError("dataflow_create_request_conflict")
         if row is not None and row["state"] == "completed":
         if row is not None and row["state"] == "completed":
+            if str(row["request_digest"]) != request_digest:
+                raise ValueError("dataflow_create_request_conflict")
             result = _object(row["result"], "dataflow create result")
             result = _object(row["result"], "dataflow create result")
             if _canonical_hash(result) != str(row["result_digest"]):
             if _canonical_hash(result) != str(row["result_digest"]):
                 raise RuntimeError("dataflow create result integrity check failed")
                 raise RuntimeError("dataflow create result integrity check failed")
@@ -271,6 +314,7 @@ class DataRuleRepository:
                 "status": "completed",
                 "status": "completed",
                 "dataflow_uid": str(row["dataflow_uid"]),
                 "dataflow_uid": str(row["dataflow_uid"]),
                 "result": copy.deepcopy(result),
                 "result": copy.deepcopy(result),
+                "request_digest": request_digest,
             }
             }
         if row is not None and row["state"] == "creating":
         if row is not None and row["state"] == "creating":
             raise ValueError("dataflow_create_in_progress")
             raise ValueError("dataflow_create_in_progress")
@@ -282,12 +326,14 @@ class DataRuleRepository:
         reservation_id: str,
         reservation_id: str,
         dataflow_uid: str,
         dataflow_uid: str,
         lease_token: str,
         lease_token: str,
+        request_digest: str,
         dataflow_node_id: int,
         dataflow_node_id: int,
         result: dict[str, Any],
         result: dict[str, Any],
     ) -> dict[str, Any]:
     ) -> dict[str, Any]:
         reservation = _uid(reservation_id, "reservation_id")
         reservation = _uid(reservation_id, "reservation_id")
         flow_uid = _uid(dataflow_uid, "dataflow_uid")
         flow_uid = _uid(dataflow_uid, "dataflow_uid")
         lease = _uid(lease_token, "lease_token")
         lease = _uid(lease_token, "lease_token")
+        request_hash = _digest(request_digest, "request_digest")
         if (
         if (
             isinstance(dataflow_node_id, bool)
             isinstance(dataflow_node_id, bool)
             or not isinstance(dataflow_node_id, int)
             or not isinstance(dataflow_node_id, int)
@@ -310,12 +356,14 @@ class DataRuleRepository:
                     "AND dataflow_uid = CAST(:dataflow_uid AS uuid) "
                     "AND dataflow_uid = CAST(:dataflow_uid AS uuid) "
                     "AND state = 'creating' "
                     "AND state = 'creating' "
                     "AND lease_token = CAST(:lease_token AS uuid) "
                     "AND lease_token = CAST(:lease_token AS uuid) "
+                    "AND request_digest = :request_digest "
                     "RETURNING result /* complete_dataflow_create */"
                     "RETURNING result /* complete_dataflow_create */"
                 ),
                 ),
                 {
                 {
                     "id": reservation,
                     "id": reservation,
                     "dataflow_uid": flow_uid,
                     "dataflow_uid": flow_uid,
                     "lease_token": lease,
                     "lease_token": lease,
+                    "request_digest": request_hash,
                     "result": _json(normalized_result),
                     "result": _json(normalized_result),
                     "result_digest": result_digest,
                     "result_digest": result_digest,
                     "dataflow_node_id": dataflow_node_id,
                     "dataflow_node_id": dataflow_node_id,
@@ -328,14 +376,19 @@ class DataRuleRepository:
             replay = (
             replay = (
                 self.session.execute(
                 self.session.execute(
                     text(
                     text(
-                        "SELECT result, result_digest "
+                        "SELECT result, result_digest, request_digest "
                         "FROM public.dataflow_draft_reservations "
                         "FROM public.dataflow_draft_reservations "
                         "WHERE id = CAST(:id AS uuid) "
                         "WHERE id = CAST(:id AS uuid) "
                         "AND dataflow_uid = CAST(:dataflow_uid AS uuid) "
                         "AND dataflow_uid = CAST(:dataflow_uid AS uuid) "
                         "AND state = 'completed' "
                         "AND state = 'completed' "
+                        "AND request_digest = :request_digest "
                         "/* replay_completed_dataflow_create */"
                         "/* replay_completed_dataflow_create */"
                     ),
                     ),
-                    {"id": reservation, "dataflow_uid": flow_uid},
+                    {
+                        "id": reservation,
+                        "dataflow_uid": flow_uid,
+                        "request_digest": request_hash,
+                    },
                 )
                 )
                 .mappings()
                 .mappings()
                 .one_or_none()
                 .one_or_none()
@@ -383,6 +436,147 @@ class DataRuleRepository:
         )
         )
         self.session.commit()
         self.session.commit()
 
 
+    def list_reconcilable_dataflow_creates(
+        self, *, limit: int = 100
+    ) -> list[dict[str, Any]]:
+        if (
+            isinstance(limit, bool)
+            or not isinstance(limit, int)
+            or limit < 1
+            or limit > 500
+        ):
+            raise ValueError("reconciliation limit is invalid")
+        rows = (
+            self.session.execute(
+                text(
+                    "SELECT id::text AS reservation_id, "
+                    "dataflow_uid::text, state, attempt_count, "
+                    "request_digest, error_code, lease_expires_at, updated_at "
+                    "FROM public.dataflow_draft_reservations "
+                    "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 (state = 'failed' OR (state = 'creating' "
+                    "AND lease_expires_at <= CURRENT_TIMESTAMP)) "
+                    "ORDER BY updated_at, id LIMIT :limit "
+                    "/* list_reconcilable_dataflow_creates */"
+                ),
+                {"limit": limit},
+            )
+            .mappings()
+            .all()
+        )
+        return [
+            {
+                "reservation_id": str(row["reservation_id"]),
+                "dataflow_uid": str(row["dataflow_uid"]),
+                "state": str(row["state"]),
+                "attempt_count": int(row["attempt_count"]),
+                "request_digest": str(row["request_digest"]),
+                "error_code": (
+                    str(row["error_code"])
+                    if row["error_code"] is not None
+                    else None
+                ),
+                "lease_expires_at": (
+                    row["lease_expires_at"].isoformat()
+                    if hasattr(row["lease_expires_at"], "isoformat")
+                    else row["lease_expires_at"]
+                ),
+                "updated_at": (
+                    row["updated_at"].isoformat()
+                    if hasattr(row["updated_at"], "isoformat")
+                    else row["updated_at"]
+                ),
+            }
+            for row in rows
+        ]
+
+    def claim_reconcilable_dataflow_creates(
+        self, *, limit: int = 100, lease_seconds: int = 60
+    ) -> list[dict[str, Any]]:
+        if (
+            isinstance(limit, bool)
+            or not isinstance(limit, int)
+            or limit < 1
+            or limit > 500
+        ):
+            raise ValueError("reconciliation limit is invalid")
+        if (
+            isinstance(lease_seconds, bool)
+            or not isinstance(lease_seconds, int)
+            or lease_seconds < 10
+            or lease_seconds > 300
+        ):
+            raise ValueError("reconciliation lease is invalid")
+        lease_token = new_governance_uid()
+        rows = (
+            self.session.execute(
+                text(
+                    "WITH candidates AS (SELECT id "
+                    "FROM public.dataflow_draft_reservations "
+                    "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 (state = 'failed' OR (state = 'creating' "
+                    "AND lease_expires_at <= CURRENT_TIMESTAMP)) "
+                    "ORDER BY updated_at, id "
+                    "FOR UPDATE SKIP LOCKED LIMIT :limit) "
+                    "UPDATE public.dataflow_draft_reservations value "
+                    "SET state = 'creating', "
+                    "lease_token = CAST(:lease_token AS uuid), "
+                    "lease_expires_at = CURRENT_TIMESTAMP + "
+                    "(:lease_seconds * INTERVAL '1 second'), "
+                    "attempt_count = attempt_count + 1, "
+                    "error_code = NULL, updated_at = CURRENT_TIMESTAMP "
+                    "FROM candidates WHERE value.id = candidates.id "
+                    "RETURNING value.id::text AS reservation_id, "
+                    "value.dataflow_uid::text, value.create_request, "
+                    "value.create_intent, value.request_digest, "
+                    "value.attempt_count "
+                    "/* claim_reconcilable_dataflow_creates */"
+                ),
+                {
+                    "limit": limit,
+                    "lease_token": lease_token,
+                    "lease_seconds": lease_seconds,
+                },
+            )
+            .mappings()
+            .all()
+        )
+        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"])
+            )
+            claimed.append(
+                {
+                    "reservation_id": str(row["reservation_id"]),
+                    "dataflow_uid": str(row["dataflow_uid"]),
+                    "lease_token": lease_token,
+                    "request_digest": request_digest,
+                    "create_intent": create_intent,
+                    "attempt": int(row["attempt_count"]),
+                    "integrity_valid": integrity_valid,
+                }
+            )
+        return claimed
+
     @staticmethod
     @staticmethod
     def candidate_hash(candidate: dict[str, Any]) -> str:
     def candidate_hash(candidate: dict[str, Any]) -> str:
         return _canonical_hash(validate_rule_candidate(candidate))
         return _canonical_hash(validate_rule_candidate(candidate))

+ 47 - 0
migrations/versions/20260724_230_dataflow_create_intent.py

@@ -0,0 +1,47 @@
+"""Bind recoverable DataFlow creates to a canonical server-owned intent."""
+
+from alembic import op
+
+revision = "20260724_230"
+down_revision = "20260724_220"
+branch_labels = None
+depends_on = None
+
+
+def upgrade() -> None:
+    op.execute(
+        """
+        ALTER TABLE public.dataflow_draft_reservations
+            ADD COLUMN create_request JSONB,
+            ADD COLUMN create_intent JSONB,
+            ADD COLUMN request_digest CHAR(64);
+
+        ALTER TABLE public.dataflow_draft_reservations
+            ADD CONSTRAINT ck_dataflow_create_intent_complete
+            CHECK (
+                (
+                    create_request IS NULL
+                    AND create_intent IS NULL
+                    AND request_digest IS NULL
+                )
+                OR (
+                    create_request IS NOT NULL
+                    AND create_intent IS NOT NULL
+                    AND request_digest ~ '^[0-9a-f]{64}$'
+                )
+            );
+
+        CREATE INDEX idx_dataflow_create_reconciliation
+            ON public.dataflow_draft_reservations(
+                state, lease_expires_at, updated_at
+            )
+            WHERE request_digest IS NOT NULL
+              AND create_intent IS NOT NULL;
+        """
+    )
+
+
+def downgrade() -> None:
+    raise RuntimeError(
+        "DataFlow create intents are forward-only and cannot downgrade"
+    )

+ 43 - 3
tests/core/data_rules/test_data_rule_repository.py

@@ -189,12 +189,20 @@ def test_dataflow_create_saga_claim_replay_and_integrity_are_closed():
         "dataflow_uid": new_governance_uid(),
         "dataflow_uid": new_governance_uid(),
         "nonce": "unguessable-test-nonce-value",
         "nonce": "unguessable-test-nonce-value",
     }
     }
+    intent = {
+        "dataflow_uid": receipt["dataflow_uid"],
+        "node": {"uid": receipt["dataflow_uid"]},
+        "tags": [],
+    }
+    create_request = {"payload": {"name_zh": "测试"}, "intent": intent}
 
 
     class AtomicSession:
     class AtomicSession:
         def __init__(self):
         def __init__(self):
             self.state = "reserved"
             self.state = "reserved"
             self.result = None
             self.result = None
             self.digest = None
             self.digest = None
+            self.request_digest = None
+            self.intent = None
             self.calls = []
             self.calls = []
 
 
         def execute(self, statement, params=None):
         def execute(self, statement, params=None):
@@ -203,8 +211,17 @@ def test_dataflow_create_saga_claim_replay_and_integrity_are_closed():
             self.calls.append((sql, values))
             self.calls.append((sql, values))
             if "begin_dataflow_create" in sql and self.state == "reserved":
             if "begin_dataflow_create" in sql and self.state == "reserved":
                 self.state = "creating"
                 self.state = "creating"
+                self.request_digest = values["request_digest"]
+                self.intent = json.loads(values["create_intent"])
                 return FakeResult(
                 return FakeResult(
-                    rows=[{"dataflow_uid": receipt["dataflow_uid"], "attempt_count": 1}]
+                    rows=[
+                        {
+                            "dataflow_uid": receipt["dataflow_uid"],
+                            "attempt_count": 1,
+                            "create_intent": self.intent,
+                            "request_digest": self.request_digest,
+                        }
+                    ]
                 )
                 )
             if "begin_dataflow_create" in sql:
             if "begin_dataflow_create" in sql:
                 return FakeResult()
                 return FakeResult()
@@ -218,6 +235,8 @@ def test_dataflow_create_saga_claim_replay_and_integrity_are_closed():
                             "result_digest": self.digest,
                             "result_digest": self.digest,
                             "lease_expires_at": None,
                             "lease_expires_at": None,
                             "expires_at": None,
                             "expires_at": None,
+                            "request_digest": self.request_digest,
+                            "create_intent": self.intent,
                         }
                         }
                     ]
                     ]
                 )
                 )
@@ -233,23 +252,44 @@ def test_dataflow_create_saga_claim_replay_and_integrity_are_closed():
 
 
     session = AtomicSession()
     session = AtomicSession()
     repository = DataRuleRepository(session)
     repository = DataRuleRepository(session)
-    claim = repository.begin_dataflow_create(receipt, actor_uid=actor)
+    claim = repository.begin_dataflow_create(
+        receipt,
+        actor_uid=actor,
+        create_request=create_request,
+        create_intent=intent,
+    )
     assert claim["status"] == "claimed"
     assert claim["status"] == "claimed"
     result = {"id": 31, "uid": receipt["dataflow_uid"]}
     result = {"id": 31, "uid": receipt["dataflow_uid"]}
     assert repository.complete_dataflow_create(
     assert repository.complete_dataflow_create(
         reservation_id=receipt["reservation_id"],
         reservation_id=receipt["reservation_id"],
         dataflow_uid=receipt["dataflow_uid"],
         dataflow_uid=receipt["dataflow_uid"],
         lease_token=claim["lease_token"],
         lease_token=claim["lease_token"],
+        request_digest=claim["request_digest"],
         dataflow_node_id=31,
         dataflow_node_id=31,
         result=result,
         result=result,
     ) == result
     ) == result
     assert repository.begin_dataflow_create(
     assert repository.begin_dataflow_create(
-        receipt, actor_uid=actor
+        receipt,
+        actor_uid=actor,
+        create_request=create_request,
+        create_intent=intent,
     ) == {
     ) == {
         "status": "completed",
         "status": "completed",
         "dataflow_uid": receipt["dataflow_uid"],
         "dataflow_uid": receipt["dataflow_uid"],
         "result": result,
         "result": result,
+        "request_digest": claim["request_digest"],
     }
     }
+    altered_request = {
+        "payload": {"name_zh": "已篡改"},
+        "intent": intent,
+    }
+    with pytest.raises(ValueError, match="request_conflict"):
+        repository.begin_dataflow_create(
+            receipt,
+            actor_uid=actor,
+            create_request=altered_request,
+            create_intent=intent,
+        )
     sql = session.calls[0][0]
     sql = session.calls[0][0]
     assert "consumed_at IS NULL" in sql
     assert "consumed_at IS NULL" in sql
     assert "lease_expires_at <= CURRENT_TIMESTAMP" in sql
     assert "lease_expires_at <= CURRENT_TIMESTAMP" in sql

+ 116 - 0
tests/integration/test_catalog_schema_snapshot_api_postgres.py

@@ -0,0 +1,116 @@
+from __future__ import annotations
+
+import os
+import uuid
+from datetime import UTC, datetime, timedelta
+
+import pytest
+from sqlalchemy import create_engine, text
+from sqlalchemy.orm import Session
+
+from app.core.common.identifiers import new_governance_uid
+from app.core.data_rules.repository import DataRuleRepository
+from app.core.data_rules.schema_resolver import SchemaResolver
+from app.core.system.tokens import decode_access_token, issue_access_token
+
+
+class StableCatalog:
+    def __init__(self, schema_ref):
+        self.schema_ref = schema_ref
+
+    def load_schema(self, schema_ref):
+        assert schema_ref == self.schema_ref
+        return {
+            "source_revision": "acceptance:revision:1",
+            "fields": [
+                {"name": "customer_id", "type": "string", "nullable": False}
+            ],
+        }
+
+
+class ResolverExercisingCatalogRepository:
+    def __init__(self, schema_ref):
+        self.schema_ref = schema_ref
+
+    def search_published_assets(self, **kwargs):
+        snapshot = kwargs["schema_resolver"].resolve(self.schema_ref)
+        return {
+            "items": [{"snapshot_id": snapshot["id"]}],
+            "total": 1,
+            "limit": kwargs["limit"],
+            "offset": kwargs["offset"],
+        }
+
+
+def test_catalog_api_commits_read_through_snapshot_and_reuses_stable_id(
+    monkeypatch,
+):
+    url = os.environ.get("DATA_RULE_POSTGRES_ACCEPTANCE_URL")
+    if not url:
+        pytest.skip("real PostgreSQL acceptance URL is not configured")
+    monkeypatch.setenv("DATABASE_URL", url)
+    from app import create_app, db
+
+    app = create_app()
+    app.config["TESTING"] = True
+    schema_ref = f"bd:acceptance-{uuid.uuid4().hex}:v1"
+    app.extensions["data_rule_repository"] = (
+        ResolverExercisingCatalogRepository(schema_ref)
+    )
+    app.extensions["data_rule_schema_resolver"] = SchemaResolver(
+        StableCatalog(schema_ref), DataRuleRepository(db.session)
+    )
+
+    def load_identity(token, *, secret):
+        claims = decode_access_token(token, secret=secret)
+        return {
+            "id": claims["sub"],
+            "username": "snapshot-acceptance",
+            "display_name": "Snapshot Acceptance",
+            "roles": claims["roles"],
+        }
+
+    monkeypatch.setattr(
+        "app.core.system.auth.load_identity_from_token", load_identity
+    )
+    token = issue_access_token(
+        user_id=new_governance_uid(),
+        roles=["viewer"],
+        secret=app.config["SECRET_KEY"],
+        now=datetime.now(UTC),
+        lifetime=timedelta(minutes=10),
+    )
+    headers = {"Authorization": f"Bearer {token}"}
+    query = (
+        f"input_schema_refs=[%22{schema_ref}%22]"
+        f"&output_schema_ref={schema_ref}"
+    )
+    engine = create_engine(url)
+    try:
+        client = app.test_client()
+        first = client.get(f"/api/rules/catalog?{query}", headers=headers)
+        second = client.get(f"/api/rules/catalog?{query}", headers=headers)
+        assert first.status_code == second.status_code == 200
+        first_id = first.get_json()["data"]["items"][0]["snapshot_id"]
+        second_id = second.get_json()["data"]["items"][0]["snapshot_id"]
+        assert first_id == second_id
+        with Session(engine) as session:
+            persisted = session.execute(
+                text(
+                    "SELECT id::text FROM public.data_schema_snapshots "
+                    "WHERE schema_ref = :schema_ref"
+                ),
+                {"schema_ref": schema_ref},
+            ).scalar_one()
+        assert persisted == first_id
+    finally:
+        with Session(engine) as session:
+            session.execute(
+                text(
+                    "DELETE FROM public.data_schema_snapshots "
+                    "WHERE schema_ref = :schema_ref"
+                ),
+                {"schema_ref": schema_ref},
+            )
+            session.commit()
+        engine.dispose()

+ 176 - 0
tests/integration/test_dataflow_create_reconciliation.py

@@ -0,0 +1,176 @@
+from __future__ import annotations
+
+import os
+
+import pytest
+from neo4j import GraphDatabase
+from sqlalchemy import create_engine, text
+from sqlalchemy.orm import Session
+
+from app.core.common.identifiers import new_governance_uid
+from app.core.data_flow.create_reconciliation import DataFlowCreateReconciler
+from app.core.data_flow.dataflows import DataFlowService
+from app.core.data_rules.repository import DataRuleRepository
+
+
+def test_real_postgres_neo4j_orphan_is_reconciled_from_stored_intent(
+    monkeypatch,
+):
+    pg_url = os.environ.get("DATA_RULE_POSTGRES_ACCEPTANCE_URL")
+    neo4j_uri = os.environ.get("DATA_RULE_NEO4J_ACCEPTANCE_URI")
+    password = os.environ.get("DATA_RULE_NEO4J_ACCEPTANCE_PASSWORD")
+    if not pg_url or not neo4j_uri or not password:
+        pytest.skip("real PostgreSQL and Neo4j acceptance are not configured")
+    user = os.environ.get("DATA_RULE_NEO4J_ACCEPTANCE_USER", "neo4j")
+    monkeypatch.setenv("DATABASE_URL", pg_url)
+    from app import create_app
+
+    app = create_app()
+    app.config.update(
+        TESTING=True,
+        NEO4J_URI=neo4j_uri,
+        NEO4J_USER=user,
+        NEO4J_PASSWORD=password,
+        NEO4J_ENCRYPTED=False,
+    )
+    pg = create_engine(pg_url)
+    graph = GraphDatabase.driver(
+        neo4j_uri, auth=(user, password), encrypted=False
+    )
+    actor = new_governance_uid()
+    reservation_id = None
+    dataflow_uid = None
+    try:
+        with app.app_context(), Session(pg) as session:
+            session.execute(
+                text(
+                    "INSERT INTO public.users "
+                    "(id, username, display_name, password_hash, status) "
+                    "VALUES (CAST(:id AS uuid), :username, "
+                    "'Reconcile Acceptance', 'not-a-login-hash', 'active')"
+                ),
+                {"id": actor, "username": f"reconcile-{actor}"},
+            )
+            session.commit()
+            repository = DataRuleRepository(session)
+            receipt = repository.reserve_dataflow_draft(actor_uid=actor)
+            reservation_id = receipt["reservation_id"]
+            dataflow_uid = receipt["dataflow_uid"]
+            session.commit()
+            node = {
+                "uid": dataflow_uid,
+                "name_zh": f"孤儿生产线-{dataflow_uid}",
+                "name_en": f"orphan-{dataflow_uid}",
+                "category": "应用类",
+                "organization": "acceptance",
+                "leader": "system",
+                "frequency": "月",
+                "describe": "reconciliation acceptance",
+                "status": "active",
+                "update_mode": "append",
+                "script_type": "governed",
+                "script_requirement": "{}",
+                "script_path": "",
+            }
+            intent = {
+                "dataflow_uid": dataflow_uid,
+                "node": node,
+                "tags": [],
+            }
+            create_request = {
+                "payload": {"name_zh": node["name_zh"]},
+                "intent": intent,
+            }
+            closed = {
+                key: receipt[key]
+                for key in ("reservation_id", "dataflow_uid", "nonce")
+            }
+            claim = repository.begin_dataflow_create(
+                closed,
+                actor_uid=actor,
+                create_request=create_request,
+                create_intent=intent,
+                lease_seconds=10,
+            )
+            session.commit()
+            DataFlowService._merge_governed_dataflow(node)
+            session.execute(
+                text(
+                    "UPDATE public.dataflow_draft_reservations "
+                    "SET lease_expires_at = CURRENT_TIMESTAMP - INTERVAL '1 second' "
+                    "WHERE id = CAST(:id AS uuid)"
+                ),
+                {"id": reservation_id},
+            )
+            session.commit()
+            monkeypatch.setattr(
+                DataFlowService,
+                "validate_governed_create_intent",
+                lambda value, *, repository: value,
+            )
+            reconciler = DataFlowCreateReconciler(repository)
+            before_preview = session.execute(
+                text(
+                    "SELECT state, attempt_count "
+                    "FROM public.dataflow_draft_reservations "
+                    "WHERE id = CAST(:id AS uuid)"
+                ),
+                {"id": reservation_id},
+            ).one()
+            dry_run = reconciler.run(dry_run=True)
+            assert dry_run["candidate_count"] == 1
+            after_preview = session.execute(
+                text(
+                    "SELECT state, attempt_count "
+                    "FROM public.dataflow_draft_reservations "
+                    "WHERE id = CAST(:id AS uuid)"
+                ),
+                {"id": reservation_id},
+            ).one()
+            assert after_preview == before_preview
+            report = reconciler.run(dry_run=False)
+            assert report["reconciled_count"] == 1
+            assert report["items"][0]["attempt"] == claim["attempt"] + 1
+            assert reconciler.run(dry_run=False)["candidate_count"] == 0
+            state = session.execute(
+                text(
+                    "SELECT state, request_digest, result_digest "
+                    "FROM public.dataflow_draft_reservations "
+                    "WHERE id = CAST(:id AS uuid)"
+                ),
+                {"id": reservation_id},
+            ).mappings().one()
+            assert state["state"] == "completed"
+            assert state["request_digest"]
+            assert state["result_digest"]
+        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
+    finally:
+        if dataflow_uid:
+            with graph.session() as graph_session:
+                graph_session.run(
+                    "MATCH (n:DataFlow {uid: $uid}) DETACH DELETE n",
+                    {"uid": dataflow_uid},
+                )
+        graph.close()
+        with Session(pg) as session:
+            if reservation_id:
+                session.execute(
+                    text(
+                        "DELETE FROM public.dataflow_draft_reservations "
+                        "WHERE id = CAST(:id AS uuid)"
+                    ),
+                    {"id": reservation_id},
+                )
+            session.execute(
+                text(
+                    "DELETE FROM public.users WHERE id = CAST(:id AS uuid)"
+                ),
+                {"id": actor},
+            )
+            session.commit()
+        pg.dispose()

+ 68 - 2
tests/integration/test_dataflow_draft_reservation_postgres.py

@@ -12,13 +12,26 @@ from app.core.common.identifiers import new_governance_uid
 from app.core.data_rules.repository import DataRuleRepository
 from app.core.data_rules.repository import DataRuleRepository
 
 
 
 
-def _claim(url, receipt, actor):
+def _request(receipt, *, name="Saga Acceptance"):
+    intent = {
+        "dataflow_uid": receipt["dataflow_uid"],
+        "node": {"uid": receipt["dataflow_uid"], "name_zh": name},
+        "tags": [],
+    }
+    return {"payload": {"name_zh": name}, "intent": intent}, intent
+
+
+def _claim(url, receipt, actor, *, name="Saga Acceptance"):
     engine = create_engine(url)
     engine = create_engine(url)
     try:
     try:
         with Session(engine) as session:
         with Session(engine) as session:
             try:
             try:
+                create_request, create_intent = _request(receipt, name=name)
                 value = DataRuleRepository(session).begin_dataflow_create(
                 value = DataRuleRepository(session).begin_dataflow_create(
-                    receipt, actor_uid=actor
+                    receipt,
+                    actor_uid=actor,
+                    create_request=create_request,
+                    create_intent=create_intent,
                 )
                 )
                 session.commit()
                 session.commit()
                 return value
                 return value
@@ -29,6 +42,19 @@ def _claim(url, receipt, actor):
         engine.dispose()
         engine.dispose()
 
 
 
 
+def _reconciliation_claim(url):
+    engine = create_engine(url)
+    try:
+        with Session(engine) as session:
+            rows = DataRuleRepository(
+                session
+            ).claim_reconcilable_dataflow_creates(limit=1)
+            session.commit()
+            return len(rows)
+    finally:
+        engine.dispose()
+
+
 def test_real_postgres_saga_fk_expiry_concurrent_claim_replay_and_lease_recovery():
 def test_real_postgres_saga_fk_expiry_concurrent_claim_replay_and_lease_recovery():
     url = os.environ.get("DATA_RULE_POSTGRES_ACCEPTANCE_URL")
     url = os.environ.get("DATA_RULE_POSTGRES_ACCEPTANCE_URL")
     if not url:
     if not url:
@@ -76,6 +102,7 @@ def test_real_postgres_saga_fk_expiry_concurrent_claim_replay_and_lease_recovery
                 reservation_id=receipt["reservation_id"],
                 reservation_id=receipt["reservation_id"],
                 dataflow_uid=receipt["dataflow_uid"],
                 dataflow_uid=receipt["dataflow_uid"],
                 lease_token=claimed[0]["lease_token"],
                 lease_token=claimed[0]["lease_token"],
+                request_digest=claimed[0]["request_digest"],
                 dataflow_node_id=73,
                 dataflow_node_id=73,
                 result=result,
                 result=result,
             ) == result
             ) == result
@@ -83,6 +110,9 @@ def test_real_postgres_saga_fk_expiry_concurrent_claim_replay_and_lease_recovery
         replay = _claim(url, closed, actor)
         replay = _claim(url, closed, actor)
         assert replay["status"] == "completed"
         assert replay["status"] == "completed"
         assert replay["result"] == result
         assert replay["result"] == result
+        assert _claim(url, closed, actor, name="Tampered")["status"] == (
+            "dataflow_create_request_conflict"
+        )
 
 
         with Session(engine) as session:
         with Session(engine) as session:
             expired = DataRuleRepository(session).reserve_dataflow_draft(
             expired = DataRuleRepository(session).reserve_dataflow_draft(
@@ -119,6 +149,8 @@ def test_real_postgres_saga_fk_expiry_concurrent_claim_replay_and_lease_recovery
                     for key in ("reservation_id", "dataflow_uid", "nonce")
                     for key in ("reservation_id", "dataflow_uid", "nonce")
                 },
                 },
                 actor_uid=actor,
                 actor_uid=actor,
+                create_request=_request(recoverable)[0],
+                create_intent=_request(recoverable)[1],
                 lease_seconds=10,
                 lease_seconds=10,
             )
             )
             session.commit()
             session.commit()
@@ -143,6 +175,40 @@ def test_real_postgres_saga_fk_expiry_concurrent_claim_replay_and_lease_recovery
         assert recovered["status"] == "claimed"
         assert recovered["status"] == "claimed"
         assert recovered["attempt"] == 2
         assert recovered["attempt"] == 2
 
 
+        with Session(engine) as session:
+            concurrent = DataRuleRepository(
+                session
+            ).reserve_dataflow_draft(actor_uid=actor)
+            created_ids.append(concurrent["reservation_id"])
+            session.commit()
+            concurrent_closed = {
+                key: concurrent[key]
+                for key in ("reservation_id", "dataflow_uid", "nonce")
+            }
+            request_value, intent_value = _request(concurrent)
+            DataRuleRepository(session).begin_dataflow_create(
+                concurrent_closed,
+                actor_uid=actor,
+                create_request=request_value,
+                create_intent=intent_value,
+                lease_seconds=10,
+            )
+            session.commit()
+            session.execute(
+                text(
+                    "UPDATE public.dataflow_draft_reservations "
+                    "SET lease_expires_at = CURRENT_TIMESTAMP - INTERVAL '1 second' "
+                    "WHERE id = CAST(:id AS uuid)"
+                ),
+                {"id": concurrent["reservation_id"]},
+            )
+            session.commit()
+        with ThreadPoolExecutor(max_workers=2) as pool:
+            reconciliation_claims = list(
+                pool.map(lambda _index: _reconciliation_claim(url), range(2))
+            )
+        assert sorted(reconciliation_claims) == [0, 1]
+
         with Session(engine) as session:
         with Session(engine) as session:
             with pytest.raises(IntegrityError):
             with pytest.raises(IntegrityError):
                 DataRuleRepository(session).reserve_dataflow_draft(
                 DataRuleRepository(session).reserve_dataflow_draft(

+ 103 - 1
tests/test_dataflow_create_saga.py

@@ -5,6 +5,7 @@ import copy
 import pytest
 import pytest
 
 
 from app.core.common.identifiers import new_governance_uid
 from app.core.common.identifiers import new_governance_uid
+from app.core.data_flow.create_reconciliation import DataFlowCreateReconciler
 from app.core.data_flow.dataflows import DataFlowService
 from app.core.data_flow.dataflows import DataFlowService
 
 
 
 
@@ -116,12 +117,21 @@ def test_completed_saga_replays_result_without_another_neo4j_write(monkeypatch):
         def load_published_assets(self, _flow):
         def load_published_assets(self, _flow):
             return {}, {}
             return {}, {}
 
 
-        def begin_dataflow_create(self, _receipt, *, actor_uid):
+        def begin_dataflow_create(
+            self,
+            _receipt,
+            *,
+            actor_uid,
+            create_request,
+            create_intent,
+        ):
             assert actor_uid
             assert actor_uid
+            assert create_request["intent"] == create_intent
             return {
             return {
                 "status": "completed",
                 "status": "completed",
                 "dataflow_uid": expected["uid"],
                 "dataflow_uid": expected["uid"],
                 "result": expected,
                 "result": expected,
+                "request_digest": "a" * 64,
             }
             }
 
 
     flow = {
     flow = {
@@ -243,3 +253,95 @@ def test_graph_failure_is_persisted_but_finalize_failure_leaves_reconcilable_lea
             repository=finalize_failure,
             repository=finalize_failure,
             actor_uid=actor,
             actor_uid=actor,
         )
         )
+
+
+def test_reconciler_dry_run_is_read_only_and_digest_failure_is_isolated(
+    monkeypatch,
+):
+    good_uid = new_governance_uid()
+    bad_uid = new_governance_uid()
+
+    class Session:
+        def __init__(self):
+            self.commits = 0
+            self.rollbacks = 0
+
+        def commit(self):
+            self.commits += 1
+
+        def rollback(self):
+            self.rollbacks += 1
+
+    class Repository:
+        def __init__(self):
+            self.session = Session()
+            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 [
+                {
+                    "reservation_id": new_governance_uid(),
+                    "dataflow_uid": bad_uid,
+                    "lease_token": new_governance_uid(),
+                    "request_digest": "a" * 64,
+                    "create_intent": {},
+                    "attempt": 2,
+                    "integrity_valid": False,
+                },
+                {
+                    "reservation_id": new_governance_uid(),
+                    "dataflow_uid": good_uid,
+                    "lease_token": new_governance_uid(),
+                    "request_digest": "b" * 64,
+                    "create_intent": {"dataflow_uid": good_uid},
+                    "attempt": 2,
+                    "integrity_valid": True,
+                },
+            ]
+
+        def commit_dataflow_create_failure(self, **kwargs):
+            self.failures.append(kwargs)
+            self.session.commit()
+
+        def complete_dataflow_create(self, **kwargs):
+            self.completions.append(kwargs)
+            return {"id": 73, "uid": kwargs["dataflow_uid"]}
+
+    repository = Repository()
+    reconciler = DataFlowCreateReconciler(repository)
+    preview = reconciler.run(dry_run=True, limit=2)
+    assert preview["candidate_count"] == 1
+    assert repository.claim_calls == 0
+    assert repository.session.commits == 0
+
+    monkeypatch.setattr(
+        DataFlowService,
+        "validate_governed_create_intent",
+        lambda value, *, repository: {
+            "node": {"uid": value["dataflow_uid"]},
+            "tags": [],
+        },
+    )
+    monkeypatch.setattr(
+        DataFlowService,
+        "_merge_governed_dataflow",
+        lambda node: (73, {"id": 73, **node}),
+    )
+    report = reconciler.run(dry_run=False, limit=2)
+    assert report["reconciled_count"] == 1
+    assert report["failed_count"] == 1
+    assert repository.failures[0]["error_code"] == (
+        "create_intent_integrity_failed"
+    )
+    assert repository.completions[0]["dataflow_uid"] == good_uid

+ 21 - 2
tests/test_legacy_governance_cutover.py

@@ -1,5 +1,6 @@
 from __future__ import annotations
 from __future__ import annotations
 
 
+import hashlib
 import json
 import json
 import uuid
 import uuid
 from datetime import UTC, datetime, timedelta
 from datetime import UTC, datetime, timedelta
@@ -39,13 +40,24 @@ class PublishedAssetRepository:
             "expires_at": "2099-01-01T00:00:00+00:00",
             "expires_at": "2099-01-01T00:00:00+00:00",
         }
         }
 
 
-    def begin_dataflow_create(self, receipt, *, actor_uid):
+    def begin_dataflow_create(
+        self, receipt, *, actor_uid, create_request, create_intent
+    ):
         key = (receipt["reservation_id"], actor_uid)
         key = (receipt["reservation_id"], actor_uid)
+        request_digest = hashlib.sha256(
+            json.dumps(
+                create_request,
+                sort_keys=True,
+                separators=(",", ":"),
+                ensure_ascii=False,
+            ).encode("utf-8")
+        ).hexdigest()
         if key in self.completed:
         if key in self.completed:
             return {
             return {
                 "status": "completed",
                 "status": "completed",
                 "dataflow_uid": receipt["dataflow_uid"],
                 "dataflow_uid": receipt["dataflow_uid"],
                 "result": self.completed[key],
                 "result": self.completed[key],
+                "request_digest": request_digest,
             }
             }
         self.consumed.add(key)
         self.consumed.add(key)
         return {
         return {
@@ -53,6 +65,8 @@ class PublishedAssetRepository:
             "dataflow_uid": receipt["dataflow_uid"],
             "dataflow_uid": receipt["dataflow_uid"],
             "lease_token": new_governance_uid(),
             "lease_token": new_governance_uid(),
             "attempt": 1,
             "attempt": 1,
+            "create_intent": create_intent,
+            "request_digest": request_digest,
         }
         }
 
 
     def commit_dataflow_create_claim(self):
     def commit_dataflow_create_claim(self):
@@ -342,7 +356,9 @@ def test_governed_dataflow_creation_never_generates_legacy_task_or_workflow(
     }
     }
     monkeypatch.setattr(
     monkeypatch.setattr(
         "app.core.data_flow.dataflows.translate_and_parse",
         "app.core.data_flow.dataflows.translate_and_parse",
-        lambda _name: ["customer_governed_line"],
+        lambda _name: pytest.fail(
+            "governed create used non-deterministic translation"
+        ),
     )
     )
     created = {}
     created = {}
     monkeypatch.setattr(
     monkeypatch.setattr(
@@ -377,6 +393,9 @@ def test_governed_dataflow_creation_never_generates_legacy_task_or_workflow(
 
 
     assert result["id"] == 31
     assert result["id"] == 31
     assert created["uid"] == flow["dataflow_uid"]
     assert created["uid"] == flow["dataflow_uid"]
+    assert created["name_en"] == (
+        f"dataflow_{flow['dataflow_uid'].replace('-', '')}"
+    )
     assert json.loads(created["script_requirement"])["dataflow_spec"] == flow
     assert json.loads(created["script_requirement"])["dataflow_spec"] == flow