From 3202ed5a5b81945123821a6f9003335d656cf701 Mon Sep 17 00:00:00 2001 From: Yucheng Zhu Date: Thu, 30 Jul 2026 18:02:50 -0700 Subject: [PATCH] fix(otel/v2): stop console-flood degrade, harden fan-out flush, confirm destination delete Three issues from review, each on the OTEL v2 path: A credential-less v2 config no longer degrades to a console exporter. OpenTelemetryV2Config folded its default `exporter="console"` into a real exporter whenever the exporter list was empty, so a generic destination with no endpoint (or a langfuse/weave/levo preset degrading via allow_missing_credentials) printed every span -- including prompt and completion content -- to stdout synchronously on the request path. A bare console-with-no-endpoint config is now left exporter-less (exports nothing); an explicit endpoint or non-console kind still folds. TenantFanOutSpanProcessor.force_flush and shutdown snapshot the processor cache with tuple() before iterating, so a concurrent on_end mutating the cache (insert / move_to_end / popitem) can no longer raise "OrderedDict mutated during iteration" mid-flush and drop the remaining destinations' buffered spans. Deleting a trace destination from the logging page now goes through the same delete confirmation modal as every other row instead of firing the credential delete immediately, so a mis-click can't irreversibly drop a destination and its stored collector secrets. --- litellm/integrations/otel/model/config.py | 10 ++++- litellm/integrations/otel/plumbing/routing.py | 11 ++++- .../integrations/otel/test_otel_v2_fan_out.py | 42 +++++++++++++++++++ .../integrations/otel/test_otel_v2_presets.py | 26 ++++++++++-- .../src/components/settings.tsx | 38 ++++++++--------- 5 files changed, 99 insertions(+), 28 deletions(-) diff --git a/litellm/integrations/otel/model/config.py b/litellm/integrations/otel/model/config.py index 7f33129c560..9ff6db2566d 100644 --- a/litellm/integrations/otel/model/config.py +++ b/litellm/integrations/otel/model/config.py @@ -244,8 +244,14 @@ class OpenTelemetryV2Config(BaseSettings): if self.endpoint and self.exporter == "console": self.exporter = "otlp_http" # When no explicit destinations are given, fold the single-destination - # shorthand into one spec so the provider always has a destination. - if not self.exporters: + # shorthand into one spec so the provider has a destination. A bare config + # whose only shorthand is the default console kind with no endpoint is the + # "nothing configured" degrade case (a preset returned no credentials, or v2 + # is enabled with no exporter set); leave it exporter-less so the provider + # exports nothing, rather than folding it into a console exporter that prints + # every span -- including prompt and completion content -- to stdout + # synchronously on the request path. + if not self.exporters and (self.endpoint or self.exporter != "console"): self.exporters = [ ExporterSpec( kind=self.exporter, diff --git a/litellm/integrations/otel/plumbing/routing.py b/litellm/integrations/otel/plumbing/routing.py index 37cbea69145..96b69863d1b 100644 --- a/litellm/integrations/otel/plumbing/routing.py +++ b/litellm/integrations/otel/plumbing/routing.py @@ -379,7 +379,11 @@ class TenantFanOutSpanProcessor(SpanProcessor): ) def shutdown(self) -> None: - for processor in self._processors.values(): + # Snapshot before iterating: ``on_end`` mutates ``self._processors`` (insert / + # move_to_end / popitem) on the span-ending thread and can run concurrently with + # this SDK-driven shutdown, so iterating the live mapping risks a + # "mutated during iteration" RuntimeError that the per-item except can't catch. + for processor in tuple(self._processors.values()): try: processor.shutdown() except Exception as exc: # noqa: BLE001 # a single processor's shutdown failure must not abort shutting down the rest @@ -388,7 +392,10 @@ class TenantFanOutSpanProcessor(SpanProcessor): def force_flush(self, timeout_millis: int = 30000) -> bool: all_ok = True - for processor in self._processors.values(): + # Snapshot before iterating (see ``shutdown``): a concurrent ``on_end`` mutating + # the processor cache must not abort the flush and drop the remaining destinations' + # buffered spans. + for processor in tuple(self._processors.values()): try: if not processor.force_flush(timeout_millis): all_ok = False diff --git a/tests/test_litellm/integrations/otel/test_otel_v2_fan_out.py b/tests/test_litellm/integrations/otel/test_otel_v2_fan_out.py index c014527e32b..607f1e82014 100644 --- a/tests/test_litellm/integrations/otel/test_otel_v2_fan_out.py +++ b/tests/test_litellm/integrations/otel/test_otel_v2_fan_out.py @@ -285,6 +285,48 @@ def test_fan_out_caches_processor_per_destination_key(monkeypatch): assert len(fan_out._processors) == 2 +def test_fan_out_flush_and_shutdown_survive_concurrent_cache_mutation(): + """force_flush/shutdown snapshot the processor cache before iterating, so a + concurrent ``on_end`` inserting or evicting a processor cannot raise a + 'mutated during iteration' RuntimeError and drop the remaining destinations' + buffered spans.""" + _provider, fan_out = _build_provider_with_fan_out("langfuse_otel", InMemorySpanExporter()) + + touched: list[str] = [] + + class _MutatingProc: + def __init__(self, name: str, mutate: bool = False) -> None: + self.name = name + self.mutate = mutate + + def _maybe_mutate(self, key: str) -> None: + if self.mutate: # simulate a concurrent on_end caching a new processor + fan_out._processors[(key, "x")] = _MutatingProc(key) + + def force_flush(self, timeout_millis: int = 30000) -> bool: + touched.append(self.name) + self._maybe_mutate("late-flush") + return True + + def shutdown(self) -> None: + touched.append("sd-" + self.name) + self._maybe_mutate("late-shutdown") + + fan_out._processors.clear() + fan_out._processors[("a", "1")] = _MutatingProc("a", mutate=True) + fan_out._processors[("b", "2")] = _MutatingProc("b") + fan_out._processors[("c", "3")] = _MutatingProc("c") + + # No RuntimeError despite the mutation mid-iteration; every original processor ran. + assert fan_out.force_flush() is True + assert {"a", "b", "c"} <= set(touched) + + touched.clear() + fan_out._processors[("a", "1")] = _MutatingProc("a", mutate=True) + fan_out.shutdown() + assert "sd-a" in touched + + 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, diff --git a/tests/test_litellm/integrations/otel/test_otel_v2_presets.py b/tests/test_litellm/integrations/otel/test_otel_v2_presets.py index c0ca48f09e7..e88580aaa6e 100644 --- a/tests/test_litellm/integrations/otel/test_otel_v2_presets.py +++ b/tests/test_litellm/integrations/otel/test_otel_v2_presets.py @@ -273,7 +273,27 @@ def test_generic_preset_needs_no_global_env_and_emits_genai(monkeypatch): monkeypatch.delenv(var, raising=False) cfg = generic_preset() # must not raise, no env needed assert "genai" in cfg.mapper_names - # no vendor (owned) exporter contributed -- only the base/global passthrough, if any - assert all(e.owner is None for e in cfg.exporters) + # A credential-less generic config must contribute NO exporter at all. In particular + # it must not degrade to the default console exporter, which prints every span + # (including prompt/completion content) to stdout synchronously on the request path. + assert cfg.exporters == [] # the degrade flag is accepted (Preset protocol) and irrelevant -- still builds - assert generic_preset(allow_missing_credentials=True) is not None + assert generic_preset(allow_missing_credentials=True).exporters == [] + + +def test_bare_and_degraded_configs_do_not_console_flood(monkeypatch): + # The "nothing configured" degrade case (bare config, or a credential-mandatory + # preset degrading with allow_missing_credentials) must export nothing rather than + # folding the default console exporter in and dumping every span to stdout. + from litellm.integrations.otel.model.config import OpenTelemetryV2Config + from litellm.integrations.otel.presets.langfuse import langfuse_preset + from litellm.integrations.otel.presets.weave import weave_preset + + assert OpenTelemetryV2Config().exporters == [] + # An explicit endpoint (a real destination) still folds into one OTLP exporter. + assert len(OpenTelemetryV2Config(endpoint="http://collector/v1/traces").exporters) == 1 + + for var in ("LANGFUSE_PUBLIC_KEY", "LANGFUSE_SECRET_KEY", "WANDB_API_KEY", "WANDB_PROJECT_ID"): + monkeypatch.delenv(var, raising=False) + assert langfuse_preset(allow_missing_credentials=True).exporters == [] + assert weave_preset(allow_missing_credentials=True).exporters == [] diff --git a/ui/litellm-dashboard/src/components/settings.tsx b/ui/litellm-dashboard/src/components/settings.tsx index 8b0d8c2d165..6a85612d405 100644 --- a/ui/litellm-dashboard/src/components/settings.tsx +++ b/ui/litellm-dashboard/src/components/settings.tsx @@ -296,17 +296,6 @@ const Settings: React.FC = ({ accessToken, userRole, userID, resolvedScope: resolveScope(c.credential_info?.access), })); - const handleDeleteDestination = async (name: string) => { - if (!accessToken) return; - try { - await credentialDeleteCall(accessToken, name); - NotificationsManager.success("Logging destination deleted"); - refetchCredentials(); - } catch (error) { - NotificationsManager.fromBackend(parseErrorMessage(error)); - } - }; - useEffect(() => { if (!accessToken) { return; @@ -633,13 +622,22 @@ const Settings: React.FC = ({ accessToken, userRole, userID, try { setIsDeletingCallback(true); - await deleteCallback(accessToken, callbackToDelete.name); - NotificationsManager.success(`Callback ${callbackToDelete.name} deleted successfully`); - - // Refresh the callbacks list - if (userID && userRole) { - const data = await getCallbacksCall(accessToken, userID, userRole); - setCallbacks(data.callbacks); + // A destination row carries a credentialName; it is a logging credential and is + // deleted (with its stored collector secrets) through the credential endpoint. A + // plain config callback is deleted through the callback endpoint. Both run only + // after the same delete confirmation, so a mis-click can't drop either instantly. + if (callbackToDelete.credentialName) { + await credentialDeleteCall(accessToken, callbackToDelete.credentialName); + NotificationsManager.success("Logging destination deleted"); + refetchCredentials(); + } else { + await deleteCallback(accessToken, callbackToDelete.name); + NotificationsManager.success(`Callback ${callbackToDelete.name} deleted successfully`); + // Refresh the callbacks list + if (userID && userRole) { + const data = await getCallbacksCall(accessToken, userID, userRole); + setCallbacks(data.callbacks); + } } setShowDeleteConfirmModal(false); @@ -682,9 +680,7 @@ const Settings: React.FC = ({ accessToken, userRole, userID, onEditAccess={(cb) => cb.credentialName && setEditAccessFor({ name: cb.credentialName, access: cb.access }) } - onDelete={(cb) => - cb.credentialName ? handleDeleteDestination(cb.credentialName) : handleDeleteCallback(cb) - } + onDelete={(cb) => handleDeleteCallback(cb)} onTest={async (cb) => { try { await serviceHealthCheck(accessToken, cb.name);