mirror of
https://github.com/BerriAI/litellm.git
synced 2026-10-10 03:28:53 +00:00
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.
This commit is contained in:
parent
1f6b80e659
commit
f5fb73f716
5 changed files with 253 additions and 173 deletions
|
|
@ -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."""
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
)
|
||||
|
|
|
|||
|
|
@ -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
|
||||
)
|
||||
|
|
|
|||
|
|
@ -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``.
|
||||
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue