from __future__ import annotations import inspect import json from contextlib import asynccontextmanager from unittest.mock import AsyncMock import pytest from src.jobs import awooop_ansible_candidate_backfill_job as candidate_job from src.services import awooop_ansible_check_mode_service as service class _MappingResult: def __init__(self, row: dict | None = None, scalar=None) -> None: self._row = row self._scalar = scalar def mappings(self) -> _MappingResult: return self def one_or_none(self) -> dict | None: return self._row def scalar_one_or_none(self): return self._scalar def scalar(self): return self._scalar class _SequenceDB: def __init__(self, *results: _MappingResult) -> None: self._results = list(results) self.statements: list[str] = [] self.parameters: list[dict] = [] async def execute(self, statement, parameters=None): self.statements.append(str(statement)) self.parameters.append(dict(parameters or {})) return self._results.pop(0) if self._results else _MappingResult() def _claim() -> service.AnsibleCheckModeClaim: return service.AnsibleCheckModeClaim( op_id="00000000-0000-0000-0000-000000000101", source_candidate_op_id="00000000-0000-0000-0000-000000000100", incident_id="INC-20260711-D037E5", catalog_id="ansible:awoooi-auto-repair-canary", playbook_path="infra/ansible/playbooks/awoooi-auto-repair-canary.yml", apply_playbook_path=( "infra/ansible/playbooks/awoooi-auto-repair-canary.yml" ), inventory_hosts=("host_121",), risk_level="medium", input_payload={ "automation_run_id": "00000000-0000-0000-0000-000000000102", "approval_id": "00000000-0000-0000-0000-000000000103", }, ) def _verified_result() -> service.AnsibleRunResult: return service.AnsibleRunResult( returncode=0, stdout="", stderr="", duration_ms=25, post_verifier_passed=True, ) def test_runtime_stage_ids_only_accepts_durable_receipts() -> None: stage_ids = service._runtime_stage_ids( { "runtime_stage_receipts": [ { "stage_id": "mcp_context", "durable_receipt": True, }, { "stage_id": "telegram_receipt", "durable_receipt": False, }, ] } ) assert stage_ids == {"mcp_context"} def test_incident_terminal_writer_normalizes_legacy_non_object_outcome() -> None: source = inspect.getsource(service._record_incident_terminal_disposition) assert "jsonb_typeof(CAST(outcome AS jsonb))" in source assert "jsonb_build_object(" in source assert "'legacy_outcome'" in source @pytest.mark.asyncio @pytest.mark.parametrize( ("proposal_executed", "expected_success"), [(True, False), (False, None)], ) async def test_incident_terminal_round_trips_negative_outcome_truth( monkeypatch: pytest.MonkeyPatch, proposal_executed: bool, expected_success: bool | None, ) -> None: terminal = { "schema_version": "ansible_incident_terminal_disposition_v1", "automation_run_id": "00000000-0000-0000-0000-000000000102", "apply_op_id": "00000000-0000-0000-0000-000000000104", "terminal_type": "no_write_replay_failed_waiting_repair_candidate", } db = _SequenceDB( _MappingResult({"proposal_executed": proposal_executed}), _MappingResult( { "incident_status": "MITIGATING", "updated_at": None, "resolved_at": None, "outcome": { "automation_terminal": terminal, "proposal_executed": proposal_executed, "execution_success": expected_success, }, } ), ) @asynccontextmanager async def fake_db_context(_project_id: str): yield db monkeypatch.setattr(service, "get_db_context", fake_db_context) result = await service._record_incident_terminal_disposition( _claim(), apply_op_id="00000000-0000-0000-0000-000000000104", project_id="awoooi", terminal_type=terminal["terminal_type"], success_terminal=False, retry_op_id=None, telegram_receipt_acknowledged=False, ) assert result is not None assert result["proposal_executed"] is proposal_executed assert result["execution_success"] is expected_success assert db.parameters[1]["execution_success"] is expected_success @pytest.mark.asyncio async def test_closure_refuses_to_resolve_when_any_receipt_is_missing( monkeypatch: pytest.MonkeyPatch, ) -> None: atomic_writer = AsyncMock() monkeypatch.setattr( service, "_read_verified_apply_closure_prerequisites", AsyncMock( return_value={ "ready": False, "missing": ["telegram_receipt"], "receipts": {}, } ), ) monkeypatch.setattr( service, "_commit_verified_incident_closure", atomic_writer, ) result = await service._finalize_verified_apply_closure( _claim(), apply_op_id="00000000-0000-0000-0000-000000000104", project_id="awoooi", ) assert result == { "status": "closure_receipts_pending", "closed": False, "missing": ["telegram_receipt"], } atomic_writer.assert_not_awaited() @pytest.mark.asyncio async def test_closure_resolves_only_after_durable_readback( monkeypatch: pytest.MonkeyPatch, ) -> None: monkeypatch.setattr( service, "_read_verified_apply_closure_prerequisites", AsyncMock(return_value={"ready": True, "receipts": {"all": True}}), ) atomic_writer = AsyncMock( return_value={ "automation_run_id": "00000000-0000-0000-0000-000000000102", "apply_op_id": "00000000-0000-0000-0000-000000000104", "incident_resolved": True, "closure_receipt_written": True, "resolved_lifecycle_written": True, } ) monkeypatch.setattr( service, "_commit_verified_incident_closure", atomic_writer, ) monkeypatch.setattr( service, "_read_incident_closure_readback", AsyncMock( return_value={ "closed": True, "missing": [], "incident_status": "RESOLVED", } ), ) result = await service._finalize_verified_apply_closure( _claim(), apply_op_id="00000000-0000-0000-0000-000000000104", project_id="awoooi", ) assert result["status"] == "controlled_apply_closed" assert result["closed"] is True atomic_writer.assert_awaited_once() assert atomic_writer.await_args.kwargs["closure_prerequisite_count"] == 1 @pytest.mark.asyncio async def test_atomic_closure_rolls_back_when_bundle_readback_is_missing( monkeypatch: pytest.MonkeyPatch, ) -> None: db = _SequenceDB(_MappingResult(), _MappingResult(None)) transaction_rolled_back = False @asynccontextmanager async def fake_db_context(_project_id: str): nonlocal transaction_rolled_back try: yield db except Exception: transaction_rolled_back = True raise monkeypatch.setattr(service, "get_db_context", fake_db_context) result = await service._commit_verified_incident_closure( _claim(), apply_op_id="00000000-0000-0000-0000-000000000104", project_id="awoooi", closure_prerequisite_count=12, ) assert result is None assert transaction_rolled_back is True atomic_statement = db.statements[1] assert "INSERT INTO alert_operation_log" in atomic_statement assert "UPDATE automation_operation_log apply" in atomic_statement assert "UPDATE incidents incident" in atomic_statement assert "WITH target_apply AS" in atomic_statement assert "WHEN upper(incident.status::text) = 'CLOSED'" in atomic_statement def test_terminal_candidate_contract_requires_verified_four_node_chain() -> None: predicate = candidate_job.verified_ansible_candidate_terminal_sql( "candidate" ) recorder_source = inspect.getsource( candidate_job.record_ansible_decision_audit ) assert "direct_terminal.status = 'dry_run'" in predicate assert "check_mode.parent_op_id = candidate.op_id" in predicate assert "apply.parent_op_id = check_mode.op_id" in predicate assert "replay.parent_op_id = apply.op_id" in predicate assert "replay.status IN ('success', 'failed')" in predicate assert "verified_no_write_terminal" in predicate assert "repository_readback_verified" in predicate assert "{detail,execution_success}" in predicate assert "{detail,telegram_receipt_acknowledged}" in predicate assert "ORDER BY latest.created_at DESC, latest.op_id DESC" in ( recorder_source ) @pytest.mark.asyncio async def test_projection_replay_sends_missing_receipt_without_reapplying( monkeypatch: pytest.MonkeyPatch, ) -> None: readback = AsyncMock( side_effect=[ {"receipts": {"telegram_receipt": False}}, {"receipts": {"telegram_receipt": True}}, ] ) telegram_sender = AsyncMock(return_value={"ok": True}) monkeypatch.setattr( service, "_finalize_controlled_approval_projection", AsyncMock(return_value=True), ) monkeypatch.setattr( service, "_append_alert_lifecycle_receipt", AsyncMock(return_value=True), ) monkeypatch.setattr( service, "_read_verified_apply_closure_prerequisites", readback, ) monkeypatch.setattr( service, "_send_controlled_apply_telegram_receipt", telegram_sender, ) monkeypatch.setattr( service, "_finalize_verified_apply_closure", AsyncMock(return_value={"status": "controlled_apply_closed", "closed": True}), ) result = await service._reconcile_verified_apply_closure_projections( _claim(), _verified_result(), apply_op_id="00000000-0000-0000-0000-000000000104", writeback={ "verification_passed": True, "verification_result": "success", "verification": True, "learning": True, }, project_id="awoooi", ) assert result["closed"] is True assert result["runtime_apply_executed"] is False telegram_sender.assert_awaited_once() def test_backfill_preserves_approval_and_replaces_terminal_skipped_candidate() -> None: proposal = candidate_job._build_backfill_proposal( { "alertname": "AwoooPAutoRepairCanaryT16", "severity": "medium", "approval_id": "00000000-0000-0000-0000-000000000103", } ) query_source = inspect.getsource( candidate_job._fetch_missing_candidate_incidents ) assert proposal["approval_id"] == "00000000-0000-0000-0000-000000000103" assert "approval_records" in query_source assert "verified_ansible_candidate_terminal_sql" in query_source assert "AND NOT ({terminal_candidate_sql})" in query_source assert "ORDER BY latest.created_at DESC, latest.op_id DESC" in query_source def test_retry_projection_only_repairs_latest_candidate_run() -> None: source = inspect.getsource( service._load_missing_retry_terminal_projection_rows ) assert "source_candidate.op_id = source_check.parent_op_id" in source assert "newer_candidate.created_at," in source assert "source_candidate.created_at," in source assert "repair_candidate_controlled_queue" in source assert "retry_receipt.receipt IS NULL" in source assert "controlled_apply_result" in source assert "LEFT JOIN LATERAL" in source assert "receipt.value ->> 'automation_run_id'" in source assert "upper(incident.status::text)" in source assert "ansible_incident_terminal_disposition_v1" in source assert "'proposal_executed'" in source assert "'execution_success'" in source assert "verification_result IN ('failed', 'timeout')" in source @pytest.mark.asyncio async def test_broker_reconciles_receipts_then_continues_claiming_fresh_work( monkeypatch: pytest.MonkeyPatch, ) -> None: retry_replayer = AsyncMock( return_value={ "scanned": 0, "replayed": 0, "check_mode_passed": 0, "check_mode_failed": 0, "runtime_stage_receipt_written": 0, "blockers": [], "error": None, } ) fresh_claimer = AsyncMock(return_value=[]) monkeypatch.setattr( service, "_expire_stale_ansible_execution_capabilities", AsyncMock(return_value=0), ) monkeypatch.setattr( service, "backfill_missing_retry_terminal_projections_once", AsyncMock( return_value={ "scanned": 0, "written": 0, "incident_receipt_written": 0, "lifecycle_written": 0, "verified": 0, "runtime_apply_executed": False, "error": None, } ), ) monkeypatch.setattr( service, "backfill_missing_auto_repair_execution_receipts_once", AsyncMock( return_value={ "scanned": 1, "written": 1, "incident_closure_written": 1, "telegram_receipt_acknowledged": 1, "error": None, } ), ) monkeypatch.setattr( service, "run_failed_apply_check_mode_replay_once", retry_replayer, ) monkeypatch.setattr(service, "_runtime_blockers", lambda: []) monkeypatch.setattr( service, "recent_ansible_transport_blockers", AsyncMock(return_value=[]), ) monkeypatch.setattr( service, "claim_stale_pending_check_modes", AsyncMock(return_value=[]), ) monkeypatch.setattr( service, "claim_catalog_drift_failed_check_modes", AsyncMock(return_value=[]), ) monkeypatch.setattr(service, "claim_pending_check_modes", fresh_claimer) result = await service.run_pending_check_modes_once(limit=1) assert result["claimed"] == 0 assert result["repair_receipt_backfill_priority_tick"] is True assert result["repair_receipt_closure_written"] == 1 retry_replayer.assert_awaited_once() fresh_claimer.assert_awaited_once() @pytest.mark.asyncio async def test_broker_errors_do_not_starve_unrelated_fresh_claims( monkeypatch: pytest.MonkeyPatch, ) -> None: monkeypatch.setattr( service, "_expire_stale_ansible_execution_capabilities", AsyncMock(return_value=0), ) monkeypatch.setattr( service, "backfill_missing_auto_repair_execution_receipts_once", AsyncMock( return_value={ "scanned": 3, "written": 0, "incident_closure_written": 0, "error": "legacy_projection_write_failed", } ), ) monkeypatch.setattr( service, "run_failed_apply_check_mode_replay_once", AsyncMock( return_value={ "scanned": 1, "replayed": 0, "blockers": ["foreign_retry_already_claimed"], "error": "foreign_retry_reconstruction_failed", } ), ) monkeypatch.setattr(service, "_runtime_blockers", lambda: []) monkeypatch.setattr( service, "recent_ansible_transport_blockers", AsyncMock(return_value=[]), ) monkeypatch.setattr( service, "claim_stale_pending_check_modes", AsyncMock(return_value=[]), ) monkeypatch.setattr( service, "claim_catalog_drift_failed_check_modes", AsyncMock(return_value=[]), ) fresh_claimer = AsyncMock(return_value=[]) monkeypatch.setattr(service, "claim_pending_check_modes", fresh_claimer) result = await service.run_pending_check_modes_once(limit=1) assert result["repair_receipt_backfill_error"] == ( "legacy_projection_write_failed" ) assert result["failed_apply_retry_error"] == ( "foreign_retry_reconstruction_failed" ) assert result["failed_apply_retry_blockers"] == [ "foreign_retry_already_claimed" ] fresh_claimer.assert_awaited_once() @pytest.mark.asyncio async def test_retry_terminal_projection_backfill_writes_and_verifies_no_apply( monkeypatch: pytest.MonkeyPatch, ) -> None: row = { "op_id": "00000000-0000-0000-0000-000000000104", "replay_op_id": "00000000-0000-0000-0000-000000000105", "replay_status": "success", "replay_output": {"returncode": 0}, "replay_dry_run_result": {"timed_out": False}, "replay_duration_ms": 18, "retry_receipt_present": False, "terminal_type": "no_write_replay_passed_waiting_repair_candidate", "telegram_receipt_acknowledged": False, } monkeypatch.setattr( service, "_load_missing_retry_terminal_projection_rows", AsyncMock(return_value=[row]), ) monkeypatch.setattr( service, "_claim_from_apply_operation_row", lambda _row: ( _claim(), service.AnsibleRunResult( returncode=1, stdout="", stderr="", duration_ms=0, ), ), ) terminal_writer = AsyncMock( return_value={ "automation_run_id": "00000000-0000-0000-0000-000000000102", "apply_op_id": row["op_id"], "retry_op_id": row["replay_op_id"], "terminal_type": row["terminal_type"], "repository_readback_verified": True, } ) receipt_writer = AsyncMock(return_value=True) lifecycle_writer = AsyncMock(return_value=True) verifier = AsyncMock(return_value=True) telegram_sender = AsyncMock(return_value=True) telegram_readback = AsyncMock(return_value=True) monkeypatch.setattr( service, "_record_incident_terminal_disposition", terminal_writer ) monkeypatch.setattr( service, "_append_runtime_stage_receipts_to_apply", receipt_writer ) monkeypatch.setattr( service, "_append_alert_lifecycle_receipt", lifecycle_writer ) monkeypatch.setattr( service, "_verify_retry_terminal_projection_readback", verifier ) monkeypatch.setattr( service, "_send_controlled_apply_telegram_receipt", telegram_sender, ) monkeypatch.setattr( service, "_retry_telegram_receipt_acknowledged", telegram_readback, ) result = await service.backfill_missing_retry_terminal_projections_once() assert result == { "scanned": 1, "written": 1, "retry_receipt_written": 1, "telegram_receipt_acknowledged": 1, "incident_receipt_written": 1, "lifecycle_written": 1, "verified": 1, "runtime_apply_executed": False, "error": None, } terminal_writer.assert_awaited_once() assert terminal_writer.await_args.kwargs["success_terminal"] is False retry_receipt = receipt_writer.await_args_list[0].kwargs["receipts"][0] assert retry_receipt["stage_id"] == "retry_or_rollback" assert retry_receipt["derived_from_durable_chain"] is True receipt = receipt_writer.await_args_list[1].kwargs["receipts"][0] assert receipt["stage_id"] == "incident_closure" assert receipt["derived_from_durable_chain"] is True assert lifecycle_writer.await_args.args[1] == "EXECUTION_COMPLETED" assert lifecycle_writer.await_args.kwargs["success"] is False telegram_sender.assert_awaited_once() telegram_readback.assert_awaited_once() @pytest.mark.asyncio async def test_broker_prioritizes_retry_terminal_projection_backfill( monkeypatch: pytest.MonkeyPatch, ) -> None: monkeypatch.setattr( service, "_expire_stale_ansible_execution_capabilities", AsyncMock(return_value=0), ) receipt_backfill = AsyncMock( return_value={ "scanned": 1, "written": 0, "incident_closure_written": 1, "telegram_receipt_acknowledged": 0, "retry_terminal_projection_scanned": 1, "retry_terminal_projection_written": 1, "retry_terminal_projection_verified": 1, "retry_terminal_projection_runtime_apply_executed": False, "error": None, } ) monkeypatch.setattr( service, "backfill_missing_auto_repair_execution_receipts_once", receipt_backfill, ) monkeypatch.setattr( service, "run_failed_apply_check_mode_replay_once", AsyncMock(return_value={"scanned": 0, "replayed": 0}), ) monkeypatch.setattr(service, "_runtime_blockers", lambda: []) monkeypatch.setattr( service, "recent_ansible_transport_blockers", AsyncMock(return_value=[]), ) monkeypatch.setattr( service, "claim_stale_pending_check_modes", AsyncMock(return_value=[]), ) monkeypatch.setattr( service, "claim_catalog_drift_failed_check_modes", AsyncMock(return_value=[]), ) fresh_claimer = AsyncMock(return_value=[]) monkeypatch.setattr(service, "claim_pending_check_modes", fresh_claimer) result = await service.run_pending_check_modes_once(limit=1) assert result["claimed"] == 0 assert result["retry_terminal_projection_verified"] == 1 assert result["retry_terminal_projection_runtime_apply_executed"] is False assert result["repair_receipt_backfill_priority_tick"] is True receipt_backfill.assert_awaited_once() fresh_claimer.assert_awaited_once() @pytest.mark.asyncio async def test_signal_worker_retry_preflight_is_query_only( monkeypatch: pytest.MonkeyPatch, ) -> None: monkeypatch.setattr( service, "_load_open_failed_apply_retry_row", AsyncMock(return_value={"op_id": "apply-op"}), ) result = await service.preflight_failed_apply_retry_queue_once( project_id="awoooi", window_hours=24, ) assert result["scanned"] == 1 assert result["replayed"] == 0 assert result["query_only"] is True assert result["runtime_apply_executed"] is False assert result["execution_owner"] == "awoooi-ansible-executor-broker" @pytest.mark.asyncio @pytest.mark.parametrize( ("recovery_count", "expected_replayed", "expected_terminalized"), [(0, 1, 0), (1, 0, 1)], ) async def test_stale_pending_retry_is_reclaimed_without_runtime_apply( monkeypatch: pytest.MonkeyPatch, recovery_count: int, expected_replayed: int, expected_terminalized: int, ) -> None: replay_op_id = "00000000-0000-0000-0000-000000000105" apply_op_id = "00000000-0000-0000-0000-000000000104" db = _SequenceDB( _MappingResult(), _MappingResult(scalar=replay_op_id), _MappingResult(scalar=replay_op_id), ) @asynccontextmanager async def fake_db_context(_project_id: str): yield db monkeypatch.setattr(service, "get_db_context", fake_db_context) monkeypatch.setattr( service, "_load_open_failed_apply_retry_row", AsyncMock( return_value={ "op_id": apply_op_id, "stale_replay_op_id": replay_op_id, "stale_replay_recovery_count": recovery_count, } ), ) monkeypatch.setattr( service, "_claim_from_apply_operation_row", lambda _row: ( _claim(), service.AnsibleRunResult( returncode=1, stdout="", stderr="failed apply", duration_ms=10, ), ), ) runtime_blockers = ( ["ansible_binary_missing"] if expected_terminalized else [] ) monkeypatch.setattr( service, "_runtime_blockers", lambda: runtime_blockers, ) transport_probe = AsyncMock( return_value=( ["ansible_transport_unavailable"] if expected_terminalized else [] ) ) monkeypatch.setattr( service, "recent_ansible_transport_blockers", transport_probe, ) monkeypatch.setattr( service, "build_ansible_check_mode_command", lambda **_kwargs: object(), ) replay_runner = AsyncMock( return_value=service.AnsibleRunResult( returncode=0, stdout="check ok", stderr="", duration_ms=15, ) ) monkeypatch.setattr(service, "_run_ansible_command", replay_runner) receipt_writer = AsyncMock(return_value=True) monkeypatch.setattr( service, "_record_retry_runtime_stage_receipt", receipt_writer, ) result = await service.run_failed_apply_check_mode_replay_once( timeout_seconds=30, ) assert result["replayed"] == expected_replayed if expected_terminalized: assert result["check_mode_passed"] == 0 assert result["check_mode_failed"] == 1 else: assert result["check_mode_passed"] == 1 assert result["check_mode_failed"] == 0 assert result["stale_retry_terminalized"] == expected_terminalized assert result["runtime_apply_executed"] == 0 assert "pg_advisory_xact_lock" in db.statements[0] assert "UPDATE automation_operation_log replay" in db.statements[1] assert "reclaimed_after_stale_pending" in db.statements[1] assert "input ->> 'retry_claim_id'" in db.statements[2] persisted_output = json.loads(db.parameters[2]["output"]) persisted_dry_run = json.loads(db.parameters[2]["dry_run_result"]) assert persisted_output["check_mode_executed"] is ( not bool(expected_terminalized) ) assert persisted_output["check_mode_replay_performed"] is ( not bool(expected_terminalized) ) assert persisted_output["runtime_apply_executed"] is False assert persisted_dry_run["check_mode_executed"] is ( not bool(expected_terminalized) ) assert persisted_dry_run["check_mode_replay_performed"] is ( not bool(expected_terminalized) ) assert persisted_dry_run["runtime_apply_executed"] is False if expected_terminalized: replay_runner.assert_not_awaited() transport_probe.assert_not_awaited() else: replay_runner.assert_awaited_once() transport_probe.assert_awaited_once() receipt_writer.assert_awaited_once() assert receipt_writer.await_args.kwargs[ "check_mode_replay_performed" ] is (not bool(expected_terminalized)) def test_candidate_worker_defaults_to_query_only_retry_preflight() -> None: defaults = candidate_job.enqueue_missing_ansible_candidates_once.__kwdefaults__ assert defaults is not None assert defaults["retry_replayer"] is service.preflight_failed_apply_retry_queue_once