release.py 6.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185
  1. """Release a governed DataFlow as an immutable data production line."""
  2. from __future__ import annotations
  3. import re
  4. from typing import Any
  5. from app.core.common.identifiers import ensure_governance_uid, new_governance_uid
  6. from app.core.data_rules.compiler import compile_rule_plan
  7. from app.core.data_rules.contracts import validate_dataflow_spec
  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. def _digest(value: Any, label: str) -> str:
  22. normalized = str(value or "")
  23. if not re.fullmatch(r"[0-9a-f]{64}", normalized):
  24. raise ValueError(f"{label} must be a sha256 hex digest")
  25. return normalized
  26. def _schema_hashes(
  27. flow: dict[str, Any],
  28. input_schema_hashes: Any,
  29. output_schema_hash: Any,
  30. ) -> tuple[dict[str, str], str]:
  31. if not isinstance(input_schema_hashes, dict):
  32. raise ValueError("input_schema_hashes must be an object")
  33. if set(input_schema_hashes) != set(flow["input_schema_refs"]):
  34. raise ValueError("input_schema_hashes must cover every input schema")
  35. inputs = {
  36. ref: _digest(input_schema_hashes[ref], f"schema hash for {ref}")
  37. for ref in sorted(input_schema_hashes)
  38. }
  39. return inputs, _digest(output_schema_hash, "output_schema_hash")
  40. class ProductionLineReleaseService:
  41. def __init__(self, repository):
  42. self.repository = repository
  43. def release(
  44. self,
  45. *,
  46. dataflow_uid: str,
  47. dataflow_spec: dict[str, Any],
  48. source_text: str,
  49. input_schema_hashes: dict[str, str],
  50. output_schema_hash: str,
  51. created_by: str,
  52. ) -> dict[str, Any]:
  53. path_uid = _uid(dataflow_uid, "dataflow_uid")
  54. actor = _uid(created_by, "created_by")
  55. flow = validate_dataflow_spec(dataflow_spec)
  56. if flow["dataflow_uid"] != path_uid:
  57. raise ValueError("dataflow spec uid does not match path dataflow_uid")
  58. source = _source(source_text)
  59. inputs, output = _schema_hashes(
  60. flow, input_schema_hashes, output_schema_hash
  61. )
  62. standards, rules = self.repository.load_published_assets(flow)
  63. version = self.repository.begin_dataflow_release(
  64. dataflow_spec=flow,
  65. source_text=source,
  66. input_schema_hashes=inputs,
  67. output_schema_hash=output,
  68. created_by=actor,
  69. )
  70. version_id = _uid(version.get("id"), "dataflow_version_id")
  71. compiled: dict[str, dict[str, Any]] = {}
  72. binding_ids: dict[str, str] = {}
  73. def add_binding(
  74. *,
  75. binding_key: str,
  76. component_id: str,
  77. component_kind: str,
  78. rule_version_id: str,
  79. stage: str,
  80. order_no: int,
  81. idempotency: dict[str, Any] | None,
  82. provenance: dict[str, str],
  83. ) -> None:
  84. rule = rules.get(rule_version_id)
  85. if not isinstance(rule, dict):
  86. raise ValueError(
  87. f"published rule version {rule_version_id} was not found"
  88. )
  89. plan = compiled.setdefault(
  90. rule_version_id, compile_rule_plan(rule)
  91. )
  92. binding_id = new_governance_uid()
  93. binding_ids[binding_key] = binding_id
  94. self.repository.persist_component_plan(
  95. dataflow_version_id=version_id,
  96. component_binding_id=binding_id,
  97. component_id=component_id,
  98. component_kind=component_kind,
  99. rule_version_id=rule_version_id,
  100. stage=stage,
  101. order_no=order_no,
  102. idempotency=idempotency,
  103. provenance=provenance,
  104. plan=plan,
  105. schema_hashes={
  106. "inputs": inputs,
  107. "output": output,
  108. },
  109. )
  110. for component in flow["components"]:
  111. if component["type"] == "standard.enforce":
  112. standard_id = component["standard_version_id"]
  113. standard = standards.get(standard_id)
  114. if not isinstance(standard, dict):
  115. raise ValueError(
  116. f"published standard version {standard_id} was not found"
  117. )
  118. for clause_index, clause in enumerate(standard["clauses"]):
  119. clause_id = str(clause["clause_id"])
  120. binding_key = f"{component['id']}:{clause_id}"
  121. add_binding(
  122. binding_key=binding_key,
  123. component_id=(
  124. f"{component['id']}__{clause_id}"[:100]
  125. ),
  126. component_kind="quality.check",
  127. rule_version_id=str(clause["rule_version_id"]),
  128. stage=component["stage"],
  129. order_no=(component["order"] * 1000) + clause_index,
  130. idempotency=None,
  131. provenance={
  132. "standard_version_id": standard_id,
  133. "clause_id": clause_id,
  134. },
  135. )
  136. continue
  137. add_binding(
  138. binding_key=component["id"],
  139. component_id=component["id"],
  140. component_kind=component["type"],
  141. rule_version_id=component["rule_version_id"],
  142. stage=component["stage"],
  143. order_no=component["order"] * 1000,
  144. idempotency=component.get("idempotency"),
  145. provenance={},
  146. )
  147. release_rules = {
  148. rule_id: {
  149. **rule,
  150. "execution_plan": {
  151. "backend": compiled[rule_id]["backend"],
  152. "plan_hash": compiled[rule_id]["plan_hash"],
  153. },
  154. }
  155. for rule_id, rule in rules.items()
  156. if rule_id in compiled
  157. }
  158. package = resolve_production_line(
  159. flow,
  160. standards,
  161. release_rules,
  162. component_binding_ids=binding_ids,
  163. )
  164. return self.repository.complete_dataflow_release(
  165. dataflow_version_id=version_id,
  166. package=package,
  167. )