feat(sre): require Agent99 RAG and MCP closure receipts
This commit is contained in:
@@ -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((
|
||||
|
||||
@@ -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(
|
||||
|
||||
@@ -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",
|
||||
})
|
||||
|
||||
@@ -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",
|
||||
},
|
||||
)
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
|
||||
Reference in New Issue
Block a user