release.py 13 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323
  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__(
  23. self,
  24. repository,
  25. *,
  26. schema_resolver=None,
  27. available_plan_backends=None,
  28. ):
  29. self.repository = repository
  30. self.schema_resolver = schema_resolver
  31. self.available_plan_backends = frozenset(
  32. available_plan_backends
  33. if available_plan_backends is not None
  34. else {"sql_pushdown", "polars_batch", "quality_check"}
  35. )
  36. def release(
  37. self,
  38. *,
  39. dataflow_uid: str,
  40. dataflow_spec: dict[str, Any],
  41. source_text: str,
  42. created_by: str,
  43. ) -> dict[str, Any]:
  44. path_uid = _uid(dataflow_uid, "dataflow_uid")
  45. actor = _uid(created_by, "created_by")
  46. flow = validate_dataflow_spec(dataflow_spec)
  47. if flow["dataflow_uid"] != path_uid:
  48. raise ValueError("dataflow spec uid does not match path dataflow_uid")
  49. source = _source(source_text)
  50. if self.schema_resolver is None:
  51. raise ValueError("trusted schema resolver is not configured")
  52. input_snapshots = {
  53. ref: self.schema_resolver.resolve(ref) for ref in flow["input_schema_refs"]
  54. }
  55. output_snapshot = self.schema_resolver.resolve(flow["output_schema_ref"])
  56. inputs = {
  57. ref: snapshot["schema_hash"]
  58. for ref, snapshot in sorted(input_snapshots.items())
  59. }
  60. output = output_snapshot["schema_hash"]
  61. standards, rules = self.repository.load_published_assets(flow)
  62. rule_backend_support: dict[str, frozenset[str]] = {}
  63. def validate_published_rule(rule_version_id: str) -> None:
  64. if rule_version_id in rule_backend_support:
  65. return
  66. rule = rules.get(rule_version_id)
  67. if not isinstance(rule, dict):
  68. raise ValueError(
  69. f"published rule version {rule_version_id} was not found"
  70. )
  71. rule_spec = read_rule_spec(rule.get("rule_spec"))
  72. rule_snapshot = self.schema_resolver.resolve(
  73. rule_spec["input_schema_ref"]
  74. )
  75. snapshot_fields = {
  76. field["name"]: field["type"]
  77. for field in rule_snapshot["fields"]
  78. }
  79. supported = validate_rule_expressions(
  80. rule_spec, snapshot_fields
  81. )
  82. if not supported:
  83. raise ValueError("rule has no supported execution backend")
  84. rule_backend_support[rule_version_id] = supported
  85. # Validate every loaded rule before allocating any release version.
  86. # The repository is expected to scope this mapping to published assets
  87. # available to the release, so no loaded expression is left unchecked.
  88. for loaded_rule_version_id in rules:
  89. validate_published_rule(loaded_rule_version_id)
  90. # Ensure referenced standards and their clauses resolve to loaded,
  91. # already-validated rule versions before beginning the release.
  92. for component in flow["components"]:
  93. if component["type"] == "standard.enforce":
  94. standard = standards.get(component["standard_version_id"])
  95. if not isinstance(standard, dict):
  96. raise ValueError(
  97. "published standard version "
  98. f"{component['standard_version_id']} was not found"
  99. )
  100. for clause in standard["clauses"]:
  101. validate_published_rule(str(clause["rule_version_id"]))
  102. else:
  103. validate_published_rule(component["rule_version_id"])
  104. compiled: dict[str, dict[str, Any]] = {}
  105. for rule_version_id, supported_backends in sorted(
  106. rule_backend_support.items()
  107. ):
  108. plan = compile_rule_plan(
  109. rules[rule_version_id],
  110. supported_backends=supported_backends,
  111. )
  112. if plan["backend"] not in self.available_plan_backends:
  113. raise ValueError(
  114. f"plan backend {plan['backend']} is not registered"
  115. )
  116. compiled[rule_version_id] = plan
  117. version = self.repository.begin_dataflow_release(
  118. dataflow_spec=flow,
  119. source_text=source,
  120. input_schema_hashes=inputs,
  121. output_schema_hash=output,
  122. created_by=actor,
  123. )
  124. version_id = _uid(version.get("id"), "dataflow_version_id")
  125. binding_ids: dict[str, str] = {}
  126. def add_binding(
  127. *,
  128. binding_key: str,
  129. component_id: str,
  130. component_kind: str,
  131. rule_version_id: str,
  132. stage: str,
  133. order_no: int,
  134. idempotency: dict[str, Any] | None,
  135. provenance: dict[str, str],
  136. ) -> None:
  137. rule = rules.get(rule_version_id)
  138. if not isinstance(rule, dict):
  139. raise ValueError(
  140. f"published rule version {rule_version_id} was not found"
  141. )
  142. plan = compiled[rule_version_id]
  143. binding_id = new_governance_uid()
  144. binding_ids[binding_key] = binding_id
  145. self.repository.persist_component_plan(
  146. dataflow_version_id=version_id,
  147. component_binding_id=binding_id,
  148. component_id=component_id,
  149. component_kind=component_kind,
  150. rule_version_id=rule_version_id,
  151. stage=stage,
  152. order_no=order_no,
  153. idempotency=idempotency,
  154. provenance=provenance,
  155. plan=plan,
  156. schema_hashes={
  157. "inputs": inputs,
  158. "output": output,
  159. },
  160. )
  161. for component in flow["components"]:
  162. if component["type"] == "standard.enforce":
  163. standard_id = component["standard_version_id"]
  164. standard = standards.get(standard_id)
  165. if not isinstance(standard, dict):
  166. raise ValueError(
  167. f"published standard version {standard_id} was not found"
  168. )
  169. for clause_index, clause in enumerate(standard["clauses"]):
  170. clause_id = str(clause["clause_id"])
  171. binding_key = f"{component['id']}:{clause_id}"
  172. add_binding(
  173. binding_key=binding_key,
  174. component_id=(f"{component['id']}__{clause_id}"[:100]),
  175. component_kind="quality.check",
  176. rule_version_id=str(clause["rule_version_id"]),
  177. stage=component["stage"],
  178. order_no=(component["order"] * 1000) + clause_index,
  179. idempotency=None,
  180. provenance={
  181. "standard_version_id": standard_id,
  182. "clause_id": clause_id,
  183. },
  184. )
  185. continue
  186. add_binding(
  187. binding_key=component["id"],
  188. component_id=component["id"],
  189. component_kind=component["type"],
  190. rule_version_id=component["rule_version_id"],
  191. stage=component["stage"],
  192. order_no=component["order"] * 1000,
  193. idempotency=component.get("idempotency"),
  194. provenance={},
  195. )
  196. release_rules = {
  197. rule_id: {
  198. **rule,
  199. "execution_plan": {
  200. "backend": compiled[rule_id]["backend"],
  201. "plan_hash": compiled[rule_id]["plan_hash"],
  202. },
  203. }
  204. for rule_id, rule in rules.items()
  205. if rule_id in compiled
  206. }
  207. package = resolve_production_line(
  208. flow,
  209. standards,
  210. release_rules,
  211. component_binding_ids=binding_ids,
  212. )
  213. return self.repository.complete_dataflow_release(
  214. dataflow_version_id=version_id,
  215. package=package,
  216. )
  217. class BoundSqlPlanService:
  218. """Compile and persist a physical plan only after deployment bindings exist.
  219. Compilation is deliberately persisted as ``compiled``. Test evidence and
  220. publication transitions remain separate lifecycle actions (Task 7), so
  221. this vertical slice never labels an unexecuted plan tested or published.
  222. """
  223. def __init__(self, repository, compiler_registry):
  224. self.repository = repository
  225. self.compiler_registry = compiler_registry
  226. def compile_and_persist(
  227. self,
  228. *,
  229. component_binding_id: str,
  230. rule_version_id: str,
  231. input_schema_snapshot_id: str,
  232. output_schema_snapshot_id: str,
  233. input_binding_id: str,
  234. output_binding_id: str,
  235. ) -> dict[str, Any]:
  236. ids = {
  237. "component_binding_id": _uid(
  238. component_binding_id, "component_binding_id"
  239. ),
  240. "rule_version_id": _uid(rule_version_id, "rule_version_id"),
  241. "input_schema_snapshot_id": _uid(
  242. input_schema_snapshot_id, "input_schema_snapshot_id"
  243. ),
  244. "output_schema_snapshot_id": _uid(
  245. output_schema_snapshot_id, "output_schema_snapshot_id"
  246. ),
  247. "input_binding_id": _uid(input_binding_id, "input_binding_id"),
  248. "output_binding_id": _uid(output_binding_id, "output_binding_id"),
  249. }
  250. context = self.repository.load_bound_compile_context(**ids)
  251. if not isinstance(context, dict):
  252. raise ValueError("canonical bound compile context was not found")
  253. try:
  254. component = context["component_binding"]
  255. rule_version = context["rule_version"]
  256. input_schema = context["input_schema"]
  257. output_schema = context["output_schema"]
  258. input_binding = context["input_binding"]
  259. output_binding = context["output_binding"]
  260. backend = context["backend"]
  261. except KeyError as exc:
  262. raise ValueError(
  263. "canonical bound compile context is incomplete"
  264. ) from exc
  265. canonical_ids = {
  266. "component_binding_id": component.get("id"),
  267. "rule_version_id": rule_version.get("id"),
  268. "input_schema_snapshot_id": input_schema.get("id"),
  269. "output_schema_snapshot_id": output_schema.get("id"),
  270. "input_binding_id": input_binding.get("id"),
  271. "output_binding_id": output_binding.get("id"),
  272. }
  273. if canonical_ids != ids or component.get(
  274. "rule_version_id"
  275. ) != ids["rule_version_id"]:
  276. raise ValueError(
  277. "canonical bound compile context identifiers do not match"
  278. )
  279. compiler = self.compiler_registry.select(
  280. rule_version.get("rule_spec"),
  281. input_binding,
  282. output_binding,
  283. )
  284. compiled = compiler.compile(
  285. rule_version=rule_version,
  286. input_schema=input_schema,
  287. output_schema=output_schema,
  288. input_binding=input_binding,
  289. output_binding=output_binding,
  290. backend=backend,
  291. )
  292. return self.repository.persist_bound_component_plan(
  293. component_binding_id=ids["component_binding_id"],
  294. rule_version_id=ids["rule_version_id"],
  295. input_binding_id=ids["input_binding_id"],
  296. output_binding_id=ids["output_binding_id"],
  297. compiled=compiled,
  298. status="compiled",
  299. )