release.py 8.5 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216
  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. rule_backend_support: dict[str, frozenset[str]] = {}
  52. def validate_published_rule(rule_version_id: str) -> None:
  53. if rule_version_id in rule_backend_support:
  54. return
  55. rule = rules.get(rule_version_id)
  56. if not isinstance(rule, dict):
  57. raise ValueError(
  58. f"published rule version {rule_version_id} was not found"
  59. )
  60. rule_spec = read_rule_spec(rule.get("rule_spec"))
  61. rule_snapshot = self.schema_resolver.resolve(
  62. rule_spec["input_schema_ref"]
  63. )
  64. snapshot_fields = {
  65. field["name"]: field["type"]
  66. for field in rule_snapshot["fields"]
  67. }
  68. supported = validate_rule_expressions(
  69. rule_spec, snapshot_fields
  70. )
  71. if not supported:
  72. raise ValueError("rule has no supported execution backend")
  73. rule_backend_support[rule_version_id] = supported
  74. # Validate every loaded rule before allocating any release version.
  75. # The repository is expected to scope this mapping to published assets
  76. # available to the release, so no loaded expression is left unchecked.
  77. for loaded_rule_version_id in rules:
  78. validate_published_rule(loaded_rule_version_id)
  79. # Ensure referenced standards and their clauses resolve to loaded,
  80. # already-validated rule versions before beginning the release.
  81. for component in flow["components"]:
  82. if component["type"] == "standard.enforce":
  83. standard = standards.get(component["standard_version_id"])
  84. if not isinstance(standard, dict):
  85. raise ValueError(
  86. "published standard version "
  87. f"{component['standard_version_id']} was not found"
  88. )
  89. for clause in standard["clauses"]:
  90. validate_published_rule(str(clause["rule_version_id"]))
  91. else:
  92. validate_published_rule(component["rule_version_id"])
  93. version = self.repository.begin_dataflow_release(
  94. dataflow_spec=flow,
  95. source_text=source,
  96. input_schema_hashes=inputs,
  97. output_schema_hash=output,
  98. created_by=actor,
  99. )
  100. version_id = _uid(version.get("id"), "dataflow_version_id")
  101. compiled: dict[str, dict[str, Any]] = {}
  102. binding_ids: dict[str, str] = {}
  103. def add_binding(
  104. *,
  105. binding_key: str,
  106. component_id: str,
  107. component_kind: str,
  108. rule_version_id: str,
  109. stage: str,
  110. order_no: int,
  111. idempotency: dict[str, Any] | None,
  112. provenance: dict[str, str],
  113. ) -> None:
  114. rule = rules.get(rule_version_id)
  115. if not isinstance(rule, dict):
  116. raise ValueError(
  117. f"published rule version {rule_version_id} was not found"
  118. )
  119. if rule_version_id not in compiled:
  120. compiled[rule_version_id] = compile_rule_plan(
  121. rule,
  122. supported_backends=rule_backend_support[rule_version_id],
  123. )
  124. plan = compiled[rule_version_id]
  125. binding_id = new_governance_uid()
  126. binding_ids[binding_key] = binding_id
  127. self.repository.persist_component_plan(
  128. dataflow_version_id=version_id,
  129. component_binding_id=binding_id,
  130. component_id=component_id,
  131. component_kind=component_kind,
  132. rule_version_id=rule_version_id,
  133. stage=stage,
  134. order_no=order_no,
  135. idempotency=idempotency,
  136. provenance=provenance,
  137. plan=plan,
  138. schema_hashes={
  139. "inputs": inputs,
  140. "output": output,
  141. },
  142. )
  143. for component in flow["components"]:
  144. if component["type"] == "standard.enforce":
  145. standard_id = component["standard_version_id"]
  146. standard = standards.get(standard_id)
  147. if not isinstance(standard, dict):
  148. raise ValueError(
  149. f"published standard version {standard_id} was not found"
  150. )
  151. for clause_index, clause in enumerate(standard["clauses"]):
  152. clause_id = str(clause["clause_id"])
  153. binding_key = f"{component['id']}:{clause_id}"
  154. add_binding(
  155. binding_key=binding_key,
  156. component_id=(f"{component['id']}__{clause_id}"[:100]),
  157. component_kind="quality.check",
  158. rule_version_id=str(clause["rule_version_id"]),
  159. stage=component["stage"],
  160. order_no=(component["order"] * 1000) + clause_index,
  161. idempotency=None,
  162. provenance={
  163. "standard_version_id": standard_id,
  164. "clause_id": clause_id,
  165. },
  166. )
  167. continue
  168. add_binding(
  169. binding_key=component["id"],
  170. component_id=component["id"],
  171. component_kind=component["type"],
  172. rule_version_id=component["rule_version_id"],
  173. stage=component["stage"],
  174. order_no=component["order"] * 1000,
  175. idempotency=component.get("idempotency"),
  176. provenance={},
  177. )
  178. release_rules = {
  179. rule_id: {
  180. **rule,
  181. "execution_plan": {
  182. "backend": compiled[rule_id]["backend"],
  183. "plan_hash": compiled[rule_id]["plan_hash"],
  184. },
  185. }
  186. for rule_id, rule in rules.items()
  187. if rule_id in compiled
  188. }
  189. package = resolve_production_line(
  190. flow,
  191. standards,
  192. release_rules,
  193. component_binding_ids=binding_ids,
  194. )
  195. return self.repository.complete_dataflow_release(
  196. dataflow_version_id=version_id,
  197. package=package,
  198. )