From f367bafaf28637a9e7084ce888ae42abfdd21f64 Mon Sep 17 00:00:00 2001 From: Imran Siddique Date: Tue, 18 Aug 2026 17:42:11 -0700 Subject: [PATCH] feat: add AGT audit entry adapters Signed-off-by: Imran Siddique --- ROADMAP.md | 2 +- docs/adapters.md | 14 ++ src/agentrust_telemetry/__init__.py | 4 + src/agentrust_telemetry/adapters/__init__.py | 3 + src/agentrust_telemetry/adapters/agt_audit.py | 177 ++++++++++++++++++ tests/test_agt_audit.py | 126 +++++++++++++ 6 files changed, 325 insertions(+), 1 deletion(-) create mode 100644 src/agentrust_telemetry/adapters/agt_audit.py create mode 100644 tests/test_agt_audit.py diff --git a/ROADMAP.md b/ROADMAP.md index 4a31886..fb418af 100644 --- a/ROADMAP.md +++ b/ROADMAP.md @@ -13,7 +13,7 @@ - Additional metric views and collector interoperability fixtures. - Additional W3C propagation carrier adapters beyond mutable string mappings. - Action lifecycle expansion if adopters need in-flight attempt telemetry. -- AGT audit-entry adapters and framework integrations. +- Framework integrations and future consolidated-core native event surfaces. - TypeScript SDK and mixed-language conformance. Roadmap items are intentions, not shipped behavior or compatibility commitments. diff --git a/docs/adapters.md b/docs/adapters.md index 6083285..0b257b9 100644 --- a/docs/adapters.md +++ b/docs/adapters.md @@ -45,6 +45,20 @@ version, and chronology against the supplied request. Individual The final chain-entry digest is retained as source-reported approval evidence; the adapter does not independently verify the chain. +Agent Mesh `AuditEntry` records, including those installed through +`agent-governance-toolkit-core`, can be mapped as policy decisions or completed +actions. Policy mapping requires a trusted bundle digest and explicit resource +classification. Action mapping requires the caller's digest of the full governed +action; `arguments_hash` is deliberately insufficient because it covers only +arguments. Free-form `data` and `resource` are not copied. + +Agent Mesh audit hash version 1.0 does not cover later-added policy decision, +policy version, argument hash, approver, timing, trace, or environment fields. +Calling the source object's `verify_hash()` therefore must not be represented as +integrity proof for those fields. The minimal Agent OS `AuditEntry` has no stable +source event ID, trace context, bundle identity, or action digest and is not +mapped automatically. + Adapters report source facts; they do not evaluate policy or prove source authenticity. Callers remain responsible for trusted bundle digests and correct action/resource classification. diff --git a/src/agentrust_telemetry/__init__.py b/src/agentrust_telemetry/__init__.py index dba9d16..7d26ae0 100644 --- a/src/agentrust_telemetry/__init__.py +++ b/src/agentrust_telemetry/__init__.py @@ -5,6 +5,8 @@ EventFactory, agt_approval_request, agt_approval_resolution, + agt_audit_action, + agt_audit_policy_decision, agt_policy_decision, agt_policy_decision_record, cedar_policy_decision, @@ -51,6 +53,8 @@ "active_context_ids", "agt_approval_request", "agt_approval_resolution", + "agt_audit_action", + "agt_audit_policy_decision", "agt_policy_decision", "agt_policy_decision_record", "cedar_policy_decision", diff --git a/src/agentrust_telemetry/adapters/__init__.py b/src/agentrust_telemetry/adapters/__init__.py index ef6b406..93facfe 100644 --- a/src/agentrust_telemetry/adapters/__init__.py +++ b/src/agentrust_telemetry/adapters/__init__.py @@ -7,6 +7,7 @@ agt_approval_resolution, agt_policy_decision_record, ) +from .agt_audit import agt_audit_action, agt_audit_policy_decision from .cedar import cedar_policy_decision from .opa import opa_decision_log @@ -17,6 +18,8 @@ "agt_policy_decision_record", "agt_approval_request", "agt_approval_resolution", + "agt_audit_action", + "agt_audit_policy_decision", "cedar_policy_decision", "opa_decision_log", ] diff --git a/src/agentrust_telemetry/adapters/agt_audit.py b/src/agentrust_telemetry/adapters/agt_audit.py new file mode 100644 index 0000000..7de5315 --- /dev/null +++ b/src/agentrust_telemetry/adapters/agt_audit.py @@ -0,0 +1,177 @@ +"""Strict adapters for Agent Mesh audit entries shipped by AGT core.""" + +from __future__ import annotations + +import uuid +from datetime import datetime, timezone +from typing import Any + +from .base import EventFactory + + +_AUDIT_NAMESPACE = uuid.UUID("398428fd-2730-498f-8ed8-6b4290668171") +_POLICY_EVENTS = {"policy_evaluation", "policy_violation"} +_ACTION_EVENTS = {"tool_invocation", "tool_blocked", "action"} + + +def agt_audit_policy_decision( + factory: EventFactory, + source: Any, + *, + run_id: str, + policy_engine_version: str, + bundle_digest: dict[str, str], + resource_type: str, + evaluation_duration_ns: int = 0, + enforcement_mode: str = "enforce", +) -> dict[str, Any]: + """Map a policy-oriented Agent Mesh AuditEntry without copying its data.""" + event_type = _required_string(_field(source, "event_type"), "event_type") + if event_type not in _POLICY_EVENTS: + raise ValueError(f"AGT audit event is not a policy event: {event_type!r}") + decision = _policy_decision(_field(source, "policy_decision")) + entry_id = _required_string(_field(source, "entry_id"), "entry_id") + policy: dict[str, Any] = { + "engine": "agt", + "engine_version": _required_string(policy_engine_version, "policy_engine_version"), + "bundle_digest": bundle_digest, + } + matched_rule = _field(source, "matched_rule") + if matched_rule is not None: + policy["policy_id"] = _required_string(matched_rule, "matched_rule") + return factory.build( + "policy.decision", + run_id=run_id, + agent_id=_required_string(_field(source, "agent_did"), "agent_did"), + event_id=_source_event_id("policy", entry_id), + time_unix_nano=_datetime_ns(_field(source, "timestamp"), "timestamp"), + trace_id=_optional_string(_field(source, "trace_id"), "trace_id"), + decision=decision, + policy=policy, + action_type=_required_string(_field(source, "action"), "action"), + resource_type=_required_string(resource_type, "resource_type"), + enforcement_mode=enforcement_mode, + evaluation_duration_ns=evaluation_duration_ns, + reason_codes=[f"agt.audit:{event_type}"], + ) + + +def agt_audit_action( + factory: EventFactory, + source: Any, + *, + run_id: str, + action_digest: dict[str, str], + action_kind: str, + operation: str, + duration_ns: int | None = None, +) -> dict[str, Any]: + """Map an action audit row using a caller-computed full action digest.""" + event_type = _required_string(_field(source, "event_type"), "event_type") + if event_type not in _ACTION_EVENTS: + raise ValueError(f"AGT audit event is not an action event: {event_type!r}") + entry_id = _required_string(_field(source, "entry_id"), "entry_id") + outcome = "denied" if event_type == "tool_blocked" else _action_outcome( + _field(source, "outcome") + ) + resolved_duration = duration_ns + if resolved_duration is None: + resolved_duration = _audit_duration(source) + target: dict[str, str] | None = None + target_did = _field(source, "target_did") + if target_did is not None: + target = {"kind": "agent", "id": _required_string(target_did, "target_did")} + return factory.build( + "action.executed", + run_id=run_id, + agent_id=_required_string(_field(source, "agent_did"), "agent_did"), + event_id=_source_event_id("action", entry_id), + time_unix_nano=_datetime_ns(_field(source, "timestamp"), "timestamp"), + trace_id=_optional_string(_field(source, "trace_id"), "trace_id"), + action_id=entry_id, + action_kind=action_kind, + action_name=_required_string(_field(source, "action"), "action"), + operation=_required_string(operation, "operation"), + outcome=outcome, + duration_ns=resolved_duration, + action_digest=action_digest, + **({"target": target} if target else {}), + ) + + +def _policy_decision(value: Any) -> str: + mapping = { + "allow": "allow", + "allowed": "allow", + "deny": "deny", + "denied": "deny", + "require_approval": "challenge", + "requires_approval": "challenge", + "review": "challenge", + "not_applicable": "not_applicable", + "error": "error", + } + if value not in mapping: + raise ValueError(f"unsupported AGT audit policy decision: {value!r}") + return mapping[value] + + +def _action_outcome(value: Any) -> str: + mapping = { + "success": "success", + "failure": "error", + "error": "error", + "denied": "denied", + "cancelled": "cancelled", + "timeout": "timeout", + } + if value not in mapping: + raise ValueError(f"unsupported AGT audit action outcome: {value!r}") + return mapping[value] + + +def _audit_duration(source: Any) -> int: + issued = _field(source, "issued_at") + completed = _field(source, "completed_at") + if issued is None or completed is None: + raise ValueError( + "AGT action audit requires duration_ns or both issued_at and completed_at" + ) + issued_ns = _datetime_ns(issued, "issued_at") + completed_ns = _datetime_ns(completed, "completed_at") + if completed_ns < issued_ns: + raise ValueError("AGT completed_at cannot predate issued_at") + return completed_ns - issued_ns + + +def _source_event_id(kind: str, source_id: str) -> str: + return str(uuid.uuid5(_AUDIT_NAMESPACE, f"{kind}:{source_id}")) + + +def _field(source: Any, name: str) -> Any: + if isinstance(source, dict): + return source.get(name) + return getattr(source, name, None) + + +def _required_string(value: Any, field: str) -> str: + if not isinstance(value, str) or not value: + raise ValueError(f"AGT audit {field} must be a non-empty string") + return value + + +def _optional_string(value: Any, field: str) -> str | None: + return None if value is None else _required_string(value, field) + + +def _datetime_ns(value: Any, field: str) -> int: + if not isinstance(value, datetime) or value.tzinfo is None: + raise ValueError(f"AGT audit {field} must be a timezone-aware datetime") + utc = value.astimezone(timezone.utc) + epoch = datetime(1970, 1, 1, tzinfo=timezone.utc) + delta = utc - epoch + return ( + delta.days * 86_400_000_000_000 + + delta.seconds * 1_000_000_000 + + delta.microseconds * 1_000 + ) diff --git a/tests/test_agt_audit.py b/tests/test_agt_audit.py new file mode 100644 index 0000000..9b60623 --- /dev/null +++ b/tests/test_agt_audit.py @@ -0,0 +1,126 @@ +import sys +import unittest +from datetime import datetime, timedelta, timezone +from pathlib import Path +from types import SimpleNamespace + + +ROOT = Path(__file__).resolve().parents[1] +sys.path.insert(0, str(ROOT / "src")) + +from agentrust_telemetry import ( # noqa: E402 + EventFactory, + SchemaValidator, + agt_audit_action, + agt_audit_policy_decision, +) + + +NOW = datetime(2026, 8, 18, 12, 0, tzinfo=timezone.utc) +DIGEST = {"algorithm": "sha256", "value": "a" * 64} +BUNDLE = {"algorithm": "sha256", "value": "b" * 64} + + +def audit_entry(**overrides): + values = { + "entry_id": "audit_123", + "timestamp": NOW, + "issued_at": NOW - timedelta(milliseconds=5), + "completed_at": NOW, + "event_type": "policy_evaluation", + "agent_did": "did:agt:agent-1", + "action": "tool.invoke", + "arguments_hash": "c" * 64, + "resource": "/customer/42", + "target_did": None, + "data": {"prompt": "must not cross"}, + "outcome": "success", + "policy_decision": "require_approval", + "matched_rule": "high-risk-tools", + "policy_version": "7", + "entry_hash": "d" * 64, + "previous_hash": "e" * 64, + "trace_id": "4bf92f3577b34da6a3ce929d0e0e4736", + } + values.update(overrides) + return SimpleNamespace(**values) + + +class AgtAuditAdapterTests(unittest.TestCase): + @classmethod + def setUpClass(cls): + cls.factory = EventFactory( + SchemaValidator(ROOT / "spec" / "schema"), + producer_name="agt-audit-tests", + producer_version="1.0.0", + ) + + def test_policy_audit_maps_without_free_form_data_or_resource(self): + event = agt_audit_policy_decision( + self.factory, + audit_entry(), + run_id="run-1", + policy_engine_version="5.0.0", + bundle_digest=BUNDLE, + resource_type="tool", + ) + self.assertEqual(event["decision"], "challenge") + self.assertEqual(event["policy"]["policy_id"], "high-risk-tools") + serialized = repr(event) + self.assertNotIn("/customer/42", serialized) + self.assertNotIn("must not cross", serialized) + self.assertNotIn("arguments_hash", serialized) + self.assertNotIn("entry_hash", serialized) + + def test_action_audit_requires_full_action_digest_and_computes_duration(self): + source = audit_entry(event_type="tool_invocation", policy_decision=None) + event = agt_audit_action( + self.factory, + source, + run_id="run-1", + action_digest=DIGEST, + action_kind="tool", + operation="invoke", + ) + self.assertEqual(event["duration_ns"], 5_000_000) + self.assertEqual(event["action_digest"], DIGEST) + self.assertNotEqual(event["action_digest"]["value"], source.arguments_hash) + + def test_blocked_action_is_denied_even_if_source_outcome_is_success(self): + event = agt_audit_action( + self.factory, + audit_entry(event_type="tool_blocked", outcome="success"), + run_id="run-1", + action_digest=DIGEST, + action_kind="tool", + operation="invoke", + ) + self.assertEqual(event["outcome"], "denied") + + def test_action_rejects_missing_or_reversed_timing(self): + with self.assertRaisesRegex(ValueError, "requires duration_ns"): + agt_audit_action( + self.factory, + audit_entry(event_type="tool_invocation", issued_at=None), + run_id="run-1", + action_digest=DIGEST, + action_kind="tool", + operation="invoke", + ) + with self.assertRaisesRegex(ValueError, "cannot predate"): + agt_audit_action( + self.factory, + audit_entry( + event_type="tool_invocation", + issued_at=NOW, + completed_at=NOW - timedelta(seconds=1), + ), + run_id="run-1", + action_digest=DIGEST, + action_kind="tool", + operation="invoke", + ) + + +if __name__ == "__main__": + unittest.main()