Files
awoooi/apps/api/tests/test_agent99_controlled_dispatch_reconciler_job.py
2026-07-19 01:08:37 +08:00

568 lines
19 KiB
Python

from __future__ import annotations
from pathlib import Path
from types import SimpleNamespace
import pytest
from src.jobs import agent99_controlled_dispatch_reconciler_job as job
from src.services.agent99_controlled_dispatch_ledger import (
build_agent99_dispatch_identity,
)
class _Ledger:
def __init__(self, receipt: dict[str, object]) -> None:
self.receipt = receipt
self.verifier_calls: list[dict[str, object]] = []
async def list_reconcilable(self, **_kwargs): # type: ignore[no-untyped-def]
return [{"identity": IDENTITY, "receipt": self.receipt}]
async def record_verifier(self, **kwargs): # type: ignore[no-untyped-def]
self.verifier_calls.append(kwargs)
return {
**self.receipt,
"post_verifier_passed": True,
"receipt_persisted": True,
"verifier": {
"status": "success",
"evidence_refs": kwargs["evidence_refs"],
},
}
IDENTITY = build_agent99_dispatch_identity(
project_id="awoooi",
incident_id="INC-20260711-11C751",
source_fingerprint="fp-cold-start",
route_id="agent99:host_recovery:Recover",
)
def test_agent99_playbook_identity_is_project_scoped() -> None:
other = build_agent99_dispatch_identity(
project_id="other-project",
incident_id=IDENTITY.incident_id,
source_fingerprint=IDENTITY.source_fingerprint,
route_id=IDENTITY.route_id,
)
assert job._playbook_id(IDENTITY) != job._playbook_id(other)
class _Context:
def __init__(self, db) -> None: # type: ignore[no-untyped-def]
self.db = db
async def __aenter__(self): # type: ignore[no-untyped-def]
return self.db
async def __aexit__(self, *_args): # type: ignore[no-untyped-def]
return False
class _Result:
def __init__(self, *, row=None, scalar=None) -> None: # type: ignore[no-untyped-def]
self.row = row
self.scalar_value = scalar
def mappings(self): # type: ignore[no-untyped-def]
return self
def first(self): # type: ignore[no-untyped-def]
return self.row
def scalar_one_or_none(self): # type: ignore[no-untyped-def]
return self.scalar_value
@pytest.mark.asyncio
async def test_agent99_rag_ack_requires_durable_embedding_readback(
monkeypatch,
) -> None:
class RagDB:
async def execute(self, _statement, _params): # type: ignore[no-untyped-def]
return _Result(row={
"entry_id": "00000000-0000-0000-0000-000000000099",
"embedding_persisted": True,
})
monkeypatch.setattr(
job,
"get_db_context",
lambda _project_id: _Context(RagDB()),
)
receipt = await job._ensure_rag_writeback(IDENTITY)
assert receipt == (
"knowledge_entries:00000000-0000-0000-0000-000000000099:"
"rag_embedding_verified"
)
@pytest.mark.asyncio
async def test_agent99_mcp_ack_is_same_run_durable_and_idempotent(
monkeypatch,
) -> None:
class McpDB:
def __init__(self) -> None:
self.row = None
async def execute(self, statement): # type: ignore[no-untyped-def]
if self.row is None:
values = statement.compile().params
self.row = SimpleNamespace(**values)
return _Result()
return _Result(scalar=self.row)
db = McpDB()
monkeypatch.setattr(
job,
"get_db_context",
lambda _project_id: _Context(db),
)
refs = {
"telegram_lifecycle_receipt_id": "telegram:1",
"km_writeback_ack_id": "km:1",
"rag_writeback_ack_id": "rag:1",
"playbook_trust_writeback_ack_id": "playbook:1",
}
first = await job._ensure_mcp_evidence_writeback(
IDENTITY,
receipt_refs=refs,
)
db.row = None
second = await job._ensure_mcp_evidence_writeback(
IDENTITY,
receipt_refs=refs,
)
expected = (
f"awooop_mcp_gateway_audit:{job._mcp_learning_call_id(IDENTITY)}:verified"
)
assert first == expected
assert second == expected
@pytest.mark.asyncio
async def test_backup_mcp_ack_requires_dr_scorecard_binding() -> None:
identity = build_agent99_dispatch_identity(
project_id="awoooi",
incident_id="INC-20260719-BACKUP-MCP",
source_fingerprint="backup-mcp-fingerprint",
route_id="agent99:backup_health:BackupCheck",
)
receipt = await job._ensure_mcp_evidence_writeback(
identity,
receipt_refs={
"telegram_lifecycle_receipt_id": "telegram:1",
"km_writeback_ack_id": "km:1",
"rag_writeback_ack_id": "rag:1",
"playbook_trust_writeback_ack_id": "playbook:1",
},
)
assert receipt is None
@pytest.mark.asyncio
async def test_backup_telegram_ack_is_no_write_lifecycle_not_apply(
monkeypatch,
) -> None:
identity = build_agent99_dispatch_identity(
project_id="awoooi",
incident_id="INC-20260719-BACKUP",
source_fingerprint="backup-fingerprint-1",
route_id="agent99:backup_health:BackupCheck",
)
delivered = []
async def deliver(payload): # type: ignore[no-untyped-def]
delivered.append(payload)
return {
"ok": True,
"provider_message_id": "9919",
"durable_outbound_acknowledged": True,
"destination_binding_verified": True,
}
class ForbiddenGateway:
async def send_controlled_apply_result_receipt(self, **_kwargs):
raise AssertionError("BackupCheck must not use controlled-apply receipt")
monkeypatch.setattr(job, "deliver_agent99_telegram_lifecycle", deliver)
monkeypatch.setattr(job, "get_telegram_gateway", lambda: ForbiddenGateway())
receipt = await job._ensure_telegram_receipt(
identity,
mode="BackupCheck",
)
assert receipt == "telegram_outbound:9919:durable_ack"
assert len(delivered) == 1
payload = delivered[0]
assert payload.run_id == str(identity.run_id)
assert payload.trace_id == identity.trace_id
assert payload.work_item_id == identity.work_item_id
assert payload.lifecycle == "verifying"
assert payload.executor_name == "Agent99 BackupCheck (read-only)"
assert payload.verifier_name == "backup_restore_readback_verifier"
assert payload.apply_status == "not_applicable"
assert payload.verifier_status == "passed"
assert payload.closure_status == "pending"
assert "沒有執行備份" in payload.impact
assert "LLM 未取得 runtime authority" in payload.action
@pytest.mark.asyncio
async def test_backup_telegram_ack_requires_durable_destination_proof(
monkeypatch,
) -> None:
identity = build_agent99_dispatch_identity(
project_id="awoooi",
incident_id="INC-20260719-BACKUP-FAIL",
source_fingerprint="backup-fingerprint-2",
route_id="agent99:backup_health:BackupCheck",
)
async def deliver(_payload): # type: ignore[no-untyped-def]
return {
"ok": False,
"provider_message_id": "9920",
"durable_outbound_acknowledged": True,
"destination_binding_verified": False,
}
monkeypatch.setattr(job, "deliver_agent99_telegram_lifecycle", deliver)
receipt = await job._ensure_telegram_receipt(
identity,
mode="BackupCheck",
)
assert receipt is None
@pytest.mark.asyncio
async def test_reconciler_consumes_authenticated_outcome_and_closes_same_run(
monkeypatch,
) -> None:
ledger = _Ledger({
"status": "dispatch_accepted_verifier_pending",
"post_verifier_passed": False,
"dispatch_receipt": {"kind": "host_recovery"},
})
finalized: list[dict[str, object]] = []
monkeypatch.setattr(job, "get_agent99_dispatch_ledger", lambda: ledger)
monkeypatch.setattr(
job,
"read_agent99_sre_outcome",
lambda _run_id: {
"automationRunId": str(IDENTITY.run_id),
"traceId": IDENTITY.trace_id,
"workItemId": IDENTITY.work_item_id,
"mode": "Recover",
"controlledApply": True,
"outcome": {
"schemaVersion": "agent99_outcome_contract_v1",
"state": "resolved",
"transportOk": True,
"verifierName": "agent99-recover-post-condition-v1",
"verifierPassed": True,
"sourceEventResolved": True,
"verifiedAt": "2026-07-11T20:00:00+08:00",
},
},
)
async def finalize(identity, *, mode, evidence_refs): # type: ignore[no-untyped-def]
finalized.append({
"identity": identity,
"mode": mode,
"evidence_refs": evidence_refs,
})
return {"runtime_closure_verified": True}
monkeypatch.setattr(job, "_finalize_learning", finalize)
result = await job.reconcile_agent99_controlled_dispatches_once()
assert result == {
"scanned": 1,
"source_not_resolved": 0,
"status_reconcile_requested": 0,
"status_reconcile_pending": 0,
"external_no_write_reconciled": 0,
"outcome_timeout_terminalized": 0,
"verifier_written": 1,
"closed": 1,
}
assert len(ledger.verifier_calls) == 1
refs = ledger.verifier_calls[0]["evidence_refs"]
assert refs["agent99_outcome_receipt_id"].endswith(str(IDENTITY.run_id))
assert "post_verifier_evidence_ref" in refs
assert "source_event_evidence_ref" in refs
assert finalized[0]["identity"].run_id == IDENTITY.run_id
assert finalized[0]["mode"] == "Recover"
@pytest.mark.asyncio
async def test_reconciler_does_not_finalize_without_outcome(monkeypatch) -> None:
ledger = _Ledger({
"status": "dispatch_delivery_unknown_reconcile_only",
"post_verifier_passed": False,
})
monkeypatch.setattr(job, "get_agent99_dispatch_ledger", lambda: ledger)
monkeypatch.setattr(job, "read_agent99_sre_outcome", lambda _run_id: None)
result = await job.reconcile_agent99_controlled_dispatches_once()
assert result == {
"scanned": 1,
"source_not_resolved": 0,
"status_reconcile_requested": 0,
"status_reconcile_pending": 0,
"external_no_write_reconciled": 0,
"outcome_timeout_terminalized": 0,
"verifier_written": 0,
"closed": 0,
}
assert ledger.verifier_calls == []
def test_windows_relay_outcome_readback_is_authenticated_and_public_safe() -> None:
source = Path(__file__).resolve().parents[3] / "agent99-sre-alert-relay.ps1"
text = source.read_text(encoding="utf-8")
assert '$context.Request.HttpMethod -eq "GET"' in text
assert "relay_auth_not_configured" in text
assert 'Headers["X-Agent99-Relay-Token"]' in text
token_gate = text.index("if (-not $expectedToken)")
get_branch = text.index('if ($context.Request.HttpMethod -eq "GET")')
post_branch = text.index('if ($context.Request.HttpMethod -ne "POST")')
assert token_gate < get_branch < post_branch
assert "storesRawEvidence = $false" in text
assert "stdout" not in text[text.index("function Get-AgentOutcomeReadback"):text.index("function Write-AgentRelayEvent")]
@pytest.mark.asyncio
async def test_exact_11c751_legacy_failure_backfills_once_without_recurrence() -> None:
row = {
"incident_id": "INC-20260711-11C751",
"project_id": "awoooi",
"alertname": "ColdStartGateBlocked",
"severity": "low",
"alert_category": "general",
"notification_type": "legacy_backfill",
"approval_id": "00000000-0000-0000-0000-000000000751",
"approval_fingerprint": "legacy-source-fingerprint-11c751",
"approval_status": "execution_failed",
}
calls: list[dict[str, object]] = []
dispatched_keys: set[tuple[str, str, str]] = set()
fetch_count = 0
async def fetcher(**_kwargs): # type: ignore[no-untyped-def]
nonlocal fetch_count
fetch_count += 1
if fetch_count == 1:
return [row]
return [{
**row,
"latest_approval_id": "00000000-0000-0000-0000-000000000752",
"latest_approval_fingerprint": "newer-approval-fingerprint-11c751",
}]
async def dispatcher(**kwargs): # type: ignore[no-untyped-def]
calls.append(kwargs)
key = (
str(kwargs["incident_id"]),
str(kwargs["fingerprint"]),
str(kwargs["route_id"]),
)
if key in dispatched_keys:
return {"status": "deduplicated", "dispatchPerformed": False}
dispatched_keys.add(key)
return {"status": "dispatched", "dispatchPerformed": True}
first = await job.backfill_legacy_agent99_dispatches_once(
fetcher=fetcher,
dispatcher=dispatcher,
)
second = await job.backfill_legacy_agent99_dispatches_once(
fetcher=fetcher,
dispatcher=dispatcher,
)
assert first == {
"scanned": 1,
"dispatch_performed": 1,
"deduplicated": 0,
"retryable": 0,
"blocked": 0,
}
assert second == {
"scanned": 1,
"dispatch_performed": 0,
"deduplicated": 1,
"retryable": 0,
"blocked": 0,
}
assert len(dispatched_keys) == 1
assert calls[0]["incident_id"] == "INC-20260711-11C751"
assert calls[0]["namespace"] == "host_recovery"
assert calls[0]["target_resource"] == "cold-start-gate"
assert calls[0]["route_id"] == "agent99:host_recovery:Recover"
assert calls[0]["fingerprint"].startswith("legacy-cold-start:")
assert calls[0]["fingerprint"] == calls[1]["fingerprint"]
assert calls[0]["approval_id"] == calls[1]["approval_id"]
assert calls[0]["labels"]["legacy_approval_fingerprint"] == (
calls[1]["labels"]["legacy_approval_fingerprint"]
)
assert calls[0]["labels"]["legacy_executor"] == "ansible"
assert calls[0]["labels"]["legacy_status"] == "execution_failed"
@pytest.mark.asyncio
async def test_learning_checkpoint_resumes_without_replaying_assets(
monkeypatch,
) -> None:
identity = build_agent99_dispatch_identity(
project_id="awoooi",
incident_id="INC-20260711-BACKUP",
source_fingerprint="backup-source-1",
route_id="agent99:backup_health:BackupCheck",
approval_id="00000000-0000-0000-0000-000000000751",
)
class CheckpointLedger:
def __init__(self) -> None:
self.refs: dict[str, str] = {}
async def record_learning_asset_acks(
self,
*,
identity, # type: ignore[no-untyped-def]
receipt_refs,
) -> dict[str, object]:
self.refs.update(receipt_refs)
return {
"status": "learning_asset_acks_recorded",
"receipt_persisted": True,
"learning_receipt_refs": dict(self.refs),
"runtime_closure_verified": False,
}
ledger = CheckpointLedger()
effects = {
"playbook": 0,
"km": 0,
"rag": 0,
"telegram": 0,
"mcp": 0,
"dr": 0,
}
effect_order: list[str] = []
terminal_calls = 0
async def write_playbook(*_args, **_kwargs): # type: ignore[no-untyped-def]
effects["playbook"] += 1
effect_order.append("playbook")
return "playbooks:PB-A99:version:1"
async def write_km(*_args, **_kwargs): # type: ignore[no-untyped-def]
effects["km"] += 1
effect_order.append("km")
return "knowledge_entries:km-1:embedding_verified"
async def write_telegram(*_args, **_kwargs): # type: ignore[no-untyped-def]
effects["telegram"] += 1
effect_order.append("telegram")
return "telegram_outbound:1:durable_ack"
async def write_rag(*_args, **_kwargs): # type: ignore[no-untyped-def]
effects["rag"] += 1
effect_order.append("rag")
return "knowledge_entries:km-1:rag_embedding_verified"
async def write_mcp(*_args, **_kwargs): # type: ignore[no-untyped-def]
effects["mcp"] += 1
effect_order.append("mcp")
assert _kwargs["receipt_refs"]["dr_scorecard_writeback_ack_id"] == (
"knowledge_entries:dr-1:dr_scorecard"
)
return "awooop_mcp_gateway_audit:call-1:verified"
async def write_dr(*_args, **_kwargs): # type: ignore[no-untyped-def]
effects["dr"] += 1
effect_order.append("dr")
return "knowledge_entries:dr-1:dr_scorecard"
async def terminal(**_kwargs): # type: ignore[no-untyped-def]
nonlocal terminal_calls
terminal_calls += 1
if terminal_calls == 1:
return {
"status": "learning_writeback_persistence_failed",
"runtime_closure_verified": False,
}
return {
"status": "closed_verified_learning_written",
"runtime_closure_verified": True,
}
monkeypatch.setattr(job, "get_agent99_dispatch_ledger", lambda: ledger)
monkeypatch.setattr(job, "_ensure_playbook_trust", write_playbook)
monkeypatch.setattr(job, "_ensure_km_writeback", write_km)
monkeypatch.setattr(job, "_ensure_rag_writeback", write_rag)
monkeypatch.setattr(job, "_ensure_telegram_receipt", write_telegram)
monkeypatch.setattr(job, "_ensure_mcp_evidence_writeback", write_mcp)
monkeypatch.setattr(job, "_ensure_dr_scorecard", write_dr)
monkeypatch.setattr(job, "record_agent99_learning_writeback", terminal)
first = await job._finalize_learning(
identity,
mode="BackupCheck",
evidence_refs={},
)
second = await job._finalize_learning(
identity,
mode="BackupCheck",
evidence_refs={},
)
assert first["runtime_closure_verified"] is False
assert second["runtime_closure_verified"] is True
assert effects == {
"playbook": 1,
"km": 1,
"rag": 1,
"telegram": 1,
"mcp": 1,
"dr": 1,
}
assert effect_order == ["playbook", "km", "rag", "telegram", "dr", "mcp"]
assert terminal_calls == 2
def test_legacy_backfill_query_is_bounded_and_requires_failure_evidence() -> None:
source = Path(job.__file__).read_text(encoding="utf-8")
fetch = source[
source.index("async def _fetch_legacy_cold_start_failures"):source.index(
"def _legacy_source_fingerprint"
)
]
assert 'LEGACY_COLD_START_INCIDENT_ID = "INC-20260711-11C751"' in source
assert "lower(approval_state.status::text) = 'execution_failed'" in fetch
assert "ORDER BY candidate.created_at ASC, candidate.id ASC" in fetch
assert "legacy.input ->> 'executor' = 'ansible'" in fetch
assert "LIMIT :limit" in fetch
assert "statement_timeout = '5000ms'" in fetch
assert "backfill_legacy_agent99_dispatches_once()" in source