413 lines
15 KiB
Python
413 lines
15 KiB
Python
"""Durable, idempotent Agent99 completion callback writeback."""
|
|
|
|
from __future__ import annotations
|
|
|
|
from hashlib import sha256
|
|
from typing import Any
|
|
|
|
from sqlalchemy import text
|
|
|
|
from src.db.base import get_db_context
|
|
from src.models.agent99_completion import Agent99CompletionCallbackRequest
|
|
from src.repositories.alert_operation_log_repository import (
|
|
get_alert_operation_log_repository,
|
|
)
|
|
from src.services.agent99_controlled_dispatch_ledger import (
|
|
parse_agent99_dispatch_identity,
|
|
read_agent99_dispatch_receipt_for_run,
|
|
)
|
|
from src.services.channel_hub import (
|
|
build_external_alert_provider_event_id,
|
|
record_external_alert_event,
|
|
)
|
|
|
|
_OPERATION_INCIDENT_ID_MAX_LENGTH = 30
|
|
_AGENT99_CLOSURE_DOMAINS = frozenset({
|
|
"control_plane_recovery",
|
|
"windows_vmware",
|
|
})
|
|
|
|
|
|
def _operation_incident_id(alert_id: str | None) -> str | None:
|
|
"""Keep external alert identity while fitting the immutable event schema."""
|
|
if not alert_id:
|
|
return None
|
|
if len(alert_id) <= _OPERATION_INCIDENT_ID_MAX_LENGTH:
|
|
return alert_id
|
|
digest = sha256(alert_id.encode("utf-8")).hexdigest().upper()[:21]
|
|
return f"INC-AG99-{digest}"
|
|
|
|
|
|
async def _load_callback_readback(
|
|
*,
|
|
project_id: str,
|
|
provider_event_id: str,
|
|
callback_id: str,
|
|
) -> dict[str, Any] | None:
|
|
async with get_db_context(project_id) as db:
|
|
result = await db.execute(
|
|
text(
|
|
"""
|
|
WITH event_readback AS (
|
|
SELECT
|
|
event_row.event_id::text AS event_id,
|
|
event_row.run_id::text AS run_id
|
|
FROM awooop_conversation_event event_row
|
|
WHERE event_row.project_id = :project_id
|
|
AND event_row.channel_type = 'internal'
|
|
AND event_row.provider_event_id = :provider_event_id
|
|
LIMIT 1
|
|
),
|
|
operation_counts AS (
|
|
SELECT
|
|
count(*) AS operation_count,
|
|
count(*) FILTER (
|
|
WHERE operation_row.event_type = 'EXECUTION_COMPLETED'
|
|
) AS execution_count,
|
|
count(*) FILTER (
|
|
WHERE operation_row.event_type = 'RESOLVED'
|
|
) AS resolved_count
|
|
FROM alert_operation_log operation_row
|
|
WHERE operation_row.context ->> 'callback_id' = :callback_id
|
|
AND operation_row.context ->> 'project_id' = :project_id
|
|
)
|
|
SELECT
|
|
event_readback.event_id,
|
|
event_readback.run_id,
|
|
operation_counts.operation_count,
|
|
operation_counts.execution_count,
|
|
operation_counts.resolved_count
|
|
FROM operation_counts
|
|
LEFT JOIN event_readback ON TRUE
|
|
"""
|
|
),
|
|
{
|
|
"project_id": project_id,
|
|
"provider_event_id": provider_event_id,
|
|
"callback_id": callback_id,
|
|
},
|
|
)
|
|
row = result.mappings().first()
|
|
return dict(row) if row else None
|
|
|
|
|
|
def _operation_context(
|
|
payload: Agent99CompletionCallbackRequest,
|
|
*,
|
|
effective_outcome_state: str,
|
|
closure_authority: dict[str, Any],
|
|
) -> dict[str, Any]:
|
|
return {
|
|
"schema_version": payload.schema_version,
|
|
"callback_id": payload.callback_id,
|
|
"project_id": payload.project_id,
|
|
"run_id": payload.run_id,
|
|
"trace_id": payload.trace_id,
|
|
"work_item_id": payload.work_item_id,
|
|
"alert_id": payload.alert_id,
|
|
"operation_incident_id": _operation_incident_id(payload.alert_id),
|
|
"correlation_key": payload.correlation_key,
|
|
"source": payload.source,
|
|
"mode": payload.mode,
|
|
"outcome_state": effective_outcome_state,
|
|
"claimed_outcome_state": payload.outcome_state,
|
|
"controlled_apply": payload.controlled_apply,
|
|
"transport_ok": payload.transport_ok,
|
|
"verifier_name": payload.verifier_name,
|
|
"verifier_passed": payload.verifier_passed,
|
|
"source_event_resolved": payload.source_event_resolved,
|
|
"source_event_resolution_policy": payload.source_event_resolution_policy,
|
|
"duration_seconds": payload.duration_seconds,
|
|
"evidence_ref": payload.evidence_ref,
|
|
"alert_kind": payload.alert_kind,
|
|
"alert_service": payload.alert_service,
|
|
"alert_host": payload.alert_host,
|
|
"telegram": payload.telegram.model_dump(mode="json"),
|
|
"problem": payload.problem.model_dump(mode="json"),
|
|
"occurred_at": payload.occurred_at.isoformat(),
|
|
"callback_authority_verified": bool(
|
|
closure_authority.get("verified") is True
|
|
),
|
|
"closure_allowed": bool(closure_authority.get("verified") is True),
|
|
"closure_blocker": str(
|
|
closure_authority.get("blocker") or "unknown"
|
|
)[:120],
|
|
"typed_domain": str(
|
|
closure_authority.get("typed_domain") or ""
|
|
)[:120],
|
|
"route_id": str(closure_authority.get("route_id") or "")[:180],
|
|
"stores_secret": False,
|
|
"raw_log_stored": False,
|
|
}
|
|
|
|
|
|
def _owner_plane(alert_kind: str | None) -> str:
|
|
value = (alert_kind or "").lower()
|
|
if any(marker in value for marker in ("security", "wazuh", "siem", "cve")):
|
|
return "iwooos"
|
|
return "awoooi"
|
|
|
|
|
|
def _blocked_closure_authority(blocker: str) -> dict[str, Any]:
|
|
return {
|
|
"required": True,
|
|
"verified": False,
|
|
"blocker": blocker,
|
|
"typed_domain": "",
|
|
"route_id": "",
|
|
}
|
|
|
|
|
|
async def _load_callback_closure_authority(
|
|
payload: Agent99CompletionCallbackRequest,
|
|
) -> dict[str, Any]:
|
|
"""Bind a terminal callback to the typed dispatch and verifier ledger."""
|
|
|
|
receipt = await read_agent99_dispatch_receipt_for_run(
|
|
project_id=payload.project_id,
|
|
run_id=payload.run_id,
|
|
)
|
|
if not isinstance(receipt, dict):
|
|
return _blocked_closure_authority("dispatch_run_not_found")
|
|
public_identity = receipt.get("identity")
|
|
if not isinstance(public_identity, dict):
|
|
return _blocked_closure_authority("dispatch_identity_missing")
|
|
try:
|
|
identity = parse_agent99_dispatch_identity(public_identity)
|
|
except ValueError:
|
|
return _blocked_closure_authority("dispatch_identity_invalid")
|
|
|
|
expected_incident_ids = {
|
|
identity.incident_id,
|
|
_operation_incident_id(identity.incident_id),
|
|
}
|
|
if (
|
|
identity.project_id != payload.project_id
|
|
or str(identity.run_id) != payload.run_id
|
|
or identity.trace_id != payload.trace_id
|
|
or identity.work_item_id != payload.work_item_id
|
|
or payload.alert_id not in expected_incident_ids
|
|
or str(receipt.get("run_id") or "") != payload.run_id
|
|
or str(receipt.get("trace_id") or "") != payload.trace_id
|
|
):
|
|
return _blocked_closure_authority("dispatch_callback_identity_mismatch")
|
|
|
|
scope = receipt.get("dispatch_scope")
|
|
if not isinstance(scope, dict):
|
|
return _blocked_closure_authority("dispatch_scope_missing")
|
|
typed_domain = str(scope.get("typed_domain") or "")
|
|
route_id = str(scope.get("route_id") or "")
|
|
expected_verifier = str(scope.get("verifier") or "")
|
|
if (
|
|
scope.get("schema_version") != "agent99_dispatch_scope_v1"
|
|
or typed_domain not in _AGENT99_CLOSURE_DOMAINS
|
|
or scope.get("executor") != "Agent99"
|
|
or route_id != identity.route_id
|
|
or not expected_verifier
|
|
):
|
|
return _blocked_closure_authority("typed_dispatch_scope_mismatch")
|
|
if (
|
|
bool(scope.get("controlled_apply_requested") is True)
|
|
!= payload.controlled_apply
|
|
or (
|
|
payload.controlled_apply
|
|
and receipt.get("controlled_apply_authorized") is not True
|
|
)
|
|
or receipt.get("dispatch_promoted") is not True
|
|
):
|
|
return _blocked_closure_authority("controlled_dispatch_not_authorized")
|
|
if (
|
|
payload.verifier_name != expected_verifier
|
|
or (
|
|
str(scope.get("suggested_mode") or "")
|
|
and str(scope.get("suggested_mode")) != payload.mode
|
|
)
|
|
):
|
|
return _blocked_closure_authority("callback_route_projection_mismatch")
|
|
|
|
verifier = receipt.get("verifier")
|
|
if not isinstance(verifier, dict):
|
|
return _blocked_closure_authority("independent_verifier_missing")
|
|
if (
|
|
verifier.get("schema_version")
|
|
!= "agent99_independent_verifier_receipt_v1"
|
|
or verifier.get("status") != "success"
|
|
or str(verifier.get("run_id") or "") != payload.run_id
|
|
or str(verifier.get("trace_id") or "") != payload.trace_id
|
|
or str(verifier.get("work_item_id") or "") != payload.work_item_id
|
|
or verifier.get("outcome_state") != "resolved"
|
|
or verifier.get("transport_ok") is not True
|
|
or verifier.get("verifier_passed") is not True
|
|
or verifier.get("source_event_resolved") is not True
|
|
or receipt.get("post_verifier_passed") is not True
|
|
or str(receipt.get("run_state") or "") not in {"waiting_tool", "completed"}
|
|
):
|
|
return _blocked_closure_authority("independent_verifier_not_verified")
|
|
if (
|
|
payload.transport_ok is not True
|
|
or payload.verifier_passed is not True
|
|
or payload.source_event_resolved is not True
|
|
):
|
|
return _blocked_closure_authority("callback_verifier_projection_mismatch")
|
|
return {
|
|
"required": True,
|
|
"verified": True,
|
|
"blocker": "none",
|
|
"typed_domain": typed_domain,
|
|
"route_id": route_id,
|
|
}
|
|
|
|
|
|
async def record_agent99_completion_callback(
|
|
payload: Agent99CompletionCallbackRequest,
|
|
) -> dict[str, Any]:
|
|
"""Write one callback to the AWOOOI/AwoooP truth chain and read it back."""
|
|
|
|
closure_authority = (
|
|
await _load_callback_closure_authority(payload)
|
|
if payload.outcome_state == "resolved"
|
|
else {
|
|
"required": False,
|
|
"verified": False,
|
|
"blocker": "not_terminal_claim",
|
|
"typed_domain": "",
|
|
"route_id": "",
|
|
}
|
|
)
|
|
terminal_resolved = bool(
|
|
payload.outcome_state == "resolved"
|
|
and closure_authority.get("verified") is True
|
|
)
|
|
stage = "resolved" if terminal_resolved else (
|
|
"verifying" if payload.outcome_state == "resolved" else payload.outcome_state
|
|
)
|
|
provider_event_id = build_external_alert_provider_event_id(
|
|
"agent99",
|
|
payload.callback_id,
|
|
stage,
|
|
)
|
|
before = await _load_callback_readback(
|
|
project_id=payload.project_id,
|
|
provider_event_id=provider_event_id,
|
|
callback_id=payload.callback_id,
|
|
)
|
|
duplicate = bool(
|
|
before
|
|
and (
|
|
before.get("event_id")
|
|
or int(before.get("operation_count") or 0) > 0
|
|
)
|
|
)
|
|
severity = (
|
|
"info"
|
|
if stage == "resolved"
|
|
else "critical"
|
|
if stage == "failed"
|
|
else "warning"
|
|
)
|
|
context = _operation_context(
|
|
payload,
|
|
effective_outcome_state=stage,
|
|
closure_authority=closure_authority,
|
|
)
|
|
operation_incident_id = _operation_incident_id(payload.alert_id)
|
|
event_id = await record_external_alert_event(
|
|
project_id=payload.project_id,
|
|
provider="agent99",
|
|
event_id=payload.callback_id,
|
|
stage=stage,
|
|
title=f"Agent99 {payload.mode} {stage}",
|
|
severity=severity,
|
|
namespace="agent99",
|
|
target_resource=payload.alert_service or payload.mode,
|
|
fingerprint=payload.correlation_key or payload.trace_id,
|
|
incident_id=payload.alert_id,
|
|
labels={
|
|
"agent": "agent99",
|
|
"mode": payload.mode,
|
|
"outcome_state": stage,
|
|
"claimed_outcome_state": payload.outcome_state,
|
|
"verifier_passed": str(payload.verifier_passed).lower(),
|
|
"source_event_resolved": str(payload.source_event_resolved).lower(),
|
|
"owner_plane": _owner_plane(payload.alert_kind),
|
|
"callback_authority_verified": str(terminal_resolved).lower(),
|
|
},
|
|
annotations={
|
|
"run_id": payload.run_id,
|
|
"trace_id": payload.trace_id,
|
|
"work_item_id": payload.work_item_id,
|
|
"evidence_ref": payload.evidence_ref,
|
|
},
|
|
payload=context,
|
|
is_duplicate=duplicate,
|
|
)
|
|
|
|
repository = get_alert_operation_log_repository()
|
|
if int((before or {}).get("execution_count") or 0) == 0:
|
|
await repository.append(
|
|
"EXECUTION_COMPLETED",
|
|
incident_id=operation_incident_id,
|
|
actor="agent99_completion_callback",
|
|
action_detail=f"{payload.mode}:{payload.outcome_state}",
|
|
success=bool(payload.verifier_passed and stage != "verifying"),
|
|
context=context,
|
|
)
|
|
if terminal_resolved and int((before or {}).get("resolved_count") or 0) == 0:
|
|
await repository.append(
|
|
"RESOLVED",
|
|
incident_id=operation_incident_id,
|
|
actor="agent99_completion_callback",
|
|
action_detail=f"{payload.mode}:verified_resolved",
|
|
success=True,
|
|
context=context,
|
|
)
|
|
|
|
readback = await _load_callback_readback(
|
|
project_id=payload.project_id,
|
|
provider_event_id=provider_event_id,
|
|
callback_id=payload.callback_id,
|
|
)
|
|
operation_count = int((readback or {}).get("operation_count") or 0)
|
|
execution_count = int((readback or {}).get("execution_count") or 0)
|
|
resolved_count = int((readback or {}).get("resolved_count") or 0)
|
|
durable_readback = bool(
|
|
readback
|
|
and str(readback.get("event_id") or "") == str(event_id)
|
|
and readback.get("run_id")
|
|
and execution_count > 0
|
|
and (not terminal_resolved or resolved_count > 0)
|
|
)
|
|
closure_pending = bool(
|
|
payload.outcome_state == "resolved" and not terminal_resolved
|
|
)
|
|
return {
|
|
"schema_version": "agent99_completion_callback_receipt_v1",
|
|
"ok": bool(durable_readback and not closure_pending),
|
|
"status": (
|
|
"closure_authority_pending"
|
|
if closure_pending
|
|
else "duplicate_verified"
|
|
if duplicate
|
|
else "recorded_verified"
|
|
),
|
|
"accepted": True,
|
|
"duplicate": duplicate,
|
|
"callback_id": payload.callback_id,
|
|
"run_id": payload.run_id,
|
|
"trace_id": payload.trace_id,
|
|
"work_item_id": payload.work_item_id,
|
|
"outcome_state": stage,
|
|
"claimed_outcome_state": payload.outcome_state,
|
|
"callback_authority_verified": terminal_resolved,
|
|
"closure_allowed": terminal_resolved,
|
|
"closure_blocker": str(closure_authority.get("blocker") or "unknown"),
|
|
"conversation_event_id": str(event_id),
|
|
"awooop_run_id": (readback or {}).get("run_id"),
|
|
"operation_receipt_count": operation_count,
|
|
"execution_receipt_count": execution_count,
|
|
"resolved_receipt_count": resolved_count,
|
|
"durable_readback": durable_readback,
|
|
"stores_secret": False,
|
|
"raw_log_stored": False,
|
|
}
|