release.py 6.6 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174
  1. """Release a governed DataFlow as an immutable data production line."""
  2. from __future__ import annotations
  3. from typing import Any
  4. from app.core.common.identifiers import ensure_governance_uid, new_governance_uid
  5. from app.core.data_rules.compiler import compile_rule_plan
  6. from app.core.data_rules.contracts import read_rule_spec, validate_dataflow_spec
  7. from app.core.data_rules.expressions import validate_rule_expressions
  8. from app.core.data_rules.production_line import resolve_production_line
  9. def _uid(value: Any, label: str) -> str:
  10. try:
  11. return ensure_governance_uid({"uid": str(value)})
  12. except ValueError as exc:
  13. raise ValueError(f"{label} must be a valid UUIDv7") from exc
  14. def _source(value: Any) -> str:
  15. if not isinstance(value, str) or not value.strip():
  16. raise ValueError("source_text is required")
  17. normalized = value.strip()
  18. if len(normalized) > 20_000:
  19. raise ValueError("source_text exceeds 20000 characters")
  20. return normalized
  21. class ProductionLineReleaseService:
  22. def __init__(self, repository, *, schema_resolver=None):
  23. self.repository = repository
  24. self.schema_resolver = schema_resolver
  25. def release(
  26. self,
  27. *,
  28. dataflow_uid: str,
  29. dataflow_spec: dict[str, Any],
  30. source_text: str,
  31. created_by: str,
  32. ) -> dict[str, Any]:
  33. path_uid = _uid(dataflow_uid, "dataflow_uid")
  34. actor = _uid(created_by, "created_by")
  35. flow = validate_dataflow_spec(dataflow_spec)
  36. if flow["dataflow_uid"] != path_uid:
  37. raise ValueError("dataflow spec uid does not match path dataflow_uid")
  38. source = _source(source_text)
  39. if self.schema_resolver is None:
  40. raise ValueError("trusted schema resolver is not configured")
  41. input_snapshots = {
  42. ref: self.schema_resolver.resolve(ref) for ref in flow["input_schema_refs"]
  43. }
  44. output_snapshot = self.schema_resolver.resolve(flow["output_schema_ref"])
  45. inputs = {
  46. ref: snapshot["schema_hash"]
  47. for ref, snapshot in sorted(input_snapshots.items())
  48. }
  49. output = output_snapshot["schema_hash"]
  50. standards, rules = self.repository.load_published_assets(flow)
  51. version = self.repository.begin_dataflow_release(
  52. dataflow_spec=flow,
  53. source_text=source,
  54. input_schema_hashes=inputs,
  55. output_schema_hash=output,
  56. created_by=actor,
  57. )
  58. version_id = _uid(version.get("id"), "dataflow_version_id")
  59. compiled: dict[str, dict[str, Any]] = {}
  60. binding_ids: dict[str, str] = {}
  61. def add_binding(
  62. *,
  63. binding_key: str,
  64. component_id: str,
  65. component_kind: str,
  66. rule_version_id: str,
  67. stage: str,
  68. order_no: int,
  69. idempotency: dict[str, Any] | None,
  70. provenance: dict[str, str],
  71. ) -> None:
  72. rule = rules.get(rule_version_id)
  73. if not isinstance(rule, dict):
  74. raise ValueError(
  75. f"published rule version {rule_version_id} was not found"
  76. )
  77. rule_spec = read_rule_spec(rule.get("rule_spec"))
  78. rule_snapshot = self.schema_resolver.resolve(
  79. rule_spec["input_schema_ref"]
  80. )
  81. snapshot_fields = {
  82. field["name"]: field["type"]
  83. for field in rule_snapshot["fields"]
  84. }
  85. validate_rule_expressions(rule_spec, snapshot_fields)
  86. plan = compiled.setdefault(rule_version_id, compile_rule_plan(rule))
  87. binding_id = new_governance_uid()
  88. binding_ids[binding_key] = binding_id
  89. self.repository.persist_component_plan(
  90. dataflow_version_id=version_id,
  91. component_binding_id=binding_id,
  92. component_id=component_id,
  93. component_kind=component_kind,
  94. rule_version_id=rule_version_id,
  95. stage=stage,
  96. order_no=order_no,
  97. idempotency=idempotency,
  98. provenance=provenance,
  99. plan=plan,
  100. schema_hashes={
  101. "inputs": inputs,
  102. "output": output,
  103. },
  104. )
  105. for component in flow["components"]:
  106. if component["type"] == "standard.enforce":
  107. standard_id = component["standard_version_id"]
  108. standard = standards.get(standard_id)
  109. if not isinstance(standard, dict):
  110. raise ValueError(
  111. f"published standard version {standard_id} was not found"
  112. )
  113. for clause_index, clause in enumerate(standard["clauses"]):
  114. clause_id = str(clause["clause_id"])
  115. binding_key = f"{component['id']}:{clause_id}"
  116. add_binding(
  117. binding_key=binding_key,
  118. component_id=(f"{component['id']}__{clause_id}"[:100]),
  119. component_kind="quality.check",
  120. rule_version_id=str(clause["rule_version_id"]),
  121. stage=component["stage"],
  122. order_no=(component["order"] * 1000) + clause_index,
  123. idempotency=None,
  124. provenance={
  125. "standard_version_id": standard_id,
  126. "clause_id": clause_id,
  127. },
  128. )
  129. continue
  130. add_binding(
  131. binding_key=component["id"],
  132. component_id=component["id"],
  133. component_kind=component["type"],
  134. rule_version_id=component["rule_version_id"],
  135. stage=component["stage"],
  136. order_no=component["order"] * 1000,
  137. idempotency=component.get("idempotency"),
  138. provenance={},
  139. )
  140. release_rules = {
  141. rule_id: {
  142. **rule,
  143. "execution_plan": {
  144. "backend": compiled[rule_id]["backend"],
  145. "plan_hash": compiled[rule_id]["plan_hash"],
  146. },
  147. }
  148. for rule_id, rule in rules.items()
  149. if rule_id in compiled
  150. }
  151. package = resolve_production_line(
  152. flow,
  153. standards,
  154. release_rules,
  155. component_binding_ids=binding_ids,
  156. )
  157. return self.repository.complete_dataflow_release(
  158. dataflow_version_id=version_id,
  159. package=package,
  160. )