test_rule_evidence.py 11 KB

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