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

1"""AgentOS OpenTelemetry - OTLP/Jaeger/Zipkin trace/metric export (v1.14.6).""" 

2 

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 

10 

11if TYPE_CHECKING: 

12 from agentos.observability.metrics import MetricsCollector 

13 

14logger = logging.getLogger("agentos.otel") 

15 

16 

17# ---- Enums ---- 

18 

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" 

26 

27 

28class OtelStatus(str, Enum): 

29 """Span status codes.""" 

30 OK = "OK" 

31 ERROR = "ERROR" 

32 

33 

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" 

41 

42 

43# ---- OtelConfig ---- 

44 

45@dataclass 

46class OtelConfig: 

47 """OpenTelemetry configuration.""" 

48 

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" 

61 

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 

72 

73 

74# ---- SpanHandle ---- 

75 

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) 

82 

83 

84class SpanHandle: 

85 """Wrapper around OTel span for attribute/event/exception API.""" 

86 

87 def __init__(self, span: Any): 

88 self._span = span 

89 

90 def set_attribute(self, key: str, value: Any) -> None: 

91 self._span.set_attribute(key, _normalize_value(value)) 

92 

93 def set_attributes(self, attrs: Dict[str, Any]) -> None: 

94 for k, v in attrs.items(): 

95 self.set_attribute(k, v) 

96 

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 }) 

101 

102 def record_exception(self, exception: Exception) -> None: 

103 self._span.record_exception(exception) 

104 

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)) 

109 

110 

111# ---- OtelTracer ---- 

112 

113class OtelTracer: 

114 """OpenTelemetry tracer with span management and W3C context propagation. 

115 

116 Usage: 

117 OtelTracer.init(OtelConfig(service_name="my-agent")) 

118 

119 with OtelTracer.span("llm_call", kind=SpanKind.CLIENT) as span: 

120 span.set_attribute("model", "gpt-4") 

121 result = llm.generate(prompt) 

122 

123 @OtelTracer.trace("process") 

124 async def process(input): ... 

125 """ 

126 

127 _config: Optional[OtelConfig] = None 

128 _tracer_provider: Any = None 

129 _initialized: bool = False 

130 

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() 

137 

138 cls._config = config 

139 if config.disabled: 

140 cls._initialized = True 

141 return 

142 

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 

156 

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 

163 

164 resource = Resource.create({ 

165 SERVICE_NAME: config.service_name, 

166 SERVICE_VERSION: config.service_version, 

167 **config.resource_attrs, 

168 }) 

169 

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 

182 

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 

202 

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) 

209 

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() 

221 

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) 

231 

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 ) 

240 

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() 

248 

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 

257 

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) 

264 

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 

277 

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 

290 

291 return async_wrapper if is_async else sync_wrapper 

292 

293 return decorator 

294 

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 

300 

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 

310 

311 

312# ---- OtelMeter ---- 

313 

314class OtelMeter: 

315 """Bridge MetricsCollector to OpenTelemetry metrics.""" 

316 

317 _meter: Any = None 

318 _instruments: Dict[str, Any] = {} 

319 

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 ) 

334 

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) 

346 

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 {}) 

356 

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 {}) 

367 

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 {}) 

380 

381 @classmethod 

382 def bridge(cls, collector: "MetricsCollector"): 

383 """Wire MetricsCollector snapshots to OtelMeter on flush.""" 

384 original_flush = collector.flush 

385 

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 

397 

398 collector.flush = _hooked_flush # type: ignore[method-assign] 

399 

400 

401# ---- OtelMiddleware ---- 

402 

403class OtelMiddleware: 

404 """W3C TraceContext propagation for multi-agent pipelines.""" 

405 

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 

416 

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 

427 

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 "" 

439 

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 ""