From bd47ceb4078829a8c4807eba9d6f2b37708b91cc Mon Sep 17 00:00:00 2001
From: Kinvectum <134240819+KAVentures@users.noreply.github.com>
Date: Sun, 27 Sep 2026 00:21:20 +0200
Subject: [PATCH 01/36] Add privacy-safe Claude Code OTel event adapter
---
shared/claude_otel_adapter.py | 386 ++++++++++++++++++++++++++++++++++
1 file changed, 386 insertions(+)
create mode 100644 shared/claude_otel_adapter.py
diff --git a/shared/claude_otel_adapter.py b/shared/claude_otel_adapter.py
new file mode 100644
index 00000000..742b9465
--- /dev/null
+++ b/shared/claude_otel_adapter.py
@@ -0,0 +1,386 @@
+from __future__ import annotations
+
+"""Translate Claude Code OTLP log events into privacy-safe OWG agent evidence.
+
+Claude Code can emit prompt/response text, tool parameters/content, raw API bodies,
+filesystem paths and identity attributes. This adapter deliberately ignores all of
+that. It reads only a small structural allowlist from documented Claude Code events.
+"""
+
+from datetime import datetime, timezone
+import hashlib
+import re
+from typing import Any, Iterator
+
+
+_SUPPORTED_EVENTS = frozenset({
+ "claude_code.user_prompt",
+ "claude_code.api_request",
+ "claude_code.api_error",
+ "claude_code.api_refusal",
+ "claude_code.tool_result",
+ "claude_code.tool_decision",
+ "claude_code.api_retries_exhausted",
+ "claude_code.subagent_completed",
+})
+_HUMAN_DECISION_SOURCES = frozenset({
+ "user_permanent",
+ "user_temporary",
+ "user_abort",
+ "user_reject",
+})
+_SAFE_LABEL = re.compile(r"^[A-Za-z][A-Za-z0-9_.:/-]{0,199}$")
+
+
+def _value(value: Any) -> Any:
+ if not isinstance(value, dict):
+ return value
+ for key in ("stringValue", "intValue", "doubleValue", "boolValue", "bytesValue"):
+ if key in value:
+ return value[key]
+ return None
+
+
+def _attrs(raw: Any) -> dict[str, Any]:
+ if isinstance(raw, dict):
+ return dict(raw)
+ out: dict[str, Any] = {}
+ if isinstance(raw, list):
+ for item in raw:
+ if not isinstance(item, dict):
+ continue
+ key = str(item.get("key") or "").strip()
+ if key:
+ out[key] = _value(item.get("value"))
+ return out
+
+
+def _text(value: Any, limit: int = 240) -> str:
+ return re.sub(r"\s+", " ", str(value or "")).strip()[:limit]
+
+
+def _safe_label(value: Any, *, default: str = "", limit: int = 160) -> str:
+ raw = _text(value, min(limit, 200))
+ if not raw or not _SAFE_LABEL.fullmatch(raw):
+ return default
+ return raw[:limit]
+
+
+def _bool(value: Any) -> bool | None:
+ if isinstance(value, bool):
+ return value
+ low = _text(value, 20).lower()
+ if low in {"true", "1", "yes"}:
+ return True
+ if low in {"false", "0", "no"}:
+ return False
+ return None
+
+
+def _int(value: Any) -> int | None:
+ try:
+ number = int(value)
+ except Exception:
+ return None
+ return number if 0 <= number <= 1_000_000_000 else None
+
+
+def _duration_seconds(value: Any) -> float:
+ try:
+ millis = float(value or 0)
+ except Exception:
+ return 0.0
+ return round(max(0.0, min(millis / 1000.0, 7 * 24 * 60 * 60)), 6)
+
+
+def _iso_from_nanos(value: Any) -> str:
+ try:
+ nanos = int(value)
+ except Exception:
+ return ""
+ try:
+ return datetime.fromtimestamp(nanos / 1_000_000_000, timezone.utc).isoformat()
+ except Exception:
+ return ""
+
+
+def _observed_at(attrs: dict[str, Any], native: dict[str, Any]) -> str:
+ timestamp = _text(attrs.get("event.timestamp"), 80)
+ if timestamp:
+ return timestamp
+ for key in ("timeUnixNano", "time_unix_nano", "observedTimeUnixNano", "observed_time_unix_nano"):
+ parsed = _iso_from_nanos(native.get(key))
+ if parsed:
+ return parsed
+ return datetime.now(timezone.utc).isoformat()
+
+
+def _hash(value: Any) -> str:
+ raw = _text(value, 500)
+ return hashlib.sha256(raw.encode("utf-8")).hexdigest()[:24] if raw else ""
+
+
+def _body_text(body: Any) -> str:
+ value = _value(body)
+ return _text(value, 120) if isinstance(value, (str, int, float, bool)) else ""
+
+
+def _event_name(record: dict[str, Any], attrs: dict[str, Any]) -> str:
+ raw = _text(attrs.get("event.name"), 120).lower()
+ if raw.startswith("claude_code."):
+ return raw
+ if raw:
+ candidate = f"claude_code.{raw}"
+ if candidate in _SUPPORTED_EVENTS:
+ return candidate
+ body = _body_text(record.get("body")).lower()
+ if body.startswith("claude_code."):
+ return body
+ return ""
+
+
+def _iter_records(payload: dict[str, Any], max_records: int) -> Iterator[tuple[dict[str, Any], dict[str, Any], dict[str, Any]]]:
+ seen = 0
+ for resource_log in payload.get("resourceLogs") or []:
+ if not isinstance(resource_log, dict):
+ continue
+ resource = resource_log.get("resource") if isinstance(resource_log.get("resource"), dict) else {}
+ resource_attrs = _attrs(resource.get("attributes"))
+ for scope_log in resource_log.get("scopeLogs") or []:
+ if not isinstance(scope_log, dict):
+ continue
+ for record in scope_log.get("logRecords") or []:
+ if not isinstance(record, dict):
+ continue
+ seen += 1
+ if seen > max_records:
+ raise ValueError(f"Claude Code OTLP payload exceeds {max_records} records")
+ yield record, _attrs(record.get("attributes")), resource_attrs
+
+ # Deterministic small form for adapter tests and custom relays.
+ for record in payload.get("records") or []:
+ if not isinstance(record, dict):
+ continue
+ seen += 1
+ if seen > max_records:
+ raise ValueError(f"Claude Code OTLP payload exceeds {max_records} records")
+ yield record, _attrs(record.get("attributes")), _attrs(payload.get("resource"))
+
+
+def _usage(attrs: dict[str, Any]) -> dict[str, int]:
+ out: dict[str, int] = {}
+ mapping = {
+ "input_tokens": "input_tokens",
+ "output_tokens": "output_tokens",
+ "cached_input_tokens": "cache_read_tokens",
+ }
+ for target, source in mapping.items():
+ amount = _int(attrs.get(source))
+ if amount is not None:
+ out[target] = amount
+ if "input_tokens" in out or "output_tokens" in out:
+ out["total_tokens"] = out.get("input_tokens", 0) + out.get("output_tokens", 0)
+ return out
+
+
+def _tool_category(name: str, attrs: dict[str, Any]) -> str:
+ low = name.lower()
+ tool_source = _text(attrs.get("tool_source"), 80).lower()
+ if tool_source in {"mcp", "sdk_host_builtin_mcp"} or low.startswith("mcp__") or low == "mcp_tool":
+ return "mcp"
+ if low in {"bash", "shell", "terminal", "computer"} or any(token in low for token in ("exec", "command", "powershell")):
+ return "shell"
+ if any(token in low for token in ("read", "write", "edit", "file", "notebook")):
+ return "filesystem"
+ if any(token in low for token in ("search", "grep", "glob", "find", "lookup")):
+ return "search"
+ if any(token in low for token in ("webfetch", "browser", "playwright", "chrome")):
+ return "browser"
+ if any(token in low for token in ("github", "git", "code")):
+ return "code"
+ return "other" if name else "none"
+
+
+def _event_id(run_id: str, event_name: str, attrs: dict[str, Any]) -> str:
+ if event_name in {"claude_code.tool_result", "claude_code.tool_decision"}:
+ discriminator = _text(attrs.get("tool_use_id"), 240) or _text(attrs.get("event.sequence"), 80)
+ elif event_name in {"claude_code.api_request", "claude_code.api_error", "claude_code.api_refusal"}:
+ discriminator = (
+ _text(attrs.get("request_id"), 240)
+ or _text(attrs.get("client_request_id"), 240)
+ or _text(attrs.get("event.sequence"), 80)
+ )
+ elif event_name == "claude_code.subagent_completed":
+ discriminator = "|".join([
+ _text(attrs.get("event.sequence"), 80),
+ _safe_label(attrs.get("agent_type"), default="subagent", limit=80),
+ ])
+ else:
+ discriminator = _text(attrs.get("event.sequence"), 80) or _text(attrs.get("event.timestamp"), 80)
+ material = f"{run_id}\x1f{event_name}\x1f{discriminator}".encode("utf-8")
+ return "claude-otel:" + hashlib.sha256(material).hexdigest()[:40]
+
+
+def _base(
+ *,
+ attrs: dict[str, Any],
+ resource_attrs: dict[str, Any],
+ native: dict[str, Any],
+ event_name: str,
+ operation: str,
+ status: str,
+ defaults: dict[str, Any],
+ tool_name: str = "",
+ span_id: str = "",
+) -> dict[str, Any] | None:
+ session_id = _text(attrs.get("session.id") or resource_attrs.get("session.id"), 128)
+ prompt_id = _text(attrs.get("prompt.id"), 128)
+ run_id = prompt_id or session_id
+ if not run_id:
+ return None
+ model = _safe_label(attrs.get("model"), limit=200)
+ return {
+ "event_id": _event_id(run_id, event_name, attrs),
+ "observed_at": _observed_at(attrs, native),
+ "organization_id": _text(defaults.get("organization_id"), 128),
+ "actor_id": _text(defaults.get("actor_id"), 128),
+ "device_id": _text(defaults.get("device_id"), 128) or "claude-local",
+ "sensor_id": "agent:claude-code-otel",
+ "session_id": session_id or run_id,
+ "agent_name": _safe_label(defaults.get("agent_name"), default="Claude-Code", limit=160),
+ "provider": "anthropic",
+ "framework": "claude-code",
+ "model": model,
+ "operation": operation,
+ "status": status,
+ "observation_level": "native_trace",
+ "run_id": run_id,
+ "trace_id": session_id or run_id,
+ "span_id": span_id,
+ "tool_name": tool_name,
+ "tool_category": _tool_category(tool_name, attrs),
+ "duration_seconds": _duration_seconds(attrs.get("duration_ms")),
+ "usage": _usage(attrs),
+ }
+
+
+def claude_otel_to_agent_events(
+ payload: dict[str, Any],
+ *,
+ defaults: dict[str, Any] | None = None,
+ max_records: int = 1000,
+) -> tuple[list[dict[str, Any]], dict[str, int]]:
+ """Project documented Claude Code log events without copying event content."""
+ if not isinstance(payload, dict):
+ raise ValueError("Claude Code OpenTelemetry payload must be an object")
+ defaults = dict(defaults or {})
+ events: list[dict[str, Any]] = []
+ seen = 0
+ ignored = 0
+
+ for native, attrs, resource_attrs in _iter_records(payload, max_records):
+ seen += 1
+ event_name = _event_name(native, attrs)
+ if event_name not in _SUPPORTED_EVENTS:
+ ignored += 1
+ continue
+
+ projected: dict[str, Any] | None = None
+ if event_name == "claude_code.user_prompt":
+ projected = _base(
+ attrs=attrs,
+ resource_attrs=resource_attrs,
+ native=native,
+ event_name=event_name,
+ operation="run_started",
+ status="running",
+ defaults=defaults,
+ )
+
+ elif event_name in {"claude_code.api_request", "claude_code.api_error", "claude_code.api_refusal"}:
+ status = "success" if event_name == "claude_code.api_request" else "error" if event_name == "claude_code.api_error" else "denied"
+ request_id = _text(attrs.get("request_id") or attrs.get("client_request_id"), 128)
+ projected = _base(
+ attrs=attrs,
+ resource_attrs=resource_attrs,
+ native=native,
+ event_name=event_name,
+ operation="model_call",
+ status=status,
+ defaults=defaults,
+ span_id=("request:" + _hash(request_id)) if request_id else "",
+ )
+
+ elif event_name == "claude_code.tool_result":
+ tool_name = _safe_label(attrs.get("tool_name"), default="unknown-tool", limit=160)
+ success = _bool(attrs.get("success"))
+ tool_id = _text(attrs.get("tool_use_id"), 128)
+ projected = _base(
+ attrs=attrs,
+ resource_attrs=resource_attrs,
+ native=native,
+ event_name=event_name,
+ operation="tool_call",
+ status="success" if success is True else "error" if success is False else "unknown",
+ defaults=defaults,
+ tool_name=tool_name,
+ span_id=("tool:" + _hash(tool_id)) if tool_id else "",
+ )
+
+ elif event_name == "claude_code.tool_decision":
+ source = _text(attrs.get("source"), 80).lower()
+ if source not in _HUMAN_DECISION_SOURCES:
+ ignored += 1
+ continue
+ decision = _text(attrs.get("decision"), 40).lower()
+ tool_name = _safe_label(attrs.get("tool_name"), default="unknown-tool", limit=160)
+ tool_id = _text(attrs.get("tool_use_id"), 128)
+ projected = _base(
+ attrs=attrs,
+ resource_attrs=resource_attrs,
+ native=native,
+ event_name=event_name,
+ operation="human_approval_received",
+ status="success" if decision == "accept" else "denied" if decision == "reject" else "unknown",
+ defaults=defaults,
+ tool_name=tool_name,
+ span_id=("tool:" + _hash(tool_id)) if tool_id else "",
+ )
+
+ elif event_name == "claude_code.subagent_completed":
+ agent_type = _safe_label(attrs.get("agent_type"), default="subagent", limit=80)
+ projected = _base(
+ attrs=attrs,
+ resource_attrs=resource_attrs,
+ native=native,
+ event_name=event_name,
+ operation="handoff",
+ status="success",
+ defaults=defaults,
+ tool_name=f"subagent:{agent_type}",
+ )
+
+ elif event_name == "claude_code.api_retries_exhausted":
+ projected = _base(
+ attrs=attrs,
+ resource_attrs=resource_attrs,
+ native=native,
+ event_name=event_name,
+ operation="error",
+ status="error",
+ defaults=defaults,
+ )
+
+ if projected is None:
+ ignored += 1
+ continue
+ events.append(projected)
+
+ unique: dict[str, dict[str, Any]] = {}
+ for event in events:
+ unique.setdefault(str(event["event_id"]), event)
+ return list(unique.values()), {
+ "records_seen": seen,
+ "records_ignored": ignored,
+ "agent_events": len(unique),
+ }
From 791b017fb71b8ca960cd084ae809be0735b7663a Mon Sep 17 00:00:00 2001
From: Kinvectum <134240819+KAVentures@users.noreply.github.com>
Date: Sun, 27 Sep 2026 00:21:36 +0200
Subject: [PATCH 02/36] Ingest Claude Code structural OTel events
---
server/agent_ingest.py | 26 ++++++++++++++++++++++++++
1 file changed, 26 insertions(+)
diff --git a/server/agent_ingest.py b/server/agent_ingest.py
index 18a11dfb..c0229eff 100644
--- a/server/agent_ingest.py
+++ b/server/agent_ingest.py
@@ -8,6 +8,7 @@
from shared.agent_evidence import AgentEvidenceError, agent_event_to_evidence
from shared.agent_ingress_validation import validate_agent_ingress_event
from shared.capture_control import filter_recordable
+from shared.claude_otel_adapter import claude_otel_to_agent_events
from shared.codex_otel_adapter import codex_otel_to_agent_events
from shared.otel_agent_adapter import otel_payload_to_agent_events
from .db import insert_events
@@ -16,6 +17,7 @@
MAX_AGENT_BATCH_BYTES = 2_000_000
MAX_OTEL_SPANS = 1000
MAX_CODEX_OTEL_RECORDS = 1000
+MAX_CLAUDE_OTEL_RECORDS = 1000
def _bounded_json_size(value: Any, *, maximum: int = MAX_AGENT_BATCH_BYTES) -> None:
@@ -102,3 +104,27 @@ def ingest_codex_otel_payload(
"projected": len(events),
"inserted": inserted,
}
+
+
+def ingest_claude_otel_payload(
+ payload: dict[str, Any],
+ *,
+ defaults: dict[str, Any] | None = None,
+) -> dict[str, int]:
+ """Ingest Claude Code OTLP logs through a strict structural allowlist."""
+ _bounded_json_size({"payload": payload, "defaults": defaults or {}})
+ projected, stats = claude_otel_to_agent_events(
+ payload,
+ defaults=defaults,
+ max_records=MAX_CLAUDE_OTEL_RECORDS,
+ )
+ if len(projected) > MAX_AGENT_EVENTS * 2:
+ raise AgentEvidenceError("Claude Code OpenTelemetry projection produced too many agent events")
+ events = _validated_events(projected)
+ inserted = _insert_recordable(events)
+ return {
+ "records_seen": int(stats.get("records_seen") or 0),
+ "records_ignored": int(stats.get("records_ignored") or 0),
+ "projected": len(events),
+ "inserted": inserted,
+ }
From 4080d86a5e5d8db56561a03805d2557319f0413e Mon Sep 17 00:00:00 2001
From: Kinvectum <134240819+KAVentures@users.noreply.github.com>
Date: Sun, 27 Sep 2026 00:21:59 +0200
Subject: [PATCH 03/36] Expose write-only Claude Code OTel ingest endpoint
---
server/agent_routes.py | 18 ++++++++++++++++++
1 file changed, 18 insertions(+)
diff --git a/server/agent_routes.py b/server/agent_routes.py
index 6eab1065..223960e0 100644
--- a/server/agent_routes.py
+++ b/server/agent_routes.py
@@ -11,6 +11,7 @@
from .agent_ingest import (
MAX_AGENT_BATCH_BYTES,
ingest_agent_payloads,
+ ingest_claude_otel_payload,
ingest_codex_otel_payload,
ingest_otel_payload,
)
@@ -34,6 +35,7 @@
AGENT_EVENT_PATH = "/agent-ingest/v1/events"
AGENT_OTEL_PATH = "/agent-ingest/v1/otel"
AGENT_CODEX_OTEL_PATH = "/agent-ingest/v1/codex-otel"
+AGENT_CLAUDE_OTEL_PATH = "/agent-ingest/v1/claude-otel"
class OTelDefaults(BaseModel):
@@ -140,6 +142,22 @@ async def ingest_codex_otel(request: Request) -> Response:
return Response(status_code=202)
+@router.post(AGENT_CLAUDE_OTEL_PATH, status_code=202)
+async def ingest_claude_otel(request: Request) -> Response:
+ """Accept Claude Code OTLP/HTTP JSON logs on a write-only local endpoint."""
+ _require_agent_write_bearer(request)
+ payload = await _read_bounded_json(request)
+ if not isinstance(payload, dict):
+ raise HTTPException(status_code=422, detail="OpenTelemetry payload must be an object")
+ defaults_raw = payload.pop("openworkgraph", {})
+ try:
+ defaults = OTelDefaults.model_validate(defaults_raw if isinstance(defaults_raw, dict) else {}).model_dump()
+ ingest_claude_otel_payload(payload, defaults=defaults)
+ except (AgentEvidenceError, ValueError) as exc:
+ raise HTTPException(status_code=422, detail=str(exc)) from exc
+ return Response(status_code=202)
+
+
@router.get("/v1/agent-workflows")
def get_agent_workflows(
request: Request,
From 524245b01669f3a1f7a7b1f3eb01f206fabdaded Mon Sep 17 00:00:00 2001
From: Kinvectum <134240819+KAVentures@users.noreply.github.com>
Date: Sun, 27 Sep 2026 00:22:48 +0200
Subject: [PATCH 04/36] Correlate Claude hooks by prompt and observe subagent
handoffs
---
shared/claude_code_adapter.py | 117 ++++++++++++++++++++++++++--------
1 file changed, 89 insertions(+), 28 deletions(-)
diff --git a/shared/claude_code_adapter.py b/shared/claude_code_adapter.py
index 2480e6bf..6fc49cf2 100644
--- a/shared/claude_code_adapter.py
+++ b/shared/claude_code_adapter.py
@@ -59,7 +59,7 @@ def _duration_seconds(payload: dict[str, Any]) -> float:
def _tool_category(name: str) -> str:
low = name.lower()
- if low.startswith("mcp__"):
+ if low.startswith("mcp__") or low == "mcp_tool":
return "mcp"
if low in {"bash", "shell", "terminal", "computer"} or any(
token in low for token in ("exec", "command", "powershell")
@@ -76,6 +76,21 @@ def _tool_category(name: str) -> str:
return "other" if name else "none"
+def _prompt_run_id(payload: dict[str, Any], session_id: str) -> str:
+ # Claude Code v2.1.196+ gives every within-turn hook the same prompt_id used
+ # by its OTel prompt.id. Prefer it so hook evidence and OTel evidence land in
+ # one execution. Older versions safely fall back to the session boundary.
+ return _text(payload.get("prompt_id"), 128) or session_id
+
+
+def _hook_agent_name(payload: dict[str, Any]) -> str:
+ agent_id = _text(payload.get("agent_id"), 128)
+ if not agent_id:
+ return "Claude Code"
+ agent_type = _safe_label(payload.get("agent_type"), default="subagent", limit=80)
+ return f"Claude Code/{agent_type}"
+
+
def _base_event(
payload: dict[str, Any],
*,
@@ -87,20 +102,26 @@ def _base_event(
event_key: str,
agent_name: str = "Claude Code",
span_id: str = "",
+ parent_span_id: str = "",
tool_name: str = "",
+ model: str = "",
) -> dict[str, Any]:
return {
- "event_id": _event_id(trace_id, event_key, span_id, run_id),
+ "event_id": _event_id(trace_id, event_key, span_id, run_id, agent_name),
"observed_at": observed_at,
+ "sensor_id": "agent:claude-code-hook",
+ "session_id": trace_id,
"agent_name": agent_name,
"provider": "anthropic",
"framework": "claude-code",
+ "model": model,
"operation": operation,
"status": status,
"observation_level": "native_trace",
"run_id": run_id,
"trace_id": trace_id,
"span_id": span_id,
+ "parent_span_id": parent_span_id,
"tool_name": tool_name,
"tool_category": _tool_category(tool_name),
"duration_seconds": _duration_seconds(payload),
@@ -112,11 +133,11 @@ def claude_hook_to_agent_events(
*,
observed_at: str | None = None,
) -> list[dict[str, Any]]:
- """Project one Claude Code hook invocation into zero or one safe events.
+ """Project one Claude Code hook invocation into privacy-safe structural events.
- Hooks that expose content but do not add reliable structural workflow signal
- (for example UserPromptSubmit, PreToolUse, MessageDisplay, and Stop) are
- intentionally ignored. Tool calls are recorded only after success/failure.
+ Content-bearing hooks remain ignored. Within-turn events prefer prompt_id, so
+ they correlate with Claude Code OTel without exposing prompt content. A
+ SubagentStart produces both the parent handoff and a child-run boundary.
"""
if not isinstance(payload, dict):
return []
@@ -124,11 +145,13 @@ def claude_hook_to_agent_events(
if hook not in _SUPPORTED_EVENTS:
return []
- session_id = _text(payload.get("session_id"), 240)
+ session_id = _text(payload.get("session_id"), 128)
if not session_id:
return []
timestamp = _text(observed_at, 80) or _now_iso()
trace_id = session_id
+ prompt_run_id = _prompt_run_id(payload, session_id)
+ agent_name = _hook_agent_name(payload)
if hook == "SessionStart":
return [_base_event(
@@ -139,6 +162,7 @@ def claude_hook_to_agent_events(
run_id=session_id,
trace_id=trace_id,
event_key=hook,
+ model=_safe_label(payload.get("model"), limit=160),
)]
if hook == "SessionEnd":
@@ -153,81 +177,118 @@ def claude_hook_to_agent_events(
)]
if hook in {"PostToolUse", "PostToolUseFailure"}:
- tool_use_id = _text(payload.get("tool_use_id"), 240)
+ tool_use_id = _text(payload.get("tool_use_id"), 128)
tool_name = _safe_label(payload.get("tool_name"), default="unknown-tool", limit=160)
return [_base_event(
payload,
operation="tool_call",
status="success" if hook == "PostToolUse" else "error",
observed_at=timestamp,
- run_id=session_id,
+ run_id=prompt_run_id,
trace_id=trace_id,
span_id=tool_use_id,
event_key=hook,
tool_name=tool_name,
+ agent_name=agent_name,
)]
if hook == "PermissionRequest":
- tool_use_id = _text(payload.get("tool_use_id"), 240)
+ tool_use_id = _text(payload.get("tool_use_id"), 128)
tool_name = _safe_label(payload.get("tool_name"), default="unknown-tool", limit=160)
return [_base_event(
payload,
operation="human_approval_requested",
status="running",
observed_at=timestamp,
- run_id=session_id,
+ run_id=prompt_run_id,
trace_id=trace_id,
span_id=tool_use_id,
event_key=hook,
tool_name=tool_name,
+ agent_name=agent_name,
)]
if hook == "PermissionDenied":
- # Claude Code emits PermissionDenied in auto mode. Do not mislabel an
- # automatic policy denial as a human approval/denial decision.
- tool_use_id = _text(payload.get("tool_use_id"), 240)
+ # PermissionDenied can be an automatic policy denial. Do not turn it into
+ # a human decision; Claude's tool_decision OTel event identifies the
+ # actual decision source when richer telemetry is connected.
+ tool_use_id = _text(payload.get("tool_use_id"), 128)
tool_name = _safe_label(payload.get("tool_name"), default="unknown-tool", limit=160)
return [_base_event(
payload,
operation="error",
status="denied",
observed_at=timestamp,
- run_id=session_id,
+ run_id=prompt_run_id,
trace_id=trace_id,
span_id=tool_use_id,
event_key=hook,
tool_name=tool_name,
+ agent_name=agent_name,
)]
- if hook in {"SubagentStart", "SubagentStop"}:
- agent_id = _text(payload.get("agent_id"), 240)
- if not agent_id:
+ if hook == "SubagentStart":
+ child_id = _text(payload.get("agent_id"), 128)
+ if not child_id:
+ return []
+ child_type = _safe_label(payload.get("agent_type"), default="subagent", limit=80)
+ child_name = f"Claude Code/{child_type}"
+ # One event belongs to the parent execution and records the delegation;
+ # the other opens the child execution under the same prompt correlation.
+ return [
+ _base_event(
+ payload,
+ operation="handoff",
+ status="running",
+ observed_at=timestamp,
+ run_id=prompt_run_id,
+ trace_id=trace_id,
+ span_id=child_id,
+ event_key="SubagentHandoff",
+ tool_name=f"subagent:{child_type}",
+ agent_name="Claude Code",
+ ),
+ _base_event(
+ payload,
+ operation="run_started",
+ status="running",
+ observed_at=timestamp,
+ run_id=prompt_run_id,
+ trace_id=trace_id,
+ span_id=child_id,
+ event_key=hook,
+ agent_name=child_name,
+ ),
+ ]
+
+ if hook == "SubagentStop":
+ child_id = _text(payload.get("agent_id"), 128)
+ if not child_id:
return []
- agent_type = _safe_label(payload.get("agent_type"), default="subagent", limit=80)
- sub_run_id = f"{session_id}:{agent_id}"
+ child_type = _safe_label(payload.get("agent_type"), default="subagent", limit=80)
return [_base_event(
payload,
- operation="run_started" if hook == "SubagentStart" else "run_finished",
- status="running" if hook == "SubagentStart" else "unknown",
+ operation="run_finished",
+ status="unknown",
observed_at=timestamp,
- run_id=sub_run_id,
+ run_id=prompt_run_id,
trace_id=trace_id,
- span_id=agent_id,
+ span_id=child_id,
event_key=hook,
- agent_name=f"Claude Code/{agent_type}",
+ agent_name=f"Claude Code/{child_type}",
)]
if hook == "StopFailure":
- prompt_id = _text(payload.get("prompt_id"), 240)
return [_base_event(
payload,
operation="error",
status="error",
observed_at=timestamp,
- run_id=session_id,
+ run_id=prompt_run_id,
trace_id=trace_id,
- span_id=prompt_id,
+ span_id=_text(payload.get("prompt_id"), 128),
event_key=hook,
+ agent_name=agent_name,
)]
return []
From 36f5f218f178f0e6e0b29ec8fdce1087789dc10a Mon Sep 17 00:00:00 2001
From: Kinvectum <134240819+KAVentures@users.noreply.github.com>
Date: Sun, 27 Sep 2026 00:24:12 +0200
Subject: [PATCH 05/36] Add capability-aware agent observability summaries
---
server/agent_observability.py | 188 ++++++++++++++++++++++++++++++++++
1 file changed, 188 insertions(+)
create mode 100644 server/agent_observability.py
diff --git a/server/agent_observability.py b/server/agent_observability.py
new file mode 100644
index 00000000..74de6491
--- /dev/null
+++ b/server/agent_observability.py
@@ -0,0 +1,188 @@
+from __future__ import annotations
+
+"""Capability-aware summaries over canonical agent evidence.
+
+Counts answer "what did OWG observe?" while capabilities answer "could the active
+adapter observe this signal at all?" Keeping those separate prevents unsupported
+signals from being rendered as misleading zeroes.
+"""
+
+from collections import Counter
+from typing import Any
+
+from .context_execution_linkage import _agent_groups, _meta, _one_execution
+
+
+_SIGNAL_KEYS = (
+ "model_call",
+ "tool_call",
+ "handoff",
+ "human_approval_requested",
+ "human_approval_received",
+ "token_usage",
+ "model_identity",
+ "duration",
+)
+
+
+def _sensor_ids(events: list[dict[str, Any]]) -> set[str]:
+ return {str(event.get("sensor_id") or "").strip() for event in events if event.get("sensor_id")}
+
+
+def _framework(events: list[dict[str, Any]]) -> str:
+ for event in events:
+ meta, _trace = _meta(event)
+ agent = meta.get("agent") if isinstance(meta.get("agent"), dict) else {}
+ value = str(agent.get("framework") or "").strip().lower()
+ if value:
+ return value
+ return ""
+
+
+def _capability(status: str, basis: str) -> dict[str, str]:
+ return {"status": status, "basis": basis}
+
+
+def _capabilities(events: list[dict[str, Any]]) -> dict[str, dict[str, str]]:
+ sensors = _sensor_ids(events)
+ framework = _framework(events)
+ caps = {key: _capability("unknown", "adapter_capability_not_declared") for key in _SIGNAL_KEYS}
+
+ if framework == "claude-code":
+ has_otel = "agent:claude-code-otel" in sensors
+ has_hooks = "agent:claude-code-hook" in sensors
+ if has_otel:
+ for key in ("model_call", "tool_call", "handoff", "human_approval_received", "token_usage", "model_identity", "duration"):
+ caps[key] = _capability("observable", "claude_code_otel_logs")
+ if has_hooks:
+ for key in ("tool_call", "handoff", "human_approval_requested", "duration"):
+ caps[key] = _capability("observable", "claude_code_hooks")
+ if not has_otel:
+ for key in ("model_call", "human_approval_received", "token_usage"):
+ caps[key] = _capability("not_observable", "claude_code_hooks_do_not_emit_signal")
+ if not has_hooks:
+ caps["human_approval_requested"] = _capability("not_observable", "claude_code_otel_reports_decision_not_prompt_open")
+ return caps
+
+ if framework == "codex":
+ for key in ("model_call", "tool_call", "handoff", "human_approval_received", "model_identity", "duration"):
+ caps[key] = _capability("observable", "codex_otel")
+ caps["human_approval_requested"] = _capability("not_observable", "codex_otel_decision_surface_has_no_request_event")
+ caps["token_usage"] = _capability("partial", "codex_usage_is_span_or_turn_dependent")
+ return caps
+
+ if framework.startswith("openai-agents"):
+ for key in ("model_call", "tool_call", "handoff", "token_usage", "model_identity", "duration"):
+ caps[key] = _capability("observable", "openai_agents_tracing_processor")
+ for key in ("human_approval_requested", "human_approval_received"):
+ caps[key] = _capability("not_observable", "openai_agents_processor_has_no_approval_lifecycle_signal")
+ return caps
+
+ if any(sensor.startswith("otel:") or sensor == "agent:otel" for sensor in sensors):
+ for key in ("model_call", "tool_call", "token_usage", "model_identity", "duration"):
+ caps[key] = _capability("observable", "generic_genai_otel_semantic_conventions")
+ for key in ("handoff", "human_approval_requested", "human_approval_received"):
+ caps[key] = _capability("not_observable", "generic_genai_otel_contract_has_no_portable_signal")
+ return caps
+
+
+def _effective_events(events: list[dict[str, Any]]) -> list[dict[str, Any]]:
+ """Prefer richer Claude OTel evidence where hooks report the same tool call.
+
+ OWG intentionally keeps both canonical rows. This read-only projection avoids
+ double-counting after a user upgrades an existing hook installation to OTel.
+ Handoff initiation remains hook-preferred when both feeds are present, because
+ the OTel subagent event is emitted on completion rather than delegation start.
+ """
+ framework = _framework(events)
+ sensors = _sensor_ids(events)
+ if framework != "claude-code" or not {"agent:claude-code-hook", "agent:claude-code-otel"}.issubset(sensors):
+ return events
+
+ hook_handoff_observed = any(
+ str(event.get("sensor_id") or "") == "agent:claude-code-hook"
+ and str(_meta(event)[0].get("operation") or "") == "handoff"
+ for event in events
+ )
+ output: list[dict[str, Any]] = []
+ for event in events:
+ meta, _trace = _meta(event)
+ operation = str(meta.get("operation") or "")
+ sensor = str(event.get("sensor_id") or "")
+ if operation == "tool_call" and sensor == "agent:claude-code-hook":
+ continue
+ if operation == "handoff" and hook_handoff_observed and sensor == "agent:claude-code-otel":
+ continue
+ output.append(event)
+ return output
+
+
+def _usage_totals(events: list[dict[str, Any]]) -> dict[str, int]:
+ totals: Counter[str] = Counter()
+ for event in events:
+ meta, _trace = _meta(event)
+ usage = meta.get("usage") if isinstance(meta.get("usage"), dict) else {}
+ for key in ("input_tokens", "output_tokens", "cached_input_tokens", "total_tokens"):
+ try:
+ value = int(usage.get(key))
+ except Exception:
+ continue
+ if value >= 0:
+ totals[key] += value
+ return dict(totals)
+
+
+def _models(events: list[dict[str, Any]]) -> list[str]:
+ values: set[str] = set()
+ for event in events:
+ meta, _trace = _meta(event)
+ agent = meta.get("agent") if isinstance(meta.get("agent"), dict) else {}
+ model = str(agent.get("model") or "").strip()
+ if model:
+ values.add(model[:200])
+ return sorted(values)
+
+
+def enrich_agent_execution_payload(payload: dict[str, Any], raw_events: list[dict[str, Any]]) -> dict[str, Any]:
+ groups: dict[str, list[dict[str, Any]]] = {}
+ for events in _agent_groups(raw_events):
+ if not events:
+ continue
+ execution = _one_execution(events)
+ groups[str(execution.get("execution_id") or "")] = events
+
+ for run in payload.get("executions") or []:
+ if not isinstance(run, dict):
+ continue
+ events = groups.get(str(run.get("execution_id") or ""), [])
+ if not events:
+ continue
+ effective = _effective_events(events)
+ counts = Counter(str(_meta(event)[0].get("operation") or "unknown") for event in effective)
+ caps = _capabilities(events)
+ sensors = sorted(_sensor_ids(events))
+ run["observed_operation_counts"] = dict(sorted(counts.items()))
+ run["signal_capabilities"] = caps
+ run["telemetry_sources"] = sensors
+ run["usage_totals"] = _usage_totals(effective)
+ run["models_observed"] = _models(effective)
+ run["telemetry_depth"] = (
+ "rich_native_events_plus_hooks"
+ if {"agent:claude-code-hook", "agent:claude-code-otel"}.issubset(set(sensors))
+ else "rich_native_events"
+ if "agent:claude-code-otel" in sensors or _framework(events) in {"codex", "openai-agents-python"}
+ else "hooks_only"
+ if "agent:claude-code-hook" in sensors
+ else "provider_neutral_structural"
+ )
+
+ payload["capability_semantics"] = {
+ "observable": "the active adapter is designed to emit this structural signal; zero means none was observed in this evidence window",
+ "partial": "the active adapter can expose the signal only on some runtime paths or span shapes",
+ "not_observable": "the active adapter does not expose this signal; do not interpret absence as zero runtime activity",
+ "unknown": "the integration did not declare whether this signal is observable",
+ }
+ return payload
+
+
+__all__ = ["enrich_agent_execution_payload"]
From f201a26fc99ff41e9f2f0547fe231c5a02d8c9e8 Mon Sep 17 00:00:00 2001
From: Kinvectum <134240819+KAVentures@users.noreply.github.com>
Date: Sun, 27 Sep 2026 00:24:20 +0200
Subject: [PATCH 06/36] Expose signal capability and rich run summaries
---
server/agent_execution_trace_routes.py | 2 ++
1 file changed, 2 insertions(+)
diff --git a/server/agent_execution_trace_routes.py b/server/agent_execution_trace_routes.py
index f2dbdcb7..c11ce89c 100644
--- a/server/agent_execution_trace_routes.py
+++ b/server/agent_execution_trace_routes.py
@@ -5,6 +5,7 @@
from fastapi import APIRouter, HTTPException, Request
from .agent_execution_traces import agent_execution_traces
+from .agent_observability import enrich_agent_execution_payload
from .agent_read_auth import agent_read_authorized
from .procedural_memory import load_recent_evidence
@@ -41,6 +42,7 @@ def get_agent_execution_traces(
limit=limit,
max_events_per_execution=max_events_per_execution,
)
+ payload = enrich_agent_execution_payload(payload, raw)
except (TypeError, ValueError) as exc:
raise HTTPException(status_code=422, detail="invalid agent-execution-trace query") from exc
return {
From 9a04c1424864387361a7ddae8f076ecf27187fd0 Mon Sep 17 00:00:00 2001
From: Kinvectum <134240819+KAVentures@users.noreply.github.com>
Date: Sun, 27 Sep 2026 00:27:40 +0200
Subject: [PATCH 07/36] Enrich Codex turn, usage and multi-agent telemetry
---
shared/codex_otel_adapter.py | 148 +++++++++++++++++++++++++----------
1 file changed, 106 insertions(+), 42 deletions(-)
diff --git a/shared/codex_otel_adapter.py b/shared/codex_otel_adapter.py
index f276e6cf..cd5fc849 100644
--- a/shared/codex_otel_adapter.py
+++ b/shared/codex_otel_adapter.py
@@ -2,10 +2,10 @@
"""Translate Codex OTLP JSON into structural OpenWorkGraph agent evidence.
-Codex can export diagnostic log records containing prompts, account identifiers,
-tool arguments, tool output, and error strings. This module intentionally reads
-only a small structural allowlist. It never copies OTLP record bodies or arbitrary
-attributes into canonical evidence.
+Codex telemetry can contain prompts, account identifiers, tool arguments/results,
+inter-agent message content and error strings. This adapter only projects a small
+structural allowlist and never copies record bodies or arbitrary attributes into
+canonical evidence.
"""
from datetime import datetime, timezone
@@ -19,6 +19,7 @@
"codex.tool_result",
"codex.tool_decision",
"codex.api_request",
+ "codex.agent_communication",
})
_SAFE_LABEL = re.compile(r"^[A-Za-z][A-Za-z0-9_.:/-]{0,199}$")
@@ -63,9 +64,10 @@ def _bool(value: Any) -> bool | None:
def _int(value: Any) -> int | None:
try:
- return int(value)
+ number = int(value)
except Exception:
return None
+ return number if 0 <= number <= 1_000_000_000 else None
def _duration_seconds(value: Any) -> float:
@@ -112,8 +114,30 @@ def _hash_part(value: Any) -> str:
return hashlib.sha256(raw.encode("utf-8")).hexdigest()[:24]
+def _event_name(attrs: dict[str, Any]) -> str:
+ raw = _text(attrs.get("event.name"), 160)
+ if raw in _SUPPORTED_EVENTS:
+ return raw
+
+ # Some Codex builds have exported a tracing call-site in event.name instead
+ # of the semantic name. Infer only from low-cardinality structural fields;
+ # never inspect record bodies, arguments, output or inter-agent content.
+ if attrs.get("communication_id") and attrs.get("kind") and attrs.get("state"):
+ return "codex.agent_communication"
+ if attrs.get("call_id") and attrs.get("decision") is not None and attrs.get("source") is not None:
+ return "codex.tool_decision"
+ if attrs.get("call_id") and attrs.get("tool_name") and attrs.get("success") is not None:
+ return "codex.tool_result"
+ if attrs.get("attempt") is not None and (
+ "http.response.status_code" in attrs or "error.message" in attrs or "endpoint" in attrs
+ ):
+ return "codex.api_request"
+ if attrs.get("provider_name") and "approval_policy" in attrs and "sandbox_policy" in attrs:
+ return "codex.conversation_starts"
+ return ""
+
+
def _event_id(conversation_id: str, event_name: str, attrs: dict[str, Any]) -> str:
- """Stable across Codex log + trace copies without exposing native IDs."""
if event_name == "codex.conversation_starts":
discriminator = "conversation-start"
elif event_name == "codex.tool_result":
@@ -134,6 +158,12 @@ def _event_id(conversation_id: str, event_name: str, attrs: dict[str, Any]) -> s
_text(attrs.get("attempt"), 40),
"" if request_hash else _text(attrs.get("event.timestamp"), 80),
])
+ elif event_name == "codex.agent_communication":
+ discriminator = "agent-communication|" + "|".join([
+ _hash_part(attrs.get("communication_id")),
+ _text(attrs.get("kind"), 40),
+ _text(attrs.get("state"), 40),
+ ])
else:
discriminator = _text(attrs.get("event.timestamp"), 80)
material = f"{conversation_id}\x1f{event_name}\x1f{discriminator}".encode("utf-8")
@@ -157,6 +187,24 @@ def _tool_category(name: str, namespace: str = "") -> str:
return "other" if name else "none"
+def _merge_span_structure(event_attrs: dict[str, Any], span_attrs: dict[str, Any]) -> dict[str, Any]:
+ out = dict(event_attrs)
+ for key in (
+ "conversation.id",
+ "thread.id",
+ "turn.id",
+ "model",
+ "gen_ai.request.model",
+ "gen_ai.usage.input_tokens",
+ "gen_ai.usage.output_tokens",
+ "gen_ai.usage.cache_read.input_tokens",
+ "codex.usage.total_tokens",
+ ):
+ if key not in out and key in span_attrs:
+ out[key] = span_attrs[key]
+ return out
+
+
def _iter_records(payload: dict[str, Any], max_records: int) -> Iterator[tuple[dict[str, Any], dict[str, Any]]]:
seen = 0
@@ -174,8 +222,6 @@ def _iter_records(payload: dict[str, Any], max_records: int) -> Iterator[tuple[d
raise ValueError(f"Codex OTLP payload exceeds {max_records} records")
yield record, _attrs(record.get("attributes"))
- # Codex also emits trace-safe telemetry as span events. Read event attributes,
- # never arbitrary span attributes, span names, or native event bodies.
for resource_span in payload.get("resourceSpans") or []:
if not isinstance(resource_span, dict):
continue
@@ -185,6 +231,7 @@ def _iter_records(payload: dict[str, Any], max_records: int) -> Iterator[tuple[d
for span in scope_span.get("spans") or []:
if not isinstance(span, dict):
continue
+ span_attrs = _attrs(span.get("attributes"))
for event in span.get("events") or []:
if not isinstance(event, dict):
continue
@@ -194,9 +241,8 @@ def _iter_records(payload: dict[str, Any], max_records: int) -> Iterator[tuple[d
synthetic = dict(event)
if "timeUnixNano" not in synthetic and span.get("endTimeUnixNano"):
synthetic["timeUnixNano"] = span.get("endTimeUnixNano")
- yield synthetic, _attrs(event.get("attributes"))
+ yield synthetic, _merge_span_structure(_attrs(event.get("attributes")), span_attrs)
- # Small simplified form for deterministic adapter tests/integrators.
for record in payload.get("records") or []:
if not isinstance(record, dict):
continue
@@ -206,6 +252,23 @@ def _iter_records(payload: dict[str, Any], max_records: int) -> Iterator[tuple[d
yield record, _attrs(record.get("attributes"))
+def _usage(attrs: dict[str, Any]) -> dict[str, int]:
+ mapping = {
+ "input_tokens": "gen_ai.usage.input_tokens",
+ "output_tokens": "gen_ai.usage.output_tokens",
+ "cached_input_tokens": "gen_ai.usage.cache_read.input_tokens",
+ "total_tokens": "codex.usage.total_tokens",
+ }
+ out: dict[str, int] = {}
+ for target, source in mapping.items():
+ amount = _int(attrs.get(source))
+ if amount is not None:
+ out[target] = amount
+ if "total_tokens" not in out and ("input_tokens" in out or "output_tokens" in out):
+ out["total_tokens"] = out.get("input_tokens", 0) + out.get("output_tokens", 0)
+ return out
+
+
def _base(
attrs: dict[str, Any],
native: dict[str, Any],
@@ -218,16 +281,18 @@ def _base(
tool_category: str = "none",
span_id: str = "",
) -> dict[str, Any] | None:
- conversation_id = _text(attrs.get("conversation.id"), 240)
+ conversation_id = _text(attrs.get("conversation.id") or attrs.get("thread.id") or defaults.get("run_id"), 128)
if not conversation_id:
return None
- model = _safe_label(attrs.get("model"), limit=200)
+ turn_id = _text(attrs.get("turn.id"), 128)
+ run_id = turn_id or conversation_id
+ model = _safe_label(attrs.get("model") or attrs.get("gen_ai.request.model"), limit=200)
return {
"event_id": _event_id(conversation_id, event_name, attrs),
"observed_at": _observed_at(attrs, native),
- "organization_id": _text(defaults.get("organization_id"), 240),
- "actor_id": _text(defaults.get("actor_id"), 240),
- "device_id": _text(defaults.get("device_id"), 240) or "codex-local",
+ "organization_id": _text(defaults.get("organization_id"), 128),
+ "actor_id": _text(defaults.get("actor_id"), 128),
+ "device_id": _text(defaults.get("device_id"), 128) or "codex-local",
"sensor_id": "agent:codex-otel",
"agent_name": _safe_label(defaults.get("agent_name"), default="Codex", limit=160),
"provider": "openai",
@@ -236,13 +301,14 @@ def _base(
"operation": operation,
"status": status,
"observation_level": "native_trace",
- "run_id": conversation_id,
+ "run_id": run_id,
"trace_id": conversation_id,
"span_id": span_id,
- "workflow_id": _text(defaults.get("workflow_id"), 240),
+ "workflow_id": _text(defaults.get("workflow_id"), 128),
"tool_name": tool_name,
"tool_category": tool_category,
"duration_seconds": _duration_seconds(attrs.get("duration_ms")),
+ "usage": _usage(attrs),
}
@@ -261,7 +327,7 @@ def codex_otel_to_agent_events(
for native, attrs in _iter_records(payload, max_records):
seen += 1
- event_name = _text(attrs.get("event.name"), 100)
+ event_name = _event_name(attrs)
if event_name not in _SUPPORTED_EVENTS:
ignored += 1
continue
@@ -269,12 +335,8 @@ def codex_otel_to_agent_events(
projected: dict[str, Any] | None = None
if event_name == "codex.conversation_starts":
projected = _base(
- attrs,
- native,
- defaults=defaults,
- event_name=event_name,
- operation="run_started",
- status="running",
+ attrs, native, defaults=defaults, event_name=event_name,
+ operation="run_started", status="running",
)
elif event_name == "codex.tool_result":
@@ -282,20 +344,15 @@ def codex_otel_to_agent_events(
namespace = _safe_label(attrs.get("tool_namespace"), limit=160)
success = _bool(attrs.get("success"))
projected = _base(
- attrs,
- native,
- defaults=defaults,
- event_name=event_name,
+ attrs, native, defaults=defaults, event_name=event_name,
operation="tool_call",
status="success" if success is True else "error" if success is False else "unknown",
tool_name=tool_name,
tool_category=_tool_category(tool_name, namespace),
- span_id=_text(attrs.get("call_id"), 240),
+ span_id=("call:" + _hash_part(attrs.get("call_id"))) if attrs.get("call_id") else "",
)
elif event_name == "codex.tool_decision":
- # Codex can resolve approvals through a user OR an automated reviewer.
- # Only explicit user-sourced decisions are human approval evidence.
source = _text(attrs.get("source"), 80).lower()
if source != "user":
ignored += 1
@@ -303,34 +360,41 @@ def codex_otel_to_agent_events(
tool_name = _safe_label(attrs.get("tool_name"), default="unknown-tool", limit=160)
namespace = _safe_label(attrs.get("tool_namespace"), limit=160)
decision = _text(attrs.get("decision"), 80).lower()
- denied = any(token in decision for token in ("deny", "denied", "reject", "cancel"))
+ denied = any(token in decision for token in ("deny", "denied", "reject", "cancel", "abort"))
approved = any(token in decision for token in ("approve", "approved", "allow"))
projected = _base(
- attrs,
- native,
- defaults=defaults,
- event_name=event_name,
+ attrs, native, defaults=defaults, event_name=event_name,
operation="human_approval_received",
status="denied" if denied else "success" if approved else "unknown",
tool_name=tool_name,
tool_category=_tool_category(tool_name, namespace),
- span_id=_text(attrs.get("call_id"), 240),
+ span_id=("call:" + _hash_part(attrs.get("call_id"))) if attrs.get("call_id") else "",
)
elif event_name == "codex.api_request":
status_code = _int(attrs.get("http.response.status_code"))
- # We only inspect whether an error exists; its text is never copied.
has_error = bool(_text(attrs.get("error.message"), 1))
success = status_code is not None and 200 <= status_code <= 299 and not has_error
projected = _base(
- attrs,
- native,
- defaults=defaults,
- event_name=event_name,
+ attrs, native, defaults=defaults, event_name=event_name,
operation="model_call",
status="success" if success else "error" if has_error or (status_code or 0) >= 400 else "unknown",
)
+ elif event_name == "codex.agent_communication":
+ # Inter-agent messages/results are communication, not necessarily a
+ # delegation. Count only a new spawn as a handoff.
+ if _text(attrs.get("state"), 40).lower() != "send" or _text(attrs.get("kind"), 40).lower() != "spawn":
+ ignored += 1
+ continue
+ communication_id = _text(attrs.get("communication_id"), 128)
+ projected = _base(
+ attrs, native, defaults=defaults, event_name=event_name,
+ operation="handoff", status="success",
+ tool_name="agent:spawn", tool_category="other",
+ span_id=("communication:" + _hash_part(communication_id)) if communication_id else "",
+ )
+
if projected is None:
ignored += 1
continue
From 1f42c11e30fc7926ffd686cf06f0e6c7ee4a7003 Mon Sep 17 00:00:00 2001
From: Kinvectum <134240819+KAVentures@users.noreply.github.com>
Date: Sun, 27 Sep 2026 00:31:26 +0200
Subject: [PATCH 08/36] Add privacy-safe Claude Code OTel configuration helper
---
adapters/claude_code_otel.py | 29 +++++++++++++++++++++++++++++
1 file changed, 29 insertions(+)
create mode 100644 adapters/claude_code_otel.py
diff --git a/adapters/claude_code_otel.py b/adapters/claude_code_otel.py
new file mode 100644
index 00000000..0ffbba4b
--- /dev/null
+++ b/adapters/claude_code_otel.py
@@ -0,0 +1,29 @@
+from __future__ import annotations
+
+"""Build the Claude Code settings env needed for structural OWG telemetry only."""
+
+
+def env_settings(*, token: str, base_url: str) -> dict[str, str]:
+ base = str(base_url or "").rstrip("/")
+ return {
+ "CLAUDE_CODE_ENABLE_TELEMETRY": "1",
+ "OTEL_LOGS_EXPORTER": "otlp",
+ "OTEL_EXPORTER_OTLP_LOGS_PROTOCOL": "http/json",
+ "OTEL_EXPORTER_OTLP_LOGS_ENDPOINT": f"{base}/agent-ingest/v1/claude-otel",
+ "OTEL_EXPORTER_OTLP_LOGS_HEADERS": f"Authorization=Bearer {token}",
+ # Claude leaves these content surfaces off by default. OWG writes the
+ # explicit zeroes as defense in depth; the ingest adapter independently
+ # allowlists structural attributes and would discard content anyway.
+ "OTEL_LOG_USER_PROMPTS": "0",
+ "OTEL_LOG_ASSISTANT_RESPONSES": "0",
+ "OTEL_LOG_TOOL_DETAILS": "0",
+ "OTEL_LOG_TOOL_CONTENT": "0",
+ "OTEL_LOG_RAW_API_BODIES": "0",
+ }
+
+
+def settings_fragment(*, token: str, base_url: str) -> dict[str, dict[str, str]]:
+ return {"env": env_settings(token=token, base_url=base_url)}
+
+
+__all__ = ["env_settings", "settings_fragment"]
From aa364477aecf7158b10b957c4e72718f91d32c5f Mon Sep 17 00:00:00 2001
From: Kinvectum <134240819+KAVentures@users.noreply.github.com>
Date: Sun, 27 Sep 2026 00:31:59 +0200
Subject: [PATCH 09/36] Make Claude one-click setup manage rich telemetry
safely
---
server/agent_config_writer.py | 117 +++++++++++++++++++++++++++++-----
1 file changed, 102 insertions(+), 15 deletions(-)
diff --git a/server/agent_config_writer.py b/server/agent_config_writer.py
index 2f21f8e2..faed9749 100644
--- a/server/agent_config_writer.py
+++ b/server/agent_config_writer.py
@@ -70,7 +70,7 @@ def _atomic_write(path: Path, text: str) -> None:
raise
-# --- Claude Code: hooks in ~/.claude/settings.json ---------------------------
+# --- Claude Code: hooks + logs-only OTel in ~/.claude/settings.json ----------
def _load_claude(path: Path) -> dict[str, Any]:
if not path.exists():
@@ -87,6 +87,9 @@ def _load_claude(path: Path) -> dict[str, Any]:
hooks = data.get("hooks")
if hooks is not None and not isinstance(hooks, dict):
raise ConfigConflict(f"'hooks' in {path} has an unexpected shape; use manual setup")
+ env = data.get("env")
+ if env is not None and not isinstance(env, dict):
+ raise ConfigConflict(f"'env' in {path} has an unexpected shape; use manual setup")
return data
@@ -125,24 +128,92 @@ def _strip_owg_hooks(data: dict[str, Any]) -> int:
return removed
-def claude_status() -> dict[str, Any]:
- path = claude_settings_path()
- try:
- data = _load_claude(path)
- except ConfigConflict as exc:
- return {"configured": False, "path": str(path), "error": str(exc)}
- configured = any(
+def _claude_hooks_configured(data: dict[str, Any]) -> bool:
+ return any(
_is_owg_handler(h)
for groups in (data.get("hooks") or {}).values() if isinstance(groups, list)
for group in groups if isinstance(group, dict)
for h in (group.get("hooks") or []) if isinstance(group.get("hooks"), list)
)
- return {"configured": configured, "path": str(path)}
-def claude_connect(fragment_factory: Callable[[], dict]) -> dict[str, Any]:
+def _claude_env_matches(data: dict[str, Any], managed_env: dict[str, str] | None) -> bool:
+ if not managed_env:
+ return True
+ existing = data.get("env") if isinstance(data.get("env"), dict) else {}
+ return all(str(existing.get(key) or "") == str(value) for key, value in managed_env.items())
+
+
+def _validate_claude_env_conflicts(data: dict[str, Any], managed_env: dict[str, str]) -> None:
+ existing = data.get("env") if isinstance(data.get("env"), dict) else {}
+ conflicts = [
+ key for key, desired in managed_env.items()
+ if key in existing and str(existing.get(key)) != str(desired)
+ ]
+ if conflicts:
+ joined = ", ".join(sorted(conflicts))
+ raise ConfigConflict(
+ f"Claude Code already has different telemetry settings for {joined}. "
+ "OpenWorkGraph will not overwrite them; use manual setup or an OTLP collector/tee."
+ )
+
+
+def _merge_claude_env(data: dict[str, Any], managed_env: dict[str, str]) -> None:
+ if not managed_env:
+ return
+ env = data.setdefault("env", {})
+ if not isinstance(env, dict):
+ raise ConfigConflict("Claude Code 'env' has an unexpected shape; use manual setup")
+ for key, value in managed_env.items():
+ env[key] = str(value)
+
+
+def _strip_matching_claude_env(data: dict[str, Any], managed_env: dict[str, str]) -> int:
+ """Remove only OWG values that still exactly match what OWG would write.
+
+ If the user changed a value after connecting, leave it untouched. This makes
+ Disconnect reversible without claiming ownership of unrelated telemetry keys.
+ """
+ env = data.get("env")
+ if not isinstance(env, dict):
+ return 0
+ removed = 0
+ for key, expected in managed_env.items():
+ if key in env and str(env.get(key)) == str(expected):
+ del env[key]
+ removed += 1
+ if not env:
+ data.pop("env", None)
+ return removed
+
+
+def claude_status(managed_env: dict[str, str] | None = None) -> dict[str, Any]:
+ path = claude_settings_path()
+ try:
+ data = _load_claude(path)
+ except ConfigConflict as exc:
+ return {"configured": False, "path": str(path), "error": str(exc)}
+ hooks_configured = _claude_hooks_configured(data)
+ telemetry_configured = _claude_env_matches(data, managed_env)
+ return {
+ "configured": hooks_configured and telemetry_configured,
+ "hooks_configured": hooks_configured,
+ "telemetry_configured": telemetry_configured,
+ "path": str(path),
+ }
+
+
+def claude_connect(
+ fragment_factory: Callable[[], dict],
+ managed_env: dict[str, str] | None = None,
+) -> dict[str, Any]:
path = claude_settings_path()
data = _load_claude(path)
+ desired_env = dict(managed_env or {})
+ # Fail before changing hooks if a user already owns a conflicting telemetry
+ # destination or privacy flag.
+ _validate_claude_env_conflicts(data, desired_env)
+
# Replace, never duplicate: older OpenWorkGraph handlers (e.g. the pre-0.89
# command/args form) are removed before the current ones are added.
_strip_owg_hooks(data)
@@ -152,22 +223,38 @@ def claude_connect(fragment_factory: Callable[[], dict]) -> dict[str, Any]:
if not isinstance(existing, list):
raise ConfigConflict(f"hooks.{event} in {path} has an unexpected shape; use manual setup")
existing.extend(groups)
+ _merge_claude_env(data, desired_env)
+
backup = _backup(path)
_atomic_write(path, json.dumps(data, indent=2, ensure_ascii=False) + "\n")
- return {"configured": True, "path": str(path), "backup": backup,
- "note": "Takes effect in new Claude Code sessions."}
+ return {
+ "configured": True,
+ "hooks_configured": True,
+ "telemetry_configured": bool(desired_env),
+ "path": str(path),
+ "backup": backup,
+ "note": "Takes effect in new Claude Code sessions.",
+ }
-def claude_disconnect() -> dict[str, Any]:
+def claude_disconnect(managed_env: dict[str, str] | None = None) -> dict[str, Any]:
path = claude_settings_path()
if not path.exists():
return {"configured": False, "path": str(path), "backup": None}
data = _load_claude(path)
- if not _strip_owg_hooks(data):
+ removed_hooks = _strip_owg_hooks(data)
+ removed_env = _strip_matching_claude_env(data, dict(managed_env or {}))
+ if not (removed_hooks or removed_env):
return {"configured": False, "path": str(path), "backup": None}
backup = _backup(path)
_atomic_write(path, json.dumps(data, indent=2, ensure_ascii=False) + "\n")
- return {"configured": False, "path": str(path), "backup": backup}
+ return {
+ "configured": False,
+ "hooks_configured": False,
+ "telemetry_configured": False,
+ "path": str(path),
+ "backup": backup,
+ }
# --- Codex: [otel] in ~/.codex/config.toml ------------------------------------
From 8208aa8d12db899e32262917fa6a837183bc1887 Mon Sep 17 00:00:00 2001
From: Kinvectum <134240819+KAVentures@users.noreply.github.com>
Date: Sun, 27 Sep 2026 00:32:40 +0200
Subject: [PATCH 10/36] Wire rich Claude telemetry into one-click agent setup
---
server/agent_dashboard_control_plane.py | 67 +++++++++++++++++++------
1 file changed, 51 insertions(+), 16 deletions(-)
diff --git a/server/agent_dashboard_control_plane.py b/server/agent_dashboard_control_plane.py
index 6b74d1cb..793be78f 100644
--- a/server/agent_dashboard_control_plane.py
+++ b/server/agent_dashboard_control_plane.py
@@ -10,16 +10,15 @@
treats an integration as active only when structural telemetry is observed.
"""
-import json
import shlex
import sys
-from pathlib import Path
from typing import Any
from fastapi import Request
from fastapi.responses import HTMLResponse, JSONResponse, Response
-from adapters.claude_code_hook import settings_fragment
+from adapters.claude_code_hook import settings_fragment as claude_hook_settings
+from adapters.claude_code_otel import env_settings as claude_otel_env
from adapters.codex_config import config_snippet
from server import agent_config_writer as writer
from server.agent_auth import ensure_agent_ingest_token
@@ -28,13 +27,28 @@
DASHBOARD_SCRIPT = ROOT / "dashboard" / "agent_control_plane.js"
+OBSERVABILITY_SCRIPT = ROOT / "dashboard" / "agent_observability_v090.js"
_SCRIPT_MARKER = ''
+_OBSERVABILITY_MARKER = ''
def _base_url(request: Request) -> str:
return str(request.base_url).rstrip("/")
+def _claude_env(request: Request) -> dict[str, str]:
+ return claude_otel_env(
+ token=ensure_agent_ingest_token(),
+ base_url=_base_url(request),
+ )
+
+
+def _claude_settings(request: Request) -> dict[str, Any]:
+ settings = dict(claude_hook_settings())
+ settings["env"] = _claude_env(request)
+ return settings
+
+
def agent_setup_payload(request: Request) -> dict[str, Any]:
"""Return reviewable setup material for the local dashboard owner.
@@ -78,10 +92,12 @@ def agent_setup_payload(request: Request) -> dict[str, Any]:
"integrations": {
"claude_code": {
"label": "Claude Code",
- "method": "native_lifecycle_hooks",
+ "method": "native_hooks_plus_otel_logs",
"command": f"cd {shlex.quote(str(ROOT))} && {shlex.quote(sys.executable)} -m adapters.claude_code_hook --print-settings",
"one_click": True,
- "settings": settings_fragment(),
+ "settings": _claude_settings(request),
+ "otel_endpoint": f"{base_url}/agent-ingest/v1/claude-otel",
+ "telemetry_depth": "rich_structural",
"events": [
"SessionStart",
"SessionEnd",
@@ -92,9 +108,18 @@ def agent_setup_payload(request: Request) -> dict[str, Any]:
"SubagentStart",
"SubagentStop",
"StopFailure",
+ "claude_code.user_prompt",
+ "claude_code.api_request",
+ "claude_code.api_error",
+ "claude_code.api_refusal",
+ "claude_code.tool_result",
+ "claude_code.tool_decision",
+ "claude_code.api_retries_exhausted",
+ "claude_code.subagent_completed",
],
- "instructions": "Click Connect to add these hooks to ~/.claude/settings.json (other settings are preserved and a backup is written), or merge the hooks object manually into ~/.claude/settings.json or an intentional project-scoped .claude/settings.json.",
+ "instructions": "Click Connect to add OpenWorkGraph's hooks and logs-only OpenTelemetry settings to ~/.claude/settings.json. Other settings are preserved, conflicting existing telemetry is never overwritten, and a backup is written. Manual setup can merge the shown hooks and env objects instead.",
"failure_mode": "fail_open_async",
+ "content_logging_enabled": False,
},
"codex": {
"label": "Codex",
@@ -105,13 +130,15 @@ def agent_setup_payload(request: Request) -> dict[str, Any]:
"endpoint": f"{base_url}/agent-ingest/v1/codex-otel",
"instructions": "Click Connect to add a managed [otel] block to ~/.codex/config.toml (refused if you already have your own [otel] settings), or merge these keys manually into the existing [otel] section. Do not create a second [otel] table.",
"logs_enabled_by_owg": False,
+ "telemetry_depth": "native_trace",
},
"openai_agents": {
"label": "OpenAI Agents SDK",
"method": "additional_tracing_processor",
"python": "from adapters.openai_agents import install_openai_agents_processor\n\ninstall_openai_agents_processor()",
- "instructions": "Register OpenWorkGraph as an additional tracing processor in the agent application. Existing SDK tracing remains in place.",
+ "instructions": "Register OpenWorkGraph as an additional tracing processor in the agent application. Existing SDK tracing remains in place. OWG projects model spans, tools, handoffs, hierarchy, usage when the SDK exposes it, timings and errors without serializing span content.",
"one_click": False,
+ "telemetry_depth": "native_trace",
},
"otel": {
"label": "Generic OpenTelemetry",
@@ -119,7 +146,8 @@ def agent_setup_payload(request: Request) -> dict[str, Any]:
"endpoint": otel_endpoint,
"posix": posix_otel,
"powershell": powershell_otel,
- "instructions": "Use the signal-specific traces endpoint exactly as shown. OpenWorkGraph currently accepts OTLP/HTTP JSON here, not protobuf or gRPC.",
+ "instructions": "Use the signal-specific traces endpoint exactly as shown. OpenWorkGraph currently accepts OTLP/HTTP JSON here, not protobuf or gRPC. Portable GenAI model/tool spans are supported; provider-specific handoff or approval signals require a native adapter.",
+ "telemetry_depth": "portable_genai",
},
"custom": {
"label": "Custom structural agent",
@@ -140,22 +168,22 @@ def get_agent_setup(request: Request) -> JSONResponse:
_ONE_CLICK = {
"claude_code": {
- "status": writer.claude_status,
- "connect": lambda request: writer.claude_connect(settings_fragment),
- "disconnect": writer.claude_disconnect,
+ "status": lambda request: writer.claude_status(_claude_env(request)),
+ "connect": lambda request: writer.claude_connect(claude_hook_settings, _claude_env(request)),
+ "disconnect": lambda request: writer.claude_disconnect(_claude_env(request)),
},
"codex": {
- "status": writer.codex_status,
+ "status": lambda request: writer.codex_status(),
"connect": lambda request: writer.codex_connect(
config_snippet(token=ensure_agent_ingest_token(), base_url=_base_url(request))
),
- "disconnect": writer.codex_disconnect,
+ "disconnect": lambda request: writer.codex_disconnect(),
},
}
-def get_agent_config_status() -> JSONResponse:
- status = {kind: ops["status"]() for kind, ops in _ONE_CLICK.items()}
+def get_agent_config_status(request: Request) -> JSONResponse:
+ status = {kind: ops["status"](request) for kind, ops in _ONE_CLICK.items()}
return JSONResponse({"integrations": status}, headers={"Cache-Control": "no-store"})
@@ -170,7 +198,7 @@ async def change_agent_config(request: Request) -> JSONResponse:
if ops is None or action not in {"connect", "disconnect"}:
return JSONResponse({"detail": "unsupported agent or action"}, status_code=400)
try:
- result = ops["connect"](request) if action == "connect" else ops["disconnect"]()
+ result = ops[action](request)
except writer.ConfigConflict as exc:
return JSONResponse({"detail": str(exc), "manual_setup_required": True}, status_code=409)
except OSError:
@@ -185,6 +213,10 @@ def agent_control_plane_script() -> Response:
return Response(DASHBOARD_SCRIPT.read_text(encoding="utf-8"), media_type="application/javascript")
+def agent_observability_script() -> Response:
+ return Response(OBSERVABILITY_SCRIPT.read_text(encoding="utf-8"), media_type="application/javascript")
+
+
async def _inject_agent_control_plane(request: Request, call_next):
response = await call_next(request)
if request.method.upper() != "GET" or request.url.path != "/" or response.status_code != 200:
@@ -202,6 +234,8 @@ async def _inject_agent_control_plane(request: Request, call_next):
return response
if _SCRIPT_MARKER not in text:
text = text.replace("