From c9929554bb6bc949b6518f617c85b3c63ae334e0 Mon Sep 17 00:00:00 2001 From: Yucheng He Date: Fri, 4 Sep 2026 15:57:35 -0700 Subject: [PATCH] fix(otel v2): identify a destination account by its credentials, not its header names Under additive the fan-out skips a destination the operator's own exporter already writes to, so the same account is not written twice. It compared header names as well as values, and one account answers to more than one spelling: the operator's Arize exporter sends space_id where a team destination sends arize-space-id, so every span landed in the operator's own space twice. The credentials are the identity. Compare those and leave the spelling to each backend. --- .../integrations/otel/plumbing/providers.py | 14 ++++++----- .../otel/test_otel_v2_destinations.py | 23 ++++++++++++++++++- 2 files changed, 30 insertions(+), 7 deletions(-) diff --git a/litellm/integrations/otel/plumbing/providers.py b/litellm/integrations/otel/plumbing/providers.py index 858cf8a1121..ef8a35801d8 100644 --- a/litellm/integrations/otel/plumbing/providers.py +++ b/litellm/integrations/otel/plumbing/providers.py @@ -218,8 +218,8 @@ _DRAIN_WORKERS: Final = 2 #: proxy open. _SHUTDOWN_DRAIN_SECONDS: Final = 5.0 -#: An exporter's account: its normalized endpoint and its credentials. -_SinkKey = tuple[str, tuple[tuple[str, str], ...]] +#: An exporter's account: its normalized endpoint and the credentials it presents. +_SinkKey = tuple[str, tuple[str, ...]] class _DrainPool: @@ -851,14 +851,16 @@ def operator_sink_keys(config: OpenTelemetryV2Config | None) -> frozenset[_SinkK def _sink_key(endpoint: str | None, headers: Mapping[str, str]) -> "_SinkKey | None": """The account an exporter writes to, or ``None`` when it has no fixed one. - Normalized on both counts that make the same account look like two: the operator's - spec carries the signal path a tenant destination leaves for the exporter to - append, and header names survive one round trip lowercased and the other not. + The credentials are the identity; the header names are only how each backend + spells them, and one account answers to more than one spelling (Arize takes the + operator's ``space_id`` and a tenant's ``arize-space-id``). The endpoint needs + normalizing too: the operator's spec carries the signal path that a tenant + destination leaves for the exporter to append. """ normalized: Final = _otlp_traces_endpoint(endpoint) if normalized is None: return None - return (normalized, tuple(sorted((name.lower(), value) for name, value in headers.items()))) + return (normalized, tuple(sorted(headers.values()))) def _attached_processors(provider: TracerProvider) -> "tuple[SpanProcessor, ...]": 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 04051e70cb9..d122f066670 100644 --- a/tests/test_litellm/integrations/otel/test_otel_v2_destinations.py +++ b/tests/test_litellm/integrations/otel/test_otel_v2_destinations.py @@ -39,6 +39,7 @@ from litellm.integrations.otel.presets.destinations import ( destination_capable_backends, destination_for, ) +from litellm.integrations.otel.presets.arize import arize_preset from litellm.integrations.otel.presets.langfuse import langfuse_preset from litellm.proxy._types import UserAPIKeyAuth from litellm.types.utils import StandardCallbackDynamicParams @@ -128,7 +129,7 @@ class TestRoutingMode: moment a team configures its own is what ``additive`` exists to prevent. """ - OPERATOR_SINK = ("https://cloud.langfuse.com/api/public/otel/v1/traces", (("authorization", "Basic op"),)) + OPERATOR_SINK = ("https://cloud.langfuse.com/api/public/otel/v1/traces", ("Basic op",)) #: What a tenant destination for that same project looks like before normalizing: #: no signal path yet, and the header name cased the way the backend writes it. SAME_ACCOUNT_ENDPOINT = "https://cloud.langfuse.com/api/public/otel" @@ -345,6 +346,26 @@ class TestRoutingMode: assert sink("pk-op", "sk-op") in operator, "a team naming the operator's own project" assert sink("pk-team", "sk-team") not in operator, "a different project on the same server" + def test_the_operators_own_arize_space_and_a_team_naming_it_are_one_account(self, monkeypatch): + """One account answers to two header names here: the operator's exporter sends + ``space_id`` and a team destination sends ``arize-space-id``. Keyed on the names, + additive would write the operator's own space twice for every request.""" + monkeypatch.setenv("ARIZE_SPACE_ID", "space-op") + monkeypatch.setenv("ARIZE_API_KEY", "key-op") + monkeypatch.delenv("ARIZE_SPACE_KEY", raising=False) + operator = operator_sink_keys(arize_preset()) + + def sink(space, api_key): + destination = destination_for( + "arize", + StandardCallbackDynamicParams(arize_space_key=space, arize_api_key=api_key), + ) + assert destination is not None + return _sink_key(destination.endpoint, destination.headers) + + assert sink("space-op", "key-op") in operator, "a team naming the operator's own space" + assert sink("space-team", "key-team") not in operator, "a different Arize space" + class TestFanOut: def test_every_span_of_the_request_reaches_the_destination_in_one_trace(self):