diff --git a/litellm/proxy/spend_tracking/budget_reservation.py b/litellm/proxy/spend_tracking/budget_reservation.py index e1cbd33ab3d..ee50153873e 100644 --- a/litellm/proxy/spend_tracking/budget_reservation.py +++ b/litellm/proxy/spend_tracking/budget_reservation.py @@ -57,11 +57,6 @@ class _BudgetCounter: window_start: datetime | None = None -_RELEASE_ON_CANCEL_RECONCILE_MAX_ATTEMPTS: Final = 2 -"""Retries of the per-reservation reconcile before falling back to the -counter-deleting invalidation, which a concurrent reservation also shares.""" - - _COUNTER_ENTITY_TYPES: Final[Mapping[str, str]] = { "Key": Litellm_EntityType.KEY.value, "Team": Litellm_EntityType.TEAM.value, @@ -407,11 +402,11 @@ async def release_budget_reservation_on_cancel( the provider. asyncio.shield keeps this running through the surrounding cancellation. - A reconcile failure is retried before falling back to - invalidate_budget_reservation_counters, since unlike that fallback, reconcile - only ever adjusts this reservation's own contribution to each counter and can't - clobber a concurrent reservation sharing it. Mirrors - release_or_invalidate_budget_reservation's fallback on the non-cancel path. + A failed reconcile is not retried: the reconcile is an additive INCRBYFLOAT, and + a failure (e.g. a timeout after Redis applied it) leaves it unknown whether the + refund landed, so re-applying it could refund twice. Instead the reserved + counters are dropped so the next read reseeds from the DB, the same fallback + release_or_invalidate_budget_reservation uses on the non-cancel path. """ if not budget_reservation or budget_reservation.get("finalized") is True: return @@ -422,14 +417,8 @@ async def release_budget_reservation_on_cancel( ) except asyncio.CancelledError: pass # a second cancellation while shielded; the reconcile keeps running detached regardless - except Exception: # noqa: BLE001 # a reconcile failure must not pin the counter; retry, then drop it directly - verbose_proxy_logger.exception( - "Failed to reconcile budget reservation on cancel; retrying before invalidating reserved counters" - ) - if await asyncio.shield( - _retry_reconcile_reservation_on_cancel(budget_reservation=budget_reservation, incurred_cost=incurred_cost) - ): - return + except Exception: # noqa: BLE001 # a reconcile failure must not pin the counter; drop it directly instead + verbose_proxy_logger.exception("Failed to reconcile budget reservation on cancel; invalidating counters") try: await invalidate_budget_reservation_counters(budget_reservation=budget_reservation) except Exception: # noqa: BLE001 # nothing left to try; the finalized stamp below keeps it from being reprocessed @@ -440,28 +429,6 @@ async def release_budget_reservation_on_cancel( budget_reservation["finalized"] = True -async def _retry_reconcile_reservation_on_cancel( - budget_reservation: dict, # mutable-ok: reconcile_budget_reservation stamps finalized on success - incurred_cost: float, -) -> bool: - """Retry the reconcile that just failed once. Safe to retry since it only - ever adjusts this reservation's own contribution, unlike the destructive - invalidation the caller falls back to once every retry fails.""" - for attempt in range(_RELEASE_ON_CANCEL_RECONCILE_MAX_ATTEMPTS): - try: - await reconcile_budget_reservation(budget_reservation=budget_reservation, actual_cost=incurred_cost) - return True - except Exception: # noqa: BLE001 # any reconcile failure is worth a retry here, not just specific ones - is_last_attempt = attempt == _RELEASE_ON_CANCEL_RECONCILE_MAX_ATTEMPTS - 1 - verbose_proxy_logger.warning( - "Retry %d/%d to reconcile budget reservation on cancel failed", - attempt + 1, - _RELEASE_ON_CANCEL_RECONCILE_MAX_ATTEMPTS, - exc_info=not is_last_attempt, - ) - return False - - async def invalidate_budget_reservation_counters( budget_reservation: dict | None, ) -> None: diff --git a/tests/test_litellm/proxy/test_budget_reservation.py b/tests/test_litellm/proxy/test_budget_reservation.py index eca6a2ea590..529cda6c83a 100644 --- a/tests/test_litellm/proxy/test_budget_reservation.py +++ b/tests/test_litellm/proxy/test_budget_reservation.py @@ -2485,20 +2485,31 @@ class _TeamMembershipFloorDb: class _FlakyPipelineRedisCache(_ExpiringRedisCache): """_ExpiringRedisCache, but the first ``fail_first_n`` calls to - async_increment_pipeline raise instead of applying, so a reconcile retry - recovering from a transient Redis failure is what's under test.""" + async_increment_pipeline time out. With ``apply_before_failing`` the + increments land before the timeout surfaces, the ambiguous case where the + caller can't tell whether its write applied; with ``fail_deletes`` the + counter delete that would otherwise clean that up fails too.""" - def __init__(self, fail_first_n: int = 0) -> None: + def __init__(self, fail_first_n: int = 0, apply_before_failing: bool = False, fail_deletes: bool = False) -> None: super().__init__() self.fail_first_n = fail_first_n + self.apply_before_failing = apply_before_failing + self.fail_deletes = fail_deletes self.pipeline_calls = 0 async def async_increment_pipeline(self, increment_list, **kwargs): self.pipeline_calls += 1 if self.pipeline_calls <= self.fail_first_n: - raise RuntimeError("redis down") + if self.apply_before_failing: + await super().async_increment_pipeline(increment_list, **kwargs) + raise TimeoutError("redis timeout") return await super().async_increment_pipeline(increment_list, **kwargs) + async def async_delete_cache(self, key: str, *args: object, **kwargs: object) -> None: + if self.fail_deletes: + raise ConnectionError("redis unreachable") + await super().async_delete_cache(key, *args, **kwargs) + @pytest.mark.asyncio async def test_reconcile_after_redis_counter_expiry_keeps_request_cost_enforced( @@ -3183,8 +3194,8 @@ async def test_release_budget_reservation_on_cancel_swallows_a_second_cancellati @pytest.mark.asyncio -async def test_release_budget_reservation_on_cancel_swallows_invalidate_failure_after_every_retry_fails(): - # If both the reconcile retries and the invalidate fallback fail (e.g. a persistent +async def test_release_budget_reservation_on_cancel_swallows_invalidate_failure_after_reconcile_fails(): + # If both the reconcile and the invalidate fallback fail (e.g. a persistent # outage), there is nothing left to try: the failure must be logged and swallowed, # not propagated, and the reservation still ends up finalized so it is not reprocessed. reservation = { @@ -3207,17 +3218,16 @@ async def test_release_budget_reservation_on_cancel_swallows_invalidate_failure_ @pytest.mark.asyncio -async def test_release_budget_reservation_on_cancel_invalidates_counter_when_reconcile_persistently_fails( +async def test_release_budget_reservation_on_cancel_invalidates_counter_when_reconcile_fails( spend_counter_state, ): """ - Regression for #30460 Path 1: if reconcile keeps failing on the cancel path - (e.g. a persistent Redis outage) across every retry, the pre-charge must - not be left stuck in the counter with nothing left to correct it. This - mirrors what release_or_invalidate_budget_reservation already does on the - non-cancel release path: fall back to invalidate_budget_reservation_counters - so the next read reseeds from the DB instead of enforcing the stale - reservation forever. + Regression for #30460 Path 1: if reconcile fails on the cancel path (e.g. a + Redis outage), the pre-charge must not be left stuck in the counter with + nothing left to correct it. This mirrors what + release_or_invalidate_budget_reservation already does on the non-cancel + release path: drop the counter so the next read reseeds from the DB instead + of enforcing the stale reservation until the TTL. """ counter_cache, _key_cache = spend_counter_state counter_key = "spend:key:key-cancel-redis-down" @@ -3234,40 +3244,30 @@ async def test_release_budget_reservation_on_cancel_invalidates_counter_when_rec await release_budget_reservation_on_cancel(reservation) + assert counter_key not in counter_cache.redis_cache.store assert counter_cache.in_memory_cache.get_cache(key=counter_key) is None assert reservation["finalized"] is True - # Only the first retry reaches the pipeline: its own failure already - # invalidates the counter, so later retries take the (here also - # unavailable, with no real DB) reseed path instead -- both exhausted - # before falling back to invalidate_budget_reservation_counters. assert counter_cache.redis_cache.pipeline_calls == 1 @pytest.mark.asyncio -async def test_release_budget_reservation_on_cancel_retries_before_invalidating( +async def test_release_budget_reservation_on_cancel_does_not_refund_twice_when_reconcile_times_out_after_applying( spend_counter_state, ): """ - Regression for the cancellation-fallback race (veria-ai finding on - budget_reservation.py): invalidate_budget_reservation_counters deletes the - whole aggregate counter unconditionally, which would also erase a - concurrent request's reservation or recorded spend sharing that same - key/user/team counter, not just this reservation's own contribution. A - reconcile failure that clears up on retry must settle through the - ordinary reconcile machinery instead of falling back to that unconditional - delete: here, the first attempt's own failure already invalidates the - counter (existing increment_spend_counters_pipeline cleanup), so the retry - settles it by reseeding the DB floor and adding this reservation's actual - cost, the same recovery path an expired counter already takes elsewhere in - this file, rather than being deleted a second time with nothing settled. + Regression for the duplicate-refund race (veria-ai finding on + budget_reservation.py): the reconcile's INCRBYFLOAT lands in Redis but the + pipeline response times out, and the counter delete meant to clean that up + fails too, so the counter still holds the already-refunded value. Re-applying + the reconcile would subtract this reservation's refund a second time and eat + into a concurrent request's spend sharing the counter. The refund must land + exactly once. """ - import litellm.proxy.proxy_server as ps - counter_cache, _key_cache = spend_counter_state - counter_key = "spend:key:key-cancel-retry-recovers" - counter_cache.redis_cache = _FlakyPipelineRedisCache(fail_first_n=1) + counter_key = "spend:key:key-cancel-timeout-after-apply" + counter_cache.redis_cache = _FlakyPipelineRedisCache(fail_first_n=99, apply_before_failing=True, fail_deletes=True) + # 2.0 of a concurrent request's spend plus this request's 3.0 reservation counter_cache.redis_cache.store[counter_key] = 5.0 - counter_cache.in_memory_cache.set_cache(key=counter_key, value=5.0) reservation = { "reserved_cost": 3.0, @@ -3276,13 +3276,12 @@ async def test_release_budget_reservation_on_cancel_retries_before_invalidating( "entries": [{"counter_key": counter_key, "reserved_cost": 3.0, "applied_adjustment": 0.0}], } - with patch.object( # test-quality-ok: the reseed reads the DB floor through a Prisma client the test has no seam for - ps.SpendCounterReseed, "from_db", AsyncMock(return_value=0.2) - ): - await release_budget_reservation_on_cancel(reservation) + await release_budget_reservation_on_cancel(reservation) + await release_budget_reservation_on_cancel(reservation) + # 5.0 - (3.0 - 0.5), applied once; a second application would leave 0.0 + assert counter_cache.redis_cache.store[counter_key] == pytest.approx(2.5) assert counter_cache.redis_cache.pipeline_calls == 1 - assert counter_cache.in_memory_cache.get_cache(key=counter_key) == pytest.approx(0.7) assert reservation["finalized"] is True