mirror of
https://github.com/BerriAI/litellm.git
synced 2026-10-10 03:28:53 +00:00
fix(otel v2): emit request metadata and identity under langfuse.observation.metadata and langfuse.trace.metadata
Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
This commit is contained in:
parent
439327ded2
commit
d964f34c06
5 changed files with 231 additions and 4 deletions
|
|
@ -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)
|
||||
),
|
||||
|
|
|
|||
|
|
@ -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.
|
||||
|
|
|
|||
|
|
@ -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 (),
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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": {}}
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue