From 81a799ec2eab773fc3b45a844e0e7e5311449d9e Mon Sep 17 00:00:00 2001 From: yucheng Date: Wed, 23 Sep 2026 07:03:10 +0000 Subject: [PATCH] fix(langfuse_otel): emit request metadata under langfuse.observation.metadata and langfuse.trace.metadata keys Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- .../integrations/langfuse/langfuse_otel.py | 33 +++ litellm/types/integrations/langfuse_otel.py | 1 + tests/integration/contracts.json | 6 + .../test_langfuse_otel_metadata.py | 189 ++++++++++++++++++ .../integrations/test_langfuse_otel.py | 85 ++++++++ 5 files changed, 314 insertions(+) create mode 100644 tests/integration/observability/test_langfuse_otel_metadata.py diff --git a/litellm/integrations/langfuse/langfuse_otel.py b/litellm/integrations/langfuse/langfuse_otel.py index a96fac32c2a..bc47070f40f 100644 --- a/litellm/integrations/langfuse/langfuse_otel.py +++ b/litellm/integrations/langfuse/langfuse_otel.py @@ -2,6 +2,7 @@ import base64 import json import os from datetime import datetime +from types import MappingProxyType from typing import TYPE_CHECKING, Any, Final, Optional from litellm._logging import verbose_logger @@ -29,6 +30,16 @@ LANGFUSE_CLOUD_US_ENDPOINT: Final = "https://us.cloud.langfuse.com/api/public/ot LANGFUSE_INGESTION_VERSION_HEADER: Final = "x-langfuse-ingestion-version" LANGFUSE_INGESTION_VERSION: Final = "4" +_TRACE_IDENTITY_FIELDS: Final = MappingProxyType( + { + "user_api_key_alias": "key_alias", + "user_api_key_user_id": "user_id", + "user_api_key_end_user_id": "end_user_id", + "user_api_key_team_id": "team_id", + "user_api_key_team_alias": "team_alias", + } +) + class LangfuseOtelLogger(OpenTelemetry): def __init__(self, config=None, *args, **kwargs): @@ -125,6 +136,27 @@ class LangfuseOtelLogger(OpenTelemetry): value = str(value) safe_set_attribute(span, enum_attr.value, value) + @staticmethod + def _set_request_metadata_attributes(span: Span, kwargs: dict) -> None: + from litellm.integrations.arize._utils import safe_set_attribute + from litellm.integrations.langfuse.langfuse import log_requester_metadata + from litellm.litellm_core_utils.redact_messages import redact_user_api_key_info + from litellm.litellm_core_utils.safe_json_dumps import safe_dumps + + standard_logging_object: Final = kwargs.get("standard_logging_object") + request_metadata: Final = ( + standard_logging_object.get("metadata") if isinstance(standard_logging_object, dict) else None + ) + if not isinstance(request_metadata, dict): + return + observation_metadata: Final = log_requester_metadata(redact_user_api_key_info(metadata=request_metadata)) + safe_set_attribute(span, LangfuseSpanAttributes.OBSERVATION_METADATA.value, safe_dumps(observation_metadata)) + trace_prefix: Final = LangfuseSpanAttributes.TRACE_METADATA.value + for source_key, target_key in _TRACE_IDENTITY_FIELDS.items(): + value = observation_metadata.get(source_key) + if value is not None: + safe_set_attribute(span, f"{trace_prefix}.{target_key}", value) + @staticmethod def _set_observation_output(span: Span, response_obj): """Helper to set observation output attributes.""" @@ -244,6 +276,7 @@ class LangfuseOtelLogger(OpenTelemetry): metadata: Final = LangfuseOtelLogger._extract_langfuse_metadata(kwargs) LangfuseOtelLogger._set_metadata_attributes(span=span, metadata=metadata) + LangfuseOtelLogger._set_request_metadata_attributes(span=span, kwargs=kwargs) messages: Final = kwargs.get("messages") if messages: diff --git a/litellm/types/integrations/langfuse_otel.py b/litellm/types/integrations/langfuse_otel.py index c58dc567cda..11f30a2bfa7 100644 --- a/litellm/types/integrations/langfuse_otel.py +++ b/litellm/types/integrations/langfuse_otel.py @@ -29,6 +29,7 @@ class LangfuseSpanAttributes(str, Enum): # ---- Observation input/output ---- OBSERVATION_INPUT = "langfuse.observation.input" OBSERVATION_OUTPUT = "langfuse.observation.output" + OBSERVATION_METADATA = "langfuse.observation.metadata" # ---- Trace-level metadata ---- TRACE_USER_ID = "user.id" diff --git a/tests/integration/contracts.json b/tests/integration/contracts.json index a8d9cf1df8b..b9c8245f51f 100644 --- a/tests/integration/contracts.json +++ b/tests/integration/contracts.json @@ -257,6 +257,12 @@ "tests/integration/observability/test_otel_text_completion_choices.py::test_otel_weave_output_keeps_text_completion_provider_fields_beside_the_synthesized_message": [ "other.observability.otel.text_completion_choices_keep_provider_fields" ], + "tests/integration/observability/test_langfuse_otel_metadata.py::test_langfuse_otel_emits_request_metadata_under_langfuse_observation_and_trace_keys": [ + "other.observability.langfuse_otel.request_metadata_on_langfuse_keys" + ], + "tests/integration/observability/test_langfuse_otel_metadata.py::test_langfuse_otel_metadata_redacts_user_api_key_fields_like_vanilla_langfuse": [ + "other.observability.langfuse_otel.request_metadata_redaction_matches_vanilla" + ], "tests/integration/observability/test_guardrail_effects.py::test_guardrail_rewrites_system_and_user_in_actual_anthropic_request": [ "other.observability.guardrails.rewrite_reaches_correct_anthropic_positions" ], diff --git a/tests/integration/observability/test_langfuse_otel_metadata.py b/tests/integration/observability/test_langfuse_otel_metadata.py new file mode 100644 index 00000000000..8f64f66508c --- /dev/null +++ b/tests/integration/observability/test_langfuse_otel_metadata.py @@ -0,0 +1,189 @@ +import json +import uuid +from collections.abc import Iterator, Mapping +from contextlib import contextmanager +from pathlib import Path +from typing import Final + +import pytest +import yaml +from integration._support.client import Gateway, eventually +from integration._support.process import owned_proxy +from integration._support.wire import Reply, Request, Wire, wire_server +from opentelemetry.proto.collector.trace.v1.trace_service_pb2 import ExportTraceServiceRequest + + +def _span_attributes(body: bytes) -> tuple[dict[str, object], ...]: + request: Final = ExportTraceServiceRequest.FromString(body) + return tuple( + {attribute.key: getattr(attribute.value, attribute.value.WhichOneof("value")) for attribute in span.attributes} + for resource in request.resource_spans + for scope in resource.scope_spans + for span in scope.spans + ) + + +@contextmanager +def _langfuse_rig( + gateway: Gateway, + tmp_path: Path, + settings: Mapping[str, object], + marker: str, +) -> Iterator[tuple[Gateway, Wire, Wire]]: + def upstream(request: Request) -> Reply: + assert request.target.endswith("/chat/completions"), request.target + return Reply( + body=json.dumps( + { + "id": f"chatcmpl-{marker}", + "object": "chat.completion", + "created": 1, + "model": "gpt-4o-mini", + "choices": [ + { + "index": 0, + "message": {"role": "assistant", "content": "langfuse-echo"}, + "finish_reason": "stop", + } + ], + "usage": {"prompt_tokens": 3, "completion_tokens": 2, "total_tokens": 5}, + } + ).encode() + ) + + def sink(_request: Request) -> Reply: + return Reply() + + with wire_server(upstream) as provider, wire_server(sink) as collector: + config: Final = yaml.safe_load(Path("tests/integration/proxy_config.yaml").read_text()) + config["litellm_settings"].update({"callbacks": ["langfuse_otel"], **settings}) + path: Final = tmp_path / "langfuse_otel.yaml" + path.write_text(yaml.safe_dump(config)) + with owned_proxy( + gateway, + tmp_path, + { + "LANGFUSE_PUBLIC_KEY": "pk-lf-integration", + "LANGFUSE_SECRET_KEY": "sk-lf-integration", + "LANGFUSE_HOST": collector.url, + }, + config=path, + workers=2, + ) as candidate: + yield candidate, provider, collector + + +def _observation_span(collector: Wire, response_id: str) -> dict[str, object]: + batches: Final = [] + + def spans() -> tuple[dict[str, object], ...]: + batches.extend(collector.drain()) + return tuple( + attributes + for batch in batches + for attributes in _span_attributes(batch.body) + if attributes.get("llm.response.id") == response_id + ) + + observed: Final = eventually(spans, lambda values: len(values) == 1, seconds=25) + return observed[0] + + +@pytest.mark.covers("other.observability.langfuse_otel.request_metadata_on_langfuse_keys") +def test_langfuse_otel_emits_request_metadata_under_langfuse_observation_and_trace_keys( + gateway: Gateway, tmp_path: Path +) -> None: + marker: Final = "lf-meta-" + uuid.uuid4().hex + with ( + _langfuse_rig(gateway, tmp_path, {}, marker) as (candidate, provider, collector), + candidate.scenario() as scenario, + ): + team_alias: Final = "lf-team-" + uuid.uuid4().hex + team_id: Final = scenario.team(team_alias=team_alias) + key_alias: Final = "lf-key-" + uuid.uuid4().hex + key: Final = scenario.key( + team_id=team_id, + key_alias=key_alias, + metadata={"spend_logs_metadata": {"ticket": "LIT-8283"}}, + ) + model: Final = scenario.model(api_base=provider.url + "/v1") + end_user: Final = "lf-end-user-" + uuid.uuid4().hex + response: Final = candidate.request( + "POST", + "/v1/chat/completions", + { + "model": model, + "messages": [{"role": "user", "content": marker}], + "user": end_user, + "cache": {"no-cache": True}, + }, + key=key, + ) + assert response.status_code == 200, response.text + assert response.json()["choices"][0]["message"]["content"] == "langfuse-echo" + + attrs: Final = _observation_span(collector, response.json()["id"]) + assert "metadata" in attrs, sorted(attrs) + assert json.loads(str(attrs["metadata"]))["user_api_key_alias"] == key_alias + assert "langfuse.observation.metadata" in attrs, sorted(attrs) + observation: Final = json.loads(str(attrs["langfuse.observation.metadata"])) + expected: Final = { + "user_api_key_alias": key_alias, + "user_api_key_team_id": team_id, + "user_api_key_team_alias": team_alias, + "user_api_key_end_user_id": end_user, + } + assert {key_: observation.get(key_) for key_ in expected} == expected, observation + assert observation["spend_logs_metadata"] == {"ticket": "LIT-8283"}, observation + assert { + attribute: attrs.get(attribute) + for attribute in ( + "langfuse.trace.metadata.key_alias", + "langfuse.trace.metadata.team_id", + "langfuse.trace.metadata.team_alias", + "langfuse.trace.metadata.end_user_id", + ) + } == { + "langfuse.trace.metadata.key_alias": key_alias, + "langfuse.trace.metadata.team_id": team_id, + "langfuse.trace.metadata.team_alias": team_alias, + "langfuse.trace.metadata.end_user_id": end_user, + }, attrs + + +@pytest.mark.covers("other.observability.langfuse_otel.request_metadata_redaction_matches_vanilla") +def test_langfuse_otel_metadata_redacts_user_api_key_fields_like_vanilla_langfuse( + gateway: Gateway, tmp_path: Path +) -> None: + marker: Final = "lf-redact-" + uuid.uuid4().hex + with ( + _langfuse_rig(gateway, tmp_path, {"redact_user_api_key_info": True}, marker) as ( + candidate, + provider, + collector, + ), + candidate.scenario() as scenario, + ): + key: Final = scenario.key( + key_alias="lf-redact-" + uuid.uuid4().hex, + metadata={"spend_logs_metadata": {"ticket": "LIT-8283"}}, + ) + model: Final = scenario.model(api_base=provider.url + "/v1") + response: Final = 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 + + attrs: Final = _observation_span(collector, 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 observation["spend_logs_metadata"] == {"ticket": "LIT-8283"}, observation + assert not [name for name in attrs if name.startswith("langfuse.trace.metadata.")], sorted(attrs) diff --git a/tests/test_litellm/integrations/test_langfuse_otel.py b/tests/test_litellm/integrations/test_langfuse_otel.py index 0a9ce55fe16..0eb2dbdbc87 100644 --- a/tests/test_litellm/integrations/test_langfuse_otel.py +++ b/tests/test_litellm/integrations/test_langfuse_otel.py @@ -4,6 +4,7 @@ from unittest.mock import MagicMock, patch import pytest +import litellm from litellm.integrations.langfuse.langfuse_otel import LangfuseOtelLogger from litellm.integrations.opentelemetry import OpenTelemetryConfig from litellm.types.llms.openai import ResponsesAPIResponse @@ -482,6 +483,90 @@ class TestLangfuseOtelIntegration: assert isinstance(config, OpenTelemetryConfig) # Endpoint assertion removed as side effect is gone + def test_request_metadata_is_emitted_under_langfuse_observation_and_trace_metadata_keys(self, monkeypatch): + monkeypatch.setattr(litellm, "redact_user_api_key_info", False) + request_metadata = { + "user_api_key_hash": "hash123", + "user_api_key_alias": "prod-key", + "user_api_key_user_id": "user-1", + "user_api_key_end_user_id": "end-user-1", + "user_api_key_team_id": "team-1", + "user_api_key_team_alias": "team-a", + "spend_logs_metadata": {"env": "prod"}, + "requester_ip_address": "10.0.0.1", + "requester_metadata": {"ticket": "LIT-8283"}, + } + kwargs = { + "litellm_params": {"metadata": {}}, + "standard_logging_object": {"metadata": request_metadata}, + } + + with patch( + "litellm.integrations.arize._utils.safe_set_attribute" + ) as mock_safe_set_attribute: + LangfuseOtelLogger._set_langfuse_specific_attributes(MagicMock(), kwargs, None) + + actual = {call.args[1]: call.args[2] for call in mock_safe_set_attribute.call_args_list} + + assert json.loads(actual["langfuse.observation.metadata"]) == {**request_metadata} + assert { + key: value for key, value in actual.items() if key.startswith("langfuse.trace.metadata.") + } == { + "langfuse.trace.metadata.key_alias": "prod-key", + "langfuse.trace.metadata.user_id": "user-1", + "langfuse.trace.metadata.end_user_id": "end-user-1", + "langfuse.trace.metadata.team_id": "team-1", + "langfuse.trace.metadata.team_alias": "team-a", + } + + def test_request_metadata_redaction_matches_vanilla_langfuse(self, monkeypatch): + monkeypatch.setattr(litellm, "redact_user_api_key_info", True) + request_metadata = { + "user_api_key_hash": "hash123", + "user_api_key_alias": "prod-key", + "user_api_key_user_id": "user-1", + "user_api_key_end_user_id": "end-user-1", + "user_api_key_team_id": "team-1", + "user_api_key_team_alias": "team-a", + "spend_logs_metadata": {"env": "prod"}, + "requester_ip_address": "10.0.0.1", + "requester_metadata": {"ticket": "LIT-8283"}, + } + kwargs = { + "litellm_params": {"metadata": {}}, + "standard_logging_object": {"metadata": request_metadata}, + } + + with patch( + "litellm.integrations.arize._utils.safe_set_attribute" + ) as mock_safe_set_attribute: + LangfuseOtelLogger._set_langfuse_specific_attributes(MagicMock(), kwargs, None) + + actual = {call.args[1]: call.args[2] for call in mock_safe_set_attribute.call_args_list} + + assert json.loads(actual["langfuse.observation.metadata"]) == { + "spend_logs_metadata": {"env": "prod"}, + "requester_ip_address": "10.0.0.1", + "requester_metadata": {"ticket": "LIT-8283"}, + } + assert not [key for key in actual if key.startswith("langfuse.trace.metadata.")] + + def test_request_metadata_keys_are_absent_without_standard_logging_object(self, monkeypatch): + monkeypatch.setattr(litellm, "redact_user_api_key_info", False) + kwargs = {"litellm_params": {"metadata": {"trace_metadata": {"k": "v"}}}} + + with patch( + "litellm.integrations.arize._utils.safe_set_attribute" + ) as mock_safe_set_attribute: + LangfuseOtelLogger._set_langfuse_specific_attributes(MagicMock(), kwargs, None) + + actual = {call.args[1]: call.args[2] for call in mock_safe_set_attribute.call_args_list} + + assert [key for key in actual if key.startswith("langfuse.trace.metadata")] == [ + "langfuse.trace.metadata" + ] + assert "langfuse.observation.metadata" not in actual + class TestLangfuseOtelKeyDynamicConfig: """Key/team-scoped Langfuse credentials must define the full export target