mirror of
https://github.com/BerriAI/litellm.git
synced 2026-10-05 02:41:56 +00:00
fix(otel/v2): deliver dynamic per-team/key/org langfuse_otel spans
Dynamic (team/key/org) langfuse_otel callbacks register on the success and failure hooks only, so log_pre_api_call never reached them and no LLM-call carrier was opened; _close_llm_call then dropped the span, so the tenant's Langfuse received nothing while the span was printed by the static console exporter. Create the span deferred at close when a real upstream call happened (payload present and not a no_upstream gate rejection). The langfuse preset now builds an otlp_http exporter even without proxy LANGFUSE_* env so per-request credentials have an OTLP destination, and dynamic routing derives the tenant OTLP endpoint from langfuse_host with per-endpoint provider caching, so a team on its own Langfuse region or host is delivered to instead of falling back to the console. Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
This commit is contained in:
parent
227b4fe5d4
commit
fc89e5f0d8
9 changed files with 304 additions and 70 deletions
|
|
@ -4,7 +4,6 @@ import os
|
|||
from datetime import datetime
|
||||
from typing import TYPE_CHECKING, Any, Optional, Union
|
||||
|
||||
from litellm._logging import verbose_logger
|
||||
from litellm.integrations.arize import _utils
|
||||
from litellm.integrations.langfuse.langfuse_otel_attributes import (
|
||||
LangfuseLLMObsOTELAttributes,
|
||||
|
|
@ -253,6 +252,22 @@ class LangfuseOtelLogger(OpenTelemetry):
|
|||
"""
|
||||
return os.environ.get("LANGFUSE_OTEL_HOST") or os.environ.get("LANGFUSE_HOST")
|
||||
|
||||
@staticmethod
|
||||
def get_langfuse_otel_endpoint(host: Optional[str] = None) -> str:
|
||||
"""
|
||||
Resolves the Langfuse OTLP traces endpoint from a host.
|
||||
|
||||
``host`` takes precedence (a per-request/team langfuse_host); otherwise the
|
||||
proxy env host is used, falling back to the US cloud endpoint. A host without
|
||||
a scheme is assumed https.
|
||||
"""
|
||||
langfuse_host = host or LangfuseOtelLogger._get_langfuse_otel_host()
|
||||
if not langfuse_host:
|
||||
return LANGFUSE_CLOUD_US_ENDPOINT
|
||||
if not langfuse_host.startswith("http"):
|
||||
langfuse_host = "https://" + langfuse_host
|
||||
return f"{langfuse_host.rstrip('/')}/api/public/otel"
|
||||
|
||||
def _create_open_telemetry_config_from_langfuse_env(self) -> OpenTelemetryConfig:
|
||||
"""
|
||||
Creates OpenTelemetryConfig from Langfuse environment variables.
|
||||
|
|
@ -267,20 +282,7 @@ class LangfuseOtelLogger(OpenTelemetry):
|
|||
# If no keys, return default from env (likely logging to console or something else)
|
||||
return OpenTelemetryConfig.from_env()
|
||||
|
||||
# Determine endpoint - default to US cloud
|
||||
langfuse_host = LangfuseOtelLogger._get_langfuse_otel_host()
|
||||
|
||||
if langfuse_host:
|
||||
# If LANGFUSE_HOST is provided, construct OTEL endpoint from it
|
||||
if not langfuse_host.startswith("http"):
|
||||
langfuse_host = "https://" + langfuse_host
|
||||
endpoint = f"{langfuse_host.rstrip('/')}/api/public/otel"
|
||||
verbose_logger.debug(f"Using Langfuse OTEL endpoint from host: {endpoint}")
|
||||
else:
|
||||
# Default to US cloud endpoint
|
||||
endpoint = LANGFUSE_CLOUD_US_ENDPOINT
|
||||
verbose_logger.debug(f"Using Langfuse US cloud endpoint: {endpoint}")
|
||||
|
||||
endpoint = LangfuseOtelLogger.get_langfuse_otel_endpoint()
|
||||
auth_header = LangfuseOtelLogger._get_langfuse_authorization_header(
|
||||
public_key=public_key, secret_key=secret_key
|
||||
)
|
||||
|
|
@ -316,20 +318,7 @@ class LangfuseOtelLogger(OpenTelemetry):
|
|||
"LANGFUSE_PUBLIC_KEY and LANGFUSE_SECRET_KEY must be set for Langfuse OpenTelemetry integration."
|
||||
)
|
||||
|
||||
# Determine endpoint - default to US cloud
|
||||
langfuse_host = LangfuseOtelLogger._get_langfuse_otel_host()
|
||||
|
||||
if langfuse_host:
|
||||
# If LANGFUSE_HOST is provided, construct OTEL endpoint from it
|
||||
if not langfuse_host.startswith("http"):
|
||||
langfuse_host = "https://" + langfuse_host
|
||||
endpoint = f"{langfuse_host.rstrip('/')}/api/public/otel"
|
||||
verbose_logger.debug(f"Using Langfuse OTEL endpoint from host: {endpoint}")
|
||||
else:
|
||||
# Default to US cloud endpoint
|
||||
endpoint = LANGFUSE_CLOUD_US_ENDPOINT
|
||||
verbose_logger.debug(f"Using Langfuse US cloud endpoint: {endpoint}")
|
||||
|
||||
endpoint = LangfuseOtelLogger.get_langfuse_otel_endpoint()
|
||||
auth_header = LangfuseOtelLogger._get_langfuse_authorization_header(
|
||||
public_key=public_key, secret_key=secret_key
|
||||
)
|
||||
|
|
|
|||
|
|
@ -375,47 +375,51 @@ class OpenTelemetryV2(CustomLogger):
|
|||
) -> Span | None:
|
||||
"""Finish the LLM-call span opened at ``pre_call`` (or create it deferred).
|
||||
|
||||
No carrier for this call id means ``pre_call`` never ran — the request was
|
||||
rejected at the gate or blocked by a pre-call guardrail before any upstream
|
||||
call — so there is nothing to record and no phantom span.
|
||||
No carrier for a genuine gate rejection (``is_no_upstream_call``) means
|
||||
``pre_call`` intentionally skipped it, so nothing is recorded. A missing
|
||||
carrier for a real call means ``pre_call`` never reached this logger — it
|
||||
happens for per-request dynamic callbacks (team/key/org), which register on
|
||||
the success hooks only, not on ``input_callback`` — so the span is created
|
||||
deferred here from the payload rather than dropped.
|
||||
"""
|
||||
call = LLMCallEvent.from_dict(kwargs)
|
||||
call_id = call.call_id
|
||||
# ``pop`` is the dedup: this method runs from both the success and failure
|
||||
# paths, and whichever fires first removes the carrier and closes the span.
|
||||
carrier = self._open_llm_calls.pop(call_id, None) if call_id else None
|
||||
if carrier is None:
|
||||
return None
|
||||
payload = call.payload
|
||||
if payload is None:
|
||||
if carrier.span is not None:
|
||||
if carrier is not None and carrier.span is not None:
|
||||
# Opened at the boundary but the payload never materialized — end
|
||||
# it (named provisionally) so it isn't leaked as an open span.
|
||||
carrier.span.end(end_time=to_ns(end_time))
|
||||
return None
|
||||
if carrier is None and call.is_no_upstream_call:
|
||||
return None
|
||||
data = LLMCallSpanData.from_standard_logging_payload(
|
||||
payload,
|
||||
capture_content=self.config.capture_span_content,
|
||||
time_to_first_chunk_seconds=call.time_to_first_chunk_seconds,
|
||||
)
|
||||
end_time_ns = to_ns(end_time)
|
||||
if carrier.span is not None:
|
||||
if carrier is not None and carrier.span is not None:
|
||||
# Born at the boundary: stamp attributes from the typed payload, set
|
||||
# status, and end it. Its parent (the server span) was captured at
|
||||
# creation from real ambient context.
|
||||
self._emitter.finish_span(SpanRole.LLM_CALL, carrier.span, data, end_time_ns=end_time_ns)
|
||||
return carrier.span
|
||||
# Deferred: ``pre_call`` saw no recordable parent, so create the span now.
|
||||
# The worker copied the request task's context, which carries the anchored
|
||||
# root span — parent to it (ambient fallback on the SDK path). Seed identity
|
||||
# Baggage so the span — and the SDK path, which has none — is labeled
|
||||
# consistently.
|
||||
# Deferred: ``pre_call`` either saw no recordable parent or never reached
|
||||
# this logger. The worker copied the request task's context, which carries
|
||||
# the anchored root span — parent to it (ambient fallback on the SDK path).
|
||||
# Seed identity Baggage so the span — and the SDK path, which has none — is
|
||||
# labeled consistently.
|
||||
start_time_ns = carrier.start_time_ns if carrier is not None else to_ns(start_time)
|
||||
parent_ctx = self._seed_identity_baggage(data.identity, data.request_model, resolve_request_span_context())
|
||||
return self._emitter.emit(
|
||||
SpanRole.LLM_CALL,
|
||||
data,
|
||||
parent_context=parent_ctx,
|
||||
start_time_ns=carrier.start_time_ns,
|
||||
start_time_ns=start_time_ns,
|
||||
end_time_ns=end_time_ns,
|
||||
tracer=self._tenant_tracers.tracer_for(self.tracer, call.dynamic_params),
|
||||
)
|
||||
|
|
|
|||
|
|
@ -16,13 +16,16 @@ from opentelemetry.trace import Tracer
|
|||
|
||||
from litellm._logging import verbose_logger
|
||||
from litellm.integrations.otel.model.config import OpenTelemetryV2Config
|
||||
from litellm.integrations.otel.presets import dynamic_otlp_headers
|
||||
from litellm.integrations.otel.presets import (
|
||||
dynamic_otlp_endpoint,
|
||||
dynamic_otlp_headers,
|
||||
)
|
||||
from litellm.integrations.otel.plumbing.providers import (
|
||||
build_tracer_provider,
|
||||
get_tracer,
|
||||
)
|
||||
|
||||
# Exporter kinds that ignore headers — never rewritten with dynamic credentials.
|
||||
# Exporter kinds that ignore endpoint/headers — never rewritten with dynamic credentials.
|
||||
_NON_OTLP_KINDS = ("console", "in_memory", "inmemory", "memory")
|
||||
|
||||
# Cap on distinct credential-scoped providers held at once. ``dynamic_params``
|
||||
|
|
@ -60,25 +63,29 @@ class TenantTracerCache:
|
|||
self._config = config
|
||||
self._callback_name = callback_name
|
||||
self._tracer_name = tracer_name
|
||||
self._providers: "OrderedDict[tuple[tuple[str, str], ...], TracerProvider]" = OrderedDict()
|
||||
self._providers: "OrderedDict[tuple[object, ...], TracerProvider]" = OrderedDict()
|
||||
|
||||
def tracer_for(self, default: Tracer, dynamic_params: Any) -> Tracer:
|
||||
"""Return the tracer for this request.
|
||||
|
||||
Use ``default`` unless the request's dynamic credentials require a
|
||||
credential-scoped tracer, in which case build (or reuse) one. The cache
|
||||
is a bounded LRU: the least-recently-used provider is flushed and shut
|
||||
down on overflow so its exporter threads don't accumulate.
|
||||
credential-scoped tracer, in which case build (or reuse) one. A tenant can
|
||||
override the exporter's auth headers (all dynamic backends) and its
|
||||
endpoint (backends whose host is tenant-scoped, e.g. Langfuse regions), so
|
||||
both feed the cache key. The cache is a bounded LRU: the least-recently-used
|
||||
provider is flushed and shut down on overflow so its exporter threads don't
|
||||
accumulate.
|
||||
"""
|
||||
headers = dynamic_otlp_headers(self._callback_name, dynamic_params)
|
||||
if not headers:
|
||||
endpoint = dynamic_otlp_endpoint(self._callback_name, dynamic_params)
|
||||
if not headers and not endpoint:
|
||||
return default
|
||||
cache_key = tuple(sorted(headers.items()))
|
||||
cache_key = (tuple(sorted((headers or {}).items())), endpoint)
|
||||
provider = self._providers.get(cache_key)
|
||||
if provider is not None:
|
||||
self._providers.move_to_end(cache_key)
|
||||
else:
|
||||
provider = build_tracer_provider(self._config_with_headers(headers))
|
||||
provider = build_tracer_provider(self._config_with_overrides(headers, endpoint))
|
||||
self._providers[cache_key] = provider
|
||||
if len(self._providers) > _MAX_CACHED_PROVIDERS:
|
||||
_, evicted = self._providers.popitem(last=False)
|
||||
|
|
@ -86,21 +93,26 @@ class TenantTracerCache:
|
|||
return get_tracer(provider, self._tracer_name)
|
||||
|
||||
def _config_with_headers(self, headers: Mapping[str, str]) -> OpenTelemetryV2Config:
|
||||
"""Clone the config, stamping ``headers`` onto the credential's own exporter.
|
||||
return self._config_with_overrides(dict(headers), None)
|
||||
|
||||
``headers`` are the per-request credentials of ``self._callback_name`` (the
|
||||
def _config_with_overrides(self, headers: Mapping[str, str] | None, endpoint: str | None) -> OpenTelemetryV2Config:
|
||||
"""Clone the config, stamping ``headers``/``endpoint`` onto the credential's own exporter.
|
||||
|
||||
The overrides are the per-request credentials of ``self._callback_name`` (the
|
||||
integration that built this cache), so they apply only to the exporter that
|
||||
integration contributed (``spec.owner``). A request that carries one
|
||||
tenant's Arize key must never rewrite the headers of a co-configured
|
||||
Langfuse or self-hosted collector exporter, which would leak that key to a
|
||||
different backend.
|
||||
"""
|
||||
header_str = ",".join(f"{key}={value}" for key, value in headers.items())
|
||||
header_update: dict[str, str] = {"headers": header_str}
|
||||
update: dict[str, str] = {
|
||||
**({"headers": ",".join(f"{key}={value}" for key, value in headers.items())} if headers else {}),
|
||||
**({"endpoint": endpoint} if endpoint else {}),
|
||||
}
|
||||
exporters = [
|
||||
(
|
||||
spec.model_copy(update=header_update)
|
||||
if spec.owner == self._callback_name and spec.kind.lower() not in _NON_OTLP_KINDS
|
||||
spec.model_copy(update=update)
|
||||
if update and spec.owner == self._callback_name and spec.kind.lower() not in _NON_OTLP_KINDS
|
||||
else spec
|
||||
)
|
||||
for spec in self._config.exporters
|
||||
|
|
|
|||
|
|
@ -14,6 +14,7 @@ from litellm.integrations.otel.presets.agentops import agentops_preset
|
|||
from litellm.integrations.otel.presets.arize import arize_dynamic_headers, arize_preset
|
||||
from litellm.integrations.otel.presets.base import Preset
|
||||
from litellm.integrations.otel.presets.langfuse import (
|
||||
langfuse_dynamic_endpoint,
|
||||
langfuse_dynamic_headers,
|
||||
langfuse_preset,
|
||||
)
|
||||
|
|
@ -45,6 +46,14 @@ DYNAMIC_HEADERS_BY_CALLBACK: dict[str, Callable[[StandardCallbackDynamicParams],
|
|||
"weave_otel": weave_dynamic_headers,
|
||||
}
|
||||
|
||||
#: Callback name → per-request OTLP endpoint builder. Only backends whose
|
||||
#: destination host is tenant-scoped appear here — Langfuse routes a team/key to
|
||||
#: its own regional/self-hosted host, so the endpoint (not just the auth header)
|
||||
#: must be resolved per request.
|
||||
DYNAMIC_ENDPOINT_BY_CALLBACK: dict[str, Callable[[StandardCallbackDynamicParams], str | None]] = {
|
||||
"langfuse_otel": langfuse_dynamic_endpoint,
|
||||
}
|
||||
|
||||
|
||||
def dynamic_otlp_headers(
|
||||
callback_name: str | None,
|
||||
|
|
@ -61,11 +70,27 @@ def dynamic_otlp_headers(
|
|||
return headers or None
|
||||
|
||||
|
||||
def dynamic_otlp_endpoint(
|
||||
callback_name: str | None,
|
||||
dynamic_params: StandardCallbackDynamicParams | None,
|
||||
) -> str | None:
|
||||
"""Per-request OTLP endpoint for ``callback_name``, or ``None`` if N/A.
|
||||
|
||||
``None`` means the exporter's configured endpoint is kept as-is.
|
||||
"""
|
||||
builder = DYNAMIC_ENDPOINT_BY_CALLBACK.get(callback_name or "")
|
||||
if builder is None or not dynamic_params:
|
||||
return None
|
||||
return builder(dynamic_params)
|
||||
|
||||
|
||||
__all__ = [
|
||||
"PRESET_BY_CALLBACK",
|
||||
"DYNAMIC_HEADERS_BY_CALLBACK",
|
||||
"DYNAMIC_ENDPOINT_BY_CALLBACK",
|
||||
"Preset",
|
||||
"dynamic_otlp_headers",
|
||||
"dynamic_otlp_endpoint",
|
||||
"agentops_preset",
|
||||
"arize_preset",
|
||||
"langfuse_preset",
|
||||
|
|
|
|||
|
|
@ -1,5 +1,7 @@
|
|||
"""Langfuse-OTEL preset."""
|
||||
|
||||
import os
|
||||
|
||||
from litellm.integrations.langfuse.langfuse_otel import (
|
||||
LangfuseOtelLogger as _V1Langfuse,
|
||||
)
|
||||
|
|
@ -12,21 +14,29 @@ from litellm.integrations.otel.presets.utils import ensure_mappers
|
|||
from litellm.types.utils import StandardCallbackDynamicParams
|
||||
|
||||
|
||||
def _langfuse_static_headers() -> str | None:
|
||||
"""OTLP auth header from proxy env keys, or ``None`` when only dynamic creds are used."""
|
||||
public_key = os.environ.get("LANGFUSE_PUBLIC_KEY")
|
||||
secret_key = os.environ.get("LANGFUSE_SECRET_KEY")
|
||||
if not public_key or not secret_key:
|
||||
return None
|
||||
auth = _V1Langfuse._get_langfuse_authorization_header(public_key=public_key, secret_key=secret_key)
|
||||
return f"Authorization={auth}"
|
||||
|
||||
|
||||
def langfuse_preset(
|
||||
*,
|
||||
config_overrides: OpenTelemetryV2Config | None = None,
|
||||
) -> OpenTelemetryV2Config:
|
||||
cfg = _V1Langfuse.get_langfuse_otel_config()
|
||||
kind = cfg.exporter if isinstance(cfg.exporter, str) else "otlp_http"
|
||||
base = config_overrides or OpenTelemetryV2Config()
|
||||
return base.model_copy(
|
||||
update={
|
||||
"exporters": [
|
||||
*base.exporters,
|
||||
ExporterSpec(
|
||||
kind=kind,
|
||||
endpoint=cfg.endpoint,
|
||||
headers=cfg.headers,
|
||||
kind="otlp_http",
|
||||
endpoint=_V1Langfuse.get_langfuse_otel_endpoint(),
|
||||
headers=_langfuse_static_headers(),
|
||||
owner=ExporterOwner.LANGFUSE_OTEL,
|
||||
),
|
||||
],
|
||||
|
|
@ -46,3 +56,9 @@ def langfuse_dynamic_headers(params: StandardCallbackDynamicParams) -> dict[str,
|
|||
)
|
||||
}
|
||||
return {}
|
||||
|
||||
|
||||
def langfuse_dynamic_endpoint(params: StandardCallbackDynamicParams) -> str | None:
|
||||
"""Per-request Langfuse OTLP endpoint from a team/key ``langfuse_host``, or ``None``."""
|
||||
host = params.get("langfuse_host")
|
||||
return _V1Langfuse.get_langfuse_otel_endpoint(host) if host else None
|
||||
|
|
|
|||
|
|
@ -8,7 +8,7 @@ sys.path.insert(0, os.path.abspath("../../../.."))
|
|||
from opentelemetry.trace import NoOpTracer
|
||||
|
||||
from litellm.integrations.otel.model.config import ExporterSpec, OpenTelemetryV2Config
|
||||
from litellm.integrations.otel.presets import dynamic_otlp_headers
|
||||
from litellm.integrations.otel.presets import dynamic_otlp_endpoint, dynamic_otlp_headers
|
||||
from litellm.integrations.otel.plumbing.routing import TenantTracerCache
|
||||
|
||||
|
||||
|
|
@ -173,3 +173,95 @@ def test_dynamic_headers_do_not_leak_to_other_owners_exporter():
|
|||
assert by_owner["arize"] == "arize-space-id=TEAMX,api_key=TEAMX_KEY"
|
||||
assert by_owner[None] == "x=base-collector"
|
||||
assert by_owner["langfuse_otel"] == "Authorization=Basic base-langfuse"
|
||||
|
||||
|
||||
# --- Langfuse routes a tenant to its own OTLP endpoint, not just auth --- #
|
||||
|
||||
|
||||
def test_langfuse_dynamic_endpoint_from_host():
|
||||
assert (
|
||||
dynamic_otlp_endpoint(
|
||||
"langfuse_otel", {"langfuse_host": "https://cloud.langfuse.com"}
|
||||
)
|
||||
== "https://cloud.langfuse.com/api/public/otel"
|
||||
)
|
||||
|
||||
|
||||
def test_langfuse_dynamic_endpoint_absent_without_host():
|
||||
"""Auth alone must not force an endpoint override; the exporter keeps its
|
||||
configured (proxy env / default) host."""
|
||||
assert (
|
||||
dynamic_otlp_endpoint(
|
||||
"langfuse_otel",
|
||||
{"langfuse_public_key": "pk", "langfuse_secret_key": "sk"},
|
||||
)
|
||||
is None
|
||||
)
|
||||
|
||||
|
||||
def test_only_langfuse_participates_in_dynamic_endpoint():
|
||||
assert dynamic_otlp_endpoint("arize", {"langfuse_host": "https://x"}) is None
|
||||
assert dynamic_otlp_endpoint("weave_otel", {"langfuse_host": "https://x"}) is None
|
||||
assert dynamic_otlp_endpoint(None, {"langfuse_host": "https://x"}) is None
|
||||
assert dynamic_otlp_endpoint("langfuse_otel", None) is None
|
||||
|
||||
|
||||
def test_langfuse_endpoint_and_auth_stamped_on_owned_exporter_only():
|
||||
"""A tenant's ``langfuse_host`` rewrites the endpoint (and auth) of the
|
||||
Langfuse exporter alone. Regression for dynamic Langfuse spans being exported
|
||||
to the static/console destination instead of the team's own Langfuse host."""
|
||||
cache = _cache(
|
||||
"langfuse_otel",
|
||||
exporters=[
|
||||
ExporterSpec(
|
||||
kind="otlp_http",
|
||||
endpoint="https://us.cloud.langfuse.com/api/public/otel",
|
||||
headers="Authorization=Basic base",
|
||||
owner="langfuse_otel",
|
||||
),
|
||||
ExporterSpec(
|
||||
kind="otlp_http",
|
||||
endpoint="http://self-hosted-collector:4318",
|
||||
headers="x=base-collector",
|
||||
owner=None,
|
||||
),
|
||||
ExporterSpec(kind="console", owner="langfuse_otel"),
|
||||
],
|
||||
)
|
||||
new_cfg = cache._config_with_overrides(
|
||||
{"Authorization": "Basic TENANT"},
|
||||
"https://cloud.langfuse.com/api/public/otel",
|
||||
)
|
||||
langfuse_otlp = next(
|
||||
e
|
||||
for e in new_cfg.exporters
|
||||
if e.owner == "langfuse_otel" and e.kind == "otlp_http"
|
||||
)
|
||||
assert langfuse_otlp.endpoint == "https://cloud.langfuse.com/api/public/otel"
|
||||
assert langfuse_otlp.headers == "Authorization=Basic TENANT"
|
||||
|
||||
collector = next(e for e in new_cfg.exporters if e.owner is None)
|
||||
assert collector.endpoint == "http://self-hosted-collector:4318"
|
||||
assert collector.headers == "x=base-collector"
|
||||
|
||||
console = next(e for e in new_cfg.exporters if e.kind == "console")
|
||||
assert console.endpoint is None and console.headers is None
|
||||
|
||||
|
||||
def test_provider_cached_per_endpoint_for_same_auth():
|
||||
"""Two tenants sharing keys but on different Langfuse regions must not share a
|
||||
provider, or the second tenant's spans export to the first's endpoint."""
|
||||
cache = _cache(
|
||||
"langfuse_otel",
|
||||
exporters=[ExporterSpec(kind="in_memory", owner="langfuse_otel")],
|
||||
)
|
||||
default = NoOpTracer()
|
||||
base = {"langfuse_public_key": "pk", "langfuse_secret_key": "sk"}
|
||||
eu = {**base, "langfuse_host": "https://cloud.langfuse.com"}
|
||||
us = {**base, "langfuse_host": "https://us.cloud.langfuse.com"}
|
||||
|
||||
cache.tracer_for(default, eu)
|
||||
cache.tracer_for(default, eu)
|
||||
assert len(cache._providers) == 1
|
||||
cache.tracer_for(default, us)
|
||||
assert len(cache._providers) == 2
|
||||
|
|
|
|||
|
|
@ -292,23 +292,41 @@ def test_missing_standard_logging_object_is_noop():
|
|||
assert exporter.get_finished_spans() == ()
|
||||
|
||||
|
||||
def test_no_span_when_pre_call_never_ran():
|
||||
def test_no_span_when_gate_rejection_marks_no_upstream_call():
|
||||
"""A request rejected before the upstream call — at the auth/budget gate, or
|
||||
blocked by a pre-call guardrail — never reaches ``pre_call``, so there is no
|
||||
carrier and the failure log produces no phantom CLIENT span. This replaces the
|
||||
old post-hoc heuristics: "did pre_call run?" is the only signal needed."""
|
||||
blocked by a pre-call guardrail — is tagged ``LITELLM_LOGGING_NO_UPSTREAM_LLM_CALL``
|
||||
(see ``ProxyLogging._handle_logging_proxy_only_error``), so the failure log
|
||||
opens no carrier and produces no phantom CLIENT span."""
|
||||
from litellm.constants import LITELLM_LOGGING_NO_UPSTREAM_LLM_CALL
|
||||
|
||||
logger, exporter = _logger()
|
||||
payload = _payload(
|
||||
status="failure",
|
||||
error_information={"error_class": "ProxyException", "error_code": "401"},
|
||||
)
|
||||
kwargs = _kwargs(payload=payload)
|
||||
kwargs[LITELLM_LOGGING_NO_UPSTREAM_LLM_CALL] = True
|
||||
# No log_pre_api_call: the call never started.
|
||||
asyncio.run(
|
||||
logger.async_log_failure_event(_kwargs(payload=payload), None, None, None)
|
||||
)
|
||||
asyncio.run(logger.async_log_failure_event(kwargs, None, None, None))
|
||||
assert exporter.get_finished_spans() == () # no phantom LLM span
|
||||
|
||||
|
||||
def test_dynamic_callback_emits_deferred_span_without_pre_call():
|
||||
"""A per-request dynamic callback (team/key/org langfuse_otel and friends)
|
||||
registers on the success/failure hooks only, so ``pre_call`` never reaches it
|
||||
and no carrier is opened. A real upstream call still happened (payload present,
|
||||
no ``no_upstream`` marker), so the span must be created deferred at close rather
|
||||
than dropped — otherwise the tenant's Langfuse receives nothing."""
|
||||
logger, exporter = _logger()
|
||||
# No log_pre_api_call for this logger: the pre_call hook was dispatched to the
|
||||
# globally-registered loggers, not this per-request one.
|
||||
asyncio.run(logger.async_log_success_event(_kwargs(), None, None, None))
|
||||
(span,) = exporter.get_finished_spans()
|
||||
assert span.name == "chat gpt-4o"
|
||||
assert span.attributes[LiteLLM.CALL_ID] == "call_1"
|
||||
assert span.status.status_code is StatusCode.UNSET
|
||||
|
||||
|
||||
def test_real_llm_failure_still_emitted():
|
||||
"""A genuine LLM failure: ``pre_call`` ran (the call was attempted), so the
|
||||
CLIENT span is opened at the boundary and closed ERROR."""
|
||||
|
|
|
|||
|
|
@ -75,6 +75,52 @@ def test_dynamic_cred_presets_tag_exporter_with_matching_owner(monkeypatch):
|
|||
)
|
||||
|
||||
|
||||
def test_langfuse_preset_builds_otlp_exporter_without_env_creds(monkeypatch):
|
||||
"""Regression: dynamic team/key/org Langfuse (no proxy ``LANGFUSE_*`` env)
|
||||
must still yield an ``otlp_http`` exporter so per-request credentials have an
|
||||
OTLP destination to be stamped onto. The preset previously required env creds
|
||||
(it raised without them), so the logger fell back to the console exporter and
|
||||
dynamic Langfuse spans were printed to the proxy instead of delivered."""
|
||||
from litellm.integrations.otel.model.config import ExporterOwner
|
||||
from litellm.integrations.otel.presets.langfuse import langfuse_preset
|
||||
|
||||
for var in (
|
||||
"LANGFUSE_PUBLIC_KEY",
|
||||
"LANGFUSE_SECRET_KEY",
|
||||
"LANGFUSE_HOST",
|
||||
"LANGFUSE_OTEL_HOST",
|
||||
):
|
||||
monkeypatch.delenv(var, raising=False)
|
||||
|
||||
cfg = langfuse_preset()
|
||||
langfuse_exporters = [
|
||||
e for e in cfg.exporters if e.owner == ExporterOwner.LANGFUSE_OTEL
|
||||
]
|
||||
assert len(langfuse_exporters) == 1
|
||||
spec = langfuse_exporters[0]
|
||||
assert spec.kind == "otlp_http"
|
||||
assert spec.endpoint == "https://us.cloud.langfuse.com/api/public/otel"
|
||||
assert spec.headers is None
|
||||
|
||||
|
||||
def test_langfuse_preset_uses_env_creds_for_static_headers(monkeypatch):
|
||||
"""The static/global path (proxy ``LANGFUSE_*`` env) keeps working: the
|
||||
exporter carries the env-derived endpoint and Basic auth header."""
|
||||
from litellm.integrations.otel.model.config import ExporterOwner
|
||||
from litellm.integrations.otel.presets.langfuse import langfuse_preset
|
||||
|
||||
monkeypatch.setenv("LANGFUSE_PUBLIC_KEY", "pk")
|
||||
monkeypatch.setenv("LANGFUSE_SECRET_KEY", "sk")
|
||||
monkeypatch.setenv("LANGFUSE_HOST", "https://cloud.langfuse.com")
|
||||
monkeypatch.delenv("LANGFUSE_OTEL_HOST", raising=False)
|
||||
|
||||
cfg = langfuse_preset()
|
||||
spec = next(e for e in cfg.exporters if e.owner == ExporterOwner.LANGFUSE_OTEL)
|
||||
assert spec.kind == "otlp_http"
|
||||
assert spec.endpoint == "https://cloud.langfuse.com/api/public/otel"
|
||||
assert spec.headers is not None and spec.headers.startswith("Authorization=Basic ")
|
||||
|
||||
|
||||
def test_agentops_exporter_mints_jwt_lazily(monkeypatch):
|
||||
pytest.importorskip("opentelemetry.exporter.otlp.proto.http.trace_exporter")
|
||||
monkeypatch.setattr(
|
||||
|
|
|
|||
|
|
@ -40,6 +40,38 @@ class TestLangfuseOtelIntegration:
|
|||
# assert os.environ.get("OTEL_EXPORTER_OTLP_ENDPOINT") == "https://us.cloud.langfuse.com/api/public/otel"
|
||||
# assert "Authorization=Basic" in os.environ.get("OTEL_EXPORTER_OTLP_HEADERS", "")
|
||||
|
||||
def test_get_langfuse_otel_endpoint_defaults_to_us_cloud(self):
|
||||
with patch.dict(os.environ, {}, clear=True):
|
||||
assert (
|
||||
LangfuseOtelLogger.get_langfuse_otel_endpoint()
|
||||
== "https://us.cloud.langfuse.com/api/public/otel"
|
||||
)
|
||||
|
||||
def test_get_langfuse_otel_endpoint_from_host_arg(self):
|
||||
assert (
|
||||
LangfuseOtelLogger.get_langfuse_otel_endpoint("https://cloud.langfuse.com")
|
||||
== "https://cloud.langfuse.com/api/public/otel"
|
||||
)
|
||||
|
||||
def test_get_langfuse_otel_endpoint_adds_https_when_scheme_missing(self):
|
||||
assert (
|
||||
LangfuseOtelLogger.get_langfuse_otel_endpoint("my-langfuse.com")
|
||||
== "https://my-langfuse.com/api/public/otel"
|
||||
)
|
||||
|
||||
def test_get_langfuse_otel_endpoint_arg_overrides_env(self):
|
||||
with patch.dict(
|
||||
os.environ, {"LANGFUSE_HOST": "https://env-host.com"}, clear=True
|
||||
):
|
||||
assert (
|
||||
LangfuseOtelLogger.get_langfuse_otel_endpoint("https://arg-host.com")
|
||||
== "https://arg-host.com/api/public/otel"
|
||||
)
|
||||
assert (
|
||||
LangfuseOtelLogger.get_langfuse_otel_endpoint()
|
||||
== "https://env-host.com/api/public/otel"
|
||||
)
|
||||
|
||||
def test_get_langfuse_otel_config_missing_keys(self):
|
||||
"""Test that ValueError is raised when required keys are missing."""
|
||||
with patch.dict(os.environ, {}, clear=True):
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue