test_rule_evidence.py 7.9 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288
  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 heartbeat(self, _rule_run_id, _lease_owner):
  48. return None
  49. def finish(self, rule_run_id, result):
  50. self.finished.append((rule_run_id, result))
  51. class Adapter:
  52. def __init__(self, result=None, error=None):
  53. self.result = result
  54. self.error = error
  55. self.calls = 0
  56. def execute(self, **_kwargs):
  57. self.calls += 1
  58. if self.error:
  59. raise self.error
  60. return dict(self.result or {})
  61. def _context():
  62. return {
  63. "correlation_id": new_governance_uid(),
  64. "dataflow_uid": new_governance_uid(),
  65. "deployment_id": new_governance_uid(),
  66. "environment": "test",
  67. "workflow_version": 3,
  68. "node_id": "customer_mobile",
  69. "task_jti": new_governance_uid(),
  70. }
  71. def test_rule_executor_records_success_and_bounded_violation_sample():
  72. from app.runner.rules import RulePlanExecutor
  73. node = _node()
  74. evidence = Evidence()
  75. adapter = Adapter(
  76. {
  77. "rows_in": 3,
  78. "rows_out": 2,
  79. "rows_rejected": 1,
  80. "rows_quarantined": 1,
  81. "commit_outcome": "committed",
  82. "_violation_sample": [
  83. {"mobile": f"13800138{index:03d}", "name": f"Person {index}"}
  84. for index in range(150)
  85. ],
  86. }
  87. )
  88. executor = RulePlanExecutor(
  89. Repository(node),
  90. adapters={"quality_check": adapter},
  91. evidence_writer=evidence,
  92. )
  93. result = executor.execute(node, {}, **_context())
  94. assert result["rows_in"] == 3
  95. assert result["rows_out"] == 2
  96. assert "_violation_sample" not in result
  97. assert evidence.started[0]["plan_hash"] == PLAN_HASH
  98. finished = evidence.finished[0][1]
  99. assert finished["status"] == "success"
  100. assert finished["sample_count"] == 100
  101. assert finished["redaction_policy"] == "rule-violation-default-v1"
  102. assert finished["violation_sample"][0] == {
  103. "mobile": "[REDACTED]",
  104. "name": "[REDACTED]",
  105. }
  106. @pytest.mark.parametrize(
  107. ("error", "status", "commit_outcome"),
  108. [
  109. (
  110. NodeExecutionError(
  111. "write failed",
  112. commit_outcome="not_committed",
  113. ),
  114. "failed",
  115. "not_committed",
  116. ),
  117. (
  118. NodeExecutionError(
  119. "commit uncertain",
  120. commit_outcome="unknown",
  121. ),
  122. "unknown",
  123. "unknown",
  124. ),
  125. ],
  126. )
  127. def test_rule_executor_finalizes_failure_and_unknown_commit(
  128. error, status, commit_outcome
  129. ):
  130. from app.runner.rules import RulePlanExecutor
  131. node = _node()
  132. evidence = Evidence()
  133. executor = RulePlanExecutor(
  134. Repository(node),
  135. adapters={"quality_check": Adapter(error=error)},
  136. evidence_writer=evidence,
  137. )
  138. with pytest.raises(NodeExecutionError):
  139. executor.execute(node, {}, **_context())
  140. assert evidence.finished[0][1]["status"] == status
  141. assert evidence.finished[0][1]["commit_outcome"] == commit_outcome
  142. assert "violation_sample" not in evidence.finished[0][1]
  143. def test_rule_executor_replays_finished_evidence_without_adapter_execution():
  144. from app.runner.rules import RulePlanExecutor
  145. node = _node()
  146. prior = {
  147. "status": "success",
  148. "rows_in": 3,
  149. "rows_out": 2,
  150. "rows_rejected": 1,
  151. "rows_quarantined": 1,
  152. "commit_outcome": "committed",
  153. "output_artifact": "minio://dataops-rules/rules/prior.parquet",
  154. }
  155. evidence = Evidence(replay=prior)
  156. adapter = Adapter()
  157. executor = RulePlanExecutor(
  158. Repository(node),
  159. adapters={"quality_check": adapter},
  160. evidence_writer=evidence,
  161. )
  162. result = executor.execute(node, {}, **_context())
  163. assert result["output_artifact"] == prior["output_artifact"]
  164. assert adapter.calls == 0
  165. assert evidence.finished == []
  166. def test_rule_executor_finalizes_cancellation():
  167. from asyncio import CancelledError
  168. from app.runner.rules import RulePlanExecutor
  169. node = _node()
  170. evidence = Evidence()
  171. executor = RulePlanExecutor(
  172. Repository(node),
  173. adapters={
  174. "quality_check": Adapter(error=CancelledError())
  175. },
  176. evidence_writer=evidence,
  177. )
  178. with pytest.raises(CancelledError):
  179. executor.execute(node, {}, **_context())
  180. assert evidence.finished[0][1]["status"] == "cancelled"
  181. def test_evidence_migration_is_forward_only_and_bounded():
  182. from pathlib import Path
  183. source = Path(
  184. "migrations/versions/"
  185. "20260723_160_rule_execution_evidence.py"
  186. ).read_text(encoding="utf-8")
  187. assert 'revision = "20260723_160"' in source
  188. assert 'down_revision = "20260723_150"' in source
  189. assert "UNIQUE (evidence_key)" in source
  190. assert "sample_count <= 100" in source
  191. assert "rule_violation_samples_rule_run_key" in source
  192. assert "duplicate rule_run_id" in source
  193. assert "forward-only" in source
  194. def test_public_result_rejects_nested_records_and_unclosed_fields():
  195. from app.runner.rule_evidence import PostgresRuleEvidenceWriter
  196. base = {
  197. "status": "success",
  198. "rows_in": 1,
  199. "rows_out": 1,
  200. "rows_rejected": 0,
  201. "rows_quarantined": 0,
  202. "commit_outcome": "committed",
  203. "timings": {"duration_ms": 1},
  204. }
  205. for public_result in (
  206. {"rows": [{"secret": "value"}]},
  207. {"records": [{"secret": "value"}]},
  208. {"arbitrary": {"nested": "value"}},
  209. {"violations": [{"step_id": "x", "count": 1, "raw": "secret"}]},
  210. ):
  211. with pytest.raises(ValueError, match="evidence safe"):
  212. PostgresRuleEvidenceWriter._validate_finish(
  213. {**base, "public_result": public_result}
  214. )
  215. def test_public_result_contract_is_backend_specific():
  216. from app.runner.rule_evidence import validate_public_rule_result
  217. with pytest.raises(ValueError, match="evidence safe"):
  218. validate_public_rule_result(
  219. {
  220. "rows_in": 1,
  221. "rows_out": 1,
  222. "artifact_ref": "minio://unexpected",
  223. },
  224. backend="sql_pushdown",
  225. )
  226. assert validate_public_rule_result(
  227. {
  228. "rows_in": 1,
  229. "rows_out": 1,
  230. "commit_outcome": "committed",
  231. "output_artifact": (
  232. "dataops-staging://"
  233. "01900000-0000-7000-8000-000000000099"
  234. ),
  235. },
  236. backend="sql_pushdown",
  237. )["rows_out"] == 1