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