release.py 6.1 KB

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