mirror of
https://github.com/BerriAI/litellm.git
synced 2026-10-06 02:48:13 +00:00
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.
This commit is contained in:
parent
f7eaeed00e
commit
3202ed5a5b
5 changed files with 99 additions and 28 deletions
|
|
@ -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,
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
|
|
|
|||
|
|
@ -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 == []
|
||||
|
|
|
|||
|
|
@ -296,17 +296,6 @@ const Settings: React.FC<SettingsPageProps> = ({ 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<SettingsPageProps> = ({ 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<SettingsPageProps> = ({ 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);
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue