Coverage for agentos/observability/otel_bridge.py: 0%

298 statements  

« prev     ^ index     » next       coverage.py v7.14.3, created at 2026-07-06 12:29 +0800

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

2 

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 

11 

12if TYPE_CHECKING: 

13 from agentos.observability.metrics import MetricsCollector 

14 

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

16 

17 

18# ---- Enums ---- 

19 

20 

21class OTelExporter(StrEnum): 

22 """OpenTelemetry exporter backend.""" 

23 

24 OTLP_HTTP = "otlp_http" 

25 OTLP_GRPC = "otlp_grpc" 

26 CONSOLE = "console" 

27 ZIPKIN = "zipkin" 

28 NONE = "none" 

29 

30 

31class OtelStatus(StrEnum): 

32 """Span status codes.""" 

33 

34 OK = "OK" 

35 ERROR = "ERROR" 

36 

37 

38class SpanKind(StrEnum): 

39 """Span kind for semantic conventions.""" 

40 

41 INTERNAL = "internal" 

42 CLIENT = "client" 

43 SERVER = "server" 

44 PRODUCER = "producer" 

45 CONSUMER = "consumer" 

46 

47 

48# ---- OtelConfig ---- 

49 

50 

51@dataclass 

52class OtelConfig: 

53 """OpenTelemetry configuration.""" 

54 

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" 

67 

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 

78 

79 

80# ---- SpanHandle ---- 

81 

82 

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) 

89 

90 

91class SpanHandle: 

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

93 

94 def __init__(self, span: Any): 

95 self._span = span 

96 

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

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

99 

100 def set_attributes(self, attrs: dict[str, Any]) -> None: 

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

102 self.set_attribute(k, v) 

103 

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 ) 

108 

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

110 self._span.record_exception(exception) 

111 

112 def set_status(self, status: OtelStatus, description: str = "") -> None: 

113 from opentelemetry import trace as otel_trace 

114 

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

117 

118 

119# ---- OtelTracer ---- 

120 

121 

122class OtelTracer: 

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

124 

125 Usage: 

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

127 

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

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

130 result = llm.generate(prompt) 

131 

132 @OtelTracer.trace("process") 

133 async def process(input): ... 

134 """ 

135 

136 _config: OtelConfig | None = None 

137 _tracer_provider: Any = None 

138 _initialized: bool = False 

139 

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

146 

147 cls._config = config 

148 if config.disabled: 

149 cls._initialized = True 

150 return 

151 

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 

166 

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 

173 

174 resource = Resource.create( 

175 { 

176 SERVICE_NAME: config.service_name, 

177 SERVICE_VERSION: config.service_version, 

178 **config.resource_attrs, 

179 } 

180 ) 

181 

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 

194 

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 ) 

201 

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 ) 

207 

208 return OTLPSpanExporter(endpoint=config.endpoint, insecure=True) 

209 elif config.exporter == OTelExporter.CONSOLE: 

210 from opentelemetry.sdk.trace.export import ConsoleSpanExporter 

211 

212 return ConsoleSpanExporter() 

213 elif config.exporter == OTelExporter.ZIPKIN: 

214 from opentelemetry.exporter.zipkin.proto.http import ZipkinExporter 

215 

216 return ZipkinExporter(endpoint=config.zipkin_endpoint) 

217 return None 

218 

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 

224 

225 return otel_trace.get_tracer(name) 

226 

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

238 

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) 

248 

249 from opentelemetry import trace as otel_trace 

250 

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 ) 

256 

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

264 

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 

273 

274 def decorator(func): 

275 nonlocal span_name 

276 import asyncio 

277 

278 if not span_name: 

279 span_name = func.__name__ 

280 is_async = asyncio.iscoroutinefunction(func) 

281 

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 

294 

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 

307 

308 return async_wrapper if is_async else sync_wrapper 

309 

310 return decorator 

311 

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 

317 

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 

327 

328 

329# ---- OtelMeter ---- 

330 

331 

332class OtelMeter: 

333 """Bridge MetricsCollector to OpenTelemetry metrics.""" 

334 

335 _meter: Any = None 

336 _instruments: dict[str, Any] = {} 

337 

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 

352 

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) 

364 

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

372 

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

381 

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

390 

391 @classmethod 

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

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

394 original_flush = collector.flush 

395 

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 

407 

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

409 

410 

411# ---- OtelMiddleware ---- 

412 

413 

414class OtelMiddleware: 

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

416 

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 

423 

424 propagate.inject(headers) 

425 except Exception: 

426 pass 

427 return headers 

428 

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 

435 

436 ctx = propagate.extract(headers) 

437 context.attach(ctx) 

438 except Exception: 

439 pass 

440 

441 @staticmethod 

442 def get_trace_id() -> str: 

443 try: 

444 from opentelemetry import trace as otel_trace 

445 

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

453 

454 @staticmethod 

455 def get_span_id() -> str: 

456 try: 

457 from opentelemetry import trace as otel_trace 

458 

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