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>
This commit is contained in:
yucheng 2026-09-23 07:03:10 +00:00
parent 40ec84caa2
commit 81a799ec2e
5 changed files with 314 additions and 0 deletions

View file

@ -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:

View file

@ -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"

View file

@ -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"
],

View file

@ -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)

View file

@ -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