test_rule_evidence.py 6.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231
  1. from __future__ import annotations
  2. import hashlib
  3. import json
  4. import pytest
  5. from app.core.common.identifiers import new_governance_uid
  6. from app.runner.nodes import NodeExecutionError
  7. PLAN = {"op": "not_null", "column": "mobile"}
  8. PLAN_HASH = hashlib.sha256(
  9. json.dumps(PLAN, sort_keys=True, separators=(",", ":")).encode()
  10. ).hexdigest()
  11. def _node():
  12. return {
  13. "id": "customer_mobile",
  14. "type": "quality.check",
  15. "purpose": "read",
  16. "config": {
  17. "component_binding_id": new_governance_uid(),
  18. "rule_version_id": new_governance_uid(),
  19. "execution_plan_hash": PLAN_HASH,
  20. },
  21. }
  22. class Repository:
  23. def __init__(self, node):
  24. self.record = {
  25. "component_binding_id": node["config"]["component_binding_id"],
  26. "rule_version_id": node["config"]["rule_version_id"],
  27. "backend": "quality_check",
  28. "plan": PLAN,
  29. "plan_hash": PLAN_HASH,
  30. "plan_status": "published",
  31. "rule_status": "published",
  32. "component_kind": "quality.check",
  33. "binding_idempotency": None,
  34. }
  35. def load(self, **_kwargs):
  36. return dict(self.record)
  37. class Evidence:
  38. def __init__(self, replay=None):
  39. self.started = []
  40. self.finished = []
  41. self._replay = replay
  42. def start(self, **values):
  43. self.started.append(values)
  44. return new_governance_uid()
  45. def replay(self, _rule_run_id):
  46. return self._replay
  47. def finish(self, rule_run_id, result):
  48. self.finished.append((rule_run_id, result))
  49. class Adapter:
  50. def __init__(self, result=None, error=None):
  51. self.result = result
  52. self.error = error
  53. self.calls = 0
  54. def execute(self, **_kwargs):
  55. self.calls += 1
  56. if self.error:
  57. raise self.error
  58. return dict(self.result or {})
  59. def _context():
  60. return {
  61. "correlation_id": new_governance_uid(),
  62. "dataflow_uid": new_governance_uid(),
  63. "workflow_version": 3,
  64. "node_id": "customer_mobile",
  65. }
  66. def test_rule_executor_records_success_and_bounded_violation_sample():
  67. from app.runner.rules import RulePlanExecutor
  68. node = _node()
  69. evidence = Evidence()
  70. adapter = Adapter(
  71. {
  72. "rows_in": 3,
  73. "rows_out": 2,
  74. "rows_rejected": 1,
  75. "rows_quarantined": 1,
  76. "commit_outcome": "committed",
  77. "_violation_sample": [
  78. {"mobile": f"13800138{index:03d}", "name": f"Person {index}"}
  79. for index in range(150)
  80. ],
  81. }
  82. )
  83. executor = RulePlanExecutor(
  84. Repository(node),
  85. adapters={"quality_check": adapter},
  86. evidence_writer=evidence,
  87. )
  88. result = executor.execute(node, {}, **_context())
  89. assert result["rows_in"] == 3
  90. assert result["rows_out"] == 2
  91. assert "_violation_sample" not in result
  92. assert evidence.started[0]["plan_hash"] == PLAN_HASH
  93. finished = evidence.finished[0][1]
  94. assert finished["status"] == "success"
  95. assert finished["sample_count"] == 100
  96. assert finished["redaction_policy"] == "rule-violation-default-v1"
  97. assert finished["violation_sample"][0] == {
  98. "mobile": "[REDACTED]",
  99. "name": "[REDACTED]",
  100. }
  101. @pytest.mark.parametrize(
  102. ("error", "status", "commit_outcome"),
  103. [
  104. (
  105. NodeExecutionError(
  106. "write failed",
  107. commit_outcome="not_committed",
  108. ),
  109. "failed",
  110. "not_committed",
  111. ),
  112. (
  113. NodeExecutionError(
  114. "commit uncertain",
  115. commit_outcome="unknown",
  116. ),
  117. "unknown",
  118. "unknown",
  119. ),
  120. ],
  121. )
  122. def test_rule_executor_finalizes_failure_and_unknown_commit(
  123. error, status, commit_outcome
  124. ):
  125. from app.runner.rules import RulePlanExecutor
  126. node = _node()
  127. evidence = Evidence()
  128. executor = RulePlanExecutor(
  129. Repository(node),
  130. adapters={"quality_check": Adapter(error=error)},
  131. evidence_writer=evidence,
  132. )
  133. with pytest.raises(NodeExecutionError):
  134. executor.execute(node, {}, **_context())
  135. assert evidence.finished[0][1]["status"] == status
  136. assert evidence.finished[0][1]["commit_outcome"] == commit_outcome
  137. assert "violation_sample" not in evidence.finished[0][1]
  138. def test_rule_executor_replays_finished_evidence_without_adapter_execution():
  139. from app.runner.rules import RulePlanExecutor
  140. node = _node()
  141. prior = {
  142. "status": "success",
  143. "rows_in": 3,
  144. "rows_out": 2,
  145. "rows_rejected": 1,
  146. "rows_quarantined": 1,
  147. "commit_outcome": "committed",
  148. "output_artifact": "minio://dataops-rules/rules/prior.parquet",
  149. }
  150. evidence = Evidence(replay=prior)
  151. adapter = Adapter()
  152. executor = RulePlanExecutor(
  153. Repository(node),
  154. adapters={"quality_check": adapter},
  155. evidence_writer=evidence,
  156. )
  157. result = executor.execute(node, {}, **_context())
  158. assert result["output_artifact"] == prior["output_artifact"]
  159. assert adapter.calls == 0
  160. assert evidence.finished == []
  161. def test_rule_executor_finalizes_cancellation():
  162. from asyncio import CancelledError
  163. from app.runner.rules import RulePlanExecutor
  164. node = _node()
  165. evidence = Evidence()
  166. executor = RulePlanExecutor(
  167. Repository(node),
  168. adapters={
  169. "quality_check": Adapter(error=CancelledError())
  170. },
  171. evidence_writer=evidence,
  172. )
  173. with pytest.raises(CancelledError):
  174. executor.execute(node, {}, **_context())
  175. assert evidence.finished[0][1]["status"] == "cancelled"
  176. def test_evidence_migration_is_forward_only_and_bounded():
  177. from pathlib import Path
  178. source = Path(
  179. "migrations/versions/"
  180. "20260723_160_rule_execution_evidence.py"
  181. ).read_text(encoding="utf-8")
  182. assert 'revision = "20260723_160"' in source
  183. assert 'down_revision = "20260723_150"' in source
  184. assert "UNIQUE (evidence_key)" in source
  185. assert "sample_count <= 100" in source
  186. assert "rule_violation_samples_rule_run_key" in source
  187. assert "forward-only" in source