diff --git a/litellm/integrations/otel/mappers/langfuse.py b/litellm/integrations/otel/mappers/langfuse.py index 9aff944cff0..40adf2363d9 100644 --- a/litellm/integrations/otel/mappers/langfuse.py +++ b/litellm/integrations/otel/mappers/langfuse.py @@ -28,6 +28,8 @@ from litellm.integrations.otel.model.payloads import ( LLMUsage, ) from litellm.integrations.otel.model.trace_controls import TraceControls +from litellm.integrations.otel.model.utils import as_str +from litellm.litellm_core_utils.safe_json_dumps import safe_dumps LANGFUSE_OBSERVATION_INPUT: Final = "langfuse.observation.input" LANGFUSE_OBSERVATION_OUTPUT: Final = "langfuse.observation.output" @@ -35,6 +37,19 @@ LANGFUSE_TRACE_NAME: Final = "langfuse.trace.name" LANGFUSE_TRACE_USER_ID: Final = "user.id" LANGFUSE_TRACE_SESSION_ID: Final = "session.id" LANGFUSE_TRACE_TAGS: Final = "langfuse.trace.tags" +LANGFUSE_OBSERVATION_METADATA: Final = "langfuse.observation.metadata" +LANGFUSE_TRACE_METADATA_PREFIX: Final = "langfuse.trace.metadata." +TRACE_IDENTITY_FIELDS: Final = ( + "user_api_key_alias", + "user_api_key_user_id", + "user_api_key_end_user_id", + "user_api_key_team_id", + "user_api_key_team_alias", +) + + +def _identity_field(name: str) -> Callable[[LLMCallSpanData], AttrValue | None]: + return lambda d: as_str(d.request_metadata.get(name)) or None class LangfuseMapper: @@ -45,6 +60,7 @@ class LangfuseMapper: "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, + **{f"{LANGFUSE_TRACE_METADATA_PREFIX}{name}": _identity_field(name) for name in TRACE_IDENTITY_FIELDS}, } # Sub-tables folded into their respective JSON blobs. @@ -64,6 +80,7 @@ class LangfuseMapper: # JSON-payload attributes: each builder returns the serialized blob or None. _BLOB_ATTRS: dict[str, Callable[[LLMCallSpanData], AttrValue | None]] = { + LANGFUSE_OBSERVATION_METADATA: lambda d: safe_dumps(dict(d.request_metadata)) if d.request_metadata else None, "langfuse.observation.model.parameters": lambda d: json_if( collect(LangfuseMapper._MODEL_PARAMS, d.request_params) ), diff --git a/litellm/integrations/otel/model/metadata.py b/litellm/integrations/otel/model/metadata.py index ede8ac99467..8dc7b2f858a 100644 --- a/litellm/integrations/otel/model/metadata.py +++ b/litellm/integrations/otel/model/metadata.py @@ -42,9 +42,11 @@ from types import MappingProxyType from typing import TYPE_CHECKING, Any, Final, cast from litellm.constants import LITELLM_LOGGING_NO_UPSTREAM_LLM_CALL, SESSION_ID_GENERATED_METADATA_KEY +from litellm.integrations.langfuse.langfuse import log_requester_metadata from litellm.integrations.otel.model.semconv import resolve_operation from litellm.integrations.otel.model.trace_controls import TraceControls, caller_trace_controls from litellm.integrations.otel.model.utils import as_str, as_str_mapping, to_seconds +from litellm.litellm_core_utils.redact_messages import redact_user_api_key_info if TYPE_CHECKING: from litellm.types.utils import StandardLoggingPayload @@ -417,6 +419,18 @@ def _model_info_id(model_info: object) -> str | None: return None +def exported_request_metadata(payload: StandardLoggingPayload) -> Mapping[str, object]: + """The request metadata as the logging callbacks export it: ``user_api_key_*`` + dropped when ``litellm.redact_user_api_key_info`` is on, header keys nested + under ``requester_metadata``, ``None`` values dropped.""" + raw_meta: Final = cast(Mapping[str, object], payload.get("metadata") or {}) + redacted: Final[Mapping[str, object]] = cast( + Mapping[str, object], redact_user_api_key_info(metadata=dict(raw_meta)) + ) + exported: Final[Mapping[str, object]] = cast(Mapping[str, object], log_requester_metadata(redacted)) + return MappingProxyType({key: value for key, value in exported.items() if value is not None}) + + def _team_metadata_dict(value: object) -> Mapping[str, object] | None: """The team's free-form metadata as a raw mapping, or ``None`` when missing or empty. diff --git a/litellm/integrations/otel/model/payloads.py b/litellm/integrations/otel/model/payloads.py index ea4ded90480..c1d421b5d52 100644 --- a/litellm/integrations/otel/model/payloads.py +++ b/litellm/integrations/otel/model/payloads.py @@ -12,7 +12,7 @@ from urllib.parse import urlsplit from typing_extensions import ReadOnly, TypedDict -from litellm.integrations.otel.model.metadata import RequestContext, RequestIdentity +from litellm.integrations.otel.model.metadata import RequestContext, RequestIdentity, exported_request_metadata from litellm.integrations.otel.model.semconv import ( GenAIOperation, GenAIOutputType, @@ -388,6 +388,7 @@ class LLMCallSpanData: response_cost: float | None server: ServerInfo | None identity: RequestIdentity + request_metadata: Mapping[str, object] = field(default_factory=lambda: cast(Mapping[str, object], {})) is_streaming: bool | None = None cost: LLMCost = field(default_factory=LLMCost) tools: tuple[ToolDefinition, ...] = () @@ -455,6 +456,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, + request_metadata=exported_request_metadata(payload), 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/integration/observability/test_langfuse_otel_metadata.py b/tests/integration/observability/test_langfuse_otel_metadata.py index 29f0f7d25ba..0be5cc29fc5 100644 --- a/tests/integration/observability/test_langfuse_otel_metadata.py +++ b/tests/integration/observability/test_langfuse_otel_metadata.py @@ -1014,8 +1014,7 @@ def test_langfuse_otel_slow_sink_does_not_deadlock_exports(gateway: Gateway, tmp ) -@pytest.mark.covers("other.observability.langfuse_otel.v2_mapper_unchanged") -def test_langfuse_otel_v2_mapper_keeps_existing_trace_metadata_keys(gateway: Gateway, tmp_path: Path) -> None: +def test_langfuse_otel_v2_mapper_emits_request_metadata_and_identity(gateway: Gateway, tmp_path: Path) -> None: marker: Final = "lf-v2-" + uuid.uuid4().hex with ( _langfuse_rig( @@ -1049,4 +1048,101 @@ def test_langfuse_otel_v2_mapper_keeps_existing_trace_metadata_keys(gateway: Gat attrs: Final = next(s for s in observed if s.get("gen_ai.response.id") == response.json()["id"]) assert attrs.get("langfuse.trace.metadata.team_id") == identity["team_id"], attrs assert attrs.get("langfuse.trace.metadata.team_alias") == identity["team_alias"], attrs - assert "langfuse.observation.metadata" not in attrs, sorted(attrs) + assert "langfuse.observation.metadata" in attrs, sorted(attrs) + observation: Final = json.loads(str(attrs["langfuse.observation.metadata"])) + expected: Final = _chat_identity_fields(identity) + assert {key: observation.get(key) for key in expected} == dict(expected), observation + assert observation["spend_logs_metadata"] == {"ticket": "LIT-8283"}, observation + assert not [key_ for key_, value in observation.items() if value is None], observation + flattened: Final = { + attribute: attrs.get(attribute) + for attribute in (f"langfuse.trace.metadata.{field}" for field in _IDENTITIES) + if attribute in attrs + } + assert flattened == {f"langfuse.trace.metadata.{field}": value for field, value in expected.items()}, attrs + + +def test_langfuse_otel_v2_mapper_redacts_user_api_key_fields(gateway: Gateway, tmp_path: Path) -> None: + marker: Final = "lf-v2-redact-" + uuid.uuid4().hex + with ( + _langfuse_rig( + gateway, + tmp_path, + {"redact_user_api_key_info": True}, + marker, + overrides={"LITELLM_OTEL_V2": "1"}, + ) as rig, + rig.candidate.scenario() as scenario, + ): + key: Final = scenario.key( + key_alias="lf-v2-redact-" + uuid.uuid4().hex, + metadata={"spend_logs_metadata": {"ticket": "LIT-8283"}}, + ) + model: Final = scenario.model(api_base=rig.provider.url + "/v1") + response: Final = rig.candidate.request( + "POST", + "/v1/chat/completions", + { + "model": model, + "messages": [{"role": "user", "content": marker}], + "cache": {"no-cache": True}, + }, + key=key, + ) + assert response.status_code == 200, response.text + watch: Final = _watched_spans(rig.collector) + observed: Final = eventually( + watch, + lambda spans: sum(s.get("gen_ai.response.id") == response.json()["id"] for s in spans) == 1, + seconds=25, + ) + attrs: Final = next(s for s in observed if s.get("gen_ai.response.id") == response.json()["id"]) + assert "langfuse.observation.metadata" in attrs, sorted(attrs) + observation: Final = json.loads(str(attrs["langfuse.observation.metadata"])) + assert not [name for name in observation if name.startswith("user_api_key")], observation + assert not [key_ for key_, value in observation.items() if value is None], observation + assert observation["spend_logs_metadata"] == {"ticket": "LIT-8283"}, observation + assert not [name for name in attrs if name.startswith("langfuse.trace.metadata.user_api_key_")], sorted(attrs) + + +def test_langfuse_otel_v2_mapper_teamless_key_emits_no_team_identity(gateway: Gateway, tmp_path: Path) -> None: + marker: Final = "lf-v2-noteam-" + uuid.uuid4().hex + with ( + _langfuse_rig( + gateway, + tmp_path, + {}, + marker, + overrides={"LITELLM_OTEL_V2": "1"}, + ) as rig, + rig.candidate.scenario() as scenario, + ): + key, identity = _scenario_identity(scenario, rig.provider, team=False) + response: Final = rig.candidate.request( + "POST", + "/v1/chat/completions", + { + "model": identity["model"], + "messages": [{"role": "user", "content": marker}], + "cache": {"no-cache": True}, + }, + key=key, + ) + assert response.status_code == 200, response.text + watch: Final = _watched_spans(rig.collector) + observed: Final = eventually( + watch, + lambda spans: sum(s.get("gen_ai.response.id") == response.json()["id"] for s in spans) == 1, + seconds=25, + ) + attrs: Final = next(s for s in observed if s.get("gen_ai.response.id") == response.json()["id"]) + assert "langfuse.observation.metadata" in attrs, sorted(attrs) + observation: Final = json.loads(str(attrs["langfuse.observation.metadata"])) + for field in ("user_api_key_team_id", "user_api_key_team_alias", "user_api_key_end_user_id"): + assert field not in observation, observation + assert not [key_ for key_, value in observation.items() if value is None], observation + assert { + attribute: attrs.get(attribute) + for attribute in (f"langfuse.trace.metadata.{field}" for field in _IDENTITIES) + if attribute in attrs + } == {"langfuse.trace.metadata.user_api_key_alias": identity["key_alias"]}, attrs 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 1e2ae24a329..8cbe4d18fdc 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 @@ -332,3 +332,101 @@ def test_resolve_mappers_composition_layers_vocabularies(): def test_resolve_mappers_rejects_unknown_name(): with pytest.raises(ValueError, match="unknown mapper name 'nope'"): resolve_mappers(["genai", "nope"]) + + +def test_langfuse_mapper_emits_request_metadata_and_identity(): + payload: Final[dict[str, object]] = { + "call_type": "acompletion", + "custom_llm_provider": "openai", + "model": "gpt-4o", + "metadata": { + "user_api_key_alias": "prod-key", + "user_api_key_user_id": "user-1", + "user_api_key_end_user_id": "end-1", + "user_api_key_team_id": "team-9", + "user_api_key_team_alias": "team nine", + "user_api_key_max_budget": None, + "spend_logs_metadata": {"ticket": "LIT-8283"}, + "requester_metadata": {"headers": {"x-tenant": "acme"}}, + }, + "response": {}, + } + attrs: Final = LangfuseMapper().map(LLMCallSpanData.from_standard_logging_payload(payload)) + + assert json.loads(attrs["langfuse.observation.metadata"]) == { + "user_api_key_alias": "prod-key", + "user_api_key_user_id": "user-1", + "user_api_key_end_user_id": "end-1", + "user_api_key_team_id": "team-9", + "user_api_key_team_alias": "team nine", + "spend_logs_metadata": {"ticket": "LIT-8283"}, + "requester_metadata": {"headers": {"x-tenant": "acme"}}, + } + assert attrs["langfuse.trace.metadata.user_api_key_alias"] == "prod-key" + assert attrs["langfuse.trace.metadata.user_api_key_user_id"] == "user-1" + assert attrs["langfuse.trace.metadata.user_api_key_end_user_id"] == "end-1" + assert attrs["langfuse.trace.metadata.user_api_key_team_id"] == "team-9" + assert attrs["langfuse.trace.metadata.user_api_key_team_alias"] == "team nine" + assert attrs["langfuse.trace.metadata.team_id"] == "team-9" + assert attrs["langfuse.trace.metadata.team_alias"] == "team nine" + assert attrs["langfuse.observation.metadata.provider"] == "openai" + + +def test_langfuse_mapper_skips_empty_identity_values(): + payload: Final[dict[str, object]] = { + "call_type": "acompletion", + "custom_llm_provider": "openai", + "model": "gpt-4o", + "metadata": { + "user_api_key_alias": None, + "user_api_key_user_id": "", + "user_api_key_end_user_id": "end-1", + }, + "response": {}, + } + attrs: Final = LangfuseMapper().map(LLMCallSpanData.from_standard_logging_payload(payload)) + + assert "langfuse.trace.metadata.user_api_key_alias" not in attrs + assert "langfuse.trace.metadata.user_api_key_user_id" not in attrs + assert attrs["langfuse.trace.metadata.user_api_key_end_user_id"] == "end-1" + + +def test_langfuse_mapper_redacts_user_api_key_fields(): + import litellm + + saved: Final = litellm.redact_user_api_key_info + try: + litellm.redact_user_api_key_info = True + payload: Final[dict[str, object]] = { + "call_type": "acompletion", + "custom_llm_provider": "openai", + "model": "gpt-4o", + "metadata": { + "user_api_key_alias": "prod-key", + "user_api_key_team_id": "team-9", + "spend_logs_metadata": {"ticket": "LIT-8283"}, + }, + "response": {}, + } + attrs: Final = LangfuseMapper().map(LLMCallSpanData.from_standard_logging_payload(payload)) + finally: + litellm.redact_user_api_key_info = saved + + exported: Final = json.loads(attrs["langfuse.observation.metadata"]) + assert not [key for key in exported if key.startswith("user_api_key")] + assert exported["spend_logs_metadata"] == {"ticket": "LIT-8283"} + assert not [key for key in attrs if key.startswith("langfuse.trace.metadata.user_api_key_")] + assert attrs["langfuse.trace.metadata.team_id"] == "team-9" + + +def test_langfuse_mapper_omits_observation_metadata_without_payload_metadata(): + payload: Final[dict[str, object]] = { + "call_type": "acompletion", + "custom_llm_provider": "openai", + "model": "gpt-4o", + "metadata": {}, + "response": {}, + } + attrs: Final = LangfuseMapper().map(LLMCallSpanData.from_standard_logging_payload(payload)) + + assert json.loads(attrs["langfuse.observation.metadata"]) == {"requester_metadata": {}}