mirror of
https://github.com/BerriAI/litellm.git
synced 2026-10-05 02:41:56 +00:00
fix(otel): map Langfuse trace user, session, name, tags in v2 mapper
This commit is contained in:
parent
4d54324515
commit
60194646fc
10 changed files with 230 additions and 23 deletions
|
|
@ -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.
|
||||
|
|
@ -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/`)
|
||||
|
||||
|
|
|
|||
|
|
@ -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.
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
||||
|
||||
|
|
|
|||
|
|
@ -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.<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 (
|
||||
|
|
@ -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
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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:
|
||||
|
|
@ -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",
|
||||
)
|
||||
|
|
|
|||
|
|
@ -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),
|
||||
)
|
||||
|
||||
|
||||
|
|
|
|||
|
|
@ -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 (),
|
||||
|
|
|
|||
|
|
@ -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"}],
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue