mirror of
https://github.com/BerriAI/litellm.git
synced 2026-10-07 02:59:05 +00:00
fix(proxy): don't retry the additive cancel-path reconcile on failure
The cancel-path reconcile applies an INCRBYFLOAT refund. If Redis applies it but the pipeline response times out (and the follow-up counter delete also fails), the reservation's applied_adjustment is never stamped, so retrying the reconcile subtracts the same refund a second time and under-counts spend shared with concurrent requests. Even when that delete succeeds, a concurrent request can reseed the counter before the retry, which then refunds against a counter that no longer holds this reservation. Drop the retry and go straight to invalidating the reserved counters (a non-additive delete, so the next read reseeds from the DB) and finalizing the reservation, the same fallback release_or_invalidate_budget_reservation uses on the non-cancel path. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
parent
168adaba87
commit
e91f3fb5da
2 changed files with 47 additions and 81 deletions
|
|
@ -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:
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
||||
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue