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

feat: resolve trusted rule schemas and datasets

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

+ 9 - 9
app/api/data_rules/routes.py

@@ -23,9 +23,9 @@ from app.core.data_rules.contracts import (
 from app.core.data_rules.production_line import resolve_production_line
 from app.core.data_rules.release import ProductionLineReleaseService
 from app.core.data_rules.repository import DataRuleRepository
+from app.core.data_rules.schema_resolver import SchemaResolver
 from app.models.result import failed, success
 
-
 _VALIDATORS = {
     "rule": (validate_rule_spec, rule_spec_hash),
     "standard": (validate_standard_spec, standard_spec_hash),
@@ -59,12 +59,16 @@ def _repository() -> DataRuleRepository:
 
 
 def _release_service() -> ProductionLineReleaseService:
-    configured = current_app.extensions.get(
-        "production_line_release_service"
-    )
+    configured = current_app.extensions.get("production_line_release_service")
     if configured is not None:
         return configured
-    return ProductionLineReleaseService(_repository())
+    repository = _repository()
+    resolver = current_app.extensions.get("data_rule_schema_resolver")
+    if resolver is None:
+        catalog = current_app.extensions.get("data_rule_metadata_catalog")
+        if catalog is not None:
+            resolver = SchemaResolver(catalog, repository)
+    return ProductionLineReleaseService(repository, schema_resolver=resolver)
 
 
 @bp.get("/capabilities")
@@ -276,16 +280,12 @@ def release_production_line(dataflow_uid: str):
             {
                 "dataflow_spec",
                 "source_text",
-                "input_schema_hashes",
-                "output_schema_hash",
             }
         )
         result = _release_service().release(
             dataflow_uid=dataflow_uid,
             dataflow_spec=body.get("dataflow_spec"),
             source_text=body.get("source_text"),
-            input_schema_hashes=body.get("input_schema_hashes"),
-            output_schema_hash=body.get("output_schema_hash"),
             created_by=g.current_user["id"],
         )
         db.session.commit()

+ 15 - 36
app/core/data_rules/release.py

@@ -2,7 +2,6 @@
 
 from __future__ import annotations
 
-import re
 from typing import Any
 
 from app.core.common.identifiers import ensure_governance_uid, new_governance_uid
@@ -27,32 +26,10 @@ def _source(value: Any) -> str:
     return normalized
 
 
-def _digest(value: Any, label: str) -> str:
-    normalized = str(value or "")
-    if not re.fullmatch(r"[0-9a-f]{64}", normalized):
-        raise ValueError(f"{label} must be a sha256 hex digest")
-    return normalized
-
-
-def _schema_hashes(
-    flow: dict[str, Any],
-    input_schema_hashes: Any,
-    output_schema_hash: Any,
-) -> tuple[dict[str, str], str]:
-    if not isinstance(input_schema_hashes, dict):
-        raise ValueError("input_schema_hashes must be an object")
-    if set(input_schema_hashes) != set(flow["input_schema_refs"]):
-        raise ValueError("input_schema_hashes must cover every input schema")
-    inputs = {
-        ref: _digest(input_schema_hashes[ref], f"schema hash for {ref}")
-        for ref in sorted(input_schema_hashes)
-    }
-    return inputs, _digest(output_schema_hash, "output_schema_hash")
-
-
 class ProductionLineReleaseService:
-    def __init__(self, repository):
+    def __init__(self, repository, *, schema_resolver=None):
         self.repository = repository
+        self.schema_resolver = schema_resolver
 
     def release(
         self,
@@ -60,8 +37,6 @@ class ProductionLineReleaseService:
         dataflow_uid: str,
         dataflow_spec: dict[str, Any],
         source_text: str,
-        input_schema_hashes: dict[str, str],
-        output_schema_hash: str,
         created_by: str,
     ) -> dict[str, Any]:
         path_uid = _uid(dataflow_uid, "dataflow_uid")
@@ -70,9 +45,17 @@ class ProductionLineReleaseService:
         if flow["dataflow_uid"] != path_uid:
             raise ValueError("dataflow spec uid does not match path dataflow_uid")
         source = _source(source_text)
-        inputs, output = _schema_hashes(
-            flow, input_schema_hashes, output_schema_hash
-        )
+        if self.schema_resolver is None:
+            raise ValueError("trusted schema resolver is not configured")
+        input_snapshots = {
+            ref: self.schema_resolver.resolve(ref) for ref in flow["input_schema_refs"]
+        }
+        output_snapshot = self.schema_resolver.resolve(flow["output_schema_ref"])
+        inputs = {
+            ref: snapshot["schema_hash"]
+            for ref, snapshot in sorted(input_snapshots.items())
+        }
+        output = output_snapshot["schema_hash"]
 
         standards, rules = self.repository.load_published_assets(flow)
         version = self.repository.begin_dataflow_release(
@@ -102,9 +85,7 @@ class ProductionLineReleaseService:
                 raise ValueError(
                     f"published rule version {rule_version_id} was not found"
                 )
-            plan = compiled.setdefault(
-                rule_version_id, compile_rule_plan(rule)
-            )
+            plan = compiled.setdefault(rule_version_id, compile_rule_plan(rule))
             binding_id = new_governance_uid()
             binding_ids[binding_key] = binding_id
             self.repository.persist_component_plan(
@@ -137,9 +118,7 @@ class ProductionLineReleaseService:
                     binding_key = f"{component['id']}:{clause_id}"
                     add_binding(
                         binding_key=binding_key,
-                        component_id=(
-                            f"{component['id']}__{clause_id}"[:100]
-                        ),
+                        component_id=(f"{component['id']}__{clause_id}"[:100]),
                         component_kind="quality.check",
                         rule_version_id=str(clause["rule_version_id"]),
                         stage=component["stage"],

+ 310 - 168
app/core/data_rules/repository.py

@@ -2,8 +2,8 @@
 
 from __future__ import annotations
 
-import json
 import hashlib
+import json
 import re
 from typing import Any
 
@@ -23,7 +23,10 @@ from app.core.data_rules.contracts import (
     validate_rule_spec,
     validate_standard_spec,
 )
-
+from app.core.data_rules.execution_contracts import (
+    validate_dataset_binding,
+    validate_schema_snapshot,
+)
 
 RULE_CATEGORIES = {
     "general",
@@ -175,15 +178,11 @@ class DataRuleRepository:
                 "id": generation_id,
                 "rule_version_id": version_id,
                 "authoring_surface": surface,
-                "source_text_hash": hashlib.sha256(
-                    source.encode("utf-8")
-                ).hexdigest(),
+                "source_text_hash": hashlib.sha256(source.encode("utf-8")).hexdigest(),
                 "model_provider": _text(
                     evidence.get("model_provider"), "model_provider", 80
                 ),
-                "model_name": _text(
-                    evidence.get("model_name"), "model_name", 120
-                ),
+                "model_name": _text(evidence.get("model_name"), "model_name", 120),
                 "prompt_version": _text(
                     evidence.get("prompt_version"), "prompt_version", 80
                 ),
@@ -197,9 +196,7 @@ class DataRuleRepository:
                 "ambiguities": _json(candidate["ambiguities"]),
                 "repair_attempts": repair_attempts,
                 "decision": decision,
-                "decision_detail": _json(
-                    {"explanation": candidate["explanation"]}
-                ),
+                "decision_detail": _json({"explanation": candidate["explanation"]}),
                 "correlation_id": correlation_id,
             },
         )
@@ -237,16 +234,20 @@ class DataRuleRepository:
         rule_uid = spec["rule_uid"]
         digest = rule_spec_hash(spec)
         self._lock(rule_uid)
-        existing = self.session.execute(
-            text(
-                "SELECT id::text AS id, version_no, status "
-                "FROM public.data_rule_versions "
-                "WHERE rule_uid = CAST(:rule_uid AS uuid) "
-                "AND spec_hash = :spec_hash "
-                "/* existing_version */"
-            ),
-            {"rule_uid": rule_uid, "spec_hash": digest},
-        ).mappings().one_or_none()
+        existing = (
+            self.session.execute(
+                text(
+                    "SELECT id::text AS id, version_no, status "
+                    "FROM public.data_rule_versions "
+                    "WHERE rule_uid = CAST(:rule_uid AS uuid) "
+                    "AND spec_hash = :spec_hash "
+                    "/* existing_version */"
+                ),
+                {"rule_uid": rule_uid, "spec_hash": digest},
+            )
+            .mappings()
+            .one_or_none()
+        )
         if existing is not None:
             return {
                 **dict(existing),
@@ -320,17 +321,21 @@ class DataRuleRepository:
     ) -> dict[str, Any]:
         version = _uid(version_id, "version_id")
         _uid(published_by, "published_by")
-        row = self.session.execute(
-            text(
-                "UPDATE public.data_rule_versions "
-                "SET status = 'published', published_at = CURRENT_TIMESTAMP "
-                "WHERE id = CAST(:version_id AS uuid) "
-                "AND status = 'validated' "
-                "RETURNING id::text AS id, rule_uid::text AS rule_uid, "
-                "version_no, status, spec_hash"
-            ),
-            {"version_id": version},
-        ).mappings().one_or_none()
+        row = (
+            self.session.execute(
+                text(
+                    "UPDATE public.data_rule_versions "
+                    "SET status = 'published', published_at = CURRENT_TIMESTAMP "
+                    "WHERE id = CAST(:version_id AS uuid) "
+                    "AND status = 'validated' "
+                    "RETURNING id::text AS id, rule_uid::text AS rule_uid, "
+                    "version_no, status, spec_hash"
+                ),
+                {"version_id": version},
+            )
+            .mappings()
+            .one_or_none()
+        )
         if row is not None:
             return dict(row)
         status = self.session.execute(
@@ -363,33 +368,39 @@ class DataRuleRepository:
             {clause["rule_version_id"] for clause in spec["clauses"]}
         )
 
-        published_rows = self.session.execute(
-            text(
-                "SELECT id::text AS id "
-                "FROM public.data_rule_versions "
-                "WHERE id = ANY(CAST(:rule_version_ids AS uuid[])) "
-                "AND status = 'published' "
-                "/* published_rule_reference */"
-            ),
-            {"rule_version_ids": rule_version_ids},
-        ).mappings().all()
+        published_rows = (
+            self.session.execute(
+                text(
+                    "SELECT id::text AS id "
+                    "FROM public.data_rule_versions "
+                    "WHERE id = ANY(CAST(:rule_version_ids AS uuid[])) "
+                    "AND status = 'published' "
+                    "/* published_rule_reference */"
+                ),
+                {"rule_version_ids": rule_version_ids},
+            )
+            .mappings()
+            .all()
+        )
         published_ids = {str(row["id"]) for row in published_rows}
         if published_ids != set(rule_version_ids):
-            raise ValueError(
-                "standard clauses require published rule versions"
-            )
+            raise ValueError("standard clauses require published rule versions")
 
         self._lock(standard_uid)
-        existing = self.session.execute(
-            text(
-                "SELECT id::text AS id, version_no, status "
-                "FROM public.data_standard_versions "
-                "WHERE standard_uid = CAST(:standard_uid AS uuid) "
-                "AND spec_hash = :spec_hash "
-                "/* existing_version */"
-            ),
-            {"standard_uid": standard_uid, "spec_hash": digest},
-        ).mappings().one_or_none()
+        existing = (
+            self.session.execute(
+                text(
+                    "SELECT id::text AS id, version_no, status "
+                    "FROM public.data_standard_versions "
+                    "WHERE standard_uid = CAST(:standard_uid AS uuid) "
+                    "AND spec_hash = :spec_hash "
+                    "/* existing_version */"
+                ),
+                {"standard_uid": standard_uid, "spec_hash": digest},
+            )
+            .mappings()
+            .one_or_none()
+        )
         if existing is not None:
             return {
                 **dict(existing),
@@ -466,9 +477,7 @@ class DataRuleRepository:
                     },
                 )
         except IntegrityError as exc:
-            raise ValueError(
-                "standard version conflicts with existing data"
-            ) from exc
+            raise ValueError("standard version conflicts with existing data") from exc
         return {
             "id": version_id,
             "standard_uid": standard_uid,
@@ -483,17 +492,21 @@ class DataRuleRepository:
     ) -> dict[str, Any]:
         version = _uid(version_id, "version_id")
         _uid(published_by, "published_by")
-        row = self.session.execute(
-            text(
-                "UPDATE public.data_standard_versions "
-                "SET status = 'published', published_at = CURRENT_TIMESTAMP "
-                "WHERE id = CAST(:version_id AS uuid) "
-                "AND status = 'validated' "
-                "RETURNING id::text AS id, standard_uid::text AS standard_uid, "
-                "version_no, status, spec_hash"
-            ),
-            {"version_id": version},
-        ).mappings().one_or_none()
+        row = (
+            self.session.execute(
+                text(
+                    "UPDATE public.data_standard_versions "
+                    "SET status = 'published', published_at = CURRENT_TIMESTAMP "
+                    "WHERE id = CAST(:version_id AS uuid) "
+                    "AND status = 'validated' "
+                    "RETURNING id::text AS id, standard_uid::text AS standard_uid, "
+                    "version_no, status, spec_hash"
+                ),
+                {"version_id": version},
+            )
+            .mappings()
+            .one_or_none()
+        )
         if row is not None:
             return dict(row)
         status = self.session.execute(
@@ -505,9 +518,7 @@ class DataRuleRepository:
             {"version_id": version},
         ).scalar_one_or_none()
         if status == "published":
-            raise ValueError(
-                "standard version is already published and immutable"
-            )
+            raise ValueError("standard version is already published and immutable")
         if status is None:
             raise ValueError("standard version was not found")
         raise ValueError("only validated standard versions may be published")
@@ -527,25 +538,29 @@ class DataRuleRepository:
         )
         standard_rows = []
         if standard_ids:
-            standard_rows = self.session.execute(
-                text(
-                    "SELECT sv.id::text AS id, sv.status, "
-                    "jsonb_agg(jsonb_build_object("
-                    "'clause_id', b.clause_id, "
-                    "'rule_version_id', b.rule_version_id::text, "
-                    "'severity', b.severity, "
-                    "'exception_policy', b.exception_policy) "
-                    "ORDER BY b.clause_id) AS clauses "
-                    "FROM public.data_standard_versions sv "
-                    "JOIN public.standard_rule_bindings b "
-                    "ON b.standard_version_id = sv.id "
-                    "WHERE sv.id = ANY(CAST(:standard_ids AS uuid[])) "
-                    "AND sv.status = 'published' "
-                    "/* published_standard_assets */ "
-                    "GROUP BY sv.id, sv.status"
-                ),
-                {"standard_ids": standard_ids},
-            ).mappings().all()
+            standard_rows = (
+                self.session.execute(
+                    text(
+                        "SELECT sv.id::text AS id, sv.status, "
+                        "jsonb_agg(jsonb_build_object("
+                        "'clause_id', b.clause_id, "
+                        "'rule_version_id', b.rule_version_id::text, "
+                        "'severity', b.severity, "
+                        "'exception_policy', b.exception_policy) "
+                        "ORDER BY b.clause_id) AS clauses "
+                        "FROM public.data_standard_versions sv "
+                        "JOIN public.standard_rule_bindings b "
+                        "ON b.standard_version_id = sv.id "
+                        "WHERE sv.id = ANY(CAST(:standard_ids AS uuid[])) "
+                        "AND sv.status = 'published' "
+                        "/* published_standard_assets */ "
+                        "GROUP BY sv.id, sv.status"
+                    ),
+                    {"standard_ids": standard_ids},
+                )
+                .mappings()
+                .all()
+            )
         standards = {
             str(row["id"]): {
                 "id": str(row["id"]),
@@ -566,22 +581,25 @@ class DataRuleRepository:
         }
         for standard in standards.values():
             rule_ids.update(
-                str(clause["rule_version_id"])
-                for clause in standard["clauses"]
+                str(clause["rule_version_id"]) for clause in standard["clauses"]
             )
         rule_rows = []
         if rule_ids:
-            rule_rows = self.session.execute(
-                text(
-                    "SELECT rv.id::text AS id, rv.status, rv.rule_spec, "
-                    "rv.spec_hash "
-                    "FROM public.data_rule_versions rv "
-                    "WHERE rv.id = ANY(CAST(:rule_ids AS uuid[])) "
-                    "AND rv.status = 'published' "
-                    "/* published_rule_assets */"
-                ),
-                {"rule_ids": sorted(rule_ids)},
-            ).mappings().all()
+            rule_rows = (
+                self.session.execute(
+                    text(
+                        "SELECT rv.id::text AS id, rv.status, rv.rule_spec, "
+                        "rv.spec_hash "
+                        "FROM public.data_rule_versions rv "
+                        "WHERE rv.id = ANY(CAST(:rule_ids AS uuid[])) "
+                        "AND rv.status = 'published' "
+                        "/* published_rule_assets */"
+                    ),
+                    {"rule_ids": sorted(rule_ids)},
+                )
+                .mappings()
+                .all()
+            )
         rules = {
             str(row["id"]): {
                 "id": str(row["id"]),
@@ -592,9 +610,7 @@ class DataRuleRepository:
             for row in rule_rows
         }
         if set(rules) != rule_ids:
-            raise ValueError(
-                "dataflow references missing or unpublished rule versions"
-            )
+            raise ValueError("dataflow references missing or unpublished rule versions")
         return standards, rules
 
     def begin_dataflow_release(
@@ -661,6 +677,125 @@ class DataRuleRepository:
             "dataflow_spec_hash": dataflow_spec_hash(flow),
         }
 
+    def find_schema_snapshot(
+        self, *, schema_ref: str, schema_hash: str
+    ) -> dict[str, Any] | None:
+        row = (
+            self.session.execute(
+                text(
+                    "SELECT id::text AS id, schema_ref, schema_hash, fields, "
+                    "source_revision FROM public.data_schema_snapshots "
+                    "WHERE schema_ref = :schema_ref AND schema_hash = :schema_hash"
+                ),
+                {
+                    "schema_ref": _text(schema_ref, "schema_ref", 500),
+                    "schema_hash": _digest(schema_hash, "schema_hash"),
+                },
+            )
+            .mappings()
+            .one_or_none()
+        )
+        if row is None:
+            return None
+        snapshot = validate_schema_snapshot(
+            {
+                "schema_ref": str(row["schema_ref"]),
+                "schema_hash": str(row["schema_hash"]),
+                "fields": _array(row["fields"], "schema snapshot fields"),
+                "source_revision": str(row["source_revision"]),
+            }
+        )
+        return {"id": str(row["id"]), **snapshot}
+
+    def persist_schema_snapshot(self, *, snapshot: dict[str, Any]) -> dict[str, Any]:
+        normalized = validate_schema_snapshot(snapshot)
+        existing = self.find_schema_snapshot(
+            schema_ref=normalized["schema_ref"], schema_hash=normalized["schema_hash"]
+        )
+        if existing is not None:
+            return existing
+        snapshot_id = new_governance_uid()
+        inserted = (
+            self.session.execute(
+                text(
+                    "INSERT INTO public.data_schema_snapshots "
+                    "(id, schema_ref, schema_hash, fields, source_revision) "
+                    "VALUES (CAST(:id AS uuid), :schema_ref, :schema_hash, "
+                    "CAST(:fields AS jsonb), :source_revision) "
+                    "ON CONFLICT (schema_ref, schema_hash) DO NOTHING "
+                    "RETURNING id::text AS id"
+                ),
+                {
+                    "id": snapshot_id,
+                    "schema_ref": normalized["schema_ref"],
+                    "schema_hash": normalized["schema_hash"],
+                    "fields": _json(normalized["fields"]),
+                    "source_revision": normalized["source_revision"],
+                },
+            )
+            .mappings()
+            .one_or_none()
+        )
+        if inserted is None:
+            existing = self.find_schema_snapshot(
+                schema_ref=normalized["schema_ref"],
+                schema_hash=normalized["schema_hash"],
+            )
+            if existing is not None:
+                return existing
+            raise ValueError("schema snapshot could not be persisted")
+        return {"id": str(inserted["id"]), **normalized}
+
+    def persist_dataset_binding(
+        self,
+        *,
+        deployment_id: str,
+        logical_ref: str,
+        binding: dict[str, Any],
+        actor_uid: str,
+    ) -> dict[str, Any]:
+        deployment = _uid(deployment_id, "deployment_id")
+        _uid(actor_uid, "actor_uid")
+        logical = _text(logical_ref, "logical_ref", 500)
+        normalized = validate_dataset_binding(binding)
+        binding_hash = _canonical_hash({"logical_ref": logical, **normalized})
+        binding_id = new_governance_uid()
+        try:
+            row = (
+                self.session.execute(
+                    text(
+                        "INSERT INTO public.dataflow_dataset_bindings "
+                        "(id, dataflow_deployment_id, logical_ref, data_source_uid, "
+                        "object_kind, object_ref, schema_snapshot_id, dialect, "
+                        "access_mode, write_mode, binding_hash) "
+                        "VALUES (CAST(:id AS uuid), CAST(:deployment_id AS uuid), "
+                        ":logical_ref, CAST(:data_source_uid AS uuid), :object_kind, "
+                        ":object_ref, CAST(:schema_snapshot_id AS uuid), :dialect, "
+                        ":access_mode, :write_mode, :binding_hash) "
+                        "RETURNING id::text AS id"
+                    ),
+                    {
+                        "id": binding_id,
+                        "deployment_id": deployment,
+                        "logical_ref": logical,
+                        **normalized,
+                        "binding_hash": binding_hash,
+                    },
+                )
+                .mappings()
+                .one()
+            )
+        except IntegrityError as exc:
+            raise ValueError(
+                "dataset binding already exists for deployment logical_ref"
+            ) from exc
+        return {
+            "id": str(row["id"]),
+            "logical_ref": logical,
+            **normalized,
+            "binding_hash": binding_hash,
+        }
+
     def persist_component_plan(
         self,
         *,
@@ -697,9 +832,7 @@ class DataRuleRepository:
             raise ValueError("compiled plan contains unsupported fields")
         if plan["backend"] not in PLAN_BACKENDS:
             raise ValueError("unsupported plan backend")
-        compiler_version = _text(
-            plan["compiler_version"], "compiler_version", 80
-        )
+        compiler_version = _text(plan["compiler_version"], "compiler_version", 80)
         plan_body = _object(plan["plan"], "compiled plan body")
         plan_hash = _digest(plan["plan_hash"], "plan_hash")
         if _canonical_hash(plan_body) != plan_hash:
@@ -730,9 +863,7 @@ class DataRuleRepository:
                 "rule_version_id": rule_id,
                 "stage": _text(stage, "stage", 30),
                 "order_no": order_no,
-                "idempotency": _json(idempotency)
-                if idempotency is not None
-                else None,
+                "idempotency": _json(idempotency) if idempotency is not None else None,
                 "provenance": _json(provenance),
             },
         )
@@ -766,22 +897,26 @@ class DataRuleRepository:
         if not isinstance(package, dict):
             raise ValueError("production-line package must be an object")
         package_hash = _digest(package.get("package_hash"), "package_hash")
-        row = self.session.execute(
-            text(
-                "UPDATE public.dataflow_versions "
-                "SET package = CAST(:package AS jsonb), "
-                "package_hash = :package_hash, status = 'released', "
-                "released_at = CURRENT_TIMESTAMP "
-                "WHERE id = CAST(:version_id AS uuid) "
-                "AND status = 'validated' "
-                "RETURNING id::text AS id, version_no, status, package_hash"
-            ),
-            {
-                "version_id": version_id,
-                "package": _json(package),
-                "package_hash": package_hash,
-            },
-        ).mappings().one_or_none()
+        row = (
+            self.session.execute(
+                text(
+                    "UPDATE public.dataflow_versions "
+                    "SET package = CAST(:package AS jsonb), "
+                    "package_hash = :package_hash, status = 'released', "
+                    "released_at = CURRENT_TIMESTAMP "
+                    "WHERE id = CAST(:version_id AS uuid) "
+                    "AND status = 'validated' "
+                    "RETURNING id::text AS id, version_no, status, package_hash"
+                ),
+                {
+                    "version_id": version_id,
+                    "package": _json(package),
+                    "package_hash": package_hash,
+                },
+            )
+            .mappings()
+            .one_or_none()
+        )
         if row is None:
             raise ValueError("dataflow release is not in validated state")
         return {**dict(row), "package": package}
@@ -805,25 +940,29 @@ class DataRuleRepository:
 
         standard_rows = []
         if standard_ids:
-            standard_rows = self.session.execute(
-                text(
-                    "SELECT sv.id::text AS id, sv.status, "
-                    "jsonb_agg(jsonb_build_object("
-                    "'clause_id', b.clause_id, "
-                    "'rule_version_id', b.rule_version_id::text, "
-                    "'severity', b.severity, "
-                    "'exception_policy', b.exception_policy) "
-                    "ORDER BY b.clause_id) AS clauses "
-                    "FROM public.data_standard_versions sv "
-                    "JOIN public.standard_rule_bindings b "
-                    "ON b.standard_version_id = sv.id "
-                    "WHERE sv.id = ANY(CAST(:standard_ids AS uuid[])) "
-                    "AND sv.status = 'published' "
-                    "/* catalog_standard_versions */ "
-                    "GROUP BY sv.id, sv.status"
-                ),
-                {"standard_ids": standard_ids},
-            ).mappings().all()
+            standard_rows = (
+                self.session.execute(
+                    text(
+                        "SELECT sv.id::text AS id, sv.status, "
+                        "jsonb_agg(jsonb_build_object("
+                        "'clause_id', b.clause_id, "
+                        "'rule_version_id', b.rule_version_id::text, "
+                        "'severity', b.severity, "
+                        "'exception_policy', b.exception_policy) "
+                        "ORDER BY b.clause_id) AS clauses "
+                        "FROM public.data_standard_versions sv "
+                        "JOIN public.standard_rule_bindings b "
+                        "ON b.standard_version_id = sv.id "
+                        "WHERE sv.id = ANY(CAST(:standard_ids AS uuid[])) "
+                        "AND sv.status = 'published' "
+                        "/* catalog_standard_versions */ "
+                        "GROUP BY sv.id, sv.status"
+                    ),
+                    {"standard_ids": standard_ids},
+                )
+                .mappings()
+                .all()
+            )
         standards = {
             str(row["id"]): {
                 "id": str(row["id"]),
@@ -840,29 +979,32 @@ class DataRuleRepository:
         rule_ids = set(direct_rule_ids)
         for standard in standards.values():
             rule_ids.update(
-                str(clause["rule_version_id"])
-                for clause in standard["clauses"]
+                str(clause["rule_version_id"]) for clause in standard["clauses"]
             )
         rule_rows = []
         if rule_ids:
-            rule_rows = self.session.execute(
-                text(
-                    "SELECT DISTINCT ON (rv.id) "
-                    "rv.id::text AS id, rv.status, rv.rule_spec, rv.spec_hash, "
-                    "p.backend, p.plan_hash "
-                    "FROM public.data_rule_versions rv "
-                    "JOIN public.dataflow_component_bindings cb "
-                    "ON cb.rule_version_id = rv.id "
-                    "JOIN public.rule_execution_plans p "
-                    "ON p.component_binding_id = cb.id "
-                    "WHERE rv.id = ANY(CAST(:rule_ids AS uuid[])) "
-                    "AND rv.status = 'published' "
-                    "AND p.status = 'published' "
-                    "/* catalog_rule_versions */ "
-                    "ORDER BY rv.id, p.created_at DESC"
-                ),
-                {"rule_ids": sorted(rule_ids)},
-            ).mappings().all()
+            rule_rows = (
+                self.session.execute(
+                    text(
+                        "SELECT DISTINCT ON (rv.id) "
+                        "rv.id::text AS id, rv.status, rv.rule_spec, rv.spec_hash, "
+                        "p.backend, p.plan_hash "
+                        "FROM public.data_rule_versions rv "
+                        "JOIN public.dataflow_component_bindings cb "
+                        "ON cb.rule_version_id = rv.id "
+                        "JOIN public.rule_execution_plans p "
+                        "ON p.component_binding_id = cb.id "
+                        "WHERE rv.id = ANY(CAST(:rule_ids AS uuid[])) "
+                        "AND rv.status = 'published' "
+                        "AND p.status = 'published' "
+                        "/* catalog_rule_versions */ "
+                        "ORDER BY rv.id, p.created_at DESC"
+                    ),
+                    {"rule_ids": sorted(rule_ids)},
+                )
+                .mappings()
+                .all()
+            )
         rules = {
             str(row["id"]): {
                 "id": str(row["id"]),

+ 221 - 0
app/core/data_rules/schema_resolver.py

@@ -0,0 +1,221 @@
+"""Trusted schema snapshots and deployment-time physical dataset bindings."""
+
+from __future__ import annotations
+
+import copy
+import re
+from typing import Any, Protocol
+
+from sqlalchemy import text
+
+from app.core.common.identifiers import ensure_governance_uid
+from app.core.data_rules.execution_contracts import (
+    canonical_schema_hash,
+    validate_dataset_binding,
+    validate_schema_snapshot,
+)
+
+
+class SchemaMetadataCatalog(Protocol):
+    """Narrow stable-reference lookup supplied by the metadata catalog."""
+
+    def load_schema(self, schema_ref: str) -> dict[str, Any] | None: ...
+
+
+class SchemaSnapshotRepository(Protocol):
+    def find_schema_snapshot(
+        self, *, schema_ref: str, schema_hash: str
+    ) -> dict[str, Any] | None: ...
+
+    def persist_schema_snapshot(
+        self, *, snapshot: dict[str, Any]
+    ) -> dict[str, Any]: ...
+
+
+_SUPPORTED_FIELD_TYPES = {
+    "binary",
+    "boolean",
+    "date",
+    "decimal",
+    "double",
+    "float",
+    "integer",
+    "json",
+    "string",
+    "timestamp",
+    "timestamptz",
+}
+_FIELD_TYPE_ALIASES = {
+    "bigint": "integer",
+    "bool": "boolean",
+    "character varying": "string",
+    "datetime": "timestamp",
+    "int": "integer",
+    "numeric": "decimal",
+    "text": "string",
+    "varchar": "string",
+}
+_IDENTIFIER = re.compile(r"^[A-Za-z_][A-Za-z0-9_]{0,127}$")
+
+
+def _uid(value: Any, label: str) -> str:
+    try:
+        return ensure_governance_uid({"uid": str(value)})
+    except ValueError as exc:
+        raise ValueError(f"{label} must be a valid UUIDv7") from exc
+
+
+def _required_string(value: Any, label: str, maximum: int) -> str:
+    if not isinstance(value, str) or not value.strip():
+        raise ValueError(f"{label} is required")
+    normalized = value.strip()
+    if len(normalized) > maximum:
+        raise ValueError(f"{label} exceeds {maximum} characters")
+    return normalized
+
+
+def _normalized_fields(fields: Any) -> list[dict[str, Any]]:
+    if not isinstance(fields, list):
+        raise ValueError("metadata schema fields must be an array")
+    normalized: list[dict[str, Any]] = []
+    for field in fields:
+        if not isinstance(field, dict):
+            raise ValueError("metadata schema field must be an object")
+        item = copy.deepcopy(field)
+        field_type = item.get("type")
+        if not isinstance(field_type, str):
+            raise ValueError("metadata schema field type is required")
+        field_type = _FIELD_TYPE_ALIASES.get(
+            field_type.strip().lower(), field_type.strip().lower()
+        )
+        if field_type not in _SUPPORTED_FIELD_TYPES:
+            raise ValueError("unsupported metadata schema field type")
+        item["type"] = field_type
+        normalized.append(item)
+    # The execution contract performs the closed-shape, duplicate-name, and
+    # precision/scale/timezone validation used by the runtime as well.
+    return validate_schema_snapshot(
+        {
+            "schema_ref": "metadata:validation",
+            "schema_hash": canonical_schema_hash(normalized),
+            "fields": normalized,
+            "source_revision": "metadata:validation",
+        }
+    )["fields"]
+
+
+class SchemaResolver:
+    """Resolve stable schema references without trusting request payloads."""
+
+    def __init__(
+        self,
+        metadata_catalog: SchemaMetadataCatalog,
+        repository: SchemaSnapshotRepository,
+    ):
+        self.metadata_catalog = metadata_catalog
+        self.repository = repository
+
+    def resolve(self, schema_ref: str) -> dict[str, Any]:
+        ref = _required_string(schema_ref, "schema_ref", 500)
+        metadata = self.metadata_catalog.load_schema(ref)
+        if not isinstance(metadata, dict):
+            raise ValueError("schema_ref was not found in the metadata catalog")
+        source_revision = _required_string(
+            metadata.get("source_revision"), "metadata source_revision", 200
+        )
+        fields = _normalized_fields(metadata.get("fields"))
+        snapshot = validate_schema_snapshot(
+            {
+                "schema_ref": ref,
+                "schema_hash": canonical_schema_hash(fields),
+                "fields": fields,
+                "source_revision": source_revision,
+            }
+        )
+        existing = self.repository.find_schema_snapshot(
+            schema_ref=ref, schema_hash=snapshot["schema_hash"]
+        )
+        if existing is not None:
+            return copy.deepcopy(existing)
+        return self.repository.persist_schema_snapshot(snapshot=snapshot)
+
+
+class DatasetBindingService:
+    """Attach credential-free physical datasets only when deployment is created."""
+
+    def __init__(self, repository, datasource_definitions, pool_manager):
+        self.repository = repository
+        self.datasource_definitions = datasource_definitions
+        self.pool_manager = pool_manager
+
+    @staticmethod
+    def _object_name(object_ref: str) -> tuple[str, str]:
+        parts = object_ref.split(".")
+        if len(parts) != 2 or not all(_IDENTIFIER.fullmatch(part) for part in parts):
+            raise ValueError("table or view object_ref must be schema.name")
+        return parts[0], parts[1]
+
+    def _verify_relation(self, binding: dict[str, Any]) -> None:
+        if binding["object_kind"] not in {"table", "view"}:
+            return
+        schema_name, object_name = self._object_name(binding["object_ref"])
+        expected_type = "BASE TABLE" if binding["object_kind"] == "table" else "VIEW"
+        statement = text(
+            "SELECT 1 FROM information_schema.tables "
+            "WHERE table_schema = :schema_name "
+            "AND table_name = :object_name "
+            "AND table_type = :table_type"
+        )
+        with self.pool_manager.connect(
+            binding["data_source_uid"], purpose="metadata_preview"
+        ) as connection:
+            result = connection.execute(
+                statement,
+                {
+                    "schema_name": schema_name,
+                    "object_name": object_name,
+                    "table_type": expected_type,
+                },
+            )
+            relation = result.first()
+        if relation is None:
+            raise ValueError("dataset table or view was not found")
+
+    def bind(
+        self,
+        deployment_id: str,
+        bindings: list[dict],
+        actor_uid: str,
+    ) -> list[dict]:
+        deployment = _uid(deployment_id, "deployment_id")
+        actor = _uid(actor_uid, "actor_uid")
+        if not isinstance(bindings, list) or not bindings:
+            raise ValueError("dataset bindings must be a non-empty array")
+        results: list[dict] = []
+        logical_refs: set[str] = set()
+        for raw_binding in bindings:
+            if not isinstance(raw_binding, dict):
+                raise ValueError("dataset binding must be an object")
+            raw = copy.deepcopy(raw_binding)
+            logical_ref = _required_string(
+                raw.pop("logical_ref", None), "logical_ref", 500
+            )
+            if logical_ref in logical_refs:
+                raise ValueError("dataset binding logical_ref values must be unique")
+            logical_refs.add(logical_ref)
+            binding = validate_dataset_binding(raw)
+            definition = self.datasource_definitions.get(binding["data_source_uid"])
+            if definition is None:
+                raise ValueError("data source was not found")
+            if getattr(definition, "status", True) is False:
+                raise ValueError("data source is inactive")
+            self._verify_relation(binding)
+            results.append(
+                self.repository.persist_dataset_binding(
+                    deployment_id=deployment,
+                    logical_ref=logical_ref,
+                    binding=binding,
+                    actor_uid=actor,
+                )
+            )
+        return results

+ 25 - 8
tests/core/data_rules/test_release.py

@@ -5,6 +5,7 @@ import re
 
 from app.core.common.identifiers import new_governance_uid
 from app.core.data_rules.contracts import rule_spec_hash
+from app.core.data_rules.execution_contracts import canonical_schema_hash
 from tests.core.data_rules.test_contracts import (
     valid_dataflow_spec,
     valid_rule_spec,
@@ -41,6 +42,24 @@ class ReleaseRepository:
         }
 
 
+class FakeSchemaResolver:
+    def __init__(self):
+        self.calls = []
+
+    def resolve(self, schema_ref):
+        self.calls.append(schema_ref)
+        fields = [
+            {"name": "customer_id", "type": "string", "nullable": False},
+        ]
+        return {
+            "id": new_governance_uid(),
+            "schema_ref": schema_ref,
+            "schema_hash": canonical_schema_hash(fields),
+            "fields": fields,
+            "source_revision": "catalog:1",
+        }
+
+
 def _published_rule(version_id, spec):
     return {
         "id": version_id,
@@ -95,20 +114,19 @@ def test_release_expands_standard_and_persists_fixed_bindings_and_plans():
         }
     }
     rules = {
-        standard_rule_id: _published_rule(
-            standard_rule_id, assertion_only_rule()
-        ),
+        standard_rule_id: _published_rule(standard_rule_id, assertion_only_rule()),
         direct_rule_id: _published_rule(direct_rule_id, valid_rule_spec()),
     }
     repository = ReleaseRepository(standards=standards, rules=rules)
     flow = valid_dataflow_spec(standard_id, direct_rule_id)
 
-    result = ProductionLineReleaseService(repository).release(
+    schema_resolver = FakeSchemaResolver()
+    result = ProductionLineReleaseService(
+        repository, schema_resolver=schema_resolver
+    ).release(
         dataflow_uid=flow["dataflow_uid"],
         dataflow_spec=flow,
         source_text="清洗客户数据后执行客户标准",
-        input_schema_hashes={"bd:customer_raw:v2": "a" * 64},
-        output_schema_hash="b" * 64,
         created_by=new_governance_uid(),
     )
 
@@ -137,6 +155,7 @@ def test_release_expands_standard_and_persists_fixed_bindings_and_plans():
     )
     assert direct_call["component_kind"] == "rule.apply"
     assert direct_call["idempotency"]["strategy"] == "partition_replace"
+    assert schema_resolver.calls == ["bd:customer_raw:v2", "bd:customer:v7"]
 
 
 def test_release_rejects_path_uid_mismatch_before_database_writes():
@@ -152,8 +171,6 @@ def test_release_rejects_path_uid_mismatch_before_database_writes():
             dataflow_uid=new_governance_uid(),
             dataflow_spec=flow,
             source_text="生产线",
-            input_schema_hashes={flow["input_schema_refs"][0]: "a" * 64},
-            output_schema_hash="b" * 64,
             created_by=new_governance_uid(),
         )
 

+ 232 - 0
tests/core/data_rules/test_schema_resolver.py

@@ -0,0 +1,232 @@
+from __future__ import annotations
+
+import copy
+
+import pytest
+
+from app.core.common.identifiers import new_governance_uid
+from app.core.data_rules.execution_contracts import canonical_schema_hash
+
+
+class FakeMetadataCatalog:
+    def __init__(self, schemas):
+        self.schemas = copy.deepcopy(schemas)
+
+    def load_schema(self, schema_ref):
+        return copy.deepcopy(self.schemas.get(schema_ref))
+
+
+class SnapshotRepository:
+    def __init__(self):
+        self.snapshots = {}
+
+    def find_schema_snapshot(self, *, schema_ref, schema_hash):
+        return copy.deepcopy(self.snapshots.get((schema_ref, schema_hash)))
+
+    def persist_schema_snapshot(self, *, snapshot):
+        saved = {"id": new_governance_uid(), **copy.deepcopy(snapshot)}
+        self.snapshots[(snapshot["schema_ref"], snapshot["schema_hash"])] = saved
+        return copy.deepcopy(saved)
+
+
+def test_schema_resolver_hashes_server_metadata_not_request_values():
+    from app.core.data_rules.schema_resolver import SchemaResolver
+
+    metadata = FakeMetadataCatalog(
+        {
+            "bd:customer:v7": {
+                "source_revision": "neo4j:42",
+                "fields": [
+                    {"name": "customer_id", "type": "string", "nullable": False},
+                    {"name": "mobile", "type": "string", "nullable": True},
+                ],
+            }
+        }
+    )
+
+    snapshot = SchemaResolver(metadata, SnapshotRepository()).resolve("bd:customer:v7")
+
+    assert snapshot["schema_hash"] == canonical_schema_hash(snapshot["fields"])
+    assert snapshot["source_revision"] == "neo4j:42"
+
+
+def test_schema_resolver_reuses_immutable_snapshot_for_the_same_schema():
+    from app.core.data_rules.schema_resolver import SchemaResolver
+
+    metadata = FakeMetadataCatalog(
+        {
+            "bd:customer:v7": {
+                "source_revision": "neo4j:42",
+                "fields": [
+                    {"name": "mobile", "type": "string", "nullable": True},
+                ],
+            }
+        }
+    )
+    repository = SnapshotRepository()
+    resolver = SchemaResolver(metadata, repository)
+
+    first = resolver.resolve("bd:customer:v7")
+    second = resolver.resolve("bd:customer:v7")
+
+    assert second == first
+    assert len(repository.snapshots) == 1
+
+
+@pytest.mark.parametrize(
+    "fields, message",
+    [
+        ([], "non-empty"),
+        (
+            [
+                {"name": "mobile", "type": "string", "nullable": True},
+                {"name": "mobile", "type": "string", "nullable": False},
+            ],
+            "unique",
+        ),
+        ([{"name": "payload", "type": "object", "nullable": False}], "unsupported"),
+    ],
+)
+def test_schema_resolver_rejects_untrusted_metadata_fields(fields, message):
+    from app.core.data_rules.schema_resolver import SchemaResolver
+
+    metadata = FakeMetadataCatalog(
+        {"bd:customer:v7": {"source_revision": "neo4j:42", "fields": fields}}
+    )
+
+    with pytest.raises(ValueError, match=message):
+        SchemaResolver(metadata, SnapshotRepository()).resolve("bd:customer:v7")
+
+
+class FakeDefinitions:
+    def __init__(self, definition):
+        self.definition = definition
+        self.requested = []
+
+    def get(self, uid):
+        self.requested.append(uid)
+        return self.definition
+
+
+class FakeConnection:
+    def __init__(self, relation_exists=True):
+        self.calls = []
+        self.relation_exists = relation_exists
+
+    def execute(self, statement, parameters):
+        self.calls.append((str(statement), parameters))
+        return FakeMetadataResult(self.relation_exists)
+
+
+class FakeMetadataResult:
+    def __init__(self, exists):
+        self.exists = exists
+
+    def first(self):
+        return (1,) if self.exists else None
+
+
+class FakePoolManager:
+    def __init__(self, relation_exists=True):
+        self.connection = FakeConnection(relation_exists)
+        self.calls = []
+
+    def connect(self, uid, purpose):
+        self.calls.append((uid, purpose))
+
+        class Context:
+            def __enter__(_self):
+                return self.connection
+
+            def __exit__(_self, *_args):
+                return False
+
+        return Context()
+
+
+class BindingRepository:
+    def __init__(self):
+        self.calls = []
+
+    def persist_dataset_binding(self, **kwargs):
+        self.calls.append(copy.deepcopy(kwargs))
+        return {"id": new_governance_uid(), **kwargs["binding"]}
+
+
+def test_dataset_binding_verifies_table_via_read_only_metadata_pool():
+    from app.core.data_rules.schema_resolver import DatasetBindingService
+
+    source_uid = new_governance_uid()
+    snapshot_uid = new_governance_uid()
+    definitions = FakeDefinitions(object())
+    pool_manager = FakePoolManager()
+    repository = BindingRepository()
+    binding = {
+        "logical_ref": "input:customers",
+        "data_source_uid": source_uid,
+        "object_kind": "table",
+        "object_ref": "crm.customers",
+        "schema_snapshot_id": snapshot_uid,
+        "access_mode": "read",
+        "dialect": "postgresql",
+        "write_mode": "append",
+    }
+
+    result = DatasetBindingService(repository, definitions, pool_manager).bind(
+        new_governance_uid(), [binding], new_governance_uid()
+    )
+
+    assert result[0]["data_source_uid"] == source_uid
+    assert definitions.requested == [source_uid]
+    assert pool_manager.calls == [(source_uid, "metadata_preview")]
+    assert "information_schema.tables" in pool_manager.connection.calls[0][0]
+    assert "credentials" not in str(repository.calls)
+
+
+def test_dataset_binding_rejects_credentials_before_accessing_data_source():
+    from app.core.data_rules.schema_resolver import DatasetBindingService
+
+    definitions = FakeDefinitions(object())
+    with pytest.raises(ValueError, match="credential|secret"):
+        DatasetBindingService(BindingRepository(), definitions, FakePoolManager()).bind(
+            new_governance_uid(),
+            [
+                {
+                    "logical_ref": "input:customers",
+                    "data_source_uid": new_governance_uid(),
+                    "object_kind": "table",
+                    "object_ref": "crm.customers",
+                    "schema_snapshot_id": new_governance_uid(),
+                    "access_mode": "read",
+                    "dialect": "postgresql",
+                    "write_mode": "append",
+                    "password": "not-allowed",
+                }
+            ],
+            new_governance_uid(),
+        )
+    assert definitions.requested == []
+
+
+def test_dataset_binding_rejects_missing_table_after_metadata_lookup():
+    from app.core.data_rules.schema_resolver import DatasetBindingService
+
+    with pytest.raises(ValueError, match="was not found"):
+        DatasetBindingService(
+            BindingRepository(), FakeDefinitions(object()), FakePoolManager(False)
+        ).bind(
+            new_governance_uid(),
+            [
+                {
+                    "logical_ref": "input:customers",
+                    "data_source_uid": new_governance_uid(),
+                    "object_kind": "table",
+                    "object_ref": "crm.customers",
+                    "schema_snapshot_id": new_governance_uid(),
+                    "access_mode": "read",
+                    "dialect": "postgresql",
+                    "write_mode": "append",
+                }
+            ],
+            new_governance_uid(),
+        )

+ 20 - 8
tests/integration/test_data_rule_control_plane.py

@@ -7,17 +7,30 @@ 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.schema_resolver import SchemaResolver
 from tests.core.data_rules.test_contracts import (
     valid_dataflow_spec,
     valid_rule_spec,
     valid_standard_spec,
 )
 
-
 pytestmark = pytest.mark.integration
 
 
+class FakeMetadataCatalog:
+    def load_schema(self, schema_ref):
+        return {
+            "source_revision": "integration:1",
+            "fields": [
+                {
+                    "name": "customer_id",
+                    "type": "string",
+                    "nullable": False,
+                }
+            ],
+        }
+
+
 @pytest.fixture()
 def database_url():
     value = os.environ.get("TEST_DATABASE_URL")
@@ -82,19 +95,18 @@ def test_postgres_rule_standard_and_production_line_release_is_atomic(
             )
 
             flow = valid_dataflow_spec(standard["id"], transform["id"])
-            released = ProductionLineReleaseService(repository).release(
+            released = ProductionLineReleaseService(
+                repository,
+                schema_resolver=SchemaResolver(FakeMetadataCatalog(), repository),
+            ).release(
                 dataflow_uid=flow["dataflow_uid"],
                 dataflow_spec=flow,
                 source_text="清洗客户数据并执行客户数据标准",
-                input_schema_hashes={"bd:customer_raw:v2": "a" * 64},
-                output_schema_hash="b" * 64,
                 created_by=actor,
             )
 
             assert released["status"] == "released"
-            assert released["package"]["standard_version_ids"] == [
-                standard["id"]
-            ]
+            assert released["package"]["standard_version_ids"] == [standard["id"]]
             assert released["package"]["rule_version_ids"] == sorted(
                 [quality["id"], transform["id"]]
             )

+ 30 - 14
tests/test_data_rule_api.py

@@ -4,9 +4,7 @@ from datetime import datetime, timedelta, timezone
 
 from app.core.common.identifiers import new_governance_uid
 from app.core.data_rules.contracts import rule_spec_hash
-from app.core.system.tokens import issue_access_token
-from app.core.system.tokens import decode_access_token
-
+from app.core.system.tokens import decode_access_token, issue_access_token
 from tests.core.data_rules.test_contracts import (
     valid_dataflow_spec,
     valid_rule_spec,
@@ -212,21 +210,15 @@ def test_production_line_resolve_preview_expands_standard_without_writing(monkey
     standard_id = new_governance_uid()
     standard_rule_id = new_governance_uid()
     direct_rule_id = new_governance_uid()
-    standard_rule = published_rule(
-        standard_rule_id, assertion_only_rule()
-    )
+    standard_rule = published_rule(standard_rule_id, assertion_only_rule())
     direct_rule = published_rule(direct_rule_id)
 
     response = client.post(
         "/api/rules/production-lines/resolve",
         json={
-            "dataflow_spec": valid_dataflow_spec(
-                standard_id, direct_rule_id
-            ),
+            "dataflow_spec": valid_dataflow_spec(standard_id, direct_rule_id),
             "standard_versions": {
-                standard_id: published_standard(
-                    standard_id, standard_rule_id
-                )
+                standard_id: published_standard(standard_id, standard_rule_id)
             },
             "rule_versions": {
                 standard_rule_id: standard_rule,
@@ -368,8 +360,6 @@ def test_dataflow_release_uses_server_assets_and_release_permission(monkeypatch)
     payload = {
         "source_text": "客户数据生产线",
         "dataflow_spec": flow,
-        "input_schema_hashes": {"bd:customer_raw:v2": "a" * 64},
-        "output_schema_hash": "b" * 64,
     }
 
     forbidden = client.post(
@@ -391,3 +381,29 @@ def test_dataflow_release_uses_server_assets_and_release_permission(monkeypatch)
     assert "standard_versions" not in service.calls[0]
     assert "rule_versions" not in service.calls[0]
     assert "component_binding_ids" not in service.calls[0]
+
+
+def test_dataflow_release_rejects_client_authored_schema_hashes(monkeypatch):
+    from app import create_app
+
+    app = create_app()
+    _use_token_identity(monkeypatch)
+    app.config["TESTING"] = True
+    service = FakeReleaseService()
+    app.extensions["production_line_release_service"] = service
+    client = app.test_client()
+    flow = valid_dataflow_spec()
+
+    response = client.post(
+        f"/api/rules/production-lines/{flow['dataflow_uid']}/release",
+        json={
+            "source_text": "客户数据生产线",
+            "dataflow_spec": flow,
+            "input_schema_hashes": {"bd:customer_raw:v2": "a" * 64},
+            "output_schema_hash": "b" * 64,
+        },
+        headers=_headers(app, "admin"),
+    )
+
+    assert response.status_code == 409
+    assert service.calls == []