fix(proxy): keep the reservation lease alive across a failed Redis EXPIRE

Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
This commit is contained in:
yassin 2026-09-08 23:24:21 +00:00
parent a558a0b6a9
commit 3cd4a768ae
3 changed files with 33 additions and 9 deletions

View file

@ -1799,16 +1799,11 @@ class RedisCache(BaseCache):
@_redis_circuit_breaker_guard
async def async_refresh_ttl(self, key: str, ttl: int | None = None) -> bool:
"""EXPIRE an existing key without touching its value. False when the key is absent or Redis failed."""
"""EXPIRE an existing key without touching its value. False when the key is absent."""
_used_ttl: Final = self.get_ttl(ttl=ttl)
if _used_ttl is None:
return False
try:
return await self._async_commands().expire(self.check_and_fix_namespace(key=key), _used_ttl)
except Exception as e:
verbose_logger.debug("Redis EXPIRE Error: %s", e)
_record_swallowed_redis_failure(self._circuit_breaker, e)
return False
return await self._async_commands().expire(self.check_and_fix_namespace(key=key), _used_ttl)
@_redis_circuit_breaker_guard
async def async_rpush(

View file

@ -3225,7 +3225,11 @@ async def increment_spend_counter(counter_key: str, increment: float):
async def refresh_spend_counter_ttl(counter_key: str) -> bool:
if spend_counter_cache.redis_cache is None:
return False
return await spend_counter_cache.redis_cache.async_refresh_ttl(key=counter_key)
try:
return await spend_counter_cache.redis_cache.async_refresh_ttl(key=counter_key)
except Exception as e:
verbose_proxy_logger.debug("spend counter TTL refresh skipped for %s: %s", counter_key, e)
return False
async def _increment_spend_counter_cache(counter_key: str, increment: float):

View file

@ -2202,11 +2202,13 @@ async def test_release_non_numeric_counter_reseeds_from_db(spend_counter_state):
class _ExpiringRedisCache:
"""In-memory stand-in for RedisCache with real wall-clock key expiry."""
def __init__(self, default_ttl: float = 60.0) -> None:
def __init__(self, default_ttl: float = 60.0, fail_first_refresh: bool = False) -> None:
self.default_ttl = default_ttl
self.store: dict[str, float] = {}
self.expires_at: dict[str, float] = {}
self.refresh_attempts = 0
self.refresh_count = 0
self.fail_first_refresh = fail_first_refresh
def _evict_expired(self, key: str) -> None:
if self.expires_at.get(key, float("inf")) <= time.monotonic():
@ -2239,6 +2241,9 @@ class _ExpiringRedisCache:
self.expires_at.pop(key, None)
async def async_refresh_ttl(self, key: str, ttl: int | None = None) -> bool:
self.refresh_attempts += 1
if self.fail_first_refresh and self.refresh_attempts == 1:
raise ConnectionError("Redis circuit breaker is open")
self._evict_expired(key)
if key not in self.store:
return False
@ -2279,6 +2284,26 @@ async def test_reservation_survives_redis_counter_ttl_while_request_in_flight(
assert await redis_cache.async_get_cache(key=counter_key) is None
@pytest.mark.asyncio
async def test_reservation_lease_keeps_renewing_after_transient_redis_failure(
spend_counter_state,
):
"""One failed EXPIRE (Redis blip, open circuit breaker) must not end renewal for the
rest of the request."""
counter_cache, key_cache = spend_counter_state
redis_cache = _ExpiringRedisCache(default_ttl=0.2, fail_first_refresh=True)
counter_cache.redis_cache = redis_cache
proxy_logging_obj = ProxyLogging(user_api_key_cache=key_cache)
valid_token = UserAPIKeyAuth(token="key-lease-blip", spend=0.0, max_budget=1.0)
reservation = await _reserve(valid_token, 0.6, key_cache, proxy_logging_obj)
assert reservation is not None
await asyncio.sleep(0.5)
assert redis_cache.refresh_attempts >= 3
await release_budget_reservation(reservation)
@pytest.mark.asyncio
async def test_reconcile_after_redis_counter_expiry_keeps_request_cost_enforced(
spend_counter_state,