mirror of
https://github.com/BerriAI/litellm.git
synced 2026-10-09 03:18:44 +00:00
feat(otel/v2): fan out proxy-internal spans to admin-resolved per-tenant destinations
LIT-3850's gen-AI fan-out only routed the LLM-call span; the FastAPI server span, auth phase, postgres lookups, and post-call cost ledger flowed only to the admin-configured backend, leaving tenant backends with an orphaned gen-AI child. Hoist resolves admin-owned destinations at the auth boundary, anchors them on a request-scoped ContextVar, and a new SpanProcessor on the main v2 TracerProvider forwards every finished span to each destination matching the v2 logger's backend. The TenantTracerCache clone providers built for the gen-AI span do not carry this processor, so each span exports exactly once per backend. The contextvar pins per-request isolation: concurrent requests with different destinations cannot leak spans to each other's backends, even on the same proxy process. Trust boundary unchanged: destinations are resolved server-side from litellm.credential_list against the caller's identity chain. A client body that smuggles otel_destinations cannot influence the resolution. Adds six regression tests pinning the contract: forwarding, backend filter, no-op when no destinations, per-(endpoint, headers) cache, per-request task isolation, and that clone providers do not carry the fan-out processor.
This commit is contained in:
parent
83a94f78c4
commit
a42efec324
7 changed files with 536 additions and 11 deletions
|
|
@ -110,7 +110,11 @@ class OpenTelemetryV2(CustomLogger):
|
|||
self.config: OpenTelemetryV2Config = config or OpenTelemetryV2Config(**kwargs)
|
||||
self.callback_name = callback_name
|
||||
self._tracer_provider: TracerProvider = (
|
||||
tracer_provider if tracer_provider is not None else build_tracer_provider(self.config)
|
||||
tracer_provider
|
||||
if tracer_provider is not None
|
||||
else build_tracer_provider(
|
||||
self.config, tenant_fan_out_owner=callback_name
|
||||
)
|
||||
)
|
||||
self.tracer: Tracer = get_tracer(self._tracer_provider, LITELLM_TRACER_NAME)
|
||||
self._metrics_recorder = self._init_metrics(meter_provider)
|
||||
|
|
|
|||
|
|
@ -27,9 +27,28 @@ _PROPAGATOR = TraceContextTextMapPropagator()
|
|||
# and is inherited by ``asyncio.create_task`` children — i.e. the async logging
|
||||
# callbacks that close the span. It is never reset: the contextvar dies with the
|
||||
# request task, so there is nothing to leak.
|
||||
_request_root_span: "ContextVar[Span | None]" = ContextVar(
|
||||
"litellm_otel_request_root_span", default=None
|
||||
)
|
||||
_request_root_span: "ContextVar[Span | None]" = ContextVar("litellm_otel_request_root_span", default=None)
|
||||
|
||||
# Per-request admin-resolved destinations. Set once at the auth boundary (the
|
||||
# earliest point a request's identity is known) and read by the global-provider
|
||||
# fan-out processor at ``on_end`` time, so every span the proxy emits for this
|
||||
# request -- the FastAPI server span, the ``auth`` phase, DB lookups, the
|
||||
# batch-write cost ledger -- ships to every per-tenant destination the admin
|
||||
# assigned. Lives on a ``ContextVar`` so it follows the request task across
|
||||
# ``asyncio.create_task`` children (the success/failure logging callbacks close
|
||||
# the LLM span in a worker copied from the request context). Request-scoped: the
|
||||
# contextvar dies with the request task, so nothing leaks across requests.
|
||||
_request_destinations: "ContextVar[tuple]" = ContextVar("litellm_otel_request_destinations", default=())
|
||||
|
||||
|
||||
def set_request_destinations(destinations: tuple) -> None:
|
||||
"""Anchor the admin-resolved destinations for this request."""
|
||||
_request_destinations.set(tuple(destinations))
|
||||
|
||||
|
||||
def request_destinations() -> tuple:
|
||||
"""Destinations the request fans out to, or empty when none were resolved."""
|
||||
return _request_destinations.get()
|
||||
|
||||
|
||||
def set_request_root_span(span: Span) -> None:
|
||||
|
|
@ -49,9 +68,7 @@ def request_root_span() -> "Span | None":
|
|||
return span if is_recordable_span(span) else None
|
||||
|
||||
|
||||
def set_request_baggage(
|
||||
values: Mapping[str, str], context: Context | None = None
|
||||
) -> Context:
|
||||
def set_request_baggage(values: Mapping[str, str], context: Context | None = None) -> Context:
|
||||
"""Return a context with ``values`` written into Baggage."""
|
||||
ctx = context
|
||||
for key, value in values.items():
|
||||
|
|
|
|||
143
litellm/integrations/otel/plumbing/fan_out.py
Normal file
143
litellm/integrations/otel/plumbing/fan_out.py
Normal file
|
|
@ -0,0 +1,143 @@
|
|||
"""Per-request span fan-out to admin-resolved destinations.
|
||||
|
||||
Attached to the main ``TracerProvider`` so every span on it (the FastAPI server
|
||||
span, the proxy's ``auth`` phase span, DB lookups, the post-call cost ledger)
|
||||
ships to every per-tenant destination the admin assigned for this request, on
|
||||
top of the provider's configured global exporters.
|
||||
|
||||
Spans emitted through the per-tenant ``TenantTracerCache`` clone providers (the
|
||||
gen-AI LLM-call span and its MCP-tool sibling) reach tenant backends through
|
||||
the clone's own exporters; the clone's provider has its own processor list and
|
||||
does NOT carry this fan-out processor, so a span is exported once per backend.
|
||||
|
||||
Identical ``(trace_id, span_id)`` deduplication on the OTLP receiver collapses
|
||||
the rare overlap into one span per backend (e.g. when the configured global
|
||||
exporter and a per-tenant destination point at the same vendor account).
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from collections import OrderedDict
|
||||
from typing import TYPE_CHECKING
|
||||
|
||||
from opentelemetry.context import Context
|
||||
from opentelemetry.sdk.trace import ReadableSpan, Span, SpanProcessor
|
||||
|
||||
from litellm._logging import verbose_logger
|
||||
from litellm.integrations.otel.model.config import ExporterSpec
|
||||
from litellm.integrations.otel.plumbing.context import request_destinations
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from litellm.integrations.otel.model.destination import OtelDestination
|
||||
|
||||
# Bound on cached per-destination processors. One processor per
|
||||
# ``(endpoint, sorted(headers))`` pair, so the working set is one entry per
|
||||
# admin-resolved tenant credential -- a real-world deployment with hundreds of
|
||||
# tenants stays well under this. The LRU shuts down the evicted processor's
|
||||
# exporter thread so the working set is reclaimed.
|
||||
_MAX_CACHED_PROCESSORS = 256
|
||||
|
||||
|
||||
def _processor_key(destination: "OtelDestination") -> tuple:
|
||||
return (destination.endpoint, tuple(sorted(destination.headers.items())))
|
||||
|
||||
|
||||
class TenantFanOutSpanProcessor(SpanProcessor):
|
||||
"""Forward each finished span to every admin-resolved per-tenant destination.
|
||||
|
||||
The destinations are looked up from a request-scoped contextvar set during
|
||||
auth, so the processor is stateless across requests and isolation across
|
||||
concurrent requests is guaranteed by Python's contextvars.
|
||||
"""
|
||||
|
||||
def __init__(self, owner_callback_name: str | None) -> None:
|
||||
self._owner = owner_callback_name
|
||||
# Built lazily so we avoid importing the providers module at class
|
||||
# definition time (which would create a circular import with routing).
|
||||
self._processors: "OrderedDict[tuple, SpanProcessor]" = OrderedDict()
|
||||
|
||||
def on_start(self, span: "Span", parent_context: Context | None = None) -> None:
|
||||
return None
|
||||
|
||||
def on_end(self, span: ReadableSpan) -> None:
|
||||
destinations = request_destinations()
|
||||
if not destinations:
|
||||
return
|
||||
for destination in destinations:
|
||||
if destination.callback_name != self._owner:
|
||||
continue
|
||||
processor = self._processor_for(destination)
|
||||
if processor is None:
|
||||
continue
|
||||
try:
|
||||
processor.on_end(span)
|
||||
except Exception as exc:
|
||||
verbose_logger.debug(
|
||||
"OTel V2 fan-out: forwarding span to %s failed: %s",
|
||||
destination.endpoint,
|
||||
exc,
|
||||
)
|
||||
|
||||
def shutdown(self) -> None:
|
||||
for processor in self._processors.values():
|
||||
try:
|
||||
processor.shutdown()
|
||||
except Exception as exc:
|
||||
verbose_logger.debug("OTel V2 fan-out: processor shutdown failed: %s", exc)
|
||||
self._processors.clear()
|
||||
|
||||
def force_flush(self, timeout_millis: int = 30000) -> bool:
|
||||
all_ok = True
|
||||
for processor in self._processors.values():
|
||||
try:
|
||||
if not processor.force_flush(timeout_millis):
|
||||
all_ok = False
|
||||
except Exception:
|
||||
all_ok = False
|
||||
return all_ok
|
||||
|
||||
def _processor_for(self, destination: "OtelDestination") -> "SpanProcessor | None":
|
||||
key = _processor_key(destination)
|
||||
cached = self._processors.get(key)
|
||||
if cached is not None:
|
||||
self._processors.move_to_end(key)
|
||||
return cached
|
||||
from litellm.integrations.otel.plumbing.providers import (
|
||||
_exporter_from_spec,
|
||||
_processor_for,
|
||||
)
|
||||
|
||||
try:
|
||||
spec = ExporterSpec(
|
||||
kind=_resolve_kind(destination),
|
||||
endpoint=destination.endpoint,
|
||||
headers=destination.header_string(),
|
||||
owner=None,
|
||||
)
|
||||
exporter = _exporter_from_spec(spec)
|
||||
processor = _processor_for(exporter, use_simple=False)
|
||||
except Exception as exc:
|
||||
verbose_logger.debug(
|
||||
"OTel V2 fan-out: failed to build processor for %s: %s",
|
||||
destination.endpoint,
|
||||
exc,
|
||||
)
|
||||
return None
|
||||
self._processors[key] = processor
|
||||
if len(self._processors) > _MAX_CACHED_PROCESSORS:
|
||||
_, evicted = self._processors.popitem(last=False)
|
||||
try:
|
||||
evicted.shutdown()
|
||||
except Exception:
|
||||
pass
|
||||
return processor
|
||||
|
||||
|
||||
# Per-backend transport: Arize speaks OTLP/gRPC, every other current preset
|
||||
# speaks OTLP/HTTP. Kept here so the fan-out picks the same transport the
|
||||
# preset's own exporter uses, mirroring ``TenantTracerCache._owned_otlp_kind``.
|
||||
_GRPC_BACKENDS = frozenset({"arize"})
|
||||
|
||||
|
||||
def _resolve_kind(destination: "OtelDestination") -> str:
|
||||
return "otlp_grpc" if destination.callback_name in _GRPC_BACKENDS else "otlp_http"
|
||||
|
|
@ -290,6 +290,7 @@ def build_tracer_provider(
|
|||
exporter: SpanExporter | None = None,
|
||||
baggage_processor: SpanProcessor | None = None,
|
||||
use_simple_processor: bool | None = None,
|
||||
tenant_fan_out_owner: str | None = None,
|
||||
) -> TracerProvider:
|
||||
"""Build the shared :class:`TracerProvider`.
|
||||
|
||||
|
|
@ -298,6 +299,13 @@ def build_tracer_provider(
|
|||
``config.exporters`` entry — this is what fans spans out to multiple
|
||||
backends. ``exporter`` and ``use_simple_processor`` are explicit overrides:
|
||||
pass a single exporter to attach exactly that one (used by tests).
|
||||
|
||||
``tenant_fan_out_owner`` — when set, attach a ``TenantFanOutSpanProcessor``
|
||||
that forwards each finished span to the admin-resolved per-tenant
|
||||
destinations whose ``callback_name`` matches the owner. Only the MAIN v2
|
||||
logger provider opts in; the per-tenant clone providers (built by
|
||||
``TenantTracerCache``) intentionally do not, so the LLM-call span exported
|
||||
through them is not also fanned out here.
|
||||
"""
|
||||
provider = TracerProvider(resource=build_resource(config))
|
||||
if baggage_processor is None:
|
||||
|
|
@ -306,6 +314,15 @@ def build_tracer_provider(
|
|||
)
|
||||
provider.add_span_processor(baggage_processor)
|
||||
|
||||
if tenant_fan_out_owner is not None:
|
||||
from litellm.integrations.otel.plumbing.fan_out import (
|
||||
TenantFanOutSpanProcessor,
|
||||
)
|
||||
|
||||
provider.add_span_processor(
|
||||
TenantFanOutSpanProcessor(owner_callback_name=tenant_fan_out_owner)
|
||||
)
|
||||
|
||||
if exporter is not None:
|
||||
provider.add_span_processor(_processor_for(exporter, use_simple_processor))
|
||||
return provider
|
||||
|
|
|
|||
|
|
@ -949,6 +949,51 @@ async def _resolve_jwt_to_virtual_key(
|
|||
return None
|
||||
|
||||
|
||||
async def _hoist_request_destinations(
|
||||
request: Request, user_api_key_dict: UserAPIKeyAuth
|
||||
) -> None:
|
||||
"""Resolve admin-owned OTEL destinations for this request and anchor them.
|
||||
|
||||
Runs after the auth builder, while we are still inside the request task, so
|
||||
the ``ContextVar`` is visible to every ``SpanProcessor.on_end`` that fires
|
||||
for spans this request opens. Stashes the same list on ``request.state`` so
|
||||
``_apply_admin_logging_exporters`` can reuse it without a second DB pass.
|
||||
|
||||
Best-effort: a resolver failure must not break the request. The contextvar
|
||||
is left at its default (empty tuple), so the fan-out processor no-ops.
|
||||
"""
|
||||
try:
|
||||
from litellm.integrations.otel.model.destination import OtelDestination
|
||||
from litellm.integrations.otel.plumbing.context import (
|
||||
set_request_destinations,
|
||||
)
|
||||
from litellm.proxy.litellm_pre_call_utils import (
|
||||
_resolve_logging_exporters,
|
||||
)
|
||||
|
||||
destinations_raw, _backends = await _resolve_logging_exporters(
|
||||
user_api_key_dict
|
||||
)
|
||||
destinations = tuple(
|
||||
OtelDestination(
|
||||
callback_name=item.get("callback_name"),
|
||||
endpoint=item.get("endpoint", ""),
|
||||
headers=item.get("headers") or {},
|
||||
)
|
||||
for item in destinations_raw
|
||||
if isinstance(item, dict) and item.get("endpoint")
|
||||
)
|
||||
set_request_destinations(destinations)
|
||||
try:
|
||||
request.state.otel_destinations = destinations_raw
|
||||
except Exception:
|
||||
pass
|
||||
except Exception as exc:
|
||||
verbose_proxy_logger.debug(
|
||||
"OTel V2: hoist destination resolution failed: %s", exc
|
||||
)
|
||||
|
||||
|
||||
def _ensure_parent_otel_span_on_request_state(request: Request) -> None:
|
||||
"""Idempotently create the OTEL SERVER span and stash it on
|
||||
``request.state.parent_otel_span``. Safe to call multiple times.
|
||||
|
|
@ -2598,6 +2643,15 @@ async def user_api_key_auth(
|
|||
)
|
||||
user_api_key_auth_obj.budget_reservation = None
|
||||
|
||||
# Admin-resolved OTEL destinations: anchor them on this request's task
|
||||
# context BEFORE downstream spans close, so the global-provider fan-out
|
||||
# processor forwards every span (server, auth, db, batch-write) to the
|
||||
# admin-assigned per-tenant backends -- not just the gen-AI span the
|
||||
# ``TenantTracerCache`` already routes. Also stashed on ``request.state``
|
||||
# so ``_apply_admin_logging_exporters`` reuses the result instead of
|
||||
# re-resolving.
|
||||
await _hoist_request_destinations(request, user_api_key_auth_obj)
|
||||
|
||||
## ENSURE DISABLE ROUTE WORKS ACROSS ALL USER AUTH FLOWS ##
|
||||
RouteChecks.should_call_route(
|
||||
route=route, valid_token=user_api_key_auth_obj, request=request
|
||||
|
|
|
|||
|
|
@ -682,7 +682,9 @@ async def _resolve_logging_exporters(
|
|||
|
||||
|
||||
async def _apply_admin_logging_exporters(
|
||||
data: dict, user_api_key_dict: UserAPIKeyAuth
|
||||
data: dict,
|
||||
user_api_key_dict: UserAPIKeyAuth,
|
||||
cached_destinations: "list | None" = None,
|
||||
) -> None:
|
||||
"""Stamp the resolved fan-out destinations onto ``data`` and activate their
|
||||
backends.
|
||||
|
|
@ -691,8 +693,23 @@ async def _apply_admin_logging_exporters(
|
|||
``all_litellm_params``, so scrubbed from the provider request body), not a
|
||||
top-level key, so an unknown field cannot leak to the provider. Default-deny
|
||||
means an identity with no assignment gets no per-tenant destination here.
|
||||
|
||||
``cached_destinations`` -- when ``user_api_key_auth`` already resolved the
|
||||
destinations on this request (the FastAPI path), reuse the result instead of
|
||||
running the resolver a second time. The SDK path passes ``None`` and the
|
||||
resolver runs here.
|
||||
"""
|
||||
destinations, backends = await _resolve_logging_exporters(user_api_key_dict)
|
||||
if cached_destinations is not None:
|
||||
destinations = list(cached_destinations)
|
||||
backends = list(
|
||||
dict.fromkeys(
|
||||
str(d["callback_name"])
|
||||
for d in destinations
|
||||
if isinstance(d, dict) and d.get("callback_name")
|
||||
)
|
||||
)
|
||||
else:
|
||||
destinations, backends = await _resolve_logging_exporters(user_api_key_dict)
|
||||
if not destinations:
|
||||
return
|
||||
proxy_metadata = data.get("litellm_metadata")
|
||||
|
|
@ -2012,8 +2029,13 @@ async def add_litellm_data_to_request(
|
|||
|
||||
# Admin-owned exporter assignment: resolve the union of exporters assigned across
|
||||
# the request's identity chain (key + team + org) into fan-out destinations and
|
||||
# activate their backends. Default-deny: an unassigned identity gets none.
|
||||
await _apply_admin_logging_exporters(data, user_api_key_dict)
|
||||
# activate their backends. Default-deny: an unassigned identity gets none. Reuse
|
||||
# the result ``user_api_key_auth`` already cached on ``request.state`` so the
|
||||
# resolver runs once per request, not twice.
|
||||
cached = getattr(getattr(request, "state", None), "otel_destinations", None)
|
||||
await _apply_admin_logging_exporters(
|
||||
data, user_api_key_dict, cached_destinations=cached
|
||||
)
|
||||
|
||||
# Add disabled callbacks from key metadata
|
||||
if (
|
||||
|
|
|
|||
268
tests/test_litellm/integrations/otel/test_otel_v2_fan_out.py
Normal file
268
tests/test_litellm/integrations/otel/test_otel_v2_fan_out.py
Normal file
|
|
@ -0,0 +1,268 @@
|
|||
"""Tests for the per-request tenant fan-out SpanProcessor (Hoist).
|
||||
|
||||
The fan-out processor lives on the main v2 ``TracerProvider`` and forwards every
|
||||
finished span to the admin-resolved per-tenant destinations carried on the
|
||||
request's ``ContextVar``. The contract under test:
|
||||
|
||||
- Spans on the main provider land at every per-tenant destination matching this
|
||||
backend (the ``owner_callback_name``).
|
||||
- When destinations is empty (an unassigned identity, the SDK path, a request
|
||||
before the resolver ran), the processor is a no-op.
|
||||
- Concurrent requests with different destinations stay isolated -- contextvars
|
||||
scope per task, so one request's tenant doesn't receive another's spans.
|
||||
- The processor caches one BatchSpanProcessor per ``(endpoint, headers)`` pair
|
||||
and skips destinations whose ``callback_name`` doesn't match its owner.
|
||||
"""
|
||||
|
||||
import asyncio
|
||||
from typing import Any
|
||||
|
||||
import pytest
|
||||
|
||||
pytest.importorskip("opentelemetry")
|
||||
|
||||
from opentelemetry.sdk.trace import TracerProvider # noqa: E402
|
||||
from opentelemetry.sdk.trace.export import SimpleSpanProcessor # noqa: E402
|
||||
from opentelemetry.sdk.trace.export.in_memory_span_exporter import ( # noqa: E402
|
||||
InMemorySpanExporter,
|
||||
)
|
||||
from opentelemetry.sdk.trace.export import SpanExporter # noqa: E402
|
||||
|
||||
from litellm.integrations.otel.model.destination import OtelDestination # noqa: E402
|
||||
from litellm.integrations.otel.plumbing.context import ( # noqa: E402
|
||||
set_request_destinations,
|
||||
_request_destinations,
|
||||
)
|
||||
from litellm.integrations.otel.plumbing.fan_out import ( # noqa: E402
|
||||
TenantFanOutSpanProcessor,
|
||||
_processor_key,
|
||||
)
|
||||
from litellm.integrations.otel.plumbing import providers # noqa: E402
|
||||
from litellm.integrations.otel.model.config import OpenTelemetryV2Config # noqa: E402
|
||||
|
||||
|
||||
@pytest.fixture(autouse=True)
|
||||
def _reset_request_destinations():
|
||||
_request_destinations.set(())
|
||||
yield
|
||||
_request_destinations.set(())
|
||||
|
||||
|
||||
def _build_provider_with_fan_out(
|
||||
owner: str, exporter: SpanExporter
|
||||
) -> tuple[TracerProvider, TenantFanOutSpanProcessor]:
|
||||
"""Build a real TracerProvider with the in-memory exporter as the configured
|
||||
backend, plus the fan-out processor wired in. The fan-out processor's
|
||||
per-destination processors are real BatchSpanProcessors, replaced in the
|
||||
test by an injected SimpleSpanProcessor against an in-memory exporter via
|
||||
``monkeypatch`` so the test can read what was forwarded."""
|
||||
cfg = OpenTelemetryV2Config(exporter="in_memory")
|
||||
provider = providers.build_tracer_provider(
|
||||
cfg, exporter=exporter, tenant_fan_out_owner=owner
|
||||
)
|
||||
fan_out = next(
|
||||
p
|
||||
for p in provider._active_span_processor._span_processors
|
||||
if isinstance(p, TenantFanOutSpanProcessor)
|
||||
)
|
||||
return provider, fan_out
|
||||
|
||||
|
||||
def test_fan_out_forwards_span_to_matching_destination(monkeypatch):
|
||||
"""A span on the main provider lands at the per-tenant destination whose
|
||||
callback_name matches the fan-out owner. Removing the fan-out branch makes
|
||||
this test fail (the tenant exporter never sees the span)."""
|
||||
main_exporter = InMemorySpanExporter()
|
||||
provider, fan_out = _build_provider_with_fan_out("langfuse_otel", main_exporter)
|
||||
|
||||
tenant_exporter = InMemorySpanExporter()
|
||||
# Swap the lazy per-destination processor build for an in-memory one so we
|
||||
# don't actually hit the network.
|
||||
def _stub_processor_for(self, destination):
|
||||
return SimpleSpanProcessor(tenant_exporter)
|
||||
|
||||
monkeypatch.setattr(
|
||||
TenantFanOutSpanProcessor, "_processor_for", _stub_processor_for
|
||||
)
|
||||
|
||||
set_request_destinations(
|
||||
(
|
||||
OtelDestination(
|
||||
callback_name="langfuse_otel",
|
||||
endpoint="https://cloud.langfuse.com/api/public/otel",
|
||||
headers={"Authorization": "Bearer pk:sk"},
|
||||
),
|
||||
)
|
||||
)
|
||||
tracer = provider.get_tracer("test")
|
||||
with tracer.start_as_current_span("auth /chat/completions"):
|
||||
pass
|
||||
|
||||
main_names = [s.name for s in main_exporter.get_finished_spans()]
|
||||
tenant_names = [s.name for s in tenant_exporter.get_finished_spans()]
|
||||
assert "auth /chat/completions" in main_names
|
||||
assert "auth /chat/completions" in tenant_names
|
||||
|
||||
|
||||
def test_fan_out_skips_destinations_for_other_backends(monkeypatch):
|
||||
"""Each fan-out processor owns ONE backend; a destination tagged for a
|
||||
different backend is skipped, so the wrong exporter never sees the span."""
|
||||
main_exporter = InMemorySpanExporter()
|
||||
_provider, _fan_out = _build_provider_with_fan_out(
|
||||
"langfuse_otel", main_exporter
|
||||
)
|
||||
|
||||
tenant_exporter = InMemorySpanExporter()
|
||||
monkeypatch.setattr(
|
||||
TenantFanOutSpanProcessor,
|
||||
"_processor_for",
|
||||
lambda self, destination: SimpleSpanProcessor(tenant_exporter),
|
||||
)
|
||||
|
||||
set_request_destinations(
|
||||
(
|
||||
OtelDestination(
|
||||
callback_name="arize",
|
||||
endpoint="https://otlp.arize.com/v1",
|
||||
headers={"api_key": "k", "space_id": "s"},
|
||||
),
|
||||
)
|
||||
)
|
||||
tracer = _provider.get_tracer("test")
|
||||
with tracer.start_as_current_span("auth"):
|
||||
pass
|
||||
|
||||
assert tenant_exporter.get_finished_spans() == ()
|
||||
|
||||
|
||||
def test_fan_out_noop_when_no_destinations(monkeypatch):
|
||||
"""An unassigned identity / pre-auth / SDK path leaves the contextvar at
|
||||
its empty-tuple default; the processor must short-circuit, NOT crash, NOT
|
||||
forward."""
|
||||
main_exporter = InMemorySpanExporter()
|
||||
provider, _ = _build_provider_with_fan_out("langfuse_otel", main_exporter)
|
||||
forwarded: list[Any] = []
|
||||
monkeypatch.setattr(
|
||||
TenantFanOutSpanProcessor,
|
||||
"_processor_for",
|
||||
lambda self, destination: forwarded.append(destination) or None,
|
||||
)
|
||||
tracer = provider.get_tracer("test")
|
||||
with tracer.start_as_current_span("auth"):
|
||||
pass
|
||||
assert forwarded == []
|
||||
|
||||
|
||||
def test_fan_out_caches_processor_per_destination_key(monkeypatch):
|
||||
"""Two requests with the same destination must share one cached
|
||||
BatchSpanProcessor; two different destinations must build two. Otherwise
|
||||
every request rebuilds the OTLP exporter (and its background thread)."""
|
||||
main_exporter = InMemorySpanExporter()
|
||||
_provider, fan_out = _build_provider_with_fan_out("langfuse_otel", main_exporter)
|
||||
|
||||
built: list[OtelDestination] = []
|
||||
|
||||
real_processor_for = TenantFanOutSpanProcessor._processor_for
|
||||
|
||||
def _spy(self, destination):
|
||||
result = real_processor_for(self, destination)
|
||||
if result is not None:
|
||||
built.append(destination)
|
||||
return result
|
||||
|
||||
# Build always returns None to keep the test offline, but we tracked
|
||||
# invocations via the spy above. Easier: stub the inner builder.
|
||||
def _stub(self, destination):
|
||||
built.append(destination)
|
||||
key = _processor_key(destination)
|
||||
if key in self._processors:
|
||||
return self._processors[key]
|
||||
proc = SimpleSpanProcessor(InMemorySpanExporter())
|
||||
self._processors[key] = proc
|
||||
return proc
|
||||
|
||||
monkeypatch.setattr(TenantFanOutSpanProcessor, "_processor_for", _stub)
|
||||
|
||||
dest_a = OtelDestination(
|
||||
callback_name="langfuse_otel", endpoint="https://a/", headers={"k": "1"}
|
||||
)
|
||||
dest_b = OtelDestination(
|
||||
callback_name="langfuse_otel", endpoint="https://b/", headers={"k": "1"}
|
||||
)
|
||||
set_request_destinations((dest_a,))
|
||||
tracer = _provider.get_tracer("test")
|
||||
with tracer.start_as_current_span("s1"):
|
||||
pass
|
||||
set_request_destinations((dest_a,))
|
||||
with tracer.start_as_current_span("s2"):
|
||||
pass
|
||||
set_request_destinations((dest_b,))
|
||||
with tracer.start_as_current_span("s3"):
|
||||
pass
|
||||
|
||||
# _stub appends every call; the cache size is what matters: two unique
|
||||
# ``(endpoint, headers)`` pairs -> two cached processors.
|
||||
assert len(fan_out._processors) == 2
|
||||
|
||||
|
||||
def test_fan_out_per_request_isolation_with_concurrent_tasks(monkeypatch):
|
||||
"""Two requests running concurrently with different destinations must each
|
||||
see only THEIR destination's spans. The contextvar scopes per asyncio task,
|
||||
so this is a pin on the contextvar approach: switching to a global variable
|
||||
would break this test."""
|
||||
main_exporter = InMemorySpanExporter()
|
||||
provider, _ = _build_provider_with_fan_out("langfuse_otel", main_exporter)
|
||||
|
||||
exporter_a = InMemorySpanExporter()
|
||||
exporter_b = InMemorySpanExporter()
|
||||
|
||||
def _stub(self, destination):
|
||||
if destination.endpoint == "https://a/":
|
||||
return SimpleSpanProcessor(exporter_a)
|
||||
return SimpleSpanProcessor(exporter_b)
|
||||
|
||||
monkeypatch.setattr(TenantFanOutSpanProcessor, "_processor_for", _stub)
|
||||
tracer = provider.get_tracer("test")
|
||||
|
||||
async def fire(label: str, endpoint: str):
|
||||
set_request_destinations(
|
||||
(
|
||||
OtelDestination(
|
||||
callback_name="langfuse_otel",
|
||||
endpoint=endpoint,
|
||||
headers={},
|
||||
),
|
||||
)
|
||||
)
|
||||
with tracer.start_as_current_span(label):
|
||||
await asyncio.sleep(0)
|
||||
|
||||
async def driver():
|
||||
await asyncio.gather(
|
||||
fire("req_a", "https://a/"),
|
||||
fire("req_b", "https://b/"),
|
||||
)
|
||||
|
||||
asyncio.run(driver())
|
||||
|
||||
names_a = {s.name for s in exporter_a.get_finished_spans()}
|
||||
names_b = {s.name for s in exporter_b.get_finished_spans()}
|
||||
assert "req_a" in names_a and "req_b" not in names_a
|
||||
assert "req_b" in names_b and "req_a" not in names_b
|
||||
|
||||
|
||||
def test_fan_out_owner_set_by_build_tracer_provider():
|
||||
"""The v2 logger opts the MAIN provider into fan-out by passing its
|
||||
callback_name; the TenantTracerCache clone providers do NOT, so the gen-AI
|
||||
span emitted through them is not also forwarded by the fan-out processor.
|
||||
Regressing this (adding the processor to clones) would double-export every
|
||||
gen-AI span."""
|
||||
from litellm.integrations.otel.logger import OpenTelemetryV2
|
||||
|
||||
logger = OpenTelemetryV2(callback_name="arize")
|
||||
main_procs = logger._tracer_provider._active_span_processor._span_processors
|
||||
assert any(isinstance(p, TenantFanOutSpanProcessor) for p in main_procs)
|
||||
# Build a tenant clone provider and confirm it has no fan-out processor.
|
||||
clone = providers.build_tracer_provider(logger.config)
|
||||
clone_procs = clone._active_span_processor._span_processors
|
||||
assert not any(isinstance(p, TenantFanOutSpanProcessor) for p in clone_procs)
|
||||
Loading…
Add table
Reference in a new issue