test_rule_evidence.py 11 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384
  1. from __future__ import annotations
  2. import hashlib
  3. import json
  4. import time
  5. import pytest
  6. from app.core.common.identifiers import new_governance_uid
  7. from app.runner.nodes import NodeExecutionError
  8. PLAN = {"op": "not_null", "column": "mobile"}
  9. PLAN_HASH = hashlib.sha256(
  10. json.dumps(PLAN, sort_keys=True, separators=(",", ":")).encode()
  11. ).hexdigest()
  12. def _node():
  13. return {
  14. "id": "customer_mobile",
  15. "type": "quality.check",
  16. "purpose": "read",
  17. "config": {
  18. "component_binding_id": new_governance_uid(),
  19. "rule_version_id": new_governance_uid(),
  20. "execution_plan_hash": PLAN_HASH,
  21. },
  22. }
  23. class Repository:
  24. def __init__(self, node):
  25. self.record = {
  26. "component_binding_id": node["config"]["component_binding_id"],
  27. "rule_version_id": node["config"]["rule_version_id"],
  28. "backend": "quality_check",
  29. "plan": PLAN,
  30. "plan_hash": PLAN_HASH,
  31. "plan_status": "published",
  32. "rule_status": "published",
  33. "component_kind": "quality.check",
  34. "binding_idempotency": None,
  35. }
  36. def load(self, **_kwargs):
  37. return dict(self.record)
  38. class Evidence:
  39. def __init__(self, replay=None):
  40. self.started = []
  41. self.finished = []
  42. self._replay = replay
  43. self.heartbeat_interval_seconds = 0.01
  44. self.heartbeats = 0
  45. def start(self, **values):
  46. self.started.append(values)
  47. return new_governance_uid()
  48. def replay(self, _rule_run_id):
  49. return self._replay
  50. def heartbeat(self, _rule_run_id, _lease_owner):
  51. self.heartbeats += 1
  52. return None
  53. def finish(self, rule_run_id, result):
  54. self.finished.append((rule_run_id, result))
  55. class Adapter:
  56. def __init__(self, result=None, error=None):
  57. self.result = result
  58. self.error = error
  59. self.calls = 0
  60. def execute(self, **_kwargs):
  61. self.calls += 1
  62. if self.error:
  63. raise self.error
  64. return dict(self.result or {})
  65. def _context():
  66. return {
  67. "correlation_id": new_governance_uid(),
  68. "dataflow_uid": new_governance_uid(),
  69. "deployment_id": new_governance_uid(),
  70. "environment": "test",
  71. "workflow_version": 3,
  72. "node_id": "customer_mobile",
  73. "task_jti": new_governance_uid(),
  74. }
  75. def test_rule_executor_records_success_and_bounded_violation_sample():
  76. from app.runner.rules import RulePlanExecutor
  77. node = _node()
  78. evidence = Evidence()
  79. adapter = Adapter(
  80. {
  81. "rows_in": 3,
  82. "rows_out": 2,
  83. "rows_rejected": 1,
  84. "rows_quarantined": 1,
  85. "commit_outcome": "committed",
  86. "_violation_sample": [
  87. {"mobile": f"13800138{index:03d}", "name": f"Person {index}"}
  88. for index in range(150)
  89. ],
  90. }
  91. )
  92. executor = RulePlanExecutor(
  93. Repository(node),
  94. adapters={"quality_check": adapter},
  95. evidence_writer=evidence,
  96. )
  97. result = executor.execute(node, {}, **_context())
  98. assert result["rows_in"] == 3
  99. assert result["rows_out"] == 2
  100. assert "_violation_sample" not in result
  101. assert evidence.started[0]["plan_hash"] == PLAN_HASH
  102. finished = evidence.finished[0][1]
  103. assert finished["status"] == "success"
  104. assert finished["sample_count"] == 100
  105. assert finished["redaction_policy"] == "rule-violation-default-v1"
  106. assert finished["violation_sample"][0] == {
  107. "mobile": "[REDACTED]",
  108. "name": "[REDACTED]",
  109. }
  110. @pytest.mark.parametrize(
  111. ("error", "status", "commit_outcome"),
  112. [
  113. (
  114. NodeExecutionError(
  115. "write failed",
  116. commit_outcome="not_committed",
  117. ),
  118. "failed",
  119. "not_committed",
  120. ),
  121. (
  122. NodeExecutionError(
  123. "commit uncertain",
  124. commit_outcome="unknown",
  125. ),
  126. "unknown",
  127. "unknown",
  128. ),
  129. ],
  130. )
  131. def test_rule_executor_finalizes_failure_and_unknown_commit(
  132. error, status, commit_outcome
  133. ):
  134. from app.runner.rules import RulePlanExecutor
  135. node = _node()
  136. evidence = Evidence()
  137. executor = RulePlanExecutor(
  138. Repository(node),
  139. adapters={"quality_check": Adapter(error=error)},
  140. evidence_writer=evidence,
  141. )
  142. with pytest.raises(NodeExecutionError):
  143. executor.execute(node, {}, **_context())
  144. assert evidence.finished[0][1]["status"] == status
  145. assert evidence.finished[0][1]["commit_outcome"] == commit_outcome
  146. assert "violation_sample" not in evidence.finished[0][1]
  147. def test_rule_executor_replays_finished_evidence_without_adapter_execution():
  148. from app.runner.rules import RulePlanExecutor
  149. node = _node()
  150. prior = {
  151. "status": "success",
  152. "rows_in": 3,
  153. "rows_out": 2,
  154. "rows_rejected": 1,
  155. "rows_quarantined": 1,
  156. "commit_outcome": "committed",
  157. "output_artifact": "minio://dataops-rules/rules/prior.parquet",
  158. }
  159. evidence = Evidence(replay=prior)
  160. adapter = Adapter()
  161. executor = RulePlanExecutor(
  162. Repository(node),
  163. adapters={"quality_check": adapter},
  164. evidence_writer=evidence,
  165. )
  166. result = executor.execute(node, {}, **_context())
  167. assert result["output_artifact"] == prior["output_artifact"]
  168. assert adapter.calls == 0
  169. assert evidence.finished == []
  170. def test_rule_executor_finalizes_cancellation():
  171. from asyncio import CancelledError
  172. from app.runner.rules import RulePlanExecutor
  173. node = _node()
  174. evidence = Evidence()
  175. executor = RulePlanExecutor(
  176. Repository(node),
  177. adapters={
  178. "quality_check": Adapter(error=CancelledError())
  179. },
  180. evidence_writer=evidence,
  181. )
  182. with pytest.raises(CancelledError):
  183. executor.execute(node, {}, **_context())
  184. assert evidence.finished[0][1]["status"] == "cancelled"
  185. def test_evidence_migration_is_forward_only_and_bounded():
  186. from pathlib import Path
  187. source = Path(
  188. "migrations/versions/"
  189. "20260723_160_rule_execution_evidence.py"
  190. ).read_text(encoding="utf-8")
  191. assert 'revision = "20260723_160"' in source
  192. assert 'down_revision = "20260723_150"' in source
  193. assert "UNIQUE (evidence_key)" in source
  194. assert "sample_count <= 100" in source
  195. assert "rule_violation_samples_rule_run_key" in source
  196. assert "duplicate rule_run_id" in source
  197. assert "forward-only" in source
  198. def test_public_result_rejects_nested_records_and_unclosed_fields():
  199. from app.runner.rule_evidence import PostgresRuleEvidenceWriter
  200. base = {
  201. "status": "success",
  202. "rows_in": 1,
  203. "rows_out": 1,
  204. "rows_rejected": 0,
  205. "rows_quarantined": 0,
  206. "commit_outcome": "committed",
  207. "timings": {"duration_ms": 1},
  208. }
  209. for public_result in (
  210. {"rows": [{"secret": "value"}]},
  211. {"records": [{"secret": "value"}]},
  212. {"arbitrary": {"nested": "value"}},
  213. {"violations": [{"step_id": "x", "count": 1, "raw": "secret"}]},
  214. ):
  215. with pytest.raises(ValueError, match="evidence safe"):
  216. PostgresRuleEvidenceWriter._validate_finish(
  217. {**base, "public_result": public_result}
  218. )
  219. def test_public_result_contract_is_backend_specific():
  220. from app.runner.rule_evidence import validate_public_rule_result
  221. with pytest.raises(ValueError, match="evidence safe"):
  222. validate_public_rule_result(
  223. {
  224. "rows_in": 1,
  225. "rows_out": 1,
  226. "artifact_ref": "minio://unexpected",
  227. },
  228. backend="sql_pushdown",
  229. )
  230. assert validate_public_rule_result(
  231. {
  232. "rows_in": 1,
  233. "rows_out": 1,
  234. "commit_outcome": "committed",
  235. "output_artifact": (
  236. "dataops-staging://"
  237. "01900000-0000-7000-8000-000000000099"
  238. ),
  239. },
  240. backend="sql_pushdown",
  241. )["rows_out"] == 1
  242. def test_rule_executor_renews_lease_while_adapter_is_running():
  243. from app.runner.rules import RulePlanExecutor
  244. class SlowAdapter(Adapter):
  245. def execute(self, **_kwargs):
  246. time.sleep(0.04)
  247. return {
  248. "rows_in": 1,
  249. "rows_out": 1,
  250. "rows_rejected": 0,
  251. "rows_quarantined": 0,
  252. "commit_outcome": "committed",
  253. }
  254. node = _node()
  255. evidence = Evidence()
  256. result = RulePlanExecutor(
  257. Repository(node),
  258. adapters={"quality_check": SlowAdapter()},
  259. evidence_writer=evidence,
  260. ).execute(node, {}, **_context())
  261. assert result["rows_out"] == 1
  262. assert evidence.heartbeats >= 2
  263. def test_heartbeat_failure_after_committed_adapter_is_unknown_committed():
  264. from app.runner.rules import RulePlanExecutor
  265. class FailingHeartbeatEvidence(Evidence):
  266. def heartbeat(self, _rule_run_id, _lease_owner):
  267. self.heartbeats += 1
  268. if self.heartbeats > 1:
  269. raise RuntimeError("postgres unavailable")
  270. class SlowCommittedAdapter(Adapter):
  271. def execute(self, **_kwargs):
  272. time.sleep(0.04)
  273. return {
  274. "rows_in": 1,
  275. "rows_out": 1,
  276. "rows_rejected": 0,
  277. "rows_quarantined": 0,
  278. "commit_outcome": "committed",
  279. }
  280. node = _node()
  281. evidence = FailingHeartbeatEvidence()
  282. executor = RulePlanExecutor(
  283. Repository(node),
  284. adapters={"quality_check": SlowCommittedAdapter()},
  285. evidence_writer=evidence,
  286. )
  287. with pytest.raises(NodeExecutionError) as error:
  288. executor.execute(node, {}, **_context())
  289. assert error.value.commit_outcome == "unknown"
  290. assert evidence.finished[-1][1]["status"] == "unknown"
  291. assert evidence.finished[-1][1]["commit_outcome"] == "committed"
  292. def test_committed_adapter_with_invalid_public_result_is_unknown_committed():
  293. from app.runner.rules import RulePlanExecutor
  294. node = _node()
  295. evidence = Evidence()
  296. executor = RulePlanExecutor(
  297. Repository(node),
  298. adapters={
  299. "quality_check": Adapter(
  300. {
  301. "rows_in": 1,
  302. "rows_out": 1,
  303. "rows_rejected": 0,
  304. "rows_quarantined": 0,
  305. "commit_outcome": "committed",
  306. "raw_records": [{"secret": "value"}],
  307. }
  308. )
  309. },
  310. evidence_writer=evidence,
  311. )
  312. with pytest.raises(NodeExecutionError) as error:
  313. executor.execute(node, {}, **_context())
  314. assert error.value.commit_outcome == "unknown"
  315. assert evidence.finished[-1][1]["status"] == "unknown"
  316. assert evidence.finished[-1][1]["commit_outcome"] == "committed"