| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387 |
- 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",
- "publication_audit_trusted": True,
- "logical_evidence_trusted": True,
- "physical_evidence_trusted": True,
- "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"
|