| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323 |
- """Release a governed DataFlow as an immutable data production line."""
- from __future__ import annotations
- from typing import Any
- from app.core.common.identifiers import ensure_governance_uid, new_governance_uid
- from app.core.data_rules.compiler import compile_rule_plan
- from app.core.data_rules.contracts import read_rule_spec, validate_dataflow_spec
- from app.core.data_rules.expressions import validate_rule_expressions
- from app.core.data_rules.production_line import resolve_production_line
- 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 _source(value: Any) -> str:
- if not isinstance(value, str) or not value.strip():
- raise ValueError("source_text is required")
- normalized = value.strip()
- if len(normalized) > 20_000:
- raise ValueError("source_text exceeds 20000 characters")
- return normalized
- class ProductionLineReleaseService:
- def __init__(
- self,
- repository,
- *,
- schema_resolver=None,
- available_plan_backends=None,
- ):
- self.repository = repository
- self.schema_resolver = schema_resolver
- self.available_plan_backends = frozenset(
- available_plan_backends
- if available_plan_backends is not None
- else {"sql_pushdown", "polars_batch", "quality_check"}
- )
- def release(
- self,
- *,
- dataflow_uid: str,
- dataflow_spec: dict[str, Any],
- source_text: str,
- created_by: str,
- ) -> dict[str, Any]:
- path_uid = _uid(dataflow_uid, "dataflow_uid")
- actor = _uid(created_by, "created_by")
- flow = validate_dataflow_spec(dataflow_spec)
- if flow["dataflow_uid"] != path_uid:
- raise ValueError("dataflow spec uid does not match path dataflow_uid")
- source = _source(source_text)
- 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)
- rule_backend_support: dict[str, frozenset[str]] = {}
- def validate_published_rule(rule_version_id: str) -> None:
- if rule_version_id in rule_backend_support:
- return
- rule = rules.get(rule_version_id)
- if not isinstance(rule, dict):
- raise ValueError(
- f"published rule version {rule_version_id} was not found"
- )
- rule_spec = read_rule_spec(rule.get("rule_spec"))
- rule_snapshot = self.schema_resolver.resolve(
- rule_spec["input_schema_ref"]
- )
- snapshot_fields = {
- field["name"]: field["type"]
- for field in rule_snapshot["fields"]
- }
- supported = validate_rule_expressions(
- rule_spec, snapshot_fields
- )
- if not supported:
- raise ValueError("rule has no supported execution backend")
- rule_backend_support[rule_version_id] = supported
- # Validate every loaded rule before allocating any release version.
- # The repository is expected to scope this mapping to published assets
- # available to the release, so no loaded expression is left unchecked.
- for loaded_rule_version_id in rules:
- validate_published_rule(loaded_rule_version_id)
- # Ensure referenced standards and their clauses resolve to loaded,
- # already-validated rule versions before beginning the release.
- for component in flow["components"]:
- if component["type"] == "standard.enforce":
- standard = standards.get(component["standard_version_id"])
- if not isinstance(standard, dict):
- raise ValueError(
- "published standard version "
- f"{component['standard_version_id']} was not found"
- )
- for clause in standard["clauses"]:
- validate_published_rule(str(clause["rule_version_id"]))
- else:
- validate_published_rule(component["rule_version_id"])
- compiled: dict[str, dict[str, Any]] = {}
- for rule_version_id, supported_backends in sorted(
- rule_backend_support.items()
- ):
- plan = compile_rule_plan(
- rules[rule_version_id],
- supported_backends=supported_backends,
- )
- if plan["backend"] not in self.available_plan_backends:
- raise ValueError(
- f"plan backend {plan['backend']} is not registered"
- )
- compiled[rule_version_id] = plan
- version = self.repository.begin_dataflow_release(
- dataflow_spec=flow,
- source_text=source,
- input_schema_hashes=inputs,
- output_schema_hash=output,
- created_by=actor,
- )
- version_id = _uid(version.get("id"), "dataflow_version_id")
- binding_ids: dict[str, str] = {}
- def add_binding(
- *,
- binding_key: str,
- component_id: str,
- component_kind: str,
- rule_version_id: str,
- stage: str,
- order_no: int,
- idempotency: dict[str, Any] | None,
- provenance: dict[str, str],
- ) -> None:
- rule = rules.get(rule_version_id)
- if not isinstance(rule, dict):
- raise ValueError(
- f"published rule version {rule_version_id} was not found"
- )
- plan = compiled[rule_version_id]
- binding_id = new_governance_uid()
- binding_ids[binding_key] = binding_id
- self.repository.persist_component_plan(
- dataflow_version_id=version_id,
- component_binding_id=binding_id,
- component_id=component_id,
- component_kind=component_kind,
- rule_version_id=rule_version_id,
- stage=stage,
- order_no=order_no,
- idempotency=idempotency,
- provenance=provenance,
- plan=plan,
- schema_hashes={
- "inputs": inputs,
- "output": output,
- },
- )
- for component in flow["components"]:
- if component["type"] == "standard.enforce":
- standard_id = component["standard_version_id"]
- standard = standards.get(standard_id)
- if not isinstance(standard, dict):
- raise ValueError(
- f"published standard version {standard_id} was not found"
- )
- for clause_index, clause in enumerate(standard["clauses"]):
- clause_id = str(clause["clause_id"])
- binding_key = f"{component['id']}:{clause_id}"
- add_binding(
- binding_key=binding_key,
- component_id=(f"{component['id']}__{clause_id}"[:100]),
- component_kind="quality.check",
- rule_version_id=str(clause["rule_version_id"]),
- stage=component["stage"],
- order_no=(component["order"] * 1000) + clause_index,
- idempotency=None,
- provenance={
- "standard_version_id": standard_id,
- "clause_id": clause_id,
- },
- )
- continue
- add_binding(
- binding_key=component["id"],
- component_id=component["id"],
- component_kind=component["type"],
- rule_version_id=component["rule_version_id"],
- stage=component["stage"],
- order_no=component["order"] * 1000,
- idempotency=component.get("idempotency"),
- provenance={},
- )
- release_rules = {
- rule_id: {
- **rule,
- "execution_plan": {
- "backend": compiled[rule_id]["backend"],
- "plan_hash": compiled[rule_id]["plan_hash"],
- },
- }
- for rule_id, rule in rules.items()
- if rule_id in compiled
- }
- package = resolve_production_line(
- flow,
- standards,
- release_rules,
- component_binding_ids=binding_ids,
- )
- return self.repository.complete_dataflow_release(
- dataflow_version_id=version_id,
- package=package,
- )
- class BoundSqlPlanService:
- """Compile and persist a physical plan only after deployment bindings exist.
- Compilation is deliberately persisted as ``compiled``. Test evidence and
- publication transitions remain separate lifecycle actions (Task 7), so
- this vertical slice never labels an unexecuted plan tested or published.
- """
- def __init__(self, repository, compiler_registry):
- self.repository = repository
- self.compiler_registry = compiler_registry
- def compile_and_persist(
- self,
- *,
- component_binding_id: str,
- rule_version_id: str,
- input_schema_snapshot_id: str,
- output_schema_snapshot_id: str,
- input_binding_id: str,
- output_binding_id: str,
- ) -> dict[str, Any]:
- ids = {
- "component_binding_id": _uid(
- component_binding_id, "component_binding_id"
- ),
- "rule_version_id": _uid(rule_version_id, "rule_version_id"),
- "input_schema_snapshot_id": _uid(
- input_schema_snapshot_id, "input_schema_snapshot_id"
- ),
- "output_schema_snapshot_id": _uid(
- output_schema_snapshot_id, "output_schema_snapshot_id"
- ),
- "input_binding_id": _uid(input_binding_id, "input_binding_id"),
- "output_binding_id": _uid(output_binding_id, "output_binding_id"),
- }
- context = self.repository.load_bound_compile_context(**ids)
- if not isinstance(context, dict):
- raise ValueError("canonical bound compile context was not found")
- try:
- component = context["component_binding"]
- rule_version = context["rule_version"]
- input_schema = context["input_schema"]
- output_schema = context["output_schema"]
- input_binding = context["input_binding"]
- output_binding = context["output_binding"]
- backend = context["backend"]
- except KeyError as exc:
- raise ValueError(
- "canonical bound compile context is incomplete"
- ) from exc
- canonical_ids = {
- "component_binding_id": component.get("id"),
- "rule_version_id": rule_version.get("id"),
- "input_schema_snapshot_id": input_schema.get("id"),
- "output_schema_snapshot_id": output_schema.get("id"),
- "input_binding_id": input_binding.get("id"),
- "output_binding_id": output_binding.get("id"),
- }
- if canonical_ids != ids or component.get(
- "rule_version_id"
- ) != ids["rule_version_id"]:
- raise ValueError(
- "canonical bound compile context identifiers do not match"
- )
- compiler = self.compiler_registry.select(
- rule_version.get("rule_spec"),
- input_binding,
- output_binding,
- )
- compiled = compiler.compile(
- rule_version=rule_version,
- input_schema=input_schema,
- output_schema=output_schema,
- input_binding=input_binding,
- output_binding=output_binding,
- backend=backend,
- )
- return self.repository.persist_bound_component_plan(
- component_binding_id=ids["component_binding_id"],
- rule_version_id=ids["rule_version_id"],
- input_binding_id=ids["input_binding_id"],
- output_binding_id=ids["output_binding_id"],
- compiled=compiled,
- status="compiled",
- )
|