From b54402392303f666a46c2a68d8f2e03f27b175ad Mon Sep 17 00:00:00 2001 From: Yucheng Zhu Date: Wed, 19 Aug 2026 14:35:34 -0700 Subject: [PATCH] fix(observability): report a swallowed Redis outage as an error, not contention `acquire_lock` caught exceptions to tell a failed attempt from losing the election, but an outage does not always raise. `RedisCache.async_set_cache` catches connection errors, records a swallowed-failure marker and returns None, and None is also how redis reports SET NX losing the race. Both arrived as `not_acquired`, so a Redis outage was published as ordinary contention: the distinction the split was added to make, lost in the most important case. The marker is the only thing that separates them, so it is compared across the attempt. It is a per-task ContextVar, so a concurrent caller's failure cannot be misread as this one's. `swallowed_redis_failure_count` exposes it rather than having callers reach for the private ContextVar. Verified against a stub that mimics RedisCache exactly: an outage now records error, losing the race still records not_acquired, and success still records acquired. --- litellm/caching/redis_cache.py | 10 +++++++ .../db_transaction_queue/pod_lock_manager.py | 11 ++++++- .../test_pod_lock_manager.py | 30 +++++++++++++++++++ 3 files changed, 50 insertions(+), 1 deletion(-) diff --git a/litellm/caching/redis_cache.py b/litellm/caching/redis_cache.py index 934ba500ef9..a820d06646f 100644 --- a/litellm/caching/redis_cache.py +++ b/litellm/caching/redis_cache.py @@ -217,6 +217,16 @@ def _record_swallowed_redis_failure(breaker: RedisCircuitBreaker, exc: BaseExcep _swallowed_redis_failures.set(_swallowed_redis_failures.get() + 1) +def swallowed_redis_failure_count() -> int: + """Redis failures this task has swallowed so far. + + Several methods catch their own connection errors and return a default, so a + caller that needs to tell "Redis did not answer" from an ordinary negative + result compares this across the call. + """ + return _swallowed_redis_failures.get() + + async def _run_under_circuit_breaker( breaker: RedisCircuitBreaker, name: str, diff --git a/litellm/proxy/db/db_transaction_queue/pod_lock_manager.py b/litellm/proxy/db/db_transaction_queue/pod_lock_manager.py index 33b71627894..9df4da3e552 100644 --- a/litellm/proxy/db/db_transaction_queue/pod_lock_manager.py +++ b/litellm/proxy/db/db_transaction_queue/pod_lock_manager.py @@ -4,7 +4,7 @@ from typing import TYPE_CHECKING, Any, Final from litellm._logging import verbose_proxy_logger from litellm._uuid import uuid -from litellm.caching.redis_cache import RedisCache +from litellm.caching.redis_cache import RedisCache, swallowed_redis_failure_count from litellm.constants import DEFAULT_CRON_JOB_LOCK_TTL_SECONDS from litellm.proxy.db.db_transaction_queue.base_update_queue import service_logger_obj from litellm.types.integrations.prometheus import LockAttemptResult @@ -73,12 +73,21 @@ end verbose_proxy_logger.debug("redis_cache is None, skipping acquire_lock") _record_lock_attempt(cronjob_id, LockAttemptResult.NO_REDIS) return None + swallowed_before: Final = swallowed_redis_failure_count() try: acquired: Final = await self._attempt_acquire_lock(cronjob_id, ttl=ttl, allow_reentrant=allow_reentrant) except Exception as e: verbose_proxy_logger.error("Error acquiring Redis lock for %s: %s", cronjob_id, e) _record_lock_attempt(cronjob_id, LockAttemptResult.ERROR) return False + # An outage does not always raise. RedisCache catches connection errors and + # returns None, which is also how redis reports SET NX losing the race, so + # only the swallowed-failure marker separates the two. Without it every + # outage is published as ordinary contention. + if swallowed_redis_failure_count() != swallowed_before: + verbose_proxy_logger.error("Redis unavailable while acquiring lock for %s", cronjob_id) + _record_lock_attempt(cronjob_id, LockAttemptResult.ERROR) + return False _record_lock_attempt(cronjob_id, LockAttemptResult.ACQUIRED if acquired else LockAttemptResult.NOT_ACQUIRED) return acquired diff --git a/tests/test_litellm/proxy/db/db_transaction_queue/test_pod_lock_manager.py b/tests/test_litellm/proxy/db/db_transaction_queue/test_pod_lock_manager.py index 031ab51688a..ebc4fd89f23 100644 --- a/tests/test_litellm/proxy/db/db_transaction_queue/test_pod_lock_manager.py +++ b/tests/test_litellm/proxy/db/db_transaction_queue/test_pod_lock_manager.py @@ -506,6 +506,36 @@ async def test_recording_the_lock_outcome_never_blocks_the_job(): assert await manager.acquire_lock(cronjob_id="db_spend_update_job") is True +@pytest.mark.asyncio +async def test_a_swallowed_redis_outage_is_not_reported_as_losing_the_election(): + """RedisCache catches connection errors and returns None rather than raising, + and None is also how redis reports SET NX losing the race. Without the + swallowed-failure marker every outage publishes as ordinary contention, + which is the distinction this metric exists to make.""" + from unittest.mock import AsyncMock, MagicMock, patch + + from litellm.caching.redis_cache import _record_swallowed_redis_failure + from litellm.proxy.db.db_transaction_queue.pod_lock_manager import PodLockManager + from litellm.types.integrations.prometheus import LockAttemptResult + + breaker = MagicMock() + + async def swallowing_set(*args, **kwargs): + _record_swallowed_redis_failure(breaker, ConnectionError("redis down")) + return None + + cache = MagicMock() + cache.async_set_cache = swallowing_set + cache.async_get_cache = AsyncMock(return_value=None) + manager = PodLockManager(redis_cache=cache) + logger = MagicMock() + + with patch("litellm.integrations.prometheus.PrometheusLogger.get_instance", return_value=logger): + assert await manager.acquire_lock(cronjob_id="db_spend_update_job") is False + + logger.record_cronjob_lock_attempt.assert_called_once_with("db_spend_update_job", LockAttemptResult.ERROR) + + @pytest.mark.asyncio async def test_a_redis_failure_is_not_reported_as_losing_the_election(): """not_acquired means another pod won. An attempt that errored is a