|
| 1 | +"""可选可观测性接入:OpenTelemetry → Langfuse / 任意 OTLP 后端。 |
| 2 | +
|
| 3 | +默认**关闭**;仅当配置了 Langfuse 或通用 OTLP 环境变量时才启用,零侵入、缺依赖优雅降级。 |
| 4 | +
|
| 5 | +启用方式(二选一): |
| 6 | + 1) Langfuse:LANGFUSE_PUBLIC_KEY + LANGFUSE_SECRET_KEY(+ LANGFUSE_BASE_URL,自建必填; |
| 7 | + 不填默认 https://cloud.langfuse.com)。自动拼 OTLP 端点 /api/public/otel/v1/traces + Basic 鉴权。 |
| 8 | + 2) 通用 OTLP:OTEL_EXPORTER_OTLP_ENDPOINT(+ 可选 OTEL_EXPORTER_OTLP_HEADERS="k=v,k2=v2")。 |
| 9 | +
|
| 10 | +采集内容: |
| 11 | + - AgnoInstrumentor(核心):agent / 工具调用 / LLM 生成 span,含 token 用量; |
| 12 | + - 可选全栈(OBS_INSTRUMENT_STACK=true):FastAPI 请求 / SQLAlchemy / httpx 出站 span。 |
| 13 | +
|
| 14 | +依赖(未装则本模块自动跳过,不影响主程序): |
| 15 | + opentelemetry-sdk / opentelemetry-exporter-otlp-proto-http / openinference-instrumentation-agno |
| 16 | + (全栈可选:opentelemetry-instrumentation-{fastapi,sqlalchemy,httpx}) |
| 17 | +
|
| 18 | +备注:MCP 路径在独立 asyncio task 内跑 agent(见 chat.mcp_bridge);asyncio.create_task 会自动 |
| 19 | +复制当前 contextvars(OTel context 即基于 contextvars),故父子 span 链路自动传播,无需特殊处理。 |
| 20 | +""" |
| 21 | + |
| 22 | +from __future__ import annotations |
| 23 | + |
| 24 | +import base64 |
| 25 | +import os |
| 26 | +from typing import TYPE_CHECKING |
| 27 | + |
| 28 | +from utils.log_util import logger |
| 29 | + |
| 30 | +if TYPE_CHECKING: |
| 31 | + from fastapi import FastAPI |
| 32 | + |
| 33 | +_initialized = False |
| 34 | + |
| 35 | + |
| 36 | +def _resolve_otlp_target() -> tuple[str, dict[str, str]] | None: |
| 37 | + """从环境变量解析 (traces_endpoint, headers);未配置返回 None。""" |
| 38 | + pub = os.environ.get('LANGFUSE_PUBLIC_KEY') |
| 39 | + sec = os.environ.get('LANGFUSE_SECRET_KEY') |
| 40 | + generic = os.environ.get('OTEL_EXPORTER_OTLP_ENDPOINT') |
| 41 | + |
| 42 | + if pub and sec: |
| 43 | + base = (os.environ.get('LANGFUSE_BASE_URL') or os.environ.get('LANGFUSE_HOST') or 'https://cloud.langfuse.com').rstrip('/') |
| 44 | + endpoint = f'{base}/api/public/otel/v1/traces' |
| 45 | + auth = base64.b64encode(f'{pub}:{sec}'.encode()).decode() |
| 46 | + return endpoint, {'Authorization': f'Basic {auth}'} |
| 47 | + |
| 48 | + if generic: |
| 49 | + endpoint = generic.rstrip('/') |
| 50 | + if not endpoint.endswith('/v1/traces'): |
| 51 | + endpoint = f'{endpoint}/v1/traces' |
| 52 | + headers: dict[str, str] = {} |
| 53 | + raw = os.environ.get('OTEL_EXPORTER_OTLP_HEADERS', '') |
| 54 | + for pair in raw.split(','): |
| 55 | + if '=' in pair: |
| 56 | + k, v = pair.split('=', 1) |
| 57 | + headers[k.strip()] = v.strip() |
| 58 | + return endpoint, headers |
| 59 | + |
| 60 | + return None |
| 61 | + |
| 62 | + |
| 63 | +def init_observability(app: FastAPI | None = None) -> bool: |
| 64 | + """按环境变量启用 OTel → Langfuse/OTLP。未配置或缺依赖则返回 False(不启用,不报错)。幂等。""" |
| 65 | + global _initialized |
| 66 | + if _initialized: |
| 67 | + return True |
| 68 | + |
| 69 | + target = _resolve_otlp_target() |
| 70 | + if target is None: |
| 71 | + return False # 未配置 → 默认关 |
| 72 | + endpoint, headers = target |
| 73 | + |
| 74 | + try: |
| 75 | + import opentelemetry.trace as ot |
| 76 | + from opentelemetry.exporter.otlp.proto.http.trace_exporter import OTLPSpanExporter |
| 77 | + from opentelemetry.sdk.resources import Resource |
| 78 | + from opentelemetry.sdk.trace import TracerProvider |
| 79 | + from opentelemetry.sdk.trace.export import BatchSpanProcessor |
| 80 | + except Exception as e: |
| 81 | + logger.warning(f'[otel] 已配置可观测性但缺 opentelemetry 依赖,跳过:{e}') |
| 82 | + return False |
| 83 | + |
| 84 | + service = os.environ.get('OTEL_SERVICE_NAME', 'ezdata') |
| 85 | + environment = os.environ.get('OTEL_DEPLOYMENT_ENVIRONMENT') or os.environ.get('APP_ENV', 'dev') |
| 86 | + provider = TracerProvider( |
| 87 | + resource=Resource.create({'service.name': service, 'deployment.environment': environment}) |
| 88 | + ) |
| 89 | + provider.add_span_processor(BatchSpanProcessor(OTLPSpanExporter(endpoint=endpoint, headers=headers or None))) |
| 90 | + ot.set_tracer_provider(provider) |
| 91 | + |
| 92 | + # 核心:agno instrumentor(agent/工具/LLM 生成 span) |
| 93 | + try: |
| 94 | + from openinference.instrumentation.agno import AgnoInstrumentor |
| 95 | + |
| 96 | + AgnoInstrumentor().instrument(tracer_provider=provider) |
| 97 | + except Exception as e: |
| 98 | + logger.warning(f'[otel] AgnoInstrumentor 未启用(缺 openinference-instrumentation-agno?):{e}') |
| 99 | + |
| 100 | + # 可选:全栈 span(默认关,OBS_INSTRUMENT_STACK=true 开) |
| 101 | + if (os.environ.get('OBS_INSTRUMENT_STACK', '') or '').lower() in ('1', 'true', 'yes'): |
| 102 | + _instrument_stack(app, provider) |
| 103 | + |
| 104 | + _initialized = True |
| 105 | + logger.info(f'[otel] 可观测性已启用 → {endpoint}(service={service}, env={environment})') |
| 106 | + return True |
| 107 | + |
| 108 | + |
| 109 | +def _instrument_stack(app: FastAPI | None, provider) -> None: |
| 110 | + """可选:FastAPI / SQLAlchemy / httpx 自动埋点。单个失败只告警、不影响其余。""" |
| 111 | + if app is not None: |
| 112 | + try: |
| 113 | + from opentelemetry.instrumentation.fastapi import FastAPIInstrumentor |
| 114 | + |
| 115 | + FastAPIInstrumentor.instrument_app(app, tracer_provider=provider) |
| 116 | + except Exception as e: |
| 117 | + logger.warning(f'[otel] FastAPI 埋点跳过:{e}') |
| 118 | + try: |
| 119 | + from opentelemetry.instrumentation.sqlalchemy import SQLAlchemyInstrumentor |
| 120 | + |
| 121 | + SQLAlchemyInstrumentor().instrument(tracer_provider=provider) |
| 122 | + except Exception as e: |
| 123 | + logger.warning(f'[otel] SQLAlchemy 埋点跳过:{e}') |
| 124 | + try: |
| 125 | + from opentelemetry.instrumentation.httpx import HTTPXClientInstrumentor |
| 126 | + |
| 127 | + HTTPXClientInstrumentor().instrument(tracer_provider=provider) |
| 128 | + except Exception as e: |
| 129 | + logger.warning(f'[otel] httpx 埋点跳过:{e}') |
0 commit comments