From 17e60dd36f709af4504014ddf3a28f398d9476d2 Mon Sep 17 00:00:00 2001 From: Your Name Date: Sun, 19 Jul 2026 00:24:44 +0800 Subject: [PATCH] feat(sre): require Agent99 RAG and MCP closure receipts --- ...nt99_controlled_dispatch_reconciler_job.py | 149 ++++++++++++++++++ .../agent99_controlled_dispatch_ledger.py | 25 +-- .../src/services/agent99_public_receipts.py | 2 + ...t99_controlled_dispatch_ledger_contract.py | 4 + ...nt99_controlled_dispatch_reconciler_job.py | 125 ++++++++++++++- 5 files changed, 294 insertions(+), 11 deletions(-) diff --git a/apps/api/src/jobs/agent99_controlled_dispatch_reconciler_job.py b/apps/api/src/jobs/agent99_controlled_dispatch_reconciler_job.py index 728d64997..2cb73c048 100644 --- a/apps/api/src/jobs/agent99_controlled_dispatch_reconciler_job.py +++ b/apps/api/src/jobs/agent99_controlled_dispatch_reconciler_job.py @@ -4,13 +4,17 @@ from __future__ import annotations import asyncio import hashlib +import json from collections.abc import Awaitable, Callable from datetime import UTC, datetime from typing import Any +from uuid import NAMESPACE_URL, uuid5 import structlog from sqlalchemy import select, text +from sqlalchemy.dialects.postgresql import insert as pg_insert +from src.db.awooop_models import AwoooPMcpGatewayAudit from src.db.base import get_db_context from src.db.models import PlaybookRecord from src.models.knowledge import ( @@ -438,6 +442,140 @@ async def _ensure_km_writeback( return None +async def _ensure_rag_writeback( + identity: Agent99DispatchIdentity, +) -> str | None: + """Read back the durable embedding created by the Agent99 KM writer.""" + + try: + async with get_db_context(identity.project_id) as db: + result = await db.execute( + text(""" + SELECT id::text AS entry_id, + embedding IS NOT NULL AS embedding_persisted + FROM knowledge_entries + WHERE project_id = :project_id + AND related_incident_id = :incident_id + AND path_type = :path_type + ORDER BY created_at DESC + LIMIT 1 + """), + { + "project_id": identity.project_id, + "incident_id": identity.incident_id, + "path_type": f"agent99:{identity.run_id}", + }, + ) + row = result.mappings().first() + if not row or row.get("embedding_persisted") is not True: + return None + return f"knowledge_entries:{row['entry_id']}:rag_embedding_verified" + except Exception as exc: + logger.warning( + "agent99_rag_writeback_readback_failed", + run_id=str(identity.run_id), + error=str(exc), + ) + return None + + +def _mcp_learning_call_id(identity: Agent99DispatchIdentity): + return uuid5( + NAMESPACE_URL, + f"agent99-mcp-learning:{identity.project_id}:{identity.run_id}", + ) + + +async def _ensure_mcp_evidence_writeback( + identity: Agent99DispatchIdentity, + *, + receipt_refs: dict[str, str], +) -> str | None: + """Publish and read back a public-safe same-run MCP evidence binding.""" + + required = { + "telegram_lifecycle_receipt_id", + "km_writeback_ack_id", + "rag_writeback_ack_id", + "playbook_trust_writeback_ack_id", + } + safe_refs = sanitize_agent99_public_receipt_refs( + receipt_refs, + allowed_keys=required, + ) + if not required.issubset(safe_refs): + return None + bindings = {key: safe_refs[key] for key in sorted(required)} + identity_hash = hashlib.sha256( + json.dumps( + identity.public_dict(), + sort_keys=True, + separators=(",", ":"), + ).encode() + ).hexdigest() + bindings_hash = hashlib.sha256( + json.dumps(bindings, sort_keys=True, separators=(",", ":")).encode() + ).hexdigest() + call_id = _mcp_learning_call_id(identity) + gate_result = { + "schema_version": "agent99_mcp_learning_evidence_v1", + "incident_id": identity.incident_id, + "route_id": identity.route_id, + "receipt_refs": bindings, + "raw_evidence_stored": False, + "secret_value_stored": False, + } + try: + async with get_db_context(identity.project_id) as db: + await db.execute( + pg_insert(AwoooPMcpGatewayAudit) + .values( + call_id=call_id, + project_id=identity.project_id, + run_id=identity.run_id, + trace_id=identity.trace_id, + agent_id="agent99_controlled_dispatch_reconciler", + tool_name="agent99_learning_evidence_publish", + input_hash=identity_hash, + output_hash=bindings_hash, + gate_result=gate_result, + result_status="success", + latency_ms=0, + ) + .on_conflict_do_nothing( + index_elements=[AwoooPMcpGatewayAudit.call_id] + ) + ) + selected = await db.execute( + select(AwoooPMcpGatewayAudit).where( + AwoooPMcpGatewayAudit.call_id == call_id, + AwoooPMcpGatewayAudit.project_id == identity.project_id, + AwoooPMcpGatewayAudit.run_id == identity.run_id, + ) + ) + row = selected.scalar_one_or_none() + if not ( + row + and str(row.trace_id or "") == identity.trace_id + and row.agent_id == "agent99_controlled_dispatch_reconciler" + and row.tool_name == "agent99_learning_evidence_publish" + and row.input_hash == identity_hash + and row.output_hash == bindings_hash + and row.result_status == "success" + and isinstance(row.gate_result, dict) + and row.gate_result == gate_result + ): + return None + return f"awooop_mcp_gateway_audit:{call_id}:verified" + except Exception as exc: + logger.warning( + "agent99_mcp_learning_evidence_writeback_failed", + run_id=str(identity.run_id), + error=str(exc), + ) + return None + + async def _ensure_telegram_receipt( identity: Agent99DispatchIdentity, *, @@ -575,10 +713,21 @@ async def _finalize_learning( "km_writeback_ack_id", lambda: _ensure_km_writeback(identity, mode=mode), ), + ( + "rag_writeback_ack_id", + lambda: _ensure_rag_writeback(identity), + ), ( "telegram_lifecycle_receipt_id", lambda: _ensure_telegram_receipt(identity, mode=mode), ), + ( + "mcp_evidence_writeback_ack_id", + lambda: _ensure_mcp_evidence_writeback( + identity, + receipt_refs=receipt_refs, + ), + ), ] if "backup_health" in identity.route_id: assets.append(( diff --git a/apps/api/src/services/agent99_controlled_dispatch_ledger.py b/apps/api/src/services/agent99_controlled_dispatch_ledger.py index c4fe253af..e59339cb6 100644 --- a/apps/api/src/services/agent99_controlled_dispatch_ledger.py +++ b/apps/api/src/services/agent99_controlled_dispatch_ledger.py @@ -498,6 +498,8 @@ def build_agent99_dispatch_receipt_envelope( "incident_closure_receipt_id", "telegram_lifecycle_receipt_id", "km_writeback_ack_id", + "rag_writeback_ack_id", + "mcp_evidence_writeback_ack_id", "playbook_trust_writeback_ack_id", ], }, @@ -1406,6 +1408,8 @@ class PostgresAgent99DispatchLedger: "incident_closure_receipt_id", "telegram_lifecycle_receipt_id", "km_writeback_ack_id", + "rag_writeback_ack_id", + "mcp_evidence_writeback_ack_id", "playbook_trust_writeback_ack_id", ], }, @@ -2068,15 +2072,17 @@ class PostgresAgent99DispatchLedger: ) -> dict[str, Any]: """Atomically close the incident, operation log, and same-run ledger. - KM, PlayBook, DR, and Telegram each provide a durable idempotent - acknowledgement first. The incident terminal state, immutable closure - event, run state, and step-3 receipt then commit in one PostgreSQL - transaction. Redis is only a post-commit projection. + KM, RAG, MCP evidence, PlayBook, DR, and Telegram each provide a + durable idempotent acknowledgement first. The incident terminal state, + immutable closure event, run state, and step-3 receipt then commit in + one PostgreSQL transaction. Redis is only a post-commit projection. """ required_refs = { "telegram_lifecycle_receipt_id", "km_writeback_ack_id", + "rag_writeback_ack_id", + "mcp_evidence_writeback_ack_id", "playbook_trust_writeback_ack_id", } if "backup_health" in identity.route_id: @@ -2354,11 +2360,12 @@ class PostgresAgent99DispatchLedger: ) -> dict[str, Any]: """Checkpoint per-run learning side effects before terminal closure. - Reconciliation may crash after KM, PlayBook, Telegram, or DR succeeds. - Persisting each public acknowledgement on the same run lets the next - tick resume at the first missing asset instead of replaying completed - side effects. The asset writers retain their own deterministic keys as - the final crash fence between side effect and this checkpoint. + Reconciliation may crash after KM, RAG, MCP evidence, PlayBook, + Telegram, or DR succeeds. Persisting each public acknowledgement on + the same run lets the next tick resume at the first missing asset + instead of replaying completed side effects. The asset writers retain + their own deterministic keys as the final crash fence between side + effect and this checkpoint. """ supplied = sanitize_agent99_public_receipt_refs( diff --git a/apps/api/src/services/agent99_public_receipts.py b/apps/api/src/services/agent99_public_receipts.py index 5851358fd..95889607f 100644 --- a/apps/api/src/services/agent99_public_receipts.py +++ b/apps/api/src/services/agent99_public_receipts.py @@ -29,6 +29,8 @@ AGENT99_LEARNING_RECEIPT_REFS = frozenset({ "incident_closure_receipt_id", "telegram_lifecycle_receipt_id", "km_writeback_ack_id", + "rag_writeback_ack_id", + "mcp_evidence_writeback_ack_id", "playbook_trust_writeback_ack_id", "dr_scorecard_writeback_ack_id", }) diff --git a/apps/api/tests/test_agent99_controlled_dispatch_ledger_contract.py b/apps/api/tests/test_agent99_controlled_dispatch_ledger_contract.py index c08a40afc..96df044b2 100644 --- a/apps/api/tests/test_agent99_controlled_dispatch_ledger_contract.py +++ b/apps/api/tests/test_agent99_controlled_dispatch_ledger_contract.py @@ -68,6 +68,8 @@ def test_accepted_without_inbox_trigger_is_delivery_unknown_not_authorized() -> "incident_closure_receipt_id", "telegram_lifecycle_receipt_id", "km_writeback_ack_id", + "rag_writeback_ack_id", + "mcp_evidence_writeback_ack_id", "playbook_trust_writeback_ack_id", ] @@ -133,6 +135,8 @@ async def test_backup_learning_requires_dr_scorecard_ack() -> None: "incident_closure_receipt_id": "incident:1", "telegram_lifecycle_receipt_id": "telegram:1", "km_writeback_ack_id": "km:1", + "rag_writeback_ack_id": "rag:1", + "mcp_evidence_writeback_ack_id": "mcp:1", "playbook_trust_writeback_ack_id": "playbook:1", }, ) diff --git a/apps/api/tests/test_agent99_controlled_dispatch_reconciler_job.py b/apps/api/tests/test_agent99_controlled_dispatch_reconciler_job.py index fc5b707a5..ad47d8827 100644 --- a/apps/api/tests/test_agent99_controlled_dispatch_reconciler_job.py +++ b/apps/api/tests/test_agent99_controlled_dispatch_reconciler_job.py @@ -1,6 +1,7 @@ from __future__ import annotations from pathlib import Path +from types import SimpleNamespace import pytest @@ -50,6 +51,102 @@ def test_agent99_playbook_identity_is_project_scoped() -> None: 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_reconciler_consumes_authenticated_outcome_and_closes_same_run( monkeypatch, @@ -262,7 +359,14 @@ async def test_learning_checkpoint_resumes_without_replaying_assets( } ledger = CheckpointLedger() - effects = {"playbook": 0, "km": 0, "telegram": 0, "dr": 0} + effects = { + "playbook": 0, + "km": 0, + "rag": 0, + "telegram": 0, + "mcp": 0, + "dr": 0, + } terminal_calls = 0 async def write_playbook(*_args, **_kwargs): # type: ignore[no-untyped-def] @@ -277,6 +381,14 @@ async def test_learning_checkpoint_resumes_without_replaying_assets( effects["telegram"] += 1 return "telegram_outbound:1:durable_ack" + async def write_rag(*_args, **_kwargs): # type: ignore[no-untyped-def] + effects["rag"] += 1 + return "knowledge_entries:km-1:rag_embedding_verified" + + async def write_mcp(*_args, **_kwargs): # type: ignore[no-untyped-def] + effects["mcp"] += 1 + return "awooop_mcp_gateway_audit:call-1:verified" + async def write_dr(*_args, **_kwargs): # type: ignore[no-untyped-def] effects["dr"] += 1 return "knowledge_entries:dr-1:dr_scorecard" @@ -297,7 +409,9 @@ async def test_learning_checkpoint_resumes_without_replaying_assets( 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) @@ -314,7 +428,14 @@ async def test_learning_checkpoint_resumes_without_replaying_assets( assert first["runtime_closure_verified"] is False assert second["runtime_closure_verified"] is True - assert effects == {"playbook": 1, "km": 1, "telegram": 1, "dr": 1} + assert effects == { + "playbook": 1, + "km": 1, + "rag": 1, + "telegram": 1, + "mcp": 1, + "dr": 1, + } assert terminal_calls == 2