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