fix(redis): log an open circuit breaker once instead of a traceback per request and count sync timeouts as timeouts

While the Redis circuit breaker is open every guarded call was refused with a bare
Exception that each swallowing catch site logged as an ERROR traceback, so a
sub-second latency blip turned into thousands of tracebacks per minute and pinned
every replica's CPU. The sync guard also recorded socket timeouts as hard
connectivity failures, so with least-busy routing the breaker opened on the first
slow replies and the timeout-only min-duration guard never applied.

Refusals now raise RedisCircuitBreakerOpenError, and the catch sites route it
through log_redis_failure, which logs a refusal at DEBUG and everything else at
the caller's level. The sync guard passes is_timeout like the async one.
This commit is contained in:
mateo-berri 2026-09-10 13:50:34 -07:00
parent 692a311efb
commit 9da3e63ae9
9 changed files with 226 additions and 27 deletions

View file

@ -10,6 +10,7 @@
import ast
import hashlib
import json
import logging
import time
import traceback
from collections.abc import Mapping
@ -32,7 +33,7 @@ from .dual_cache import DualCache # noqa: F401
from .gcs_cache import GCSCache
from .in_memory_cache import InMemoryCache
from .qdrant_semantic_cache import QdrantSemanticCache
from .redis_cache import RedisCache
from .redis_cache import RedisCache, log_redis_failure
from .redis_cluster_cache import RedisClusterCache
from .redis_semantic_cache import RedisSemanticCache
from .s3_cache import S3Cache
@ -678,7 +679,7 @@ class Cache:
cache_key, cached_data, kwargs = self._add_cache_logic(result=result, **kwargs)
self.cache.set_cache(cache_key, cached_data, **kwargs)
except Exception as e:
verbose_logger.exception("LiteLLM Cache: Excepton add_cache: %s", e)
log_redis_failure(verbose_logger, logging.ERROR, "LiteLLM Cache: exception in add_cache", e)
async def async_add_cache(self, result, dynamic_cache_object: BaseCache | None = None, **kwargs):
"""
@ -697,7 +698,7 @@ class Cache:
else:
await self.cache.async_set_cache(cache_key, cached_data, **kwargs)
except Exception as e:
verbose_logger.exception("LiteLLM Cache: Excepton add_cache: %s", e)
log_redis_failure(verbose_logger, logging.ERROR, "LiteLLM Cache: exception in add_cache", e)
def _convert_to_cached_embedding(
self,
@ -876,7 +877,7 @@ class Cache:
else:
await self.cache.async_set_cache_pipeline(cache_list=cache_list, **kwargs)
except Exception as e:
verbose_logger.exception("LiteLLM Cache: Excepton add_cache: %s", e)
log_redis_failure(verbose_logger, logging.ERROR, "LiteLLM Cache: exception in add_cache", e)
def should_use_cache(self, **kwargs):
"""

View file

@ -8,8 +8,8 @@ Has 4 primary methods:
- async_get_cache
"""
import logging
import time
import traceback
from collections.abc import Sequence
from threading import Lock
from typing import TYPE_CHECKING, Any, Final
@ -23,7 +23,7 @@ from litellm.constants import DEFAULT_MAX_REDIS_BATCH_CACHE_SIZE
from .base_cache import BaseCache
from .in_memory_cache import InMemoryCache
from .redis_cache import RedisCache
from .redis_cache import RedisCache, log_redis_failure
if TYPE_CHECKING:
from opentelemetry.trace import Span as _Span
@ -177,8 +177,8 @@ class DualCache(BaseCache):
print_verbose(f"get cache: cache result: {result}")
return result
except Exception:
verbose_logger.error(traceback.format_exc())
except Exception as e:
log_redis_failure(verbose_logger, logging.ERROR, "LiteLLM Cache: exception in get_cache", e)
def batch_get_cache(
self,
@ -217,8 +217,8 @@ class DualCache(BaseCache):
return list( # mutable-ok: public list contract
redis_result.get(key) if value is None else value for key, value in zip(keys, result)
)
except Exception:
verbose_logger.error(traceback.format_exc())
except Exception as e:
log_redis_failure(verbose_logger, logging.ERROR, "LiteLLM Cache: exception in batch_get_cache", e)
async def async_get_cache(
self,
@ -250,8 +250,8 @@ class DualCache(BaseCache):
print_verbose(f"get cache: cache result: {result}")
return result
except Exception:
verbose_logger.error(traceback.format_exc())
except Exception as e:
log_redis_failure(verbose_logger, logging.ERROR, "LiteLLM Cache: exception in async_get_cache", e)
def _reserve_redis_batch_keys(
self,
@ -339,8 +339,8 @@ class DualCache(BaseCache):
await self.in_memory_cache.async_set_cache(key, value, **self._backfill_kwargs(kwargs))
return result
except Exception:
verbose_logger.error(traceback.format_exc())
except Exception as e:
log_redis_failure(verbose_logger, logging.ERROR, "LiteLLM Cache: exception in async_batch_get_cache", e)
async def async_set_cache(self, key, value, local_only: bool = False, **kwargs):
print_verbose(f"async set cache: cache key: {key}; local_only: {local_only}; value: {value}")
@ -353,7 +353,7 @@ class DualCache(BaseCache):
if self.redis_cache is not None and local_only is False:
await self.redis_cache.async_set_cache(key, value, **kwargs)
except Exception as e:
verbose_logger.exception("LiteLLM Cache: Excepton async add_cache: %s", e)
log_redis_failure(verbose_logger, logging.ERROR, "LiteLLM Cache: exception in async add_cache", e)
# async_batch_set_cache
async def async_set_cache_pipeline(self, cache_list: list, local_only: bool = False, **kwargs):
@ -372,7 +372,7 @@ class DualCache(BaseCache):
cache_list=cache_list, ttl=kwargs.pop("ttl", None), **kwargs
)
except Exception as e:
verbose_logger.exception("LiteLLM Cache: Excepton async add_cache: %s", e)
log_redis_failure(verbose_logger, logging.ERROR, "LiteLLM Cache: exception in async add_cache", e)
async def async_increment_cache(
self,
@ -439,8 +439,10 @@ class DualCache(BaseCache):
return result
except Exception as e:
verbose_logger.warning(
"Redis async_increment_cache_pipeline failed, falling back to in-memory result: %s",
log_redis_failure(
verbose_logger,
logging.WARNING,
"Redis async_increment_cache_pipeline failed, falling back to in-memory result",
e,
)
return result

View file

@ -14,6 +14,7 @@ import functools
import hashlib
import inspect
import json
import logging
import time
from collections.abc import Awaitable, Callable, Sequence
from contextvars import ContextVar
@ -391,10 +392,26 @@ def _record_swallowed_redis_failure(breaker: RedisCircuitBreaker, exc: BaseExcep
_swallowed_redis_failures.set(_swallowed_redis_failures.get() + 1)
class RedisCircuitBreakerOpenError(Exception):
"""Raised in place of a Redis call while the circuit breaker is open."""
def log_redis_failure(logger: logging.Logger, level: int, message: str, exc: BaseException) -> None:
"""Log a Redis failure the caller is about to swallow.
An open breaker refuses every call until Redis recovers and announced itself once when it
opened, so the calls it refuses are logged at debug instead of once per request at ``level``.
"""
if isinstance(exc, RedisCircuitBreakerOpenError):
logger.debug("%s: %s", message, exc)
return
logger.log(level, "%s: %s", message, exc, exc_info=exc if level >= logging.ERROR else None)
def _enter_circuit_breaker(breaker: RedisCircuitBreaker, name: str) -> int:
"""Reject the call if the breaker is open, else return the swallowed-failure count to compare against."""
if breaker.is_open():
raise Exception(f"Redis circuit breaker is open — skipping {name}")
raise RedisCircuitBreakerOpenError(f"Redis circuit breaker is open — skipping {name}")
return _swallowed_redis_failures.get()
@ -440,7 +457,7 @@ def _run_under_circuit_breaker_sync(
result: Final = call()
except Exception as e:
if _is_redis_health_failure(e):
breaker.record_failure()
breaker.record_failure(is_timeout=_is_redis_timeout_failure(e))
raise
_exit_circuit_breaker(breaker, swallowed_before)
return result

View file

@ -6,6 +6,7 @@ This is currently in development and not yet ready for production.
import asyncio
import binascii
import logging
import os
import uuid
from collections.abc import Awaitable, Callable, Mapping, Sequence, Set
@ -26,6 +27,7 @@ from typing_extensions import NotRequired, ReadOnly
from litellm import DualCache
from litellm._logging import verbose_proxy_logger
from litellm.caching.redis_cache import log_redis_failure
from litellm.constants import DYNAMIC_RATE_LIMIT_ERROR_THRESHOLD_PER_MINUTE, INTERNAL_CALL_ORIGIN_METADATA_KEY
from litellm.integrations.custom_logger import CustomLogger
from litellm.litellm_core_utils.prompt_templates.common_utils import (
@ -3856,7 +3858,9 @@ class _PROXY_MaxParallelRequestsHandler_v3(CustomLogger):
)
except Exception as e:
verbose_proxy_logger.warning("TTL preservation failed, falling back to regular pipeline: %s", e)
log_redis_failure(
verbose_proxy_logger, logging.WARNING, "TTL preservation failed, falling back to regular pipeline", e
)
# Fallback to regular pipeline on error
await self.internal_usage_cache.dual_cache.async_increment_cache_pipeline(
increment_list=pipeline_operations,

View file

@ -1,3 +1,4 @@
import logging
from collections.abc import Mapping, Sequence
from typing import Final
@ -6,6 +7,7 @@ from typing_extensions import ReadOnly, TypedDict
from litellm._logging import verbose_router_logger
from litellm.caching.caching import DualCache
from litellm.caching.redis_cache import log_redis_failure
from litellm.integrations.custom_logger import CustomLogger
IN_FLIGHT_COUNT_TTL_SECONDS: Final = 60 * 60
@ -87,16 +89,22 @@ def _least_busy(
def _warn_unreadable(model_group: str, error: Exception) -> None:
verbose_router_logger.warning(
"least-busy routing could not read the shared in-flight counts for %s, "
"falling back to this worker's own counts: %s",
model_group,
log_redis_failure(
verbose_router_logger,
logging.WARNING,
f"least-busy routing could not read the shared in-flight counts for {model_group}, "
"falling back to this worker's own counts",
error,
)
def _warn_unwritable(key: str, error: Exception) -> None:
verbose_router_logger.warning("least-busy routing could not update the in-flight count under %s: %s", key, error)
log_redis_failure(
verbose_router_logger,
logging.WARNING,
f"least-busy routing could not update the in-flight count under {key}",
error,
)
class LeastBusyLoggingHandler(CustomLogger):

View file

@ -1,4 +1,5 @@
import asyncio
import logging
import time
import uuid
from unittest.mock import AsyncMock, MagicMock, patch
@ -7,7 +8,8 @@ import pytest
from litellm.caching.dual_cache import DualCache
from litellm.caching.in_memory_cache import InMemoryCache
from litellm.caching.redis_cache import RedisCache
from litellm.caching.redis_cache import RedisCache, _redis_circuit_breaker_guard, _redis_circuit_breaker_guard_sync
from litellm.types.caching import RedisPipelineIncrementOperation
@pytest.mark.asyncio
@ -576,3 +578,102 @@ async def test_dual_cache_late_attach_redis_wires_writes_and_ttl_async():
assert mock_redis.async_set_cache.call_args[0][:2] == (key_after, val_after)
assert in_memory.get_cache(key_after) == val_after
class _OpenBreakerRedis:
"""A RedisCache whose breaker is open, so every guarded call is refused before it starts."""
def __init__(self) -> None:
from litellm.caching.redis_cache import RedisCircuitBreaker
self._circuit_breaker = RedisCircuitBreaker(failure_threshold=3, recovery_timeout=60)
for _ in range(3):
self._circuit_breaker.record_failure()
@_redis_circuit_breaker_guard
async def async_get_cache(self, key, **kwargs):
raise AssertionError("never reached")
@_redis_circuit_breaker_guard
async def async_batch_get_cache(self, key_list, **kwargs):
raise AssertionError("never reached")
@_redis_circuit_breaker_guard
async def async_set_cache(self, key, value, **kwargs):
raise AssertionError("never reached")
@_redis_circuit_breaker_guard
async def async_set_cache_pipeline(self, cache_list, **kwargs):
raise AssertionError("never reached")
@_redis_circuit_breaker_guard
async def async_increment_pipeline(self, increment_list, **kwargs):
raise AssertionError("never reached")
@_redis_circuit_breaker_guard_sync
def get_cache(self, key, **kwargs):
raise AssertionError("never reached")
@_redis_circuit_breaker_guard_sync
def batch_get_cache(self, key_list, **kwargs):
raise AssertionError("never reached")
@pytest.mark.asyncio
@pytest.mark.parametrize(
"call",
[
lambda cache: cache.async_get_cache("k"),
lambda cache: cache.async_batch_get_cache(["k1", "k2"]),
lambda cache: cache.async_set_cache("k", "v"),
lambda cache: cache.async_set_cache_pipeline([("k", "v")]),
lambda cache: cache.async_increment_cache_pipeline(
increment_list=[RedisPipelineIncrementOperation(key="k", increment_value=1.0, ttl=60)]
),
],
ids=["get", "batch_get", "set", "set_pipeline", "increment_pipeline"],
)
async def test_an_open_circuit_breaker_is_not_an_error_per_request(caplog, call):
"""While the breaker is open every request is refused by design, and the breaker already
said so once when it opened; logging each refusal as an ERROR traceback was the storm that
pinned every worker's CPU during a Redis latency blip."""
cache = DualCache(in_memory_cache=InMemoryCache(), redis_cache=_OpenBreakerRedis()) # pyright: ignore[reportArgumentType] # duck-typed Redis double
caplog.clear()
with caplog.at_level(logging.DEBUG, logger="LiteLLM"):
await call(cache)
assert [record.levelno for record in caplog.records if record.levelno >= logging.WARNING] == []
assert any("circuit breaker is open" in record.getMessage() for record in caplog.records)
@pytest.mark.parametrize(
"call",
[lambda cache: cache.get_cache("k"), lambda cache: cache.batch_get_cache(["k1", "k2"])],
ids=["get", "batch_get"],
)
def test_an_open_circuit_breaker_is_not_an_error_per_sync_request(caplog, call):
cache = DualCache(in_memory_cache=InMemoryCache(), redis_cache=_OpenBreakerRedis()) # pyright: ignore[reportArgumentType] # duck-typed Redis double
caplog.clear()
with caplog.at_level(logging.DEBUG, logger="LiteLLM"):
call(cache)
assert [record.levelno for record in caplog.records if record.levelno >= logging.WARNING] == []
assert any("circuit breaker is open" in record.getMessage() for record in caplog.records)
@pytest.mark.asyncio
async def test_a_real_redis_failure_still_logs_an_error(caplog):
class _BrokenRedis:
async def async_get_cache(self, key, **kwargs):
raise ConnectionError("redis is down")
cache = DualCache(in_memory_cache=InMemoryCache(), redis_cache=_BrokenRedis()) # pyright: ignore[reportArgumentType] # duck-typed Redis double
with caplog.at_level(logging.DEBUG, logger="LiteLLM"):
assert await cache.async_get_cache("k") is None
errors = [record for record in caplog.records if record.levelno == logging.ERROR]
assert [record.getMessage() for record in errors] == ["LiteLLM Cache: exception in async_get_cache: redis is down"]
assert errors[0].exc_info is not None

View file

@ -1013,3 +1013,25 @@ async def test_breaker_metrics_track_state_and_failure_class():
breaker.record_success()
assert sample("litellm_redis_circuit_breaker_state", {"state": "open"}) == open_gauge_before
assert sample("litellm_redis_circuit_breaker_state", {"state": "closed"}) == closed_gauge_before + 1
def test_sync_guard_counts_a_timeout_as_a_timeout():
"""A sync Redis timeout must wait out the timeout-only min duration exactly like the async guard.
Recording it as a hard connectivity failure opened the breaker on the fifth slow reply,
which is how a latency blip took the shared cache out for every worker.
"""
from redis.exceptions import TimeoutError as RedisTimeoutError
from litellm.caching.redis_cache import RedisCircuitBreaker, _run_under_circuit_breaker_sync
breaker = RedisCircuitBreaker(failure_threshold=3, recovery_timeout=60, timeout_min_duration=5.0)
def timing_out_call() -> str:
raise RedisTimeoutError("read timed out")
for _ in range(6):
with pytest.raises(RedisTimeoutError):
_run_under_circuit_breaker_sync(breaker, "op", timing_out_call)
assert breaker.is_open() is False

View file

@ -6236,3 +6236,24 @@ async def test_post_call_success_hook_leaves_raw_provider_dict_untouched():
)
assert response == {"id": "msg_123", "type": "message", "role": "assistant", "content": []}
@pytest.mark.asyncio
async def test_an_open_circuit_breaker_falls_back_to_the_pipeline_without_a_warning(caplog):
from litellm.caching.redis_cache import RedisCircuitBreakerOpenError
from litellm.types.caching import RedisPipelineIncrementOperation
handler = _PROXY_MaxParallelRequestsHandler(internal_usage_cache=InternalUsageCache(DualCache()))
async def refused_script(keys, args):
raise RedisCircuitBreakerOpenError("Redis circuit breaker is open")
handler.token_increment_script = refused_script
with caplog.at_level(logging.DEBUG, logger="LiteLLM Proxy"):
await handler.async_increment_tokens_with_ttl_preservation(
pipeline_operations=[RedisPipelineIncrementOperation(key="quiet_key", increment_value=10.0, ttl=60)]
)
assert await handler.internal_usage_cache.dual_cache.async_get_cache("quiet_key") == 10.0
assert [record.getMessage() for record in caplog.records if record.levelno >= logging.WARNING] == []

View file

@ -1,9 +1,11 @@
import logging
from typing import Final
import pytest
from litellm.caching.caching import DualCache
from litellm.caching.in_memory_cache import InMemoryCache
from litellm.caching.redis_cache import RedisCircuitBreakerOpenError
from litellm.router_strategy.least_busy import IN_FLIGHT_COUNT_TTL_SECONDS, LeastBusyLoggingHandler
GROUP: Final = "least-busy-group"
@ -185,3 +187,24 @@ def test_calls_without_a_deployment_are_ignored() -> None:
worker.log_pre_api_call(model="m", messages=[], kwargs={})
assert shared.counts == {}
class OpenBreakerRedis(SharedRedisCounters):
def increment_with_floor(self, key: str, value: int, ttl: int) -> int:
raise RedisCircuitBreakerOpenError("Redis circuit breaker is open")
def batch_get_counts(self, key_list: list[str]) -> tuple[int | None, ...]:
raise RedisCircuitBreakerOpenError("Redis circuit breaker is open")
@pytest.mark.asyncio
async def test_an_open_circuit_breaker_falls_back_without_a_warning_per_request(caplog: pytest.LogCaptureFixture) -> None:
worker: Final = _worker(OpenBreakerRedis())
with caplog.at_level(logging.DEBUG, logger="LiteLLM Router"):
worker.log_pre_api_call(model="m", messages=[], kwargs=_call_kwargs("dep-a"))
picked: Final = worker.get_available_deployments(GROUP, HEALTHY)
assert picked is DEPLOYMENT_B
assert [record.getMessage() for record in caplog.records if record.levelno >= logging.WARNING] == []
assert sum("circuit breaker is open" in record.getMessage() for record in caplog.records) == 2