mirror of
https://github.com/BerriAI/litellm.git
synced 2026-10-10 03:28:53 +00:00
fix(otel v2): validate tenant destinations and match each backend's own endpoint
Review round on the tenant destination routing. - A key/team Langfuse host is user-supplied input, so it goes through the proxy's SSRF guard. A private address is refused, the operator keeps the trace, and the warning names user_url_allowed_hosts. The operator's own LANGFUSE_HOST is not checked. - Arize and Weave destinations now resolve their endpoint and transport through the backend's own config, so an ARIZE_HTTP_ENDPOINT collector and a self-hosted WANDB_HOST are honoured instead of the cloud default. - A half-configured backend no longer resolves: several dynamic header builders gate each credential separately, so an api key with no space id produced a non-empty but unusable header set that suppressed the operator's exporter. - A callback_type of "failure" no longer takes over the trace. The destination is resolved during auth, before the outcome is known. - The fan-out cache evicts without shutting the processor down, matching ArizePhoenixLogger: a concurrent on_end may still hold it. - The stdout placeholder is identified by what it does rather than by equality with an import-time default, so an operator's OTEL_EXPORTER_OTLP_* collector survives the credential-less path.
This commit is contained in:
parent
bcba86e426
commit
1f6b80e659
8 changed files with 390 additions and 69 deletions
|
|
@ -1,10 +1,7 @@
|
|||
"""The resolved OTLP destination a request's traces export to.
|
||||
|
||||
A destination is a backend-agnostic target: an endpoint plus the auth headers the
|
||||
exporter sends. The proxy builds one per backend from the key or team logging
|
||||
config resolved at auth, and the fan-out span processor exports the request's
|
||||
spans through it. Every OTEL backend reduces to this shape; the per-backend field
|
||||
mapping lives in ``litellm.integrations.otel.presets.destinations``.
|
||||
Backend-agnostic on purpose: every OTEL backend reduces to an endpoint plus auth
|
||||
headers. The per-backend field mapping lives in ``presets.destinations``.
|
||||
"""
|
||||
|
||||
from collections.abc import Mapping
|
||||
|
|
@ -24,10 +21,8 @@ class OtelDestination(BaseModel):
|
|||
protocol: str | None = Field(
|
||||
default=None,
|
||||
description=(
|
||||
"OTLP transport for this endpoint (``otlp_http`` / ``otlp_grpc``). The "
|
||||
"backend's intrinsic default is used when unset. A backend whose own cloud "
|
||||
"endpoint is gRPC can still be pointed at an HTTP collector, which the "
|
||||
"scheme alone cannot express: Arize's own ``https://otlp.arize.com/v1`` is gRPC."
|
||||
"OTLP transport, defaulting to the backend's own. Not derivable from the "
|
||||
"scheme: Arize's ``https://otlp.arize.com/v1`` is gRPC."
|
||||
),
|
||||
)
|
||||
|
||||
|
|
@ -42,8 +37,7 @@ class OtelDestination(BaseModel):
|
|||
return ",".join(f"{key}={quote(value, safe='')}" for key, value in self.headers.items())
|
||||
|
||||
def cache_key(self) -> tuple[str, tuple[tuple[str, str], ...], tuple[tuple[str, str], ...], str | None]:
|
||||
"""Identity for processor reuse: two requests naming the same destination
|
||||
must share one exporter rather than minting a connection pool each."""
|
||||
"""Identity for processor reuse, so one destination means one exporter."""
|
||||
return (
|
||||
self.endpoint,
|
||||
tuple(sorted(self.headers.items())),
|
||||
|
|
|
|||
|
|
@ -242,11 +242,9 @@ def _with_destination_resource(span: ReadableSpan, destination: "OtelDestination
|
|||
class TenantFanOutSpanProcessor(SpanProcessor):
|
||||
"""Export every finished span to each destination this request resolved.
|
||||
|
||||
Destinations ride a request-scoped ``ContextVar`` set during auth, so the
|
||||
processor keeps no per-request state and concurrent requests stay isolated.
|
||||
Every span is forwarded, the gen-AI span included: the tenant's account gets the
|
||||
tree the operator's would have received, still parented, because the forwarded
|
||||
view keeps the original span's trace and parent ids.
|
||||
Destinations ride a request-scoped ``ContextVar`` set during auth, so concurrent
|
||||
requests stay isolated. The forwarded view keeps the original trace and parent
|
||||
ids, so the tenant gets the same tree the operator would have received.
|
||||
"""
|
||||
|
||||
def __init__(
|
||||
|
|
@ -314,13 +312,11 @@ class TenantFanOutSpanProcessor(SpanProcessor):
|
|||
_shutdown_quietly(built)
|
||||
return existing
|
||||
self._processors[key] = built
|
||||
evicted: Final = (
|
||||
self._processors.popitem(last=False)[1]
|
||||
if len(self._processors) > _MAX_CACHED_DESTINATION_PROCESSORS
|
||||
else None
|
||||
)
|
||||
if evicted is not None:
|
||||
_shutdown_quietly(evicted)
|
||||
if len(self._processors) > _MAX_CACHED_DESTINATION_PROCESSORS:
|
||||
# Evict without shutting down: another thread may be inside ``on_end``
|
||||
# holding the victim, and a shut-down BatchSpanProcessor drops spans
|
||||
# silently. Same rule as ArizePhoenixLogger's per-project cache.
|
||||
self._processors.popitem(last=False)
|
||||
return built
|
||||
|
||||
|
||||
|
|
@ -350,11 +346,8 @@ class _OverriddenBackendFilter(SpanProcessor):
|
|||
"""Hold a span back from ``owner``'s operator-level exporter when the request
|
||||
pointed ``owner`` at a tenant's own account.
|
||||
|
||||
A team destination is an override rather than an addition, and a ``SpanProcessor``
|
||||
cannot veto its siblings (``SynchronousMultiSpanProcessor.on_end`` ignores return
|
||||
values), so suppression has to wrap the exporter's own processor. Dropping here
|
||||
also keeps the span out of ``BatchSpanProcessor``'s bounded queue instead of
|
||||
filling it with spans that will never ship.
|
||||
Wrapping is the only place this works: ``SynchronousMultiSpanProcessor.on_end``
|
||||
ignores return values, so a sibling processor can never veto the export.
|
||||
"""
|
||||
|
||||
def __init__(self, inner: SpanProcessor, owner: str) -> None:
|
||||
|
|
|
|||
|
|
@ -1,11 +1,8 @@
|
|||
"""Map a key's or team's callback vars to the OTLP destination its traces export to.
|
||||
|
||||
The auth path resolves one destination per backend the caller configured, and the
|
||||
fan-out span processor exports the whole request through it. Header building is
|
||||
delegated to each preset's existing ``*_dynamic_headers`` builder, so a destination
|
||||
authenticates exactly the way the per-request tracer route already did; only the
|
||||
endpoint needs a per-backend rule, because a backend's host is either fixed, taken
|
||||
from a region table, or named by the tenant alongside its own key pair.
|
||||
Header building is delegated to each preset's existing ``*_dynamic_headers`` builder,
|
||||
so a destination authenticates exactly the way the per-request tracer route already
|
||||
did; only the endpoint and transport need a per-backend rule.
|
||||
"""
|
||||
|
||||
import os
|
||||
|
|
@ -13,44 +10,63 @@ from collections.abc import Callable, Mapping
|
|||
from types import MappingProxyType
|
||||
from typing import Final
|
||||
|
||||
from litellm._logging import verbose_logger
|
||||
from litellm.integrations.otel.model.destination import OtelDestination
|
||||
from litellm.litellm_core_utils.url_utils import SSRFError, assert_public_url
|
||||
from litellm.types.utils import StandardCallbackDynamicParams
|
||||
|
||||
#: gRPC is Arize's own transport; an explicitly named HTTP collector overrides it.
|
||||
_ARIZE_GRPC_ENDPOINT: Final = "https://otlp.arize.com/v1"
|
||||
|
||||
|
||||
def _langfuse_endpoint(params: StandardCallbackDynamicParams) -> str | None:
|
||||
"""The tenant's own Langfuse host, else the operator's, else Langfuse US cloud.
|
||||
|
||||
Falling back to the operator's host is safe and is what V1 does: the tenant's
|
||||
own key pair still selects its own project, and a self-hosted deployment where
|
||||
every team lives on one Langfuse server is the common shape.
|
||||
A host the tenant named goes through the proxy's SSRF guard first. Anyone who can
|
||||
mint a key can write it, so without the check it points the exporter, and the
|
||||
tenant credentials it carries, at any address the proxy can reach. The operator's
|
||||
own ``LANGFUSE_HOST`` is not checked: an internal collector is a normal
|
||||
deployment and the operator is the one configuring it.
|
||||
"""
|
||||
from litellm.integrations.langfuse.langfuse_otel import (
|
||||
LANGFUSE_CLOUD_US_ENDPOINT,
|
||||
LangfuseOtelLogger,
|
||||
)
|
||||
|
||||
host: Final = params.get("langfuse_host") or LangfuseOtelLogger._get_langfuse_otel_host() # pyright: ignore[reportPrivateUsage] # reuse the backend's own env host resolver rather than duplicating it
|
||||
tenant_host: Final = params.get("langfuse_host") or None
|
||||
host: Final = tenant_host or LangfuseOtelLogger._get_langfuse_otel_host() # pyright: ignore[reportPrivateUsage] # reuse the backend's own env host resolver rather than duplicating it
|
||||
if not host:
|
||||
return LANGFUSE_CLOUD_US_ENDPOINT
|
||||
normalized: Final = host if host.startswith("http") else f"https://{host}"
|
||||
return f"{normalized.rstrip('/')}/api/public/otel"
|
||||
endpoint: Final = f"{normalized.rstrip('/')}/api/public/otel"
|
||||
if tenant_host is None:
|
||||
return endpoint
|
||||
try:
|
||||
assert_public_url(endpoint)
|
||||
except SSRFError as exc:
|
||||
verbose_logger.warning(
|
||||
"OTel V2: not exporting to key/team Langfuse host '%s' (%s). "
|
||||
"Add it to general_settings.user_url_allowed_hosts to permit it",
|
||||
host,
|
||||
exc,
|
||||
)
|
||||
return None
|
||||
return endpoint
|
||||
|
||||
|
||||
def _arize_endpoint(params: StandardCallbackDynamicParams) -> str | None:
|
||||
return os.environ.get("ARIZE_ENDPOINT") or _ARIZE_GRPC_ENDPOINT
|
||||
from litellm.integrations.arize.arize import ArizeLogger
|
||||
|
||||
return ArizeLogger.get_arize_config().endpoint
|
||||
|
||||
|
||||
def _arize_protocol(params: StandardCallbackDynamicParams) -> str | None:
|
||||
return "otlp_http" if os.environ.get("ARIZE_HTTP_ENDPOINT") and not os.environ.get("ARIZE_ENDPOINT") else None
|
||||
from litellm.integrations.arize.arize import ArizeLogger
|
||||
|
||||
return ArizeLogger.get_arize_config().protocol
|
||||
|
||||
|
||||
def _weave_endpoint(params: StandardCallbackDynamicParams) -> str | None:
|
||||
from litellm.integrations.weave.weave_otel import WEAVE_BASE_URL, WEAVE_OTEL_ENDPOINT
|
||||
from litellm.integrations.weave.weave_otel import weave_otel_endpoint
|
||||
|
||||
return WEAVE_BASE_URL + WEAVE_OTEL_ENDPOINT
|
||||
return weave_otel_endpoint(os.environ.get("WANDB_HOST"))
|
||||
|
||||
|
||||
def _newrelic_endpoint(params: StandardCallbackDynamicParams) -> str | None:
|
||||
|
|
@ -78,6 +94,19 @@ _PROTOCOL_BY_CALLBACK: Final[Mapping[str, Callable[[StandardCallbackDynamicParam
|
|||
}
|
||||
)
|
||||
|
||||
#: Headers a destination must carry to authenticate. Several dynamic-header builders
|
||||
#: gate each credential independently, so a half-configured backend yields a non-empty
|
||||
#: but unusable header set; accepting it would suppress the operator's own exporter and
|
||||
#: send the request's whole trace where it cannot be stored.
|
||||
_REQUIRED_HEADERS_BY_CALLBACK: Final[Mapping[str, frozenset[str]]] = MappingProxyType(
|
||||
{
|
||||
"langfuse_otel": frozenset({"Authorization"}),
|
||||
"arize": frozenset({"arize-space-id", "api_key"}),
|
||||
"weave_otel": frozenset({"Authorization", "project_id"}),
|
||||
"newrelic": frozenset({"api-key"}),
|
||||
}
|
||||
)
|
||||
|
||||
_NO_ATTRS: Final[Mapping[str, str]] = MappingProxyType({})
|
||||
|
||||
|
||||
|
|
@ -92,9 +121,7 @@ def destination_for(callback_name: str, params: StandardCallbackDynamicParams) -
|
|||
"""The destination ``params`` names for ``callback_name``, or ``None``.
|
||||
|
||||
``None`` means the caller configured nothing usable for this backend, so the
|
||||
request keeps the operator's global exporters. A partial config (a host with
|
||||
no key pair) resolves to ``None`` rather than to the operator's endpoint with
|
||||
the tenant's host, which would post the operator's credentials elsewhere.
|
||||
request keeps the operator's global exporters.
|
||||
"""
|
||||
from litellm.integrations.otel.presets import DYNAMIC_HEADERS_BY_CALLBACK
|
||||
|
||||
|
|
@ -103,7 +130,7 @@ def destination_for(callback_name: str, params: StandardCallbackDynamicParams) -
|
|||
if header_builder is None or endpoint_builder is None:
|
||||
return None
|
||||
headers: Final = header_builder(params)
|
||||
if not headers:
|
||||
if not _REQUIRED_HEADERS_BY_CALLBACK.get(callback_name, frozenset()) <= frozenset(headers):
|
||||
return None
|
||||
endpoint: Final = endpoint_builder(params)
|
||||
if not endpoint:
|
||||
|
|
@ -111,7 +138,7 @@ def destination_for(callback_name: str, params: StandardCallbackDynamicParams) -
|
|||
protocol_builder: Final = _PROTOCOL_BY_CALLBACK.get(callback_name)
|
||||
return OtelDestination(
|
||||
endpoint=endpoint,
|
||||
headers=MappingProxyType(dict(headers)),
|
||||
headers=MappingProxyType(dict(headers)), # mutable-ok: MappingProxyType needs a concrete mapping to wrap
|
||||
resource_attributes=_NO_ATTRS,
|
||||
callback_name=callback_name,
|
||||
protocol=protocol_builder(params) if protocol_builder is not None else None,
|
||||
|
|
|
|||
|
|
@ -5,9 +5,6 @@ from typing import Final
|
|||
|
||||
from litellm.integrations.otel.model.config import ExporterOwner, ExporterSpec
|
||||
|
||||
#: What ``OpenTelemetryV2Config._normalize`` folds in when no destination is configured.
|
||||
_DEFAULT_SHORTHAND_EXPORTER: Final = ExporterSpec()
|
||||
|
||||
|
||||
def ensure_mappers(mapper_names: Iterable[str], *names: str) -> list[str]:
|
||||
"""Return ``mapper_names`` with each of ``names`` appended if not already present.
|
||||
|
|
@ -35,6 +32,17 @@ def credential_gated_exporters(
|
|||
override filter still recognises which backend this provider speaks for.
|
||||
"""
|
||||
return (
|
||||
*(spec for spec in exporters if spec != _DEFAULT_SHORTHAND_EXPORTER),
|
||||
*(spec for spec in exporters if not _prints_to_stdout(spec)),
|
||||
ExporterSpec(owner=owner, requires_headers=True),
|
||||
)
|
||||
|
||||
|
||||
def _prints_to_stdout(spec: "ExporterSpec") -> bool:
|
||||
"""Whether ``spec`` is the placeholder ``_normalize`` folds in for an empty list.
|
||||
|
||||
Identified by what it does rather than by equality with a default instance:
|
||||
``OpenTelemetryV2Config`` reads the standard ``OTEL_EXPORTER_OTLP_*`` env vars, so
|
||||
the shorthand it synthesizes is a real operator destination whenever any of them
|
||||
is set, and only a console exporter with no endpoint prints every span.
|
||||
"""
|
||||
return spec.kind == "console" and spec.endpoint is None
|
||||
|
|
|
|||
|
|
@ -117,6 +117,14 @@ def _get_weave_authorization_header(api_key: str) -> str:
|
|||
return f"Basic {auth_header}"
|
||||
|
||||
|
||||
def weave_otel_endpoint(host: str | None) -> str:
|
||||
"""The OTLP traces endpoint for a self-managed ``host``, else Weave cloud."""
|
||||
if not host:
|
||||
return WEAVE_BASE_URL + WEAVE_OTEL_ENDPOINT
|
||||
normalized: Final = host if host.startswith("http") else f"https://{host}"
|
||||
return normalized.rstrip("/") + WEAVE_OTEL_ENDPOINT
|
||||
|
||||
|
||||
def get_weave_otel_config() -> WeaveOtelConfig:
|
||||
"""
|
||||
Retrieves the Weave OpenTelemetry configuration based on environment variables.
|
||||
|
|
@ -134,7 +142,6 @@ def get_weave_otel_config() -> WeaveOtelConfig:
|
|||
"""
|
||||
api_key: Final = os.getenv("WANDB_API_KEY")
|
||||
project_id: Final = os.getenv("WANDB_PROJECT_ID")
|
||||
host = os.getenv("WANDB_HOST")
|
||||
|
||||
if not api_key:
|
||||
raise ValueError("WANDB_API_KEY must be set for Weave OpenTelemetry integration.")
|
||||
|
|
@ -144,15 +151,8 @@ def get_weave_otel_config() -> WeaveOtelConfig:
|
|||
"WANDB_PROJECT_ID must be set for Weave OpenTelemetry integration. Format: <entity>/<project_name>"
|
||||
)
|
||||
|
||||
if host:
|
||||
if not host.startswith("http"):
|
||||
host = "https://" + host
|
||||
# Self-managed instances use a different path
|
||||
endpoint = host.rstrip("/") + WEAVE_OTEL_ENDPOINT
|
||||
verbose_logger.debug("Using Weave OTEL endpoint from host: %s", endpoint)
|
||||
else:
|
||||
endpoint = WEAVE_BASE_URL + WEAVE_OTEL_ENDPOINT
|
||||
verbose_logger.debug("Using Weave cloud endpoint: %s", endpoint)
|
||||
endpoint: Final = weave_otel_endpoint(os.getenv("WANDB_HOST"))
|
||||
verbose_logger.debug("Using Weave OTEL endpoint: %s", endpoint)
|
||||
|
||||
# Weave uses Basic auth with format: api:<WANDB_API_KEY>
|
||||
auth_header: Final = _get_weave_authorization_header(api_key=api_key)
|
||||
|
|
|
|||
|
|
@ -20,6 +20,7 @@ Admins can opt out via two ``litellm`` globals (wired from proxy config):
|
|||
"""
|
||||
|
||||
import socket
|
||||
from functools import lru_cache
|
||||
from ipaddress import ip_address, ip_network
|
||||
from typing import Any, Final, Protocol
|
||||
from urllib.parse import quote, urlparse, urlunparse
|
||||
|
|
@ -363,6 +364,54 @@ def validate_url(url: str) -> tuple[str, str]:
|
|||
return rewritten, host_header
|
||||
|
||||
|
||||
def assert_public_url(url: str) -> None:
|
||||
"""Raise ``SSRFError`` unless ``url``'s host resolves only to public addresses.
|
||||
|
||||
The validation half of :func:`validate_url`, for callers that must keep the
|
||||
original hostname on the wire (TLS SNI, vendor-side Host routing) and so cannot
|
||||
use its IP-rewriting form. It honours the same ``litellm.user_url_validation``
|
||||
master switch and ``litellm.user_url_allowed_hosts`` allowlist. Because the
|
||||
caller still connects by name, this rejects a host that resolves somewhere
|
||||
private; it does not close a DNS rebind between the check and the connection.
|
||||
"""
|
||||
if not getattr(litellm, "user_url_validation", True):
|
||||
return
|
||||
rejection: Final = _public_host_rejection(url, tuple(getattr(litellm, "user_url_allowed_hosts", None) or ()))
|
||||
if rejection is not None:
|
||||
raise SSRFError(rejection)
|
||||
|
||||
|
||||
@lru_cache(maxsize=512)
|
||||
def _public_host_rejection(url: str, allowed_hosts: tuple[str, ...]) -> str | None:
|
||||
"""Why ``url`` is not safe to reach, or ``None``.
|
||||
|
||||
A verdict rather than an exception so both outcomes are cached: callers check
|
||||
the same handful of destinations on every request and ``getaddrinfo`` blocks.
|
||||
``allowed_hosts`` is part of the key so a config reload takes effect.
|
||||
"""
|
||||
parsed: Final = urlparse(url)
|
||||
if parsed.scheme not in _ALLOWED_SCHEMES:
|
||||
return f"URL scheme '{parsed.scheme}' is not allowed"
|
||||
|
||||
hostname: Final = parsed.hostname
|
||||
if not hostname:
|
||||
return "URL has no hostname"
|
||||
|
||||
effective_port: Final = parsed.port if parsed.port is not None else _default_port_for_scheme(parsed.scheme)
|
||||
if _is_host_allowlisted(hostname, effective_port):
|
||||
return None
|
||||
|
||||
try:
|
||||
addrinfo: Final = socket.getaddrinfo(hostname, effective_port, proto=socket.IPPROTO_TCP)
|
||||
except socket.gaierror as e:
|
||||
return f"DNS resolution failed for '{hostname}': {e}"
|
||||
|
||||
blocked: Final = tuple(
|
||||
address for address in (_sockaddr_host(info[4]) for info in addrinfo) if _is_blocked_ip(address)
|
||||
)
|
||||
return f"'{hostname}' resolves to a non-public address ({blocked[0]})" if blocked else None
|
||||
|
||||
|
||||
def assert_same_origin(candidate_url: str, expected_url: str) -> None:
|
||||
"""Verify ``candidate_url`` shares scheme, host, and port with ``expected_url``.
|
||||
|
||||
|
|
|
|||
|
|
@ -997,6 +997,12 @@ def resolve_tenant_otel_destinations(
|
|||
backend to two accounts. Returns empty when OTEL V2 is off, when neither level
|
||||
named a destination-capable backend, or when the config is incomplete, and the
|
||||
request then keeps the operator's own exporters.
|
||||
|
||||
A ``failure``-only entry is skipped: a destination is resolved during auth, before
|
||||
the request has an outcome, so honouring the filter would mean holding every span
|
||||
back until the call finishes. Those entries keep today's behaviour instead, where
|
||||
the tenant's credentials reach the backend through per-request tracer routing and
|
||||
the operator's exporter is left alone.
|
||||
"""
|
||||
from litellm.integrations.otel.model.config import is_otel_v2_enabled
|
||||
from litellm.integrations.otel.presets.destinations import destination_for
|
||||
|
|
@ -1012,6 +1018,7 @@ def resolve_tenant_otel_destinations(
|
|||
destination
|
||||
for item in entries
|
||||
if (callback := _get_validated_callback_metadata(item=item, source="otel-destination")) is not None
|
||||
if callback.callback_type != "failure"
|
||||
if (destination := destination_for(callback.callback_name, _tenant_otel_params(callback.callback_vars)))
|
||||
is not None
|
||||
)
|
||||
|
|
|
|||
|
|
@ -3,6 +3,7 @@
|
|||
import contextvars
|
||||
from collections.abc import Mapping
|
||||
|
||||
import litellm
|
||||
import pytest
|
||||
from opentelemetry.sdk.trace import TracerProvider
|
||||
from opentelemetry.sdk.trace.export import SimpleSpanProcessor
|
||||
|
|
@ -41,6 +42,20 @@ LANGFUSE_DEST = OtelDestination(
|
|||
)
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def allow_test_hosts(monkeypatch):
|
||||
"""Hosts named by these fixtures do not resolve, and a tenant-supplied host now
|
||||
goes through the SSRF guard. Allowlist them so the resolution tests stay about
|
||||
resolution; ``TestTenantHostSsrfGuard`` covers the guard itself."""
|
||||
from litellm.litellm_core_utils.url_utils import _public_host_rejection
|
||||
|
||||
monkeypatch.setattr(litellm, "user_url_validation", True, raising=False)
|
||||
monkeypatch.setattr(litellm, "user_url_allowed_hosts", ["team.local", "key.local", "x"], raising=False)
|
||||
_public_host_rejection.cache_clear()
|
||||
yield
|
||||
_public_host_rejection.cache_clear()
|
||||
|
||||
|
||||
def in_fresh_context(fn, *args):
|
||||
"""Run ``fn`` in its own context so one test's destinations never leak."""
|
||||
return contextvars.copy_context().run(fn, *args)
|
||||
|
|
@ -231,6 +246,7 @@ class TestRouting:
|
|||
assert route.provider is None
|
||||
|
||||
|
||||
@pytest.mark.usefixtures("allow_test_hosts")
|
||||
class TestDestinationResolution:
|
||||
def test_a_langfuse_key_pair_and_host_become_a_destination(self, monkeypatch):
|
||||
monkeypatch.setenv("LITELLM_OTEL_V2", "true")
|
||||
|
|
@ -312,12 +328,30 @@ class TestDestinationResolution:
|
|||
assert parse_headers(destination.header_string())["authorization"] == destination.headers["Authorization"]
|
||||
|
||||
|
||||
#: Anything that makes ``OpenTelemetryV2Config`` synthesize a real operator destination.
|
||||
_OTEL_SHORTHAND_ENV = (
|
||||
"OTEL_ENDPOINT",
|
||||
"OTEL_HEADERS",
|
||||
"OTEL_EXPORTER",
|
||||
"OTEL_EXPORTER_OTLP_ENDPOINT",
|
||||
"OTEL_EXPORTER_OTLP_HEADERS",
|
||||
"OTEL_EXPORTER_OTLP_PROTOCOL",
|
||||
)
|
||||
|
||||
|
||||
def credential_less_proxy(monkeypatch) -> None:
|
||||
"""An operator with no Langfuse account and no generic OTLP collector."""
|
||||
for name in ("LANGFUSE_PUBLIC_KEY", "LANGFUSE_SECRET_KEY", *_OTEL_SHORTHAND_ENV):
|
||||
monkeypatch.delenv(name, raising=False)
|
||||
with pytest.raises(ValueError, match="LANGFUSE_PUBLIC_KEY"):
|
||||
langfuse_preset()
|
||||
|
||||
|
||||
class TestPresetDegradation:
|
||||
def test_a_credential_less_langfuse_exports_nowhere_instead_of_to_the_console(self, monkeypatch, capsys):
|
||||
"""``_normalize`` folds a console exporter in for an empty list, which would
|
||||
print every span on a proxy whose teams bring their own credentials."""
|
||||
monkeypatch.delenv("LANGFUSE_PUBLIC_KEY", raising=False)
|
||||
monkeypatch.delenv("LANGFUSE_SECRET_KEY", raising=False)
|
||||
credential_less_proxy(monkeypatch)
|
||||
|
||||
config = langfuse_preset(allow_missing_credentials=True)
|
||||
provider = build_tracer_provider(config, tenant_overrides=True)
|
||||
|
|
@ -338,9 +372,8 @@ class TestPresetDegradation:
|
|||
def test_a_credential_less_proxy_still_builds_the_v2_logger(self, monkeypatch):
|
||||
from litellm.litellm_core_utils.litellm_logging import _maybe_construct_otel_v2
|
||||
|
||||
credential_less_proxy(monkeypatch)
|
||||
monkeypatch.setenv("LITELLM_OTEL_V2", "true")
|
||||
monkeypatch.delenv("LANGFUSE_PUBLIC_KEY", raising=False)
|
||||
monkeypatch.delenv("LANGFUSE_SECRET_KEY", raising=False)
|
||||
|
||||
is_otel_v2_enabled.cache_clear()
|
||||
logger = _maybe_construct_otel_v2("langfuse_otel", [])
|
||||
|
|
@ -358,3 +391,213 @@ class TestContextIsolation:
|
|||
|
||||
assert in_fresh_context(first) == frozenset({"langfuse_otel"})
|
||||
assert in_fresh_context(request_destinations) == ()
|
||||
|
||||
|
||||
class TestOperatorShorthandSurvivesDegradation:
|
||||
def test_a_generic_otlp_collector_keeps_receiving_when_langfuse_has_no_credentials(self, monkeypatch):
|
||||
"""Only the stdout placeholder is dropped. An operator who set the standard
|
||||
OTLP env vars configured a real destination and must keep it."""
|
||||
monkeypatch.delenv("LANGFUSE_PUBLIC_KEY", raising=False)
|
||||
monkeypatch.delenv("LANGFUSE_SECRET_KEY", raising=False)
|
||||
monkeypatch.setenv("OTEL_EXPORTER_OTLP_ENDPOINT", "http://collector.local:4318")
|
||||
|
||||
config = langfuse_preset(allow_missing_credentials=True)
|
||||
|
||||
assert [spec.endpoint for spec in config.exporters] == ["http://collector.local:4318", None]
|
||||
assert [spec.kind for spec in config.exporters] == ["otlp_http", "console"]
|
||||
|
||||
def test_the_stdout_placeholder_is_still_dropped_when_it_is_the_only_exporter(self, monkeypatch):
|
||||
credential_less_proxy(monkeypatch)
|
||||
|
||||
config = langfuse_preset(allow_missing_credentials=True)
|
||||
|
||||
assert all(spec.requires_headers and not spec.headers for spec in config.exporters)
|
||||
|
||||
|
||||
class TestBackendEndpointParity:
|
||||
def test_arize_follows_its_own_http_endpoint_instead_of_the_grpc_default(self, monkeypatch):
|
||||
monkeypatch.delenv("ARIZE_ENDPOINT", raising=False)
|
||||
monkeypatch.setenv("ARIZE_HTTP_ENDPOINT", "https://otlp.arize.com/v1/traces")
|
||||
|
||||
destination = destination_for("arize", {"arize_space_id": "s", "arize_api_key": "k"})
|
||||
|
||||
assert destination.endpoint == "https://otlp.arize.com/v1/traces"
|
||||
assert destination.protocol == "otlp_http"
|
||||
|
||||
def test_arize_uses_grpc_when_nothing_is_configured(self, monkeypatch):
|
||||
monkeypatch.delenv("ARIZE_ENDPOINT", raising=False)
|
||||
monkeypatch.delenv("ARIZE_HTTP_ENDPOINT", raising=False)
|
||||
|
||||
destination = destination_for("arize", {"arize_space_id": "s", "arize_api_key": "k"})
|
||||
|
||||
assert destination.endpoint == "https://otlp.arize.com/v1"
|
||||
assert destination.protocol == "otlp_grpc"
|
||||
|
||||
def test_weave_follows_a_self_hosted_wandb_host(self, monkeypatch):
|
||||
monkeypatch.setenv("WANDB_HOST", "weave.internal.example")
|
||||
|
||||
destination = destination_for("weave_otel", {"wandb_api_key": "k", "weave_project_id": "e/p"})
|
||||
|
||||
assert destination.endpoint == "https://weave.internal.example/otel/v1/traces"
|
||||
|
||||
def test_weave_uses_the_cloud_endpoint_without_a_host(self, monkeypatch):
|
||||
monkeypatch.delenv("WANDB_HOST", raising=False)
|
||||
|
||||
destination = destination_for("weave_otel", {"wandb_api_key": "k", "weave_project_id": "e/p"})
|
||||
|
||||
assert destination.endpoint == "https://trace.wandb.ai/otel/v1/traces"
|
||||
|
||||
|
||||
class TestIncompleteCredentials:
|
||||
"""Half a credential set builds a non-empty but unusable header dict. Accepting it
|
||||
would suppress the operator's exporter and send the trace where it cannot land."""
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
"callback_name,callback_vars",
|
||||
[
|
||||
("arize", {"arize_api_key": "k"}),
|
||||
("arize", {"arize_space_id": "s"}),
|
||||
("weave_otel", {"wandb_api_key": "k"}),
|
||||
("weave_otel", {"weave_project_id": "e/p"}),
|
||||
("langfuse_otel", {"langfuse_public_key": "pk"}),
|
||||
],
|
||||
)
|
||||
def test_a_partial_credential_set_resolves_to_nothing(self, callback_name, callback_vars):
|
||||
assert destination_for(callback_name, callback_vars) is None
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
"callback_name,callback_vars",
|
||||
[
|
||||
("arize", {"arize_space_id": "s", "arize_api_key": "k"}),
|
||||
("weave_otel", {"wandb_api_key": "k", "weave_project_id": "e/p"}),
|
||||
("newrelic", {"newrelic_api_key": "k"}),
|
||||
],
|
||||
)
|
||||
def test_a_complete_credential_set_resolves(self, callback_name, callback_vars):
|
||||
assert destination_for(callback_name, callback_vars) is not None
|
||||
|
||||
|
||||
@pytest.mark.usefixtures("allow_test_hosts")
|
||||
class TestCallbackTypeFilter:
|
||||
@staticmethod
|
||||
def _auth(callback_type: str | None) -> UserAPIKeyAuth:
|
||||
return UserAPIKeyAuth(
|
||||
team_metadata={
|
||||
"logging": [
|
||||
{
|
||||
"callback_name": "langfuse_otel",
|
||||
"callback_type": callback_type,
|
||||
"callback_vars": {
|
||||
"langfuse_public_key": "pk",
|
||||
"langfuse_secret_key": "sk",
|
||||
"langfuse_host": "http://team.local",
|
||||
},
|
||||
}
|
||||
]
|
||||
}
|
||||
)
|
||||
|
||||
@pytest.mark.parametrize("callback_type", ["success", "success_and_failure", None])
|
||||
def test_an_entry_that_wants_success_traces_gets_the_whole_trace(self, monkeypatch, callback_type):
|
||||
monkeypatch.setenv("LITELLM_OTEL_V2", "true")
|
||||
is_otel_v2_enabled.cache_clear()
|
||||
|
||||
assert resolve_tenant_otel_destinations(self._auth(callback_type)) != ()
|
||||
|
||||
def test_a_failure_only_entry_does_not_take_over_the_trace(self, monkeypatch):
|
||||
monkeypatch.setenv("LITELLM_OTEL_V2", "true")
|
||||
is_otel_v2_enabled.cache_clear()
|
||||
|
||||
assert resolve_tenant_otel_destinations(self._auth("failure")) == ()
|
||||
|
||||
|
||||
class TestEvictionSafety:
|
||||
def test_evicting_a_processor_does_not_shut_it_down(self):
|
||||
"""``on_end`` hands the caller a processor and then releases the lock, so a
|
||||
concurrent eviction that shut it down would silently drop that span."""
|
||||
from litellm.integrations.otel.plumbing.providers import _MAX_CACHED_DESTINATION_PROCESSORS
|
||||
|
||||
class Recording(SimpleSpanProcessor):
|
||||
def __init__(self):
|
||||
super().__init__(InMemorySpanExporter())
|
||||
self.shutdown_calls = 0
|
||||
|
||||
def shutdown(self):
|
||||
self.shutdown_calls += 1
|
||||
|
||||
built = []
|
||||
|
||||
def factory(_destination):
|
||||
processor = Recording()
|
||||
built.append(processor)
|
||||
return processor
|
||||
|
||||
fan_out = TenantFanOutSpanProcessor(processor_factory=factory)
|
||||
for index in range(_MAX_CACHED_DESTINATION_PROCESSORS + 1):
|
||||
fan_out._processor_for(LANGFUSE_DEST.model_copy(update={"endpoint": f"http://d{index}/otel"}))
|
||||
|
||||
assert len(built) == _MAX_CACHED_DESTINATION_PROCESSORS + 1
|
||||
assert built[0].shutdown_calls == 0
|
||||
assert len(fan_out._processors) == _MAX_CACHED_DESTINATION_PROCESSORS
|
||||
|
||||
|
||||
class TestTenantHostSsrfGuard:
|
||||
"""Anyone who can mint a key can write ``langfuse_host``, so a tenant-named host
|
||||
is a user-supplied URL and goes through the proxy's SSRF guard."""
|
||||
|
||||
@staticmethod
|
||||
def _reset() -> None:
|
||||
from litellm.litellm_core_utils.url_utils import _public_host_rejection
|
||||
|
||||
_public_host_rejection.cache_clear()
|
||||
|
||||
@pytest.mark.parametrize("host", ["http://127.0.0.1:9111", "http://169.254.169.254", "http://10.0.0.5:3000"])
|
||||
def test_a_tenant_host_on_a_private_address_resolves_to_nothing(self, monkeypatch, host):
|
||||
monkeypatch.setattr(litellm, "user_url_allowed_hosts", [], raising=False)
|
||||
monkeypatch.setattr(litellm, "user_url_validation", True, raising=False)
|
||||
self._reset()
|
||||
|
||||
assert (
|
||||
destination_for(
|
||||
"langfuse_otel",
|
||||
{"langfuse_public_key": "pk", "langfuse_secret_key": "sk", "langfuse_host": host},
|
||||
)
|
||||
is None
|
||||
)
|
||||
|
||||
def test_the_operator_can_allowlist_its_teams_internal_langfuse(self, monkeypatch):
|
||||
monkeypatch.setattr(litellm, "user_url_allowed_hosts", ["127.0.0.1:9111"], raising=False)
|
||||
monkeypatch.setattr(litellm, "user_url_validation", True, raising=False)
|
||||
self._reset()
|
||||
|
||||
destination = destination_for(
|
||||
"langfuse_otel",
|
||||
{"langfuse_public_key": "pk", "langfuse_secret_key": "sk", "langfuse_host": "http://127.0.0.1:9111"},
|
||||
)
|
||||
|
||||
assert destination.endpoint == "http://127.0.0.1:9111/api/public/otel"
|
||||
|
||||
def test_the_master_switch_still_turns_the_guard_off(self, monkeypatch):
|
||||
monkeypatch.setattr(litellm, "user_url_allowed_hosts", [], raising=False)
|
||||
monkeypatch.setattr(litellm, "user_url_validation", False, raising=False)
|
||||
self._reset()
|
||||
|
||||
assert (
|
||||
destination_for(
|
||||
"langfuse_otel",
|
||||
{"langfuse_public_key": "pk", "langfuse_secret_key": "sk", "langfuse_host": "http://127.0.0.1:9111"},
|
||||
)
|
||||
is not None
|
||||
)
|
||||
|
||||
def test_the_operators_own_internal_host_is_never_blocked(self, monkeypatch):
|
||||
"""The operator configures ``LANGFUSE_HOST`` themselves, so an internal
|
||||
collector there is a deployment choice rather than caller-supplied input."""
|
||||
monkeypatch.setattr(litellm, "user_url_allowed_hosts", [], raising=False)
|
||||
monkeypatch.setattr(litellm, "user_url_validation", True, raising=False)
|
||||
monkeypatch.setenv("LANGFUSE_HOST", "http://127.0.0.1:9111")
|
||||
self._reset()
|
||||
|
||||
destination = destination_for("langfuse_otel", {"langfuse_public_key": "pk", "langfuse_secret_key": "sk"})
|
||||
|
||||
assert destination.endpoint == "http://127.0.0.1:9111/api/public/otel"
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue