test_rules.py 11 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347
  1. from __future__ import annotations
  2. import copy
  3. import hashlib
  4. import json
  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(
  11. PLAN,
  12. sort_keys=True,
  13. separators=(",", ":"),
  14. ensure_ascii=False,
  15. ).encode("utf-8")
  16. ).hexdigest()
  17. def rule_node(node_type="quality.check"):
  18. node = {
  19. "id": "customer_mobile",
  20. "type": node_type,
  21. "purpose": "read" if node_type == "quality.check" else "write",
  22. "config": {
  23. "component_binding_id": new_governance_uid(),
  24. "rule_version_id": new_governance_uid(),
  25. "execution_plan_hash": PLAN_HASH,
  26. },
  27. }
  28. if node_type == "rule.apply":
  29. node["idempotency"] = {
  30. "strategy": "upsert",
  31. "key": "customer_id",
  32. }
  33. return node
  34. class Repository:
  35. def __init__(self, record=None):
  36. self.record = record
  37. self.calls = []
  38. def load(self, **kwargs):
  39. self.calls.append(kwargs)
  40. return copy.deepcopy(self.record)
  41. class Adapter:
  42. def __init__(self):
  43. self.calls = []
  44. def execute(self, *, plan, node, parameters, write_authorized):
  45. self.calls.append(
  46. {
  47. "plan": plan,
  48. "node": node,
  49. "parameters": parameters,
  50. "write_authorized": write_authorized,
  51. }
  52. )
  53. return {"rows_rejected": 3}
  54. def published_record(node, **overrides):
  55. value = {
  56. "component_binding_id": node["config"]["component_binding_id"],
  57. "rule_version_id": node["config"]["rule_version_id"],
  58. "backend": "quality_check",
  59. "plan": PLAN,
  60. "plan_hash": PLAN_HASH,
  61. "plan_status": "published",
  62. "rule_status": "published",
  63. "publication_audit_trusted": True,
  64. "logical_evidence_trusted": True,
  65. "physical_evidence_trusted": True,
  66. "component_kind": node["type"],
  67. "binding_idempotency": node.get("idempotency"),
  68. }
  69. value.update(overrides)
  70. return value
  71. def test_rule_executor_loads_only_published_plan_by_fixed_identifiers():
  72. from app.runner.rules import RulePlanExecutor
  73. node = rule_node()
  74. adapter = Adapter()
  75. repository = Repository(published_record(node))
  76. executor = RulePlanExecutor(
  77. repository,
  78. adapters={"quality_check": adapter},
  79. )
  80. result = executor.execute(node, {"partition": "2026-07-23"})
  81. assert result["rows_rejected"] == 3
  82. assert result["rule_version_id"] == node["config"]["rule_version_id"]
  83. assert repository.calls == [
  84. {
  85. "component_binding_id": node["config"]["component_binding_id"],
  86. "rule_version_id": node["config"]["rule_version_id"],
  87. "plan_hash": PLAN_HASH,
  88. }
  89. ]
  90. assert adapter.calls[0]["parameters"] == {"partition": "2026-07-23"}
  91. @pytest.mark.parametrize(
  92. "record",
  93. [
  94. None,
  95. {"plan_status": "compiled"},
  96. {"plan_status": "tested"},
  97. {"plan_status": "revoked"},
  98. {"rule_status": "deprecated"},
  99. {"plan_hash": "b" * 64},
  100. {"publication_audit_trusted": False},
  101. {"logical_evidence_trusted": False},
  102. {"physical_evidence_trusted": False},
  103. ],
  104. )
  105. def test_rule_executor_fails_closed_for_missing_revoked_or_mismatched_plan(record):
  106. from app.runner.rules import RulePlanExecutor
  107. node = rule_node()
  108. base = published_record(node)
  109. if record is not None:
  110. base.update(record)
  111. record = base
  112. executor = RulePlanExecutor(
  113. Repository(record),
  114. adapters={"quality_check": Adapter()},
  115. )
  116. with pytest.raises(NodeExecutionError):
  117. executor.execute(node, {})
  118. def test_rule_executor_rejects_inline_plan_or_unregistered_backend():
  119. from app.runner.rules import RulePlanExecutor
  120. node = rule_node()
  121. node["config"]["plan"] = {"op": "bypass"}
  122. executor = RulePlanExecutor(
  123. Repository(published_record(node)),
  124. adapters={},
  125. )
  126. with pytest.raises(NodeExecutionError):
  127. executor.execute(node, {})
  128. clean = rule_node()
  129. with pytest.raises(NodeExecutionError):
  130. RulePlanExecutor(
  131. Repository(published_record(clean, backend="generated_python")),
  132. adapters={},
  133. ).execute(clean, {})
  134. def test_mutating_rule_requires_governed_write_authorization_and_idempotency():
  135. from app.runner.rules import RulePlanExecutor
  136. node = rule_node("rule.apply")
  137. record = published_record(node, backend="sql_pushdown")
  138. executor = RulePlanExecutor(
  139. Repository(record),
  140. adapters={"sql_pushdown": Adapter()},
  141. )
  142. with pytest.raises(NodeExecutionError):
  143. executor.execute(node, {}, write_authorized=False)
  144. del node["idempotency"]
  145. with pytest.raises(NodeExecutionError):
  146. executor.execute(node, {}, write_authorized=True)
  147. def test_quality_node_dispatches_sql_plan_to_read_only_quality_adapter():
  148. from app.core.data_rules.compilers.sql import COMPILER_VERSION
  149. from app.runner.rules import RulePlanExecutor
  150. from tests.runner.test_rule_sql import sql_plan
  151. plan, plan_hash = sql_plan()
  152. node = rule_node()
  153. node["config"]["rule_version_id"] = plan["rule_version_id"]
  154. node["config"]["execution_plan_hash"] = plan_hash
  155. quality_adapter = Adapter()
  156. record = published_record(
  157. node,
  158. backend="sql_pushdown",
  159. compiler_version=COMPILER_VERSION,
  160. plan=plan,
  161. plan_hash=plan_hash,
  162. schema_hashes={
  163. "rule_spec_hash": plan["rule_spec_hash"],
  164. "input_schema_snapshot_id": plan[
  165. "input_schema_snapshot_id"
  166. ],
  167. "input_schema_hash": plan["input_schema_hash"],
  168. "output_schema_snapshot_id": plan[
  169. "output_schema_snapshot_id"
  170. ],
  171. "output_schema_hash": plan["output_schema_hash"],
  172. },
  173. canonical_rule_spec_hash=plan["rule_spec_hash"],
  174. canonical_input_schema_snapshot_id=plan[
  175. "input_schema_snapshot_id"
  176. ],
  177. canonical_input_schema_hash=plan["input_schema_hash"],
  178. canonical_output_schema_snapshot_id=plan[
  179. "output_schema_snapshot_id"
  180. ],
  181. canonical_output_schema_hash=plan["output_schema_hash"],
  182. canonical_input_data_source_uid=plan["data_source_uid"],
  183. canonical_output_data_source_uid=plan["data_source_uid"],
  184. canonical_input_object_ref="raw.customer",
  185. canonical_output_object_ref="clean.customer",
  186. canonical_input_dialect="postgresql",
  187. canonical_output_dialect="postgresql",
  188. canonical_input_access_mode="read",
  189. canonical_input_write_mode=None,
  190. canonical_output_access_mode="write",
  191. canonical_output_write_mode="append",
  192. )
  193. repository = Repository(record)
  194. class Evidence:
  195. heartbeat_interval_seconds = 10
  196. def __init__(self):
  197. self.resolved = []
  198. def start(self, **_kwargs):
  199. return new_governance_uid()
  200. def replay(self, _rule_run_id):
  201. return None
  202. def heartbeat(self, _rule_run_id, _lease_owner):
  203. return None
  204. def finish(self, _rule_run_id, _result):
  205. return None
  206. def resolve_sql_staging(self, receipt, **kwargs):
  207. self.resolved.append((receipt, kwargs))
  208. evidence = Evidence()
  209. executor = RulePlanExecutor(
  210. repository,
  211. adapters={"quality_check": quality_adapter},
  212. evidence_writer=evidence,
  213. )
  214. context = {
  215. "correlation_id": new_governance_uid(),
  216. "dataflow_uid": new_governance_uid(),
  217. "deployment_id": new_governance_uid(),
  218. "environment": "test",
  219. "workflow_version": 1,
  220. "node_id": node["id"],
  221. "task_jti": new_governance_uid(),
  222. }
  223. receipt = (
  224. "dataops-staging://01900000-0000-7000-8000-000000000099"
  225. )
  226. result = executor.execute(
  227. node,
  228. {"input_artifact": receipt},
  229. **context,
  230. )
  231. assert result["execution_plan_hash"] == plan_hash
  232. assert quality_adapter.calls[0]["write_authorized"] is False
  233. assert quality_adapter.calls[0]["parameters"] == {}
  234. assert evidence.resolved == [
  235. (
  236. receipt,
  237. {
  238. "deployment_id": context["deployment_id"],
  239. "correlation_id": context["correlation_id"],
  240. "input_binding_id": plan["input_binding_id"],
  241. },
  242. )
  243. ]
  244. @pytest.mark.parametrize(
  245. ("field", "value"),
  246. [
  247. ("canonical_input_access_mode", "write"),
  248. ("canonical_input_write_mode", "append"),
  249. ("canonical_output_access_mode", "read"),
  250. ("canonical_output_write_mode", None),
  251. ],
  252. )
  253. def test_sql_plan_rechecks_live_binding_access_before_execution(field, value):
  254. from app.core.data_rules.compilers.sql import COMPILER_VERSION
  255. from app.runner.rules import RulePlanExecutor
  256. from tests.runner.test_rule_sql import sql_plan
  257. plan, plan_hash = sql_plan()
  258. node = rule_node("rule.apply")
  259. node["config"]["rule_version_id"] = plan["rule_version_id"]
  260. node["config"]["execution_plan_hash"] = plan_hash
  261. values = {
  262. "backend": "sql_pushdown",
  263. "compiler_version": COMPILER_VERSION,
  264. "plan": plan,
  265. "plan_hash": plan_hash,
  266. "schema_hashes": {
  267. "rule_spec_hash": plan["rule_spec_hash"],
  268. "input_schema_snapshot_id": plan["input_schema_snapshot_id"],
  269. "input_schema_hash": plan["input_schema_hash"],
  270. "output_schema_snapshot_id": plan["output_schema_snapshot_id"],
  271. "output_schema_hash": plan["output_schema_hash"],
  272. },
  273. "canonical_rule_spec_hash": plan["rule_spec_hash"],
  274. "canonical_input_schema_snapshot_id": plan["input_schema_snapshot_id"],
  275. "canonical_input_schema_hash": plan["input_schema_hash"],
  276. "canonical_output_schema_snapshot_id": plan["output_schema_snapshot_id"],
  277. "canonical_output_schema_hash": plan["output_schema_hash"],
  278. "canonical_input_data_source_uid": plan["data_source_uid"],
  279. "canonical_output_data_source_uid": plan["data_source_uid"],
  280. "canonical_input_object_ref": "raw.customer",
  281. "canonical_output_object_ref": "clean.customer",
  282. "canonical_input_dialect": "postgresql",
  283. "canonical_output_dialect": "postgresql",
  284. "canonical_input_access_mode": "read",
  285. "canonical_input_write_mode": None,
  286. "canonical_output_access_mode": "write",
  287. "canonical_output_write_mode": "append",
  288. }
  289. values[field] = value
  290. record = published_record(node, **values)
  291. with pytest.raises(
  292. NodeExecutionError,
  293. match="canonical attestation",
  294. ):
  295. RulePlanExecutor(
  296. Repository(record),
  297. adapters={"sql_pushdown": Adapter()},
  298. ).execute(node, {}, write_authorized=True)