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.
This commit is contained in:
Yucheng Zhu 2026-08-19 14:35:34 -07:00
parent 2ba08e68c1
commit b544023923
3 changed files with 50 additions and 1 deletions

View file

@ -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,

View file

@ -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

View file

@ -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