Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 7 additions & 1 deletion litellm/integrations/otel/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -179,6 +179,10 @@ nothing here imports outside it:
`config.yaml` — the latter reach the config through the logger's constructor
kwargs. `baggage_team_metadata_keys` is empty by default, so none of a team's
free-form metadata is promoted until each sub-key is explicitly allowlisted.
`langfuse_trace_metadata_keys` (`LITELLM_OTEL_LANGFUSE_TRACE_METADATA_KEYS`) is
the same kind of allowlist for the caller's own request metadata, which the
`langfuse` mapper stamps as `langfuse.trace.metadata.<key>`; also empty by
default.
- [`baggage.py`](./model/baggage.py) — the single definition of which request-identity
values are promoted into Baggage (so child spans inherit them) and under which
attribute keys.
Expand All @@ -198,7 +202,9 @@ nothing here imports outside it:
- `legacy` — an additional vocabulary using the older semconv-ai / Traceloop
attribute key names, for backends that read those.
- `openinference`, `langfuse`, `weave`, `langtrace` — vendor vocabularies.
- `resolve_mappers(names)` turns config names into mapper instances.
- `resolve_mappers(names, config)` turns config names into mapper instances,
handing each the config so a vocabulary with operator-configurable behaviour
(the Langfuse trace-metadata allowlist) reads it from the same place.

### Plumbing (`plumbing/`)

Expand Down
2 changes: 1 addition & 1 deletion litellm/integrations/otel/emitter.py
Original file line number Diff line number Diff line change
Expand Up @@ -122,7 +122,7 @@ def __init__(
# The mapper chain is the sole source of span attributes. When not
# passed in, resolve it from the config so there's one source of truth.
self._mappers: list[AttributeMapper] = (
list(mappers) if mappers is not None else resolve_mappers(config.mapper_names)
list(mappers) if mappers is not None else resolve_mappers(config.mapper_names, config)
)
# Bounded LRU (ordered by insertion / most-recent touch). Storing keys
# only — the value is unused — so it behaves like a capped set.
Expand Down
2 changes: 1 addition & 1 deletion litellm/integrations/otel/logger.py
Original file line number Diff line number Diff line change
Expand Up @@ -153,7 +153,7 @@ def __init__(
self._emitter = SpanEmitter(
self.tracer,
self.config,
mappers=resolve_mappers(self.config.mapper_names),
mappers=resolve_mappers(self.config.mapper_names, self.config),
event_recorder=self._init_events(logger_provider),
)
self._tenant_tracers = TenantTracerCache(self.config, callback_name, LITELLM_TRACER_NAME)
Expand Down
27 changes: 17 additions & 10 deletions litellm/integrations/otel/mappers/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -19,26 +19,33 @@
from litellm.integrations.otel.mappers.legacy import LegacyMapper
from litellm.integrations.otel.mappers.openinference import OpenInferenceMapper
from litellm.integrations.otel.mappers.weave import WeaveMapper
from litellm.integrations.otel.model.config import OpenTelemetryV2Config

# Registry keyed by ``config.mapper_names`` entries.
_MAPPER_BY_NAME: dict[str, Callable[[], AttributeMapper]] = {
"genai": GenAIMapper,
"legacy": LegacyMapper,
"openinference": OpenInferenceMapper,
"langfuse": LangfuseMapper,
"weave": WeaveMapper,
"langtrace": LangtraceMapper,
_MAPPER_BY_NAME: dict[str, Callable[[OpenTelemetryV2Config | None], AttributeMapper]] = {
"genai": lambda _config: GenAIMapper(),
"legacy": lambda _config: LegacyMapper(),
"openinference": lambda _config: OpenInferenceMapper(),
"langfuse": lambda config: LangfuseMapper(
trace_metadata_keys=config.langfuse_trace_metadata_keys if config else ()
),
"weave": lambda _config: WeaveMapper(),
"langtrace": lambda _config: LangtraceMapper(),
}


def resolve_mappers(names: Iterable[str]) -> list[AttributeMapper]:
"""Resolve mapper names to instances. Unknown names raise ``ValueError``."""
def resolve_mappers(names: Iterable[str], config: OpenTelemetryV2Config | None = None) -> list[AttributeMapper]:
"""Resolve mapper names to instances. Unknown names raise ``ValueError``.

``config`` is optional so a caller that only wants a vocabulary's default
behaviour (tests, ad-hoc mapping) needn't build one.
"""
out: list[AttributeMapper] = []
for name in names:
factory = _MAPPER_BY_NAME.get(name)
if factory is None:
raise ValueError(f"unknown mapper name {name!r}; known: {sorted(_MAPPER_BY_NAME)}")
out.append(factory())
out.append(factory(config))
return out


Expand Down
45 changes: 37 additions & 8 deletions litellm/integrations/otel/mappers/langfuse.py
Original file line number Diff line number Diff line change
Expand Up @@ -6,11 +6,18 @@

Every attribute is declared as a ``key -> extractor`` table entry (one callable
per mapping operation): ``_LLM_CALL_ATTRS`` for scalars and ``_BLOB_ATTRS`` for
the JSON-serialized payloads. ``_llm_call`` just applies both tables.
the JSON-serialized payloads. ``_llm_call`` applies both tables, plus the
caller's allowlisted metadata (``langfuse.trace.metadata.<key>``), which is
keyed per deployment and so can't live in a class-level table.

The trace-level controls (``user.id``, ``session.id``, ``langfuse.trace.name``,
``langfuse.trace.tags``) ride the generation span rather than a separate trace
span: Langfuse derives a trace from whichever observation carries them, which is
why they are repeated on every observation of the request.
"""

import json
from typing import Callable
from typing import Callable, Iterable

from litellm.integrations.otel.mappers.base import AttributeMap, AttrValue, SpanData
from litellm.integrations.otel.mappers.utils import (
Expand All @@ -26,14 +33,24 @@
)


TRACE_METADATA_PREFIX = "langfuse.trace.metadata."


class LangfuseMapper:
def __init__(self, trace_metadata_keys: Iterable[str] = ()) -> None:
self._trace_metadata_keys = frozenset(trace_metadata_keys)

_LLM_CALL_ATTRS: dict[str, Callable[[LLMCallSpanData], AttrValue | None]] = {
"langfuse.observation.type": lambda d: "generation",
"langfuse.observation.model.name": lambda d: d.request_model or None,
"langfuse.observation.metadata.provider": lambda d: d.provider or None,
"langfuse.observation.id": lambda d: d.identity.call_id or None,
"langfuse.trace.metadata.team_id": lambda d: d.identity.team_id or None,
"langfuse.trace.metadata.team_alias": lambda d: d.identity.team_alias or None,
"user.id": lambda d: d.annotations.user_id or d.identity.end_user or None,
"session.id": lambda d: d.annotations.session_id or None,
"langfuse.trace.name": lambda d: d.annotations.trace_name or None,
"langfuse.trace.tags": lambda d: list(d.annotations.tags) or None,
f"{TRACE_METADATA_PREFIX}team_id": lambda d: d.identity.team_id or None,
f"{TRACE_METADATA_PREFIX}team_alias": lambda d: d.identity.team_alias or None,
}

# Sub-tables folded into their respective JSON blobs.
Expand Down Expand Up @@ -71,9 +88,21 @@ def map(self, data: SpanData) -> AttributeMap:
case _:
return {}

@classmethod
def _llm_call(cls, data: LLMCallSpanData) -> AttributeMap:
def _llm_call(self, data: LLMCallSpanData) -> AttributeMap:
return {
**collect(self._LLM_CALL_ATTRS, data),
**collect(self._BLOB_ATTRS, data),
**self._trace_metadata(data),
}

def _trace_metadata(self, data: LLMCallSpanData) -> AttributeMap:
"""The caller's metadata, restricted to the operator's allowlist.

Empty unless a deployment allowlists keys, so a request can never push
arbitrary metadata of its own into the backend.
"""
return {
**collect(cls._LLM_CALL_ATTRS, data),
**collect(cls._BLOB_ATTRS, data),
f"{TRACE_METADATA_PREFIX}{key}": value
for key, value in data.annotations.requester_metadata.items()
if key in self._trace_metadata_keys
Comment on lines 104 to +107

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 Caller metadata overrides team identity

When team_id or team_alias is allowlisted, this final metadata merge replaces the corresponding proxy-authoritative attribute with the caller's requester_metadata value, causing the Langfuse trace to be attributed to the wrong team.

}
18 changes: 18 additions & 0 deletions litellm/integrations/otel/model/config.py
Original file line number Diff line number Diff line change
Expand Up @@ -203,6 +203,23 @@ class OpenTelemetryV2Config(BaseSettings):
),
)

langfuse_trace_metadata_keys: Annotated[tuple[str, ...], NoDecode] = Field(
default_factory=tuple,
validation_alias=AliasChoices(
"langfuse_trace_metadata_keys",
"LITELLM_OTEL_LANGFUSE_TRACE_METADATA_KEYS",
),
description=(
"Request-metadata keys the ``langfuse`` mapper stamps as "
"``langfuse.trace.metadata.<key>``. Empty by default so none of a "
"caller's free-form metadata reaches Langfuse until explicitly "
"allowlisted. Configure via the "
"``LITELLM_OTEL_LANGFUSE_TRACE_METADATA_KEYS`` env var "
"(comma-separated) or ``callback_settings.otel."
"langfuse_trace_metadata_keys`` in config.yaml (a YAML list)."
),
)

@field_validator("capture_message_content", mode="before")
@classmethod
def _normalize_capture_message_content(cls, value: object) -> object:
Expand All @@ -221,6 +238,7 @@ def _normalize_capture_message_content(cls, value: object) -> object:
"baggage_promoted_keys",
"baggage_metadata_keys",
"baggage_team_metadata_keys",
"langfuse_trace_metadata_keys",
"mapper_names",
mode="before",
)
Expand Down
48 changes: 47 additions & 1 deletion litellm/integrations/otel/model/metadata.py
Original file line number Diff line number Diff line change
Expand Up @@ -41,7 +41,7 @@

from litellm.constants import LITELLM_LOGGING_NO_UPSTREAM_LLM_CALL
from litellm.integrations.otel.model.semconv import resolve_operation
from litellm.integrations.otel.model.utils import as_str, to_seconds
from litellm.integrations.otel.model.utils import as_str, as_str_tuple, to_seconds

if TYPE_CHECKING:
from litellm.types.utils import StandardLoggingPayload
Expand Down Expand Up @@ -122,6 +122,50 @@ def from_user_api_key_auth(cls, auth: object) -> "RequestIdentity":
)


@dataclass(frozen=True)
class RequestAnnotations:
"""The caller-supplied trace annotations of a request, parsed once.

These are the request's own labels — the conversation/session it belongs to,
the name and end-user it should be attributed to, its tags, and whatever
free-form metadata the caller sent — as opposed to the proxy-authoritative
:class:`RequestIdentity`.

On the proxy the caller's ``metadata`` is snapshotted verbatim under
``metadata.requester_metadata`` (``StandardLoggingMetadata`` drops every key
it doesn't declare), so that snapshot is the only place the named controls
survive. ``tags`` instead comes from the payload's ``request_tags``, which
already merges request, key/team, and header tags.

``requester_metadata`` is carried raw (scalars only, stringified) and is
filtered to an operator allowlist by whoever stamps it, so an unconfigured
deployment never puts a caller's metadata on a span.
"""

session_id: str | None = None
trace_name: str | None = None
user_id: str | None = None
tags: tuple[str, ...] = ()
requester_metadata: Mapping[str, str] = field(default_factory=dict)

@classmethod
def from_payload(cls, payload: StandardLoggingPayload) -> RequestAnnotations:
metadata = payload.get("metadata")
requester = metadata.get("requester_metadata") if metadata else None
requester_meta: Mapping[str, object] = requester if isinstance(requester, Mapping) else {}
return cls(
session_id=as_str(requester_meta.get("session_id")),
trace_name=as_str(requester_meta.get("trace_name")),
user_id=as_str(requester_meta.get("trace_user_id")),
Comment on lines +153 to +159

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 SDK trace controls are dropped

When an SDK call supplies trace_user_id, session_id, or trace_name directly in its metadata, this parser reads only the proxy-created requester_metadata snapshot, causing the Langfuse span to omit user.id, session.id, and langfuse.trace.name.

Knowledge Base Used: Logging & Observability Integrations

tags=as_str_tuple(payload.get("request_tags")) or (),
requester_metadata={
key: str(value)
for key, value in requester_meta.items()
if isinstance(key, str) and isinstance(value, (str, bool, int, float))
},
)


@dataclass(frozen=True)
class RequestContext:
"""The fully-resolved view of a closed request, parsed once from the payload.
Expand All @@ -137,6 +181,7 @@ class RequestContext:
model_id: str | None
api_base: str | None
identity: RequestIdentity
annotations: RequestAnnotations = field(default_factory=RequestAnnotations)

@property
def provider_model(self) -> str | None:
Expand All @@ -160,6 +205,7 @@ def from_standard_logging_payload(cls, payload: "StandardLoggingPayload") -> "Re
model_id=as_str(payload.get("model_id")) or _model_info_id(raw_meta.get("model_info")),
api_base=as_str(payload.get("api_base")) or as_str(hidden.get("api_base")),
identity=RequestIdentity.from_payload(payload),
annotations=RequestAnnotations.from_payload(payload),
)


Expand Down
4 changes: 4 additions & 0 deletions litellm/integrations/otel/model/payloads.py
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@
from urllib.parse import urlsplit

from litellm.integrations.otel.model.metadata import (
RequestAnnotations,
RequestContext,
RequestIdentity,
)
Expand All @@ -30,6 +31,7 @@
# :mod:`metadata`; re-exported here so existing ``model.payloads`` imports keep
# resolving it.
__all__ = [
"RequestAnnotations",
"RequestContext",
"RequestIdentity",
"GuardrailSpanData",
Expand Down Expand Up @@ -297,6 +299,7 @@ class LLMCallSpanData:
response_cost: float | None
server: ServerInfo | None
identity: RequestIdentity
annotations: RequestAnnotations = field(default_factory=RequestAnnotations)
is_streaming: bool | None = None
cost: LLMCost = field(default_factory=LLMCost)
tools: tuple[ToolDefinition, ...] = ()
Expand Down Expand Up @@ -351,6 +354,7 @@ def from_standard_logging_payload(
cost=LLMCost.from_breakdown(cast("Mapping[str, object] | None", payload.get("cost_breakdown"))),
server=ServerInfo.from_api_base(context.api_base),
identity=context.identity,
annotations=context.annotations,
is_streaming=as_bool(payload.get("stream")),
tools=_extract_tools(params),
messages_in=_dicts(payload.get("messages")) if capture_content else (),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -549,6 +549,44 @@ def test_request_context_prefers_explicit_dispatched_model():
assert ctx.provider_model == "azure/my-deployment"


def test_request_annotations_parse_caller_trace_controls():
"""The caller's trace controls survive only in the ``requester_metadata``
snapshot (``StandardLoggingMetadata`` drops undeclared keys), and tags come
from the already-merged ``request_tags``."""
payload = _sample_payload(
metadata={
"requester_metadata": {
"trace_user_id": "user-1",
"session_id": "session-1",
"trace_name": "chat-request",
"tags": ["test"],
"environment": "staging",
}
},
request_tags=["test", "user_agent:curl"],
)
annotations = LLMCallSpanData.from_standard_logging_payload(payload).annotations
assert annotations.user_id == "user-1"
assert annotations.session_id == "session-1"
assert annotations.trace_name == "chat-request"
assert annotations.tags == ("test", "user_agent:curl")
assert annotations.requester_metadata == {
"trace_user_id": "user-1",
"session_id": "session-1",
"trace_name": "chat-request",
"environment": "staging",
}


def test_request_annotations_empty_without_caller_metadata():
annotations = LLMCallSpanData.from_standard_logging_payload(_sample_payload()).annotations
assert annotations.user_id is None
assert annotations.session_id is None
assert annotations.trace_name is None
assert annotations.tags == ()
assert annotations.requester_metadata == {}


def test_content_capture_opt_in_retains_bodies():
payload = _sample_payload(
messages=[{"role": "user", "content": "secret prompt"}],
Expand Down
Loading
Loading