From f5fb73f7164061d7899f6f1239faaf6663235450 Mon Sep 17 00:00:00 2001 From: Yucheng He Date: Thu, 3 Sep 2026 16:37:45 -0700 Subject: [PATCH] refactor(otel v2): reuse the proxy's own destination allowlist for tenant hosts A tenant-supplied Langfuse host is the same threat as a URL-valued `model`, so it now goes through `is_url_destination_allowed_by_host` against `provider_url_destination_allowed_hosts` instead of a second, DNS-based check of its own. The DNS lookup would have blocked the asyncio auth path on a hostname the caller picked, and its cached verdicts could blackhole a real host after one resolver blip. Evicting a destination processor now retires it to drain rather than shutting it down, since `on_end` hands a processor back and exports outside the lock. The retirees are capped so they cannot accumulate a thread each. `credential_gated_exporters` tells the synthesized stdout placeholder from a real exporter by transport rather than by the literal kind `console`, so an unrecognized kind is not mistaken for a configured collector, and an exporter the operator did configure survives. That also stops a weave test's env writes from making this look like a real OTLP exporter later in the same CI worker. --- .../integrations/otel/plumbing/providers.py | 28 ++- .../integrations/otel/presets/destinations.py | 107 ++++----- litellm/integrations/otel/presets/utils.py | 27 ++- litellm/litellm_core_utils/url_utils.py | 49 ---- .../otel/test_otel_v2_destinations.py | 215 +++++++++++++----- 5 files changed, 253 insertions(+), 173 deletions(-) diff --git a/litellm/integrations/otel/plumbing/providers.py b/litellm/integrations/otel/plumbing/providers.py index 232827e3d47..a9dc0798c07 100644 --- a/litellm/integrations/otel/plumbing/providers.py +++ b/litellm/integrations/otel/plumbing/providers.py @@ -207,6 +207,7 @@ def _processor_for(exporter: SpanExporter, use_simple: bool | None) -> SpanProce #: Distinct tenant destinations whose exporters stay alive. Each holds a connection #: pool and a batch thread, so the cache is bounded and evicts least-recently-used. _MAX_CACHED_DESTINATION_PROCESSORS: Final = 32 +_MAX_RETIRED_DESTINATION_PROCESSORS: Final = 8 class _ResourceWrappedReadableSpan(ReadableSpan): @@ -254,6 +255,7 @@ class TenantFanOutSpanProcessor(SpanProcessor): self._lock: Final = threading.Lock() self._build: Final = processor_factory if processor_factory is not None else _destination_processor self._processors: OrderedDict[object, SpanProcessor] = OrderedDict() # mutable-ok: bounded LRU + self._retired: OrderedDict[object, SpanProcessor] = OrderedDict() # mutable-ok: bounded drain list def on_start(self, span: SDKSpan, parent_context: Context | None = None) -> None: return None @@ -279,6 +281,7 @@ class TenantFanOutSpanProcessor(SpanProcessor): verbose_logger.debug("OTel V2 fan-out: processor shutdown failed: %s", exc) with self._lock: self._processors.clear() + self._retired.clear() def force_flush(self, timeout_millis: int = 30000) -> bool: results: Final = tuple(self._flush_one(processor, timeout_millis) for processor in self._snapshot()) @@ -286,7 +289,7 @@ class TenantFanOutSpanProcessor(SpanProcessor): def _snapshot(self) -> tuple[SpanProcessor, ...]: with self._lock: - return tuple(self._processors.values()) + return (*self._processors.values(), *self._retired.values()) @staticmethod def _flush_one(processor: SpanProcessor, timeout_millis: int) -> bool: @@ -312,13 +315,26 @@ class TenantFanOutSpanProcessor(SpanProcessor): _shutdown_quietly(built) return existing self._processors[key] = built - 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) + overflowed: Final = self._retired_on_overflow_locked() + if overflowed is not None: + _shutdown_quietly(overflowed) return built + def _retired_on_overflow_locked(self) -> SpanProcessor | None: + """Drop the LRU processor past the cap; return one only once it is safe to close. + + ``on_end`` hands a processor back and then exports outside the lock, so shutting + an evicted one down there loses that span. Evictions retire to drain instead, and + the retirees are capped so they cannot accumulate a thread each. + """ + if len(self._processors) <= _MAX_CACHED_DESTINATION_PROCESSORS: + return None + _, evicted = self._processors.popitem(last=False) + self._retired[id(evicted)] = evicted + if len(self._retired) <= _MAX_RETIRED_DESTINATION_PROCESSORS: + return None + return self._retired.popitem(last=False)[1] + def _destination_processor(destination: "OtelDestination") -> SpanProcessor | None: """A batching OTLP processor aimed at ``destination``, or ``None`` if unbuildable.""" diff --git a/litellm/integrations/otel/presets/destinations.py b/litellm/integrations/otel/presets/destinations.py index 2589f5dbf3e..815e8575a76 100644 --- a/litellm/integrations/otel/presets/destinations.py +++ b/litellm/integrations/otel/presets/destinations.py @@ -7,23 +7,39 @@ did; only the endpoint and transport need a per-backend rule. import os from collections.abc import Callable, Mapping +from functools import lru_cache from types import MappingProxyType from typing import Final +import litellm 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.litellm_core_utils.url_utils import is_url_destination_allowed_by_host from litellm.types.utils import StandardCallbackDynamicParams +#: An endpoint plus the OTLP transport to reach it with, or ``None`` when the backend +#: names no destination. The transport is ``None`` where the backend has only one. +_Destination = tuple[str, str | None] -def _langfuse_endpoint(params: StandardCallbackDynamicParams) -> str | None: + +@lru_cache(maxsize=128) +def _warn_host_not_allowlisted(host: str) -> None: + """Cached so one misconfigured team logs once rather than once per request.""" + verbose_logger.warning( + "OTel V2: not exporting to key/team Langfuse host '%s'. Add it to " + "litellm_settings.provider_url_destination_allowed_hosts to permit it", + host, + ) + + +def _langfuse_destination(params: StandardCallbackDynamicParams) -> "_Destination | None": """The tenant's own Langfuse host, else the operator's, else Langfuse US cloud. - 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. + A host the tenant named has to be allowlisted by the operator, the same way a + URL-valued ``model`` is: anyone who can mint a key can write it, and it becomes an + endpoint the proxy posts the request's whole trace to, carrying the tenant's own + credentials. The operator's own ``LANGFUSE_HOST`` is not checked, since an internal + collector there is a deployment choice. """ from litellm.integrations.langfuse.langfuse_otel import ( LANGFUSE_CLOUD_US_ENDPOINT, @@ -33,65 +49,50 @@ def _langfuse_endpoint(params: StandardCallbackDynamicParams) -> str | None: 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 + return (LANGFUSE_CLOUD_US_ENDPOINT, None) normalized: Final = host if host.startswith("http") else f"https://{host}" 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 (endpoint, None) + if not is_url_destination_allowed_by_host(endpoint, litellm.provider_url_destination_allowed_hosts): + _warn_host_not_allowlisted(host) return None - return endpoint + return (endpoint, None) -def _arize_endpoint(params: StandardCallbackDynamicParams) -> str | None: +def _arize_destination(params: StandardCallbackDynamicParams) -> "_Destination | None": from litellm.integrations.arize.arize import ArizeLogger - return ArizeLogger.get_arize_config().endpoint + config: Final = ArizeLogger.get_arize_config() + return (config.endpoint, config.protocol) -def _arize_protocol(params: StandardCallbackDynamicParams) -> str | None: - from litellm.integrations.arize.arize import ArizeLogger - - return ArizeLogger.get_arize_config().protocol - - -def _weave_endpoint(params: StandardCallbackDynamicParams) -> str | None: +def _weave_destination(params: StandardCallbackDynamicParams) -> "_Destination | None": from litellm.integrations.weave.weave_otel import weave_otel_endpoint - return weave_otel_endpoint(os.environ.get("WANDB_HOST")) + return (weave_otel_endpoint(os.environ.get("WANDB_HOST")), None) -def _newrelic_endpoint(params: StandardCallbackDynamicParams) -> str | None: +def _newrelic_destination(params: StandardCallbackDynamicParams) -> "_Destination | None": from litellm.integrations.otel.presets.newrelic import newrelic_dynamic_endpoint - return newrelic_dynamic_endpoint(params) + endpoint: Final = newrelic_dynamic_endpoint(params) + return (endpoint, None) if endpoint else None -#: Callback name -> endpoint resolver. A backend is destination-capable exactly +#: Callback name -> destination resolver. A backend is destination-capable exactly #: when it appears here AND in ``DYNAMIC_HEADERS_BY_CALLBACK``: without a header #: builder the destination would carry no tenant credentials, and the exporter #: would post the tenant's traffic to the operator's account. -_ENDPOINT_BY_CALLBACK: Final[Mapping[str, Callable[[StandardCallbackDynamicParams], str | None]]] = MappingProxyType( - { - "langfuse_otel": _langfuse_endpoint, - "arize": _arize_endpoint, - "weave_otel": _weave_endpoint, - "newrelic": _newrelic_endpoint, - } -) - -_PROTOCOL_BY_CALLBACK: Final[Mapping[str, Callable[[StandardCallbackDynamicParams], str | None]]] = MappingProxyType( - { - "arize": _arize_protocol, - } +_DESTINATION_BY_CALLBACK: Final[Mapping[str, Callable[[StandardCallbackDynamicParams], "_Destination | None"]]] = ( + MappingProxyType( + { + "langfuse_otel": _langfuse_destination, + "arize": _arize_destination, + "weave_otel": _weave_destination, + "newrelic": _newrelic_destination, + } + ) ) #: Headers a destination must carry to authenticate. Several dynamic-header builders @@ -114,7 +115,7 @@ def destination_capable_backends() -> frozenset[str]: """Backends a key or team can point at its own account.""" from litellm.integrations.otel.presets import DYNAMIC_HEADERS_BY_CALLBACK - return frozenset(_ENDPOINT_BY_CALLBACK) & frozenset(DYNAMIC_HEADERS_BY_CALLBACK) + return frozenset(_DESTINATION_BY_CALLBACK) & frozenset(DYNAMIC_HEADERS_BY_CALLBACK) def destination_for(callback_name: str, params: StandardCallbackDynamicParams) -> OtelDestination | None: @@ -126,20 +127,20 @@ def destination_for(callback_name: str, params: StandardCallbackDynamicParams) - from litellm.integrations.otel.presets import DYNAMIC_HEADERS_BY_CALLBACK header_builder: Final = DYNAMIC_HEADERS_BY_CALLBACK.get(callback_name) - endpoint_builder: Final = _ENDPOINT_BY_CALLBACK.get(callback_name) - if header_builder is None or endpoint_builder is None: + destination_builder: Final = _DESTINATION_BY_CALLBACK.get(callback_name) + if header_builder is None or destination_builder is None: return None headers: Final = header_builder(params) - if not _REQUIRED_HEADERS_BY_CALLBACK.get(callback_name, frozenset()) <= frozenset(headers): + if not headers or not _REQUIRED_HEADERS_BY_CALLBACK[callback_name] <= frozenset(headers): return None - endpoint: Final = endpoint_builder(params) - if not endpoint: + resolved: Final = destination_builder(params) + if resolved is None: return None - protocol_builder: Final = _PROTOCOL_BY_CALLBACK.get(callback_name) + endpoint, protocol = resolved return OtelDestination( endpoint=endpoint, 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, + protocol=protocol, ) diff --git a/litellm/integrations/otel/presets/utils.py b/litellm/integrations/otel/presets/utils.py index 002d5894e64..71ab2057fda 100644 --- a/litellm/integrations/otel/presets/utils.py +++ b/litellm/integrations/otel/presets/utils.py @@ -32,17 +32,30 @@ def credential_gated_exporters( override filter still recognises which backend this provider speaks for. """ return ( - *(spec for spec in exporters if not _prints_to_stdout(spec)), + *(spec for spec in exporters if not _is_stdout_placeholder(spec)), ExporterSpec(owner=owner, requires_headers=True), ) -def _prints_to_stdout(spec: "ExporterSpec") -> bool: +#: The fields ``OpenTelemetryV2Config._normalize`` fills the synthesized spec from. +_SHORTHAND_FIELDS: Final = frozenset({"kind", "endpoint", "headers"}) + + +def _is_stdout_placeholder(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. + Two conditions. It must have nowhere to send a span, which is what + ``exporter_transport`` answers: an unrecognized or misspelled kind falls back to the + console exporter, so comparing against the literal ``"console"`` would miss it. And + every non-shorthand field must still be at its default, which is what says the + operator did not ask for it: an exporter they configured survives, and so does the + gated spec this module appends, which would otherwise eat itself when one preset + layers onto another. """ - return spec.kind == "console" and spec.endpoint is None + from litellm.integrations.otel.plumbing.providers import exporter_transport + + return ( + exporter_transport(spec.kind) == "headerless" + and spec.endpoint is None + and spec.model_dump(exclude_defaults=True).keys() <= _SHORTHAND_FIELDS + ) diff --git a/litellm/litellm_core_utils/url_utils.py b/litellm/litellm_core_utils/url_utils.py index cee248b8384..1e43117933d 100644 --- a/litellm/litellm_core_utils/url_utils.py +++ b/litellm/litellm_core_utils/url_utils.py @@ -20,7 +20,6 @@ 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 @@ -364,54 +363,6 @@ 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/tests/test_litellm/integrations/otel/test_otel_v2_destinations.py b/tests/test_litellm/integrations/otel/test_otel_v2_destinations.py index c64ef05ce77..0e6c940fb67 100644 --- a/tests/test_litellm/integrations/otel/test_otel_v2_destinations.py +++ b/tests/test_litellm/integrations/otel/test_otel_v2_destinations.py @@ -44,16 +44,12 @@ 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() + """A tenant-supplied host must be allowlisted by the operator. Allowlist the ones + these fixtures name so the resolution tests stay about resolution; + ``TestTenantHostSsrfGuard`` covers the guard itself.""" + monkeypatch.setattr( + litellm, "provider_url_destination_allowed_hosts", ["team.local", "key.local", "x"], raising=False + ) def in_fresh_context(fn, *args): @@ -512,9 +508,44 @@ class TestCallbackTypeFilter: 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.""" + def test_an_evicted_processor_is_retired_rather_than_shut_down(self): + """``on_end`` hands a processor back and exports outside the lock, so shutting + an evicted one down there loses that span. Retirees are capped so they cannot + accumulate a thread each.""" + from litellm.integrations.otel.plumbing.providers import ( + _MAX_CACHED_DESTINATION_PROCESSORS, + _MAX_RETIRED_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): + built.append(Recording()) + return built[-1] + + fan_out = TenantFanOutSpanProcessor(processor_factory=factory) + for index in range(_MAX_CACHED_DESTINATION_PROCESSORS + _MAX_RETIRED_DESTINATION_PROCESSORS): + fan_out._processor_for(LANGFUSE_DEST.model_copy(update={"endpoint": f"http://d{index}/otel"})) + + assert [p.shutdown_calls for p in built] == [0] * len(built) + assert len(fan_out._processors) == _MAX_CACHED_DESTINATION_PROCESSORS + + for index in range(2): + fan_out._processor_for(LANGFUSE_DEST.model_copy(update={"endpoint": f"http://late{index}/otel"})) + + assert [p.shutdown_calls for p in built[:2]] == [1, 1] + assert built[2].shutdown_calls == 0 + assert len(fan_out._processors) == _MAX_CACHED_DESTINATION_PROCESSORS + + def test_a_retired_processor_is_still_flushed_and_closed_on_shutdown(self): from litellm.integrations.otel.plumbing.providers import _MAX_CACHED_DESTINATION_PROCESSORS class Recording(SimpleSpanProcessor): @@ -528,76 +559,144 @@ class TestEvictionSafety: built = [] def factory(_destination): - processor = Recording() - built.append(processor) - return processor + built.append(Recording()) + return built[-1] 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"})) + fan_out.shutdown() - assert len(built) == _MAX_CACHED_DESTINATION_PROCESSORS + 1 - assert built[0].shutdown_calls == 0 - assert len(fan_out._processors) == _MAX_CACHED_DESTINATION_PROCESSORS + assert built[0].shutdown_calls == 1 + + +class TestCredentialGatedExporters: + def test_layering_a_second_preset_does_not_eat_the_first_gated_exporter(self, monkeypatch): + """``base.Preset`` advertises ``config_overrides`` layering, and the gated spec + is itself a console exporter with no endpoint.""" + credential_less_proxy(monkeypatch) + from litellm.integrations.otel.presets.utils import credential_gated_exporters + + once = credential_gated_exporters((), ExporterOwner.LANGFUSE_OTEL) + twice = credential_gated_exporters(once, ExporterOwner.WEAVE_OTEL) + + assert [spec.owner for spec in twice] == [ExporterOwner.LANGFUSE_OTEL, ExporterOwner.WEAVE_OTEL] + + def test_an_exporter_the_operator_configured_survives(self): + from litellm.integrations.otel.presets.utils import credential_gated_exporters + + operator_console = ExporterSpec(kind="console", use_simple_processor=True) + + kept = credential_gated_exporters((operator_console,), ExporterOwner.LANGFUSE_OTEL) + + assert kept[0] == operator_console + + def test_an_otlp_exporter_on_its_default_endpoint_survives(self): + """``OTEL_EXPORTER=otlp_http`` with no endpoint is a real collector on the SDK's + default port, not the placeholder, so the transport is what tells them apart.""" + from litellm.integrations.otel.presets.utils import credential_gated_exporters + + operator_otlp = ExporterSpec(kind="otlp_http", endpoint=None, headers=None) + + kept = credential_gated_exporters((operator_otlp,), ExporterOwner.LANGFUSE_OTEL) + + assert kept[0] == operator_otlp + + def test_the_synthesized_stdout_placeholder_is_dropped(self): + from litellm.integrations.otel.presets.utils import credential_gated_exporters + + placeholder = ExporterSpec(kind="console", endpoint=None, headers=None) + + kept = credential_gated_exporters((placeholder,), ExporterOwner.LANGFUSE_OTEL) + + assert [spec.owner for spec in kept] == [ExporterOwner.LANGFUSE_OTEL] 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.""" + """Anyone who can mint a key can write ``langfuse_host``, so the host it names has + to be one the operator approved.""" + + @pytest.fixture(autouse=True) + def _guard_on(self, monkeypatch): + from litellm.integrations.otel.presets.destinations import _warn_host_not_allowlisted + + monkeypatch.setattr(litellm, "provider_url_destination_allowed_hosts", [], raising=False) + _warn_host_not_allowlisted.cache_clear() + yield + _warn_host_not_allowlisted.cache_clear() @staticmethod - def _reset() -> None: - from litellm.litellm_core_utils.url_utils import _public_host_rejection + def _langfuse(host: str) -> Mapping[str, str]: + return {"langfuse_public_key": "pk", "langfuse_secret_key": "sk", "langfuse_host": host} - _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", + "https://collector.example.com", + "https://langfuse.corp:99999", + "ftp://collector.example.com", + ], + ) + def test_a_host_the_operator_never_approved_resolves_to_nothing(self, host): + assert destination_for("langfuse_otel", self._langfuse(host)) is None - @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() + def test_userinfo_naming_an_allowlisted_host_does_not_smuggle_a_second_one(self, monkeypatch): + """``https://allowed@10.0.0.5`` reads as the allowlisted host to the eye and + posts to 10.0.0.5 on the wire.""" + monkeypatch.setattr(litellm, "provider_url_destination_allowed_hosts", ["collector.example.com"], raising=False) - assert ( - destination_for( - "langfuse_otel", - {"langfuse_public_key": "pk", "langfuse_secret_key": "sk", "langfuse_host": host}, - ) - is None + assert destination_for("langfuse_otel", self._langfuse("https://collector.example.com@10.0.0.5")) is None + + def test_a_malformed_host_does_not_take_the_other_backends_with_it(self, monkeypatch): + """``urlparse(...).port`` raises a bare ValueError, which would escape + ``destination_for`` and kill the whole resolution.""" + monkeypatch.setenv("LITELLM_OTEL_V2", "true") + monkeypatch.setenv("NEW_RELIC_OTEL_ENDPOINT", "https://otlp.nr-data.net") + monkeypatch.setattr(litellm, "provider_url_destination_allowed_hosts", ["collector.example.com"], raising=False) + is_otel_v2_enabled.cache_clear() + auth = UserAPIKeyAuth( + token="hashed", + team_metadata={ + "logging": [ + {"callback_name": "langfuse_otel", "callback_vars": self._langfuse("https://lf.corp:99999")}, + {"callback_name": "newrelic", "callback_vars": {"newrelic_api_key": "nr"}}, + ] + }, ) + assert [d.callback_name for d in resolve_tenant_otel_destinations(auth)] == ["newrelic"] + 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() + monkeypatch.setattr(litellm, "provider_url_destination_allowed_hosts", ["127.0.0.1:9111"], raising=False) - destination = destination_for( - "langfuse_otel", - {"langfuse_public_key": "pk", "langfuse_secret_key": "sk", "langfuse_host": "http://127.0.0.1:9111"}, - ) + destination = destination_for("langfuse_otel", self._langfuse("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" + + def test_an_allowlisted_host_is_taken_without_resolving_it(self, monkeypatch): + """The check runs on the asyncio auth path, so it must not block on a name the + caller chose. ``.invalid`` never resolves, and it is still accepted.""" + monkeypatch.setattr(litellm, "provider_url_destination_allowed_hosts", ["lf.invalid"], raising=False) + + destination = destination_for("langfuse_otel", self._langfuse("https://lf.invalid")) + + assert destination.endpoint == "https://lf.invalid/api/public/otel" + + def test_a_rejected_host_is_warned_about_once(self, caplog): + with caplog.at_level("WARNING", logger="LiteLLM"): + for _ in range(3): + destination_for("langfuse_otel", self._langfuse("http://10.0.0.5:3000")) + + assert sum("provider_url_destination_allowed_hosts" in record.message for record in caplog.records) == 1