Skip to content
20 changes: 15 additions & 5 deletions litellm/integrations/otel/langfuse_logger.py
Original file line number Diff line number Diff line change
@@ -1,6 +1,9 @@
from collections.abc import AsyncGenerator, AsyncIterator, Callable, Mapping
from typing import TYPE_CHECKING, Final

from opentelemetry.sdk.trace import ReadableSpan
from opentelemetry.trace import Span

from litellm._logging import verbose_logger
from litellm.integrations.otel.logger import OpenTelemetryV2
from litellm.integrations.otel.mappers.langfuse import (
Expand All @@ -19,15 +22,22 @@

class LangfuseOpenTelemetryV2(OpenTelemetryV2):
"""Names the trace from the request. Langfuse reads ``langfuse.trace.name`` off the root observation,
and the proxy's root span is still recording when the LLM call starts."""
and the proxy's root span is still recording when the LLM call starts. Only a root created by this
logger's own provider is stamped; any other root is exported to destinations that are not Langfuse."""

def log_pre_api_call(self, model: str, messages: object, kwargs: Mapping[str, object]) -> None:
root: Final = request_root_span()
root: Final = self._owned_recording_root()
name: Final = caller_trace_name(kwargs)
if root is not None and root.is_recording() and name is not None:
if root is not None and name is not None:
root.set_attribute(LANGFUSE_TRACE_NAME, name)
super().log_pre_api_call(model, messages, kwargs)

def _owned_recording_root(self) -> Span | None:
root: Final = request_root_span()
if root is None or not root.is_recording() or not isinstance(root, ReadableSpan):
return None
return root if root.resource is self.tracer_provider.resource else None


class LangfuseContentOpenTelemetryV2(LangfuseOpenTelemetryV2):
"""Stamps the request's input and output on the root observation while it is still recording.
Expand Down Expand Up @@ -59,8 +69,8 @@ async def async_post_call_streaming_iterator_hook(
self._stamp_root_io(request_data, lambda: stream_output(tuple(relayed), request_data))

def _stamp_root_io(self, data: Mapping[str, object], render_output: Callable[[], str | None]) -> None:
root: Final = request_root_span()
if root is None or not root.is_recording():
root: Final = self._owned_recording_root()
if root is None:
return
try:
output: Final = render_output()
Expand Down
34 changes: 27 additions & 7 deletions litellm/integrations/otel/logger.py
Original file line number Diff line number Diff line change
Expand Up @@ -206,6 +206,11 @@ def tracer_provider(self) -> TracerProvider:
"""The provider this logger emits through, read-only to its callers."""
return self._tracer_provider

@property
def serves_generic_collector(self) -> bool:
"""Every exporter is the operator's ``OTEL_*`` destination, as for ``otel`` and mapper-only presets."""
return all(spec.owner is None for spec in self.config.exporters)

def _init_metrics(self, meter_provider: "MeterProvider | None") -> "GenAIMetricRecorder | None":
"""Create the six GenAI histograms when metrics are enabled, else ``None``.

Expand Down Expand Up @@ -259,9 +264,19 @@ def _init_otel_logger_on_litellm_proxy(self) -> None:
self._register_in_callback_list(litellm._async_failure_callback)
except Exception:
pass
if getattr(proxy_server, "open_telemetry_logger", None) is None:
if self._outranks_for_proxy_slot(getattr(proxy_server, "open_telemetry_logger", None)):
setattr(proxy_server, "open_telemetry_logger", self)

def _outranks_for_proxy_slot(self, holder: object) -> bool:
"""The collector's logger owns the slot; a vendor preset holds it only until that logger is built."""
if holder is None:
return True
return (
self.serves_generic_collector
and isinstance(holder, OpenTelemetryV2)
and not holder.serves_generic_collector
)

# ====================================================================== #
# LLM-call callbacks — the span is opened at the ``pre_call`` boundary and
# closed here. See ``log_pre_api_call``.
Expand Down Expand Up @@ -842,10 +857,12 @@ def select_global_otel_v2_logger(
``proxy_server.open_telemetry_logger``), and every other v2 entry point —
guardrail, identity seeding, phase spans — already routes through that same
``registered`` owner. Reuse it here too so the global provider has one source
of truth instead of a second, independently-derived guess; this is the logger
a preset (arize, langfuse, …) folds the ``OTEL_*`` base exporter and its own
exporter into, so the FastAPI server span and the gen-ai spans share one
provider and one trace.
of truth instead of a second, independently-derived guess. Each logger exports
only to the destination its callback owns (``otel`` serves ``OTEL_*``, a preset
serves its own backend), so the FastAPI server span and the request root land
at the ``otel`` callback's collector when one is configured and at the first
preset's backend otherwise, whatever the callback order; a preset stamps that
root only when its own provider created it.

Fall back to ``in_memory_loggers`` for the SDK path, where no proxy global is
set (selecting from there, not ``service_callback``, which a preset logger does
Expand All @@ -855,8 +872,11 @@ def select_global_otel_v2_logger(
"""
if registered is not None:
return registered
existing: Final = next((cb for cb in in_memory_loggers if isinstance(cb, OpenTelemetryV2)), None)
return existing if existing is not None else OpenTelemetryV2()
v2_loggers: Final = tuple(cb for cb in in_memory_loggers if isinstance(cb, OpenTelemetryV2))
generic: Final = next((cb for cb in v2_loggers if cb.serves_generic_collector), None)
if generic is not None:
return generic
return v2_loggers[0] if v2_loggers else OpenTelemetryV2()


def publish_global_otel_v2_provider(
Expand Down
1 change: 0 additions & 1 deletion litellm/integrations/otel/presets/agentops.py
Original file line number Diff line number Diff line change
Expand Up @@ -53,7 +53,6 @@ def agentops_preset(
return base.model_copy(
update={
"exporters": [
*base.exporters,
ExporterSpec(
kind=_AGENTOPS_EXPORTER_KIND,
endpoint=_AGENTOPS_ENDPOINT,
Expand Down
1 change: 0 additions & 1 deletion litellm/integrations/otel/presets/arize.py
Original file line number Diff line number Diff line change
Expand Up @@ -35,7 +35,6 @@ def arize_preset(
return base.model_copy(
update={
"exporters": [
*base.exporters,
ExporterSpec(
kind=arize_cfg.protocol or "otlp_grpc",
endpoint=arize_cfg.endpoint or "https://otlp.arize.com/v1",
Expand Down
8 changes: 2 additions & 6 deletions litellm/integrations/otel/presets/langfuse.py
Original file line number Diff line number Diff line change
Expand Up @@ -10,10 +10,7 @@
ExporterSpec,
OpenTelemetryV2Config,
)
from litellm.integrations.otel.presets.utils import (
credential_gated_exporters,
ensure_mappers,
)
from litellm.integrations.otel.presets.utils import ensure_mappers
from litellm.types.utils import StandardCallbackDynamicParams


Expand All @@ -31,15 +28,14 @@ def langfuse_preset(
raise
return base.model_copy(
update={ # mutable-ok: pydantic model_copy takes a plain update mapping
"exporters": credential_gated_exporters(base.exporters, ExporterOwner.LANGFUSE_OTEL),
"exporters": [ExporterSpec(owner=ExporterOwner.LANGFUSE_OTEL, requires_headers=True)],
"mapper_names": mappers,
}
)
kind: Final = cfg.exporter if isinstance(cfg.exporter, str) else "otlp_http"
return base.model_copy(
update={
"exporters": [
*base.exporters,
ExporterSpec(
kind=kind,
endpoint=cfg.endpoint,
Expand Down
1 change: 0 additions & 1 deletion litellm/integrations/otel/presets/levo.py
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,6 @@ def levo_preset(
return base.model_copy(
update={
"exporters": [
*base.exporters,
ExporterSpec(
kind="otlp_http",
endpoint=cfg.endpoint,
Expand Down
1 change: 0 additions & 1 deletion litellm/integrations/otel/presets/newrelic.py
Original file line number Diff line number Diff line change
Expand Up @@ -54,7 +54,6 @@ def newrelic_preset(
return base.model_copy(
update={
"exporters": [
*base.exporters,
ExporterSpec(
kind="otlp_http",
endpoint=endpoint,
Expand Down
1 change: 0 additions & 1 deletion litellm/integrations/otel/presets/phoenix.py
Original file line number Diff line number Diff line change
Expand Up @@ -69,7 +69,6 @@ def phoenix_preset(
return base.model_copy(
update={
"exporters": [
*base.exporters,
ExporterSpec(
kind=cfg.protocol if hasattr(cfg, "protocol") else "otlp_http",
endpoint=cfg.endpoint,
Expand Down
31 changes: 0 additions & 31 deletions litellm/integrations/otel/presets/utils.py
Original file line number Diff line number Diff line change
Expand Up @@ -3,8 +3,6 @@
from collections.abc import Iterable
from typing import Final

from litellm.integrations.otel.model.config import ExporterOwner, ExporterSpec


def ensure_mappers(mapper_names: Iterable[str], *names: str) -> list[str]:
"""Return ``mapper_names`` with each of ``names`` appended if not already present.
Expand All @@ -17,32 +15,3 @@ def ensure_mappers(mapper_names: Iterable[str], *names: str) -> list[str]:
if name not in result:
result.append(name)
return result


def credential_gated_exporters(
exporters: "Iterable[ExporterSpec]", owner: "ExporterOwner"
) -> "tuple[ExporterSpec, ...]":
"""``exporters`` with the operator's destination replaced by a header-gated one.

Used when a credential-mandatory backend is asked to build without the operator's
own credentials, so only key/team destinations receive spans. Two things have to
happen for that to mean "export nowhere": the placeholder console spec that
``OpenTelemetryV2Config`` folds in for an empty exporter list is dropped, or every
span would be printed to stdout, and the gated spec keeps the owner so the
override filter still recognises which backend this provider speaks for.
"""
return (
*(spec for spec in exporters if not is_unconfigured_placeholder(spec)),
ExporterSpec(owner=owner, requires_headers=True),
)


def is_unconfigured_placeholder(spec: "ExporterSpec") -> bool:
"""Whether ``spec`` is the one ``_normalize`` folds in when nothing was configured.

No field set is what says the operator asked for nothing: an exporter they did
configure survives, even ``OTEL_EXPORTER=console`` whose value matches the default,
and so does the gated spec this module appends, which would otherwise eat itself
when one preset layers onto another.
"""
return not spec.model_fields_set
12 changes: 4 additions & 8 deletions litellm/integrations/otel/presets/weave.py
Original file line number Diff line number Diff line change
Expand Up @@ -7,13 +7,10 @@
ExporterSpec,
OpenTelemetryV2Config,
)
from litellm.integrations.otel.presets.utils import (
credential_gated_exporters,
ensure_mappers,
)
from litellm.integrations.otel.presets.utils import ensure_mappers
from litellm.integrations.weave.weave_otel import (
_get_weave_authorization_header,
get_weave_otel_config,
read_weave_otel_config,
)
from litellm.types.utils import StandardCallbackDynamicParams

Expand All @@ -26,20 +23,19 @@ def weave_preset(
base: Final = config_overrides or OpenTelemetryV2Config()
mappers: Final = ensure_mappers(base.mapper_names, "openinference", "weave")
try:
weave_cfg: Final = get_weave_otel_config()
weave_cfg: Final = read_weave_otel_config()
except Exception:
if not allow_missing_credentials:
raise
return base.model_copy(
update={ # mutable-ok: pydantic model_copy takes a plain update mapping
"exporters": credential_gated_exporters(base.exporters, ExporterOwner.WEAVE_OTEL),
"exporters": [ExporterSpec(owner=ExporterOwner.WEAVE_OTEL, requires_headers=True)],
"mapper_names": mappers,
}
)
return base.model_copy(
update={
"exporters": [
*base.exporters,
ExporterSpec(
kind=weave_cfg.protocol or "otlp_http",
endpoint=weave_cfg.endpoint,
Expand Down
25 changes: 10 additions & 15 deletions litellm/integrations/weave/weave_otel.py
Original file line number Diff line number Diff line change
Expand Up @@ -125,17 +125,8 @@ def weave_otel_endpoint(host: str | None) -> str:
return normalized.rstrip("/") + WEAVE_OTEL_ENDPOINT


def get_weave_otel_config() -> WeaveOtelConfig:
"""
Retrieves the Weave OpenTelemetry configuration based on environment variables.

Environment Variables:
WANDB_API_KEY: Required. W&B API key for authentication.
WANDB_PROJECT_ID: Required. Project ID in format <entity>/<project_name>.
WANDB_HOST: Optional. Custom Weave host URL. Defaults to cloud endpoint.

Returns:
WeaveOtelConfig: A Pydantic model containing Weave OTEL configuration.
def read_weave_otel_config() -> WeaveOtelConfig:
"""Weave OTLP settings from ``WANDB_API_KEY``, ``WANDB_PROJECT_ID`` and optional ``WANDB_HOST``.

Raises:
ValueError: If required environment variables are missing.
Expand All @@ -158,10 +149,6 @@ def get_weave_otel_config() -> WeaveOtelConfig:
auth_header: Final = _get_weave_authorization_header(api_key=api_key)
otlp_auth_headers: Final = f"Authorization={auth_header},project_id={project_id}"

# Set standard OTEL environment variables
os.environ["OTEL_EXPORTER_OTLP_ENDPOINT"] = endpoint
os.environ["OTEL_EXPORTER_OTLP_HEADERS"] = otlp_auth_headers

return WeaveOtelConfig(
otlp_auth_headers=otlp_auth_headers,
endpoint=endpoint,
Expand All @@ -170,6 +157,14 @@ def get_weave_otel_config() -> WeaveOtelConfig:
)


def get_weave_otel_config() -> WeaveOtelConfig:
"""``read_weave_otel_config`` plus the v1 side effect of publishing it as the process-wide ``OTEL_*`` env."""
config: Final = read_weave_otel_config()
os.environ["OTEL_EXPORTER_OTLP_ENDPOINT"] = config.endpoint
os.environ["OTEL_EXPORTER_OTLP_HEADERS"] = config.otlp_auth_headers
return config


def set_weave_otel_attributes(span: Span, kwargs: Mapping[str, object], response_obj: object):
"""
Sets OpenTelemetry span attributes for Weave observability.
Expand Down
Loading
Loading