Skip to content

Commit ca6acf8

Browse files
Add generic AGT governance event bridge (#13)
## What changed - add a structurally compatible AGT governance-event batch sink without a mandatory AGT dependency - bind the installed Agent OS export-result enum through an optional compatibility factory - normalize AGT policy decisions with caller-owned run, policy bundle, and resource classification - pre-normalize complete batches before first emission - exclude free-form reason, resource, and arbitrary attributes by default - document upstream Agent OS deprecation and partial-delivery boundaries ## Why AGT adopters need a direct path into AgentTrust telemetry without exposing AGT internals as the public telemetry API or weakening evidence and privacy guarantees. ## Validation - 62 tests and 16 schema subtests pass - actual local AGT GovernanceEventSink protocol and SinkExportResult compatibility probe passes - 8 bundled schemas match normative bytes - all conformance fixtures pass - source compilation and diff check pass - sdist and wheel build successfully - installed-wheel smoke test passes ## Limits The included mapper covers policy events only. AGT approval-protocol and legacy AuditEntry mappings remain separate roadmap work because their assurance and correlation fields differ. Destination failures after earlier accepted emissions can still produce partial batches; durable evidence callbacks must be idempotent by event ID. Signed-off-by: Imran Siddique <imran.siddique@opaque.co>
1 parent 7c9bfb8 commit ca6acf8

6 files changed

Lines changed: 385 additions & 3 deletions

File tree

‎ROADMAP.md‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -13,7 +13,7 @@
1313
- Additional metric views and collector interoperability fixtures.
1414
- Additional W3C propagation carrier adapters beyond mutable string mappings.
1515
- Action lifecycle expansion if adopters need in-flight attempt telemetry.
16-
- AGT governance-event/audit bridge and framework integration adapters.
16+
- AGT approval-protocol/audit-entry adapters and framework integrations.
1717
- TypeScript SDK and mixed-language conformance.
1818

1919
Roadmap items are intentions, not shipped behavior or compatibility commitments.

‎docs/adapters.md‎

Lines changed: 15 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,21 @@ IDs and caller-normalized diagnostic codes. Cedar evaluation errors are recorded
1919
without rewriting the final decision, matching Cedar's skip-on-error semantics.
2020
Free-form error messages are rejected to prevent accidental content leakage.
2121

22+
The optional AGT bridge implements the batch sink shape used by Agent OS without
23+
making AGT a core dependency. The generic constructor accepts a runtime's result
24+
sentinels and is the durable integration boundary. `from_agent_os` is a legacy
25+
compatibility convenience that binds the installed runtime's actual export-result
26+
enum; upstream currently deprecates `agent-os-kernel` in favor of
27+
`agent-governance-toolkit-core`. A caller-supplied mapper keeps run identity and
28+
source classification explicit. The included policy mapper accepts only AGT
29+
policy events and requires a trusted policy-engine version, bundle digest, and
30+
resource type. It does not copy free-form reason, resource, or attribute values.
31+
32+
The bridge normalizes the whole batch before emitting its first event. This
33+
prevents a malformed later source event from causing mapping-time partial
34+
delivery. Destination failures can still occur after earlier events were
35+
accepted, so durable evidence callbacks must remain idempotent by event ID.
36+
2237
Adapters report source facts; they do not evaluate policy or prove source
2338
authenticity. Callers remain responsible for trusted bundle digests and correct
2439
action/resource classification.

‎src/agentrust_telemetry/__init__.py‎

Lines changed: 9 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,12 @@
11
"""AgentTrust governance telemetry reference SDK."""
22

3-
from .adapters import EventFactory, cedar_policy_decision, opa_decision_log
3+
from .adapters import (
4+
AgtGovernanceEventSink,
5+
EventFactory,
6+
agt_policy_decision,
7+
cedar_policy_decision,
8+
opa_decision_log,
9+
)
410
from .client import EmitResult, TelemetryClient
511
from .context import ContextIds, active_context_ids
612
from .errors import (
@@ -21,6 +27,7 @@
2127
__all__ = [
2228
"ContextIds",
2329
"ContextMismatchError",
30+
"AgtGovernanceEventSink",
2431
"EmitResult",
2532
"EventValidationError",
2633
"EventFactory",
@@ -39,6 +46,7 @@
3946
"TraceConfiguration",
4047
"TraceFinalizationError",
4148
"active_context_ids",
49+
"agt_policy_decision",
4250
"cedar_policy_decision",
4351
"extract_context",
4452
"inject_context",
Lines changed: 8 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,14 @@
11
"""Source adapters for normalized AgentTrust telemetry events."""
22

33
from .base import EventFactory
4+
from .agt import AgtGovernanceEventSink, agt_policy_decision
45
from .cedar import cedar_policy_decision
56
from .opa import opa_decision_log
67

7-
__all__ = ["EventFactory", "cedar_policy_decision", "opa_decision_log"]
8+
__all__ = [
9+
"AgtGovernanceEventSink",
10+
"EventFactory",
11+
"agt_policy_decision",
12+
"cedar_policy_decision",
13+
"opa_decision_log",
14+
]
Lines changed: 207 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,207 @@
1+
"""Optional bridge from AGT governance events to AgentTrust telemetry."""
2+
3+
from __future__ import annotations
4+
5+
import math
6+
import re
7+
import uuid
8+
from datetime import datetime, timezone
9+
from enum import Enum
10+
from typing import Any, Callable, Iterable, Protocol, Sequence
11+
12+
from .base import EventFactory
13+
14+
15+
_IDENTIFIER = re.compile(r"[A-Za-z0-9_.:-]{1,128}")
16+
_TIMESTAMP = re.compile(
17+
r"^(?P<date>\d{4}-\d{2}-\d{2}T\d{2}:\d{2}:\d{2})(?:\.(?P<fraction>\d{1,9}))?(?:Z|\+00:00)$"
18+
)
19+
20+
21+
class TelemetryEmitter(Protocol):
22+
def emit(self, event: dict[str, Any]) -> Any: ...
23+
24+
25+
AgtEventMapper = Callable[[Any], Iterable[dict[str, Any]]]
26+
27+
28+
class AgtGovernanceEventSink:
29+
"""AGT-compatible batch sink without a mandatory AGT dependency.
30+
31+
Construct directly with the source runtime's result sentinels, or use
32+
:meth:`from_agent_os` when ``agent-os`` is installed.
33+
"""
34+
35+
def __init__(
36+
self,
37+
client: TelemetryEmitter,
38+
mapper: AgtEventMapper,
39+
*,
40+
success_result: Any,
41+
failure_result: Any,
42+
) -> None:
43+
self._client = client
44+
self._mapper = mapper
45+
self._success = success_result
46+
self._failure = failure_result
47+
48+
@classmethod
49+
def from_agent_os(
50+
cls,
51+
client: TelemetryEmitter,
52+
mapper: AgtEventMapper,
53+
) -> "AgtGovernanceEventSink":
54+
try:
55+
from agent_os.event_sink import SinkExportResult
56+
except ImportError as exc:
57+
raise ImportError(
58+
"AgtGovernanceEventSink.from_agent_os requires the agent-os package"
59+
) from exc
60+
return cls(
61+
client,
62+
mapper,
63+
success_result=SinkExportResult.SUCCESS,
64+
failure_result=SinkExportResult.FAILURE,
65+
)
66+
67+
def emit(self, events: Sequence[Any]) -> Any:
68+
"""Normalize then emit a batch, returning the configured AGT result."""
69+
try:
70+
normalized = [item for source in events for item in self._mapper(source)]
71+
for event in normalized:
72+
result = self._client.emit(event)
73+
if not getattr(result, "accepted", False):
74+
return self._failure
75+
if getattr(result, "projection_errors", ()):
76+
return self._failure
77+
return self._success
78+
except Exception:
79+
return self._failure
80+
81+
def shutdown(self, timeout_ms: int = 5000) -> bool:
82+
return True
83+
84+
def force_flush(self, timeout_ms: int = 30000) -> bool:
85+
return True
86+
87+
88+
def agt_policy_decision(
89+
factory: EventFactory,
90+
source: Any,
91+
*,
92+
run_id: str,
93+
policy_engine_version: str,
94+
bundle_digest: dict[str, str],
95+
resource_type: str | None = None,
96+
enforcement_mode: str = "enforce",
97+
) -> dict[str, Any]:
98+
"""Normalize one AGT policy event without copying free-form source content."""
99+
kind = _enum_value(_field(source, "kind"))
100+
if kind not in {"policy_check", "policy_violation"}:
101+
raise ValueError(f"AGT event kind is not a policy decision: {kind!r}")
102+
decision = _decision(_field(source, "decision"))
103+
agent_id = _required_string(_field(source, "agent_id"), "agent_id")
104+
action_type = _required_string(_field(source, "action"), "action")
105+
attributes = _field(source, "attributes", {})
106+
if not isinstance(attributes, dict):
107+
raise ValueError("AGT attributes must be an object")
108+
resolved_resource_type = resource_type or attributes.get("resource_type")
109+
resolved_resource_type = _required_string(resolved_resource_type, "resource_type")
110+
event_id = _event_id(_field(source, "event_id"))
111+
reason_codes = _reason_codes(attributes.get("reason_codes", []))
112+
latency_ms = _field(source, "latency_ms", 0.0)
113+
if (
114+
not isinstance(latency_ms, (int, float))
115+
or isinstance(latency_ms, bool)
116+
or not math.isfinite(latency_ms)
117+
or latency_ms < 0
118+
):
119+
raise ValueError("AGT latency_ms must be a finite non-negative number")
120+
policy: dict[str, Any] = {
121+
"engine": "agt",
122+
"engine_version": _required_string(policy_engine_version, "policy_engine_version"),
123+
"bundle_digest": bundle_digest,
124+
}
125+
policy_name = _field(source, "policy_name")
126+
if policy_name is not None:
127+
policy["policy_id"] = _required_string(policy_name, "policy_name")
128+
return factory.build(
129+
"policy.decision",
130+
run_id=run_id,
131+
agent_id=agent_id,
132+
event_id=event_id,
133+
time_unix_nano=_timestamp_ns(_field(source, "occurred_at")),
134+
trace_id=_field(source, "trace_id"),
135+
span_id=_field(source, "span_id"),
136+
decision=decision,
137+
policy=policy,
138+
action_type=action_type,
139+
resource_type=resolved_resource_type,
140+
enforcement_mode=enforcement_mode,
141+
evaluation_duration_ns=round(latency_ms * 1_000_000),
142+
reason_codes=reason_codes,
143+
)
144+
145+
146+
def _field(source: Any, name: str, default: Any = None) -> Any:
147+
if isinstance(source, dict):
148+
return source.get(name, default)
149+
return getattr(source, name, default)
150+
151+
152+
def _enum_value(value: Any) -> Any:
153+
return value.value if isinstance(value, Enum) else value
154+
155+
156+
def _decision(value: Any) -> str:
157+
normalized = _enum_value(value)
158+
mapping = {
159+
"allow": "allow",
160+
"allowed": "allow",
161+
"deny": "deny",
162+
"denied": "deny",
163+
"block": "deny",
164+
"blocked": "deny",
165+
"require_approval": "challenge",
166+
"requires_approval": "challenge",
167+
"review": "challenge",
168+
}
169+
if normalized not in mapping:
170+
raise ValueError(f"unsupported AGT policy decision: {normalized!r}")
171+
return mapping[normalized]
172+
173+
174+
def _required_string(value: Any, field: str) -> str:
175+
if not isinstance(value, str) or not value:
176+
raise ValueError(f"AGT {field} must be a non-empty string")
177+
return value
178+
179+
180+
def _reason_codes(values: Any) -> list[str]:
181+
if not isinstance(values, list):
182+
raise ValueError("AGT reason_codes must be an array")
183+
if len(values) > 32 or any(
184+
not isinstance(value, str) or _IDENTIFIER.fullmatch(value) is None
185+
for value in values
186+
):
187+
raise ValueError("AGT reason_codes must contain at most 32 identifiers")
188+
if len(values) != len(set(values)):
189+
raise ValueError("AGT reason_codes must be unique")
190+
return list(values)
191+
192+
193+
def _event_id(value: Any) -> str:
194+
try:
195+
return str(uuid.UUID(_required_string(value, "event_id")))
196+
except (ValueError, AttributeError) as exc:
197+
raise ValueError("AGT event_id must be a UUID") from exc
198+
199+
200+
def _timestamp_ns(value: Any) -> int:
201+
value = _required_string(value, "occurred_at")
202+
match = _TIMESTAMP.fullmatch(value)
203+
if match is None:
204+
raise ValueError("AGT occurred_at must be an RFC 3339 UTC timestamp")
205+
base = datetime.fromisoformat(match.group("date")).replace(tzinfo=timezone.utc)
206+
fraction = (match.group("fraction") or "").ljust(9, "0")
207+
return int(base.timestamp()) * 1_000_000_000 + int(fraction or "0")

0 commit comments

Comments
 (0)