from __future__ import annotations import hashlib import json import time import pytest from app.core.common.identifiers import new_governance_uid from app.runner.nodes import NodeExecutionError PLAN = {"op": "not_null", "column": "mobile"} PLAN_HASH = hashlib.sha256( json.dumps(PLAN, sort_keys=True, separators=(",", ":")).encode() ).hexdigest() def _node(): return { "id": "customer_mobile", "type": "quality.check", "purpose": "read", "config": { "component_binding_id": new_governance_uid(), "rule_version_id": new_governance_uid(), "execution_plan_hash": PLAN_HASH, }, } class Repository: def __init__(self, node): self.record = { "component_binding_id": node["config"]["component_binding_id"], "rule_version_id": node["config"]["rule_version_id"], "backend": "quality_check", "plan": PLAN, "plan_hash": PLAN_HASH, "plan_status": "published", "rule_status": "published", "component_kind": "quality.check", "binding_idempotency": None, } def load(self, **_kwargs): return dict(self.record) class Evidence: def __init__(self, replay=None): self.started = [] self.finished = [] self._replay = replay self.heartbeat_interval_seconds = 0.01 self.heartbeats = 0 def start(self, **values): self.started.append(values) return new_governance_uid() def replay(self, _rule_run_id): return self._replay def heartbeat(self, _rule_run_id, _lease_owner): self.heartbeats += 1 return None def finish(self, rule_run_id, result): self.finished.append((rule_run_id, result)) class Adapter: def __init__(self, result=None, error=None): self.result = result self.error = error self.calls = 0 def execute(self, **_kwargs): self.calls += 1 if self.error: raise self.error return dict(self.result or {}) def _context(): return { "correlation_id": new_governance_uid(), "dataflow_uid": new_governance_uid(), "deployment_id": new_governance_uid(), "environment": "test", "workflow_version": 3, "node_id": "customer_mobile", "task_jti": new_governance_uid(), } def test_rule_executor_records_success_and_bounded_violation_sample(): from app.runner.rules import RulePlanExecutor node = _node() evidence = Evidence() adapter = Adapter( { "rows_in": 3, "rows_out": 2, "rows_rejected": 1, "rows_quarantined": 1, "commit_outcome": "committed", "_violation_sample": [ {"mobile": f"13800138{index:03d}", "name": f"Person {index}"} for index in range(150) ], } ) executor = RulePlanExecutor( Repository(node), adapters={"quality_check": adapter}, evidence_writer=evidence, ) result = executor.execute(node, {}, **_context()) assert result["rows_in"] == 3 assert result["rows_out"] == 2 assert "_violation_sample" not in result assert evidence.started[0]["plan_hash"] == PLAN_HASH finished = evidence.finished[0][1] assert finished["status"] == "success" assert finished["sample_count"] == 100 assert finished["redaction_policy"] == "rule-violation-default-v1" assert finished["violation_sample"][0] == { "mobile": "[REDACTED]", "name": "[REDACTED]", } @pytest.mark.parametrize( ("error", "status", "commit_outcome"), [ ( NodeExecutionError( "write failed", commit_outcome="not_committed", ), "failed", "not_committed", ), ( NodeExecutionError( "commit uncertain", commit_outcome="unknown", ), "unknown", "unknown", ), ], ) def test_rule_executor_finalizes_failure_and_unknown_commit( error, status, commit_outcome ): from app.runner.rules import RulePlanExecutor node = _node() evidence = Evidence() executor = RulePlanExecutor( Repository(node), adapters={"quality_check": Adapter(error=error)}, evidence_writer=evidence, ) with pytest.raises(NodeExecutionError): executor.execute(node, {}, **_context()) assert evidence.finished[0][1]["status"] == status assert evidence.finished[0][1]["commit_outcome"] == commit_outcome assert "violation_sample" not in evidence.finished[0][1] def test_rule_executor_replays_finished_evidence_without_adapter_execution(): from app.runner.rules import RulePlanExecutor node = _node() prior = { "status": "success", "rows_in": 3, "rows_out": 2, "rows_rejected": 1, "rows_quarantined": 1, "commit_outcome": "committed", "output_artifact": "minio://dataops-rules/rules/prior.parquet", } evidence = Evidence(replay=prior) adapter = Adapter() executor = RulePlanExecutor( Repository(node), adapters={"quality_check": adapter}, evidence_writer=evidence, ) result = executor.execute(node, {}, **_context()) assert result["output_artifact"] == prior["output_artifact"] assert adapter.calls == 0 assert evidence.finished == [] def test_rule_executor_finalizes_cancellation(): from asyncio import CancelledError from app.runner.rules import RulePlanExecutor node = _node() evidence = Evidence() executor = RulePlanExecutor( Repository(node), adapters={ "quality_check": Adapter(error=CancelledError()) }, evidence_writer=evidence, ) with pytest.raises(CancelledError): executor.execute(node, {}, **_context()) assert evidence.finished[0][1]["status"] == "cancelled" def test_evidence_migration_is_forward_only_and_bounded(): from pathlib import Path source = Path( "migrations/versions/" "20260723_160_rule_execution_evidence.py" ).read_text(encoding="utf-8") assert 'revision = "20260723_160"' in source assert 'down_revision = "20260723_150"' in source assert "UNIQUE (evidence_key)" in source assert "sample_count <= 100" in source assert "rule_violation_samples_rule_run_key" in source assert "duplicate rule_run_id" in source assert "forward-only" in source def test_public_result_rejects_nested_records_and_unclosed_fields(): from app.runner.rule_evidence import PostgresRuleEvidenceWriter base = { "status": "success", "rows_in": 1, "rows_out": 1, "rows_rejected": 0, "rows_quarantined": 0, "commit_outcome": "committed", "timings": {"duration_ms": 1}, } for public_result in ( {"rows": [{"secret": "value"}]}, {"records": [{"secret": "value"}]}, {"arbitrary": {"nested": "value"}}, {"violations": [{"step_id": "x", "count": 1, "raw": "secret"}]}, ): with pytest.raises(ValueError, match="evidence safe"): PostgresRuleEvidenceWriter._validate_finish( {**base, "public_result": public_result} ) def test_public_result_contract_is_backend_specific(): from app.runner.rule_evidence import validate_public_rule_result with pytest.raises(ValueError, match="evidence safe"): validate_public_rule_result( { "rows_in": 1, "rows_out": 1, "artifact_ref": "minio://unexpected", }, backend="sql_pushdown", ) assert validate_public_rule_result( { "rows_in": 1, "rows_out": 1, "commit_outcome": "committed", "output_artifact": ( "dataops-staging://" "01900000-0000-7000-8000-000000000099" ), }, backend="sql_pushdown", )["rows_out"] == 1 def test_rule_executor_renews_lease_while_adapter_is_running(): from app.runner.rules import RulePlanExecutor class SlowAdapter(Adapter): def execute(self, **_kwargs): time.sleep(0.04) return { "rows_in": 1, "rows_out": 1, "rows_rejected": 0, "rows_quarantined": 0, "commit_outcome": "committed", } node = _node() evidence = Evidence() result = RulePlanExecutor( Repository(node), adapters={"quality_check": SlowAdapter()}, evidence_writer=evidence, ).execute(node, {}, **_context()) assert result["rows_out"] == 1 assert evidence.heartbeats >= 2 def test_heartbeat_failure_after_committed_adapter_is_unknown_committed(): from app.runner.rules import RulePlanExecutor class FailingHeartbeatEvidence(Evidence): def heartbeat(self, _rule_run_id, _lease_owner): self.heartbeats += 1 if self.heartbeats > 1: raise RuntimeError("postgres unavailable") class SlowCommittedAdapter(Adapter): def execute(self, **_kwargs): time.sleep(0.04) return { "rows_in": 1, "rows_out": 1, "rows_rejected": 0, "rows_quarantined": 0, "commit_outcome": "committed", } node = _node() evidence = FailingHeartbeatEvidence() executor = RulePlanExecutor( Repository(node), adapters={"quality_check": SlowCommittedAdapter()}, evidence_writer=evidence, ) with pytest.raises(NodeExecutionError) as error: executor.execute(node, {}, **_context()) assert error.value.commit_outcome == "unknown" assert evidence.finished[-1][1]["status"] == "unknown" assert evidence.finished[-1][1]["commit_outcome"] == "committed" def test_committed_adapter_with_invalid_public_result_is_unknown_committed(): from app.runner.rules import RulePlanExecutor node = _node() evidence = Evidence() executor = RulePlanExecutor( Repository(node), adapters={ "quality_check": Adapter( { "rows_in": 1, "rows_out": 1, "rows_rejected": 0, "rows_quarantined": 0, "commit_outcome": "committed", "raw_records": [{"secret": "value"}], } ) }, evidence_writer=evidence, ) with pytest.raises(NodeExecutionError) as error: executor.execute(node, {}, **_context()) assert error.value.commit_outcome == "unknown" assert evidence.finished[-1][1]["status"] == "unknown" assert evidence.finished[-1][1]["commit_outcome"] == "committed"