Coverage for agentos/observability/otel_bridge.py: 0%
298 statements
« prev ^ index » next coverage.py v7.14.3, created at 2026-07-08 01:44 +0800
« prev ^ index » next coverage.py v7.14.3, created at 2026-07-08 01:44 +0800
1"""AgentOS OpenTelemetry - OTLP/Jaeger/Zipkin trace/metric export (v1.14.6)."""
3import logging
4import os
5from collections.abc import Callable
6from contextlib import asynccontextmanager, contextmanager
7from dataclasses import dataclass, field
8from enum import StrEnum
9from functools import wraps
10from typing import TYPE_CHECKING, Any
12if TYPE_CHECKING:
13 from agentos.observability.metrics import MetricsCollector
15logger = logging.getLogger("agentos.otel")
18# ---- Enums ----
21class OTelExporter(StrEnum):
22 """OpenTelemetry exporter backend."""
24 OTLP_HTTP = "otlp_http"
25 OTLP_GRPC = "otlp_grpc"
26 CONSOLE = "console"
27 ZIPKIN = "zipkin"
28 NONE = "none"
31class OtelStatus(StrEnum):
32 """Span status codes."""
34 OK = "OK"
35 ERROR = "ERROR"
38class SpanKind(StrEnum):
39 """Span kind for semantic conventions."""
41 INTERNAL = "internal"
42 CLIENT = "client"
43 SERVER = "server"
44 PRODUCER = "producer"
45 CONSUMER = "consumer"
48# ---- OtelConfig ----
51@dataclass
52class OtelConfig:
53 """OpenTelemetry configuration."""
55 service_name: str = "agentos"
56 service_version: str = ""
57 exporter: OTelExporter = OTelExporter.CONSOLE
58 endpoint: str = "http://localhost:4318/v1/traces"
59 metrics_endpoint: str = "http://localhost:4318/v1/metrics"
60 resource_attrs: dict[str, str] = field(default_factory=dict)
61 sample_rate: float = 1.0
62 batch_timeout_ms: int = 5000
63 max_span_attributes: int = 128
64 disabled: bool = False
65 api_key: str = ""
66 zipkin_endpoint: str = "http://localhost:9411/api/v2/spans"
68 def with_env_overrides(self) -> "OtelConfig":
69 if v := os.getenv("OTEL_SERVICE_NAME"):
70 self.service_name = v
71 if v := os.getenv("OTEL_EXPORTER_OTLP_ENDPOINT"):
72 self.endpoint = v
73 if v := os.getenv("OTEL_EXPORTER_ZIPKIN_ENDPOINT"):
74 self.zipkin_endpoint = v
75 if os.getenv("OTEL_SDK_DISABLED"):
76 self.disabled = True
77 return self
80# ---- SpanHandle ----
83def _normalize_value(v: Any) -> Any:
84 if isinstance(v, (str, int, float, bool)):
85 return v
86 if isinstance(v, (list, tuple)):
87 return [_normalize_value(x) for x in v]
88 return str(v)
91class SpanHandle:
92 """Wrapper around OTel span for attribute/event/exception API."""
94 def __init__(self, span: Any):
95 self._span = span
97 def set_attribute(self, key: str, value: Any) -> None:
98 self._span.set_attribute(key, _normalize_value(value))
100 def set_attributes(self, attrs: dict[str, Any]) -> None:
101 for k, v in attrs.items():
102 self.set_attribute(k, v)
104 def add_event(self, name: str, attributes: dict[str, Any] | None = None) -> None:
105 self._span.add_event(
106 name, attributes={k: _normalize_value(v) for k, v in (attributes or {}).items()}
107 )
109 def record_exception(self, exception: Exception) -> None:
110 self._span.record_exception(exception)
112 def set_status(self, status: OtelStatus, description: str = "") -> None:
113 from opentelemetry import trace as otel_trace
115 code = otel_trace.StatusCode.OK if status == OtelStatus.OK else otel_trace.StatusCode.ERROR
116 self._span.set_status(otel_trace.Status(code, description))
119# ---- OtelTracer ----
122class OtelTracer:
123 """OpenTelemetry tracer with span management and W3C context propagation.
125 Usage:
126 OtelTracer.init(OtelConfig(service_name="my-agent"))
128 with OtelTracer.span("llm_call", kind=SpanKind.CLIENT) as span:
129 span.set_attribute("model", "gpt-4")
130 result = llm.generate(prompt)
132 @OtelTracer.trace("process")
133 async def process(input): ...
134 """
136 _config: OtelConfig | None = None
137 _tracer_provider: Any = None
138 _initialized: bool = False
140 @classmethod
141 def init(cls, config: OtelConfig | None = None) -> None:
142 if config is None:
143 config = OtelConfig().with_env_overrides()
144 else:
145 config = config.with_env_overrides()
147 cls._config = config
148 if config.disabled:
149 cls._initialized = True
150 return
152 try:
153 cls._init_sdk(config)
154 cls._initialized = True
155 logger.info(
156 "OtelTracer initialized: service=%s exporter=%s",
157 config.service_name,
158 config.exporter.value,
159 )
160 except ImportError:
161 logger.warning("opentelemetry packages not installed - noop tracer")
162 cls._initialized = True
163 except Exception as e:
164 logger.error("OTel init failed: %s - noop", e)
165 cls._initialized = True
167 @classmethod
168 def _init_sdk(cls, config: OtelConfig):
169 from opentelemetry.sdk.resources import SERVICE_NAME, SERVICE_VERSION, Resource
170 from opentelemetry.sdk.trace import TracerProvider
171 from opentelemetry.sdk.trace.export import BatchSpanProcessor
172 from opentelemetry.trace import set_tracer_provider
174 resource = Resource.create(
175 {
176 SERVICE_NAME: config.service_name,
177 SERVICE_VERSION: config.service_version,
178 **config.resource_attrs,
179 }
180 )
182 provider = TracerProvider(resource=resource)
183 exporter = cls._build_exporter(config)
184 if exporter:
185 provider.add_span_processor(
186 BatchSpanProcessor(
187 exporter,
188 schedule_delay_millis=config.batch_timeout_ms,
189 max_export_batch_size=512,
190 )
191 )
192 set_tracer_provider(provider)
193 cls._tracer_provider = provider
195 @classmethod
196 def _build_exporter(cls, config: OtelConfig):
197 if config.exporter == OTelExporter.OTLP_HTTP:
198 from opentelemetry.exporter.otlp.proto.http.trace_exporter import (
199 OTLPSpanExporter,
200 )
202 return OTLPSpanExporter(endpoint=config.endpoint)
203 elif config.exporter == OTelExporter.OTLP_GRPC:
204 from opentelemetry.exporter.otlp.proto.grpc.trace_exporter import (
205 OTLPSpanExporter,
206 )
208 return OTLPSpanExporter(endpoint=config.endpoint, insecure=True)
209 elif config.exporter == OTelExporter.CONSOLE:
210 from opentelemetry.sdk.trace.export import ConsoleSpanExporter
212 return ConsoleSpanExporter()
213 elif config.exporter == OTelExporter.ZIPKIN:
214 from opentelemetry.exporter.zipkin.proto.http import ZipkinExporter
216 return ZipkinExporter(endpoint=config.zipkin_endpoint)
217 return None
219 @classmethod
220 def get_tracer(cls, name: str = "agentos") -> Any:
221 if not cls._initialized:
222 cls.init()
223 from opentelemetry import trace as otel_trace
225 return otel_trace.get_tracer(name)
227 @classmethod
228 @contextmanager
229 def span(
230 cls,
231 name: str,
232 kind: SpanKind = SpanKind.INTERNAL,
233 attributes: dict[str, Any] | None = None,
234 parent: Any = None,
235 ):
236 if not cls._initialized:
237 cls.init()
239 tracer = cls.get_tracer()
240 kind_map = {
241 SpanKind.INTERNAL: 0,
242 SpanKind.CLIENT: 3,
243 SpanKind.SERVER: 2,
244 SpanKind.PRODUCER: 4,
245 SpanKind.CONSUMER: 5,
246 }
247 sk = kind_map.get(kind, 0)
249 from opentelemetry import trace as otel_trace
251 span = tracer.start_span(
252 name,
253 kind=otel_trace.SpanKind(sk),
254 attributes={k: _normalize_value(v) for k, v in (attributes or {}).items()},
255 )
257 try:
258 yield SpanHandle(span)
259 except Exception:
260 span.set_status(otel_trace.Status(otel_trace.StatusCode.ERROR))
261 raise
262 finally:
263 span.end()
265 @classmethod
266 def trace(
267 cls,
268 name: str = "",
269 kind: SpanKind = SpanKind.INTERNAL,
270 extract_attrs: Callable | None = None,
271 ):
272 span_name = name
274 def decorator(func):
275 nonlocal span_name
276 import asyncio
278 if not span_name:
279 span_name = func.__name__
280 is_async = asyncio.iscoroutinefunction(func)
282 @wraps(func)
283 async def async_wrapper(*args, **kwargs):
284 attrs = extract_attrs(*args, **kwargs) if extract_attrs else {}
285 with cls.span(span_name, kind=kind, attributes=attrs) as span:
286 try:
287 result = await func(*args, **kwargs)
288 span.set_attribute("status", "ok")
289 return result
290 except Exception as e:
291 span.record_exception(e)
292 span.set_status(OtelStatus.ERROR)
293 raise
295 @wraps(func)
296 def sync_wrapper(*args, **kwargs):
297 attrs = extract_attrs(*args, **kwargs) if extract_attrs else {}
298 with cls.span(span_name, kind=kind, attributes=attrs) as span:
299 try:
300 result = func(*args, **kwargs)
301 span.set_attribute("status", "ok")
302 return result
303 except Exception as e:
304 span.record_exception(e)
305 span.set_status(OtelStatus.ERROR)
306 raise
308 return async_wrapper if is_async else sync_wrapper
310 return decorator
312 @classmethod
313 @asynccontextmanager
314 async def async_span(cls, name: str, kind: SpanKind = SpanKind.INTERNAL, **attrs):
315 with cls.span(name, kind=kind, attributes=attrs) as span:
316 yield span
318 @classmethod
319 def shutdown(cls):
320 if cls._tracer_provider:
321 try:
322 cls._tracer_provider.shutdown()
323 except Exception:
324 pass
325 cls._tracer_provider = None
326 cls._initialized = False
329# ---- OtelMeter ----
332class OtelMeter:
333 """Bridge MetricsCollector to OpenTelemetry metrics."""
335 _meter: Any = None
336 _instruments: dict[str, Any] = {}
338 @classmethod
339 def init(cls, config: OtelConfig | None = None):
340 if config is None:
341 config = OtelConfig().with_env_overrides()
342 if config.disabled:
343 return
344 try:
345 from opentelemetry import metrics
346 from opentelemetry.exporter.otlp.proto.http.metric_exporter import (
347 OTLPMetricExporter,
348 )
349 from opentelemetry.sdk.metrics import MeterProvider
350 from opentelemetry.sdk.metrics.export import PeriodicExportingMetricReader
351 from opentelemetry.sdk.resources import SERVICE_NAME, Resource
353 resource = Resource.create({SERVICE_NAME: config.service_name})
354 exporter = OTLPMetricExporter(endpoint=config.metrics_endpoint)
355 reader = PeriodicExportingMetricReader(exporter)
356 provider = MeterProvider(resource=resource, metric_readers=[reader])
357 metrics.set_meter_provider(provider)
358 cls._meter = metrics.get_meter(config.service_name, config.service_version)
359 logger.info("OtelMeter initialized: endpoint=%s", config.metrics_endpoint)
360 except ImportError:
361 logger.warning("opentelemetry metrics packages not installed")
362 except Exception as e:
363 logger.error("OtelMeter init failed: %s", e)
365 @classmethod
366 def record_counter(cls, name: str, value: float, attrs: dict[str, str] | None = None):
367 if cls._meter is None:
368 return
369 if name not in cls._instruments:
370 cls._instruments[name] = cls._meter.create_counter(name, description="")
371 cls._instruments[name].add(value, attributes=attrs or {})
373 @classmethod
374 def record_histogram(cls, name: str, value: float, attrs: dict[str, str] | None = None):
375 if cls._meter is None:
376 return
377 key = f"hist_{name}"
378 if key not in cls._instruments:
379 cls._instruments[key] = cls._meter.create_histogram(name, description="")
380 cls._instruments[key].record(value, attributes=attrs or {})
382 @classmethod
383 def record_gauge(cls, name: str, value: float, attrs: dict[str, str] | None = None):
384 if cls._meter is None:
385 return
386 key = f"gauge_{name}"
387 if key not in cls._instruments:
388 cls._instruments[key] = cls._meter.create_up_down_counter(name, description="")
389 cls._instruments[key].add(value, attributes=attrs or {})
391 @classmethod
392 def bridge(cls, collector: "MetricsCollector"):
393 """Wire MetricsCollector snapshots to OtelMeter on flush."""
394 original_flush = collector.flush
396 def _hooked_flush():
397 snapshots = original_flush()
398 for snap in snapshots:
399 attrs = dict(snap.tags) if hasattr(snap, "tags") else {}
400 if snap.type == "counter":
401 cls.record_counter(snap.name, snap.value, attrs)
402 elif snap.type == "histogram":
403 cls.record_histogram(snap.name, snap.value, attrs)
404 elif snap.type == "gauge":
405 cls.record_gauge(snap.name, snap.value, attrs)
406 return snapshots
408 collector.flush = _hooked_flush # type: ignore[method-assign]
411# ---- OtelMiddleware ----
414class OtelMiddleware:
415 """W3C TraceContext propagation for multi-agent pipelines."""
417 @staticmethod
418 def inject_context(headers: dict[str, str] | None = None) -> dict[str, str]:
419 if headers is None:
420 headers = {}
421 try:
422 from opentelemetry import propagate
424 propagate.inject(headers)
425 except Exception:
426 pass
427 return headers
429 @staticmethod
430 def extract_context(headers: dict[str, str] | None = None) -> None:
431 if headers is None:
432 return
433 try:
434 from opentelemetry import context, propagate
436 ctx = propagate.extract(headers)
437 context.attach(ctx)
438 except Exception:
439 pass
441 @staticmethod
442 def get_trace_id() -> str:
443 try:
444 from opentelemetry import trace as otel_trace
446 span = otel_trace.get_current_span()
447 ctx = span.get_span_context()
448 if ctx.is_valid:
449 return format(ctx.trace_id, "032x")
450 except Exception:
451 pass
452 return ""
454 @staticmethod
455 def get_span_id() -> str:
456 try:
457 from opentelemetry import trace as otel_trace
459 span = otel_trace.get_current_span()
460 ctx = span.get_span_context()
461 if ctx.is_valid:
462 return format(ctx.span_id, "016x")
463 except Exception:
464 pass
465 return ""