From 1f6b80e659623cc636b7bcc23356c72bf396ce65 Mon Sep 17 00:00:00 2001 From: Yucheng He Date: Thu, 3 Sep 2026 15:50:06 -0700 Subject: [PATCH] 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. --- .../integrations/otel/model/destination.py | 16 +- .../integrations/otel/plumbing/providers.py | 27 +- .../integrations/otel/presets/destinations.py | 73 +++-- litellm/integrations/otel/presets/utils.py | 16 +- litellm/integrations/weave/weave_otel.py | 20 +- litellm/litellm_core_utils/url_utils.py | 49 ++++ litellm/proxy/litellm_pre_call_utils.py | 7 + .../otel/test_otel_v2_destinations.py | 251 +++++++++++++++++- 8 files changed, 390 insertions(+), 69 deletions(-) diff --git a/litellm/integrations/otel/model/destination.py b/litellm/integrations/otel/model/destination.py index b88a240c966..299253cac77 100644 --- a/litellm/integrations/otel/model/destination.py +++ b/litellm/integrations/otel/model/destination.py @@ -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())), diff --git a/litellm/integrations/otel/plumbing/providers.py b/litellm/integrations/otel/plumbing/providers.py index 00198234f54..232827e3d47 100644 --- a/litellm/integrations/otel/plumbing/providers.py +++ b/litellm/integrations/otel/plumbing/providers.py @@ -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: diff --git a/litellm/integrations/otel/presets/destinations.py b/litellm/integrations/otel/presets/destinations.py index 966d8f4e1f4..2589f5dbf3e 100644 --- a/litellm/integrations/otel/presets/destinations.py +++ b/litellm/integrations/otel/presets/destinations.py @@ -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, diff --git a/litellm/integrations/otel/presets/utils.py b/litellm/integrations/otel/presets/utils.py index ef6e80ff33e..002d5894e64 100644 --- a/litellm/integrations/otel/presets/utils.py +++ b/litellm/integrations/otel/presets/utils.py @@ -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 diff --git a/litellm/integrations/weave/weave_otel.py b/litellm/integrations/weave/weave_otel.py index f2cc64a9ba2..50289263f38 100644 --- a/litellm/integrations/weave/weave_otel.py +++ b/litellm/integrations/weave/weave_otel.py @@ -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: /" ) - 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: auth_header: Final = _get_weave_authorization_header(api_key=api_key) diff --git a/litellm/litellm_core_utils/url_utils.py b/litellm/litellm_core_utils/url_utils.py index 1e43117933d..cee248b8384 100644 --- a/litellm/litellm_core_utils/url_utils.py +++ b/litellm/litellm_core_utils/url_utils.py @@ -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``. diff --git a/litellm/proxy/litellm_pre_call_utils.py b/litellm/proxy/litellm_pre_call_utils.py index 6069ed9a8c6..dc78e636e37 100644 --- a/litellm/proxy/litellm_pre_call_utils.py +++ b/litellm/proxy/litellm_pre_call_utils.py @@ -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 ) diff --git a/tests/test_litellm/integrations/otel/test_otel_v2_destinations.py b/tests/test_litellm/integrations/otel/test_otel_v2_destinations.py index ee67160533f..c64ef05ce77 100644 --- a/tests/test_litellm/integrations/otel/test_otel_v2_destinations.py +++ b/tests/test_litellm/integrations/otel/test_otel_v2_destinations.py @@ -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"