diff --git a/litellm/integrations/otel/README.md b/litellm/integrations/otel/README.md index 3038bdb90b2..93e37b59a6d 100644 --- a/litellm/integrations/otel/README.md +++ b/litellm/integrations/otel/README.md @@ -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.`; 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. @@ -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/`) diff --git a/litellm/integrations/otel/emitter.py b/litellm/integrations/otel/emitter.py index 8651cf586cd..322cb2ce498 100644 --- a/litellm/integrations/otel/emitter.py +++ b/litellm/integrations/otel/emitter.py @@ -122,7 +122,7 @@ class SpanEmitter: # 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. diff --git a/litellm/integrations/otel/logger.py b/litellm/integrations/otel/logger.py index b33973f0676..ac719cbad02 100644 --- a/litellm/integrations/otel/logger.py +++ b/litellm/integrations/otel/logger.py @@ -153,7 +153,7 @@ class OpenTelemetryV2(CustomLogger): 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) diff --git a/litellm/integrations/otel/mappers/__init__.py b/litellm/integrations/otel/mappers/__init__.py index b0c1d7019db..960303f7d24 100644 --- a/litellm/integrations/otel/mappers/__init__.py +++ b/litellm/integrations/otel/mappers/__init__.py @@ -19,26 +19,35 @@ from litellm.integrations.otel.mappers.langtrace import LangtraceMapper 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, +# Registry keyed by ``config.mapper_names`` entries. Every factory takes the +# config so a vocabulary that has operator-configurable behaviour (the Langfuse +# metadata allowlist) reads it from the same source of truth as the rest. +_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 diff --git a/litellm/integrations/otel/mappers/langfuse.py b/litellm/integrations/otel/mappers/langfuse.py index 79f8f618eff..98528d3b884 100644 --- a/litellm/integrations/otel/mappers/langfuse.py +++ b/litellm/integrations/otel/mappers/langfuse.py @@ -6,11 +6,18 @@ Langfuse ingests OTLP spans and reads from its own vendor namespace 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.``), 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 ( @@ -26,14 +33,24 @@ from litellm.integrations.otel.model.payloads import ( ) +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. @@ -71,9 +88,21 @@ class LangfuseMapper: case _: return {} - @classmethod - def _llm_call(cls, data: LLMCallSpanData) -> AttributeMap: + def _llm_call(self, data: LLMCallSpanData) -> AttributeMap: return { - **collect(cls._LLM_CALL_ATTRS, data), - **collect(cls._BLOB_ATTRS, data), + **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 { + f"{TRACE_METADATA_PREFIX}{key}": value + for key, value in data.annotations.requester_metadata.items() + if key in self._trace_metadata_keys } diff --git a/litellm/integrations/otel/model/config.py b/litellm/integrations/otel/model/config.py index 7f33129c560..cf36123f69d 100644 --- a/litellm/integrations/otel/model/config.py +++ b/litellm/integrations/otel/model/config.py @@ -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.``. 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: @@ -221,6 +238,7 @@ class OpenTelemetryV2Config(BaseSettings): "baggage_promoted_keys", "baggage_metadata_keys", "baggage_team_metadata_keys", + "langfuse_trace_metadata_keys", "mapper_names", mode="before", ) diff --git a/litellm/integrations/otel/model/metadata.py b/litellm/integrations/otel/model/metadata.py index 7ff4f540908..38ae8fc8661 100644 --- a/litellm/integrations/otel/model/metadata.py +++ b/litellm/integrations/otel/model/metadata.py @@ -41,7 +41,7 @@ from typing import TYPE_CHECKING, Any, Mapping, cast 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 @@ -122,6 +122,50 @@ class 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")), + 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. @@ -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: @@ -160,6 +205,7 @@ class RequestContext: 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), ) diff --git a/litellm/integrations/otel/model/payloads.py b/litellm/integrations/otel/model/payloads.py index 4a8f01858b5..c91d9c8c921 100644 --- a/litellm/integrations/otel/model/payloads.py +++ b/litellm/integrations/otel/model/payloads.py @@ -9,6 +9,7 @@ from typing import TYPE_CHECKING, ClassVar, Mapping, cast from urllib.parse import urlsplit from litellm.integrations.otel.model.metadata import ( + RequestAnnotations, RequestContext, RequestIdentity, ) @@ -30,6 +31,7 @@ from litellm.integrations.otel.model.utils import ( # :mod:`metadata`; re-exported here so existing ``model.payloads`` imports keep # resolving it. __all__ = [ + "RequestAnnotations", "RequestContext", "RequestIdentity", "GuardrailSpanData", @@ -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, ...] = () @@ -351,6 +354,7 @@ class LLMCallSpanData: 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 (), diff --git a/tests/test_litellm/integrations/otel/test_otel_v2_sources_of_truth.py b/tests/test_litellm/integrations/otel/test_otel_v2_sources_of_truth.py index 71be28ea485..987eb01c4f6 100644 --- a/tests/test_litellm/integrations/otel/test_otel_v2_sources_of_truth.py +++ b/tests/test_litellm/integrations/otel/test_otel_v2_sources_of_truth.py @@ -549,6 +549,45 @@ 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") + # scalars only: the ``tags`` list is not a span-safe value here + 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"}], diff --git a/tests/test_litellm/integrations/otel/test_otel_v2_vendor_mappers.py b/tests/test_litellm/integrations/otel/test_otel_v2_vendor_mappers.py index 94cb79f53b8..8c3ebfef54b 100644 --- a/tests/test_litellm/integrations/otel/test_otel_v2_vendor_mappers.py +++ b/tests/test_litellm/integrations/otel/test_otel_v2_vendor_mappers.py @@ -18,10 +18,12 @@ from litellm.integrations.otel.mappers import ( WeaveMapper, resolve_mappers, ) +from litellm.integrations.otel.model.config import OpenTelemetryV2Config from litellm.integrations.otel.model.payloads import ( LLMCallSpanData, LLMRequestParams, LLMUsage, + RequestAnnotations, RequestIdentity, ServerInfo, ToolDefinition, @@ -134,6 +136,60 @@ def test_langfuse_mapper_observation_attrs(): assert attrs["langfuse.trace.metadata.team_id"] == "t1" +def test_langfuse_mapper_maps_trace_controls(): + """The trace controls LiteLLM's Langfuse contract accepts (trace user, + session, name, tags) land on the generation span under the current official + Langfuse OTLP keys, without any content capture.""" + attrs = LangfuseMapper().map( + _llm_call( + annotations=RequestAnnotations( + user_id="user-1", + session_id="session-1", + trace_name="chat-request", + tags=("test", "prod"), + ), + messages_in=(), + choices_out=(), + ) + ) + assert attrs["user.id"] == "user-1" + assert attrs["session.id"] == "session-1" + assert attrs["langfuse.trace.name"] == "chat-request" + assert attrs["langfuse.trace.tags"] == ["test", "prod"] + + +def test_langfuse_mapper_trace_user_falls_back_to_end_user(): + attrs = LangfuseMapper().map(_llm_call(identity=RequestIdentity(call_id="c1", end_user="end-user-9"))) + assert attrs["user.id"] == "end-user-9" + + +def test_langfuse_mapper_omits_absent_trace_controls(): + attrs = LangfuseMapper().map(_llm_call()) + for key in ("user.id", "session.id", "langfuse.trace.name", "langfuse.trace.tags"): + assert key not in attrs + + +def test_langfuse_mapper_trace_metadata_needs_an_allowlist(): + """A caller's free-form metadata reaches Langfuse only for keys the operator + allowlisted, so a request can't push arbitrary metadata into the backend.""" + data = _llm_call( + annotations=RequestAnnotations(requester_metadata={"environment": "staging", "internal_note": "secret"}) + ) + assert "langfuse.trace.metadata.environment" not in LangfuseMapper().map(data) + attrs = LangfuseMapper(trace_metadata_keys=["environment"]).map(data) + assert attrs["langfuse.trace.metadata.environment"] == "staging" + assert "langfuse.trace.metadata.internal_note" not in attrs + + +def test_resolve_mappers_passes_langfuse_allowlist_from_config(): + config = OpenTelemetryV2Config(mapper_names=["langfuse"], langfuse_trace_metadata_keys=["environment"]) + data = _llm_call(annotations=RequestAnnotations(requester_metadata={"environment": "staging"})) + union: dict = {} + for mapper in resolve_mappers(config.mapper_names, config): + union.update(mapper.map(data)) + assert union["langfuse.trace.metadata.environment"] == "staging" + + def test_langfuse_mapper_skips_when_no_messages(): data = _llm_call(messages_in=(), choices_out=()) attrs = LangfuseMapper().map(data)