mirror of
https://github.com/BerriAI/litellm.git
synced 2026-10-09 03:18:44 +00:00
fix(pod_lock): release cron lock by matching async_set_cache JSON encoding (#30600)
acquire_lock stores the pod_id through async_set_cache, which JSON-encodes the value, so Redis holds the quoted string "<pod_id>". release_lock's Lua compare-and-delete compared the raw pod_id, so the equality check never matched and the lock was never deleted; it only cleared on TTL expiry. That stalled the spend-update drain whenever the leader pod restarted, letting the litellm_daily_*_spend_update_buffer lists grow unbounded in Redis. Compare against json.dumps(self.pod_id) so the release matches the stored value. The GET+DEL fallback already round-trips through async_get_cache and is unaffected. Co-authored-by: Claude <noreply@anthropic.com>
This commit is contained in:
parent
ba29657d09
commit
187b205b34
2 changed files with 78 additions and 2 deletions
|
|
@ -1,4 +1,5 @@
|
|||
import asyncio
|
||||
import json
|
||||
from litellm._uuid import uuid
|
||||
from typing import TYPE_CHECKING, Any, Optional
|
||||
|
||||
|
|
@ -167,8 +168,11 @@ end
|
|||
self._release_lock_script = script_register(
|
||||
self._COMPARE_AND_DELETE_LOCK_SCRIPT
|
||||
)
|
||||
# acquire_lock stores the pod_id via async_set_cache, which
|
||||
# JSON-encodes the value; compare against the same encoding so
|
||||
# the Lua equality check matches and the lock is released
|
||||
result = await self._release_lock_script(
|
||||
keys=[lock_key], args=[self.pod_id]
|
||||
keys=[lock_key], args=[json.dumps(self.pod_id)]
|
||||
)
|
||||
return int(result or 0)
|
||||
except Exception:
|
||||
|
|
|
|||
|
|
@ -327,7 +327,7 @@ async def test_release_lock_uses_atomic_compare_delete_script_when_available(
|
|||
PodLockManager._COMPARE_AND_DELETE_LOCK_SCRIPT
|
||||
)
|
||||
script_callable.assert_called_once_with(
|
||||
keys=[lock_key], args=[pod_lock_manager.pod_id]
|
||||
keys=[lock_key], args=[json.dumps(pod_lock_manager.pod_id)]
|
||||
)
|
||||
mock_redis.async_get_cache.assert_not_called()
|
||||
mock_redis.async_delete_cache.assert_not_called()
|
||||
|
|
@ -364,6 +364,78 @@ async def test_release_lock_lua_path_emits_released_event(pod_lock_manager, mock
|
|||
)
|
||||
|
||||
|
||||
class FakeRedisLockStore:
|
||||
"""
|
||||
Minimal stand-in that mirrors how RedisCache actually stores values:
|
||||
async_set_cache JSON-encodes the value, and the compare-and-delete Lua
|
||||
script compares against the raw stored bytes. This is what exposes the
|
||||
quoted-vs-raw mismatch that a value-agnostic mock cannot catch.
|
||||
"""
|
||||
|
||||
def __init__(self):
|
||||
self.store: dict = {}
|
||||
|
||||
async def async_set_cache(self, key, value, nx=False, ttl=None, **kwargs):
|
||||
if nx and key in self.store:
|
||||
return None
|
||||
self.store[key] = json.dumps(value)
|
||||
return True
|
||||
|
||||
async def async_get_cache(self, key, **kwargs):
|
||||
raw = self.store.get(key)
|
||||
return json.loads(raw) if raw is not None else None
|
||||
|
||||
async def async_delete_cache(self, key, **kwargs):
|
||||
return 1 if self.store.pop(key, None) is not None else 0
|
||||
|
||||
def async_register_script(self, script):
|
||||
async def _run(keys, args):
|
||||
key = keys[0]
|
||||
if self.store.get(key) == args[0]:
|
||||
del self.store[key]
|
||||
return 1
|
||||
return 0
|
||||
|
||||
return _run
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_release_lock_deletes_lock_held_by_same_pod():
|
||||
"""
|
||||
Regression: acquire_lock stores the pod_id JSON-encoded, so release_lock's
|
||||
Lua compare-and-delete must use the same encoding or the comparison never
|
||||
matches and the lock leaks until its TTL expires (stalling the spend-update
|
||||
drain and growing the Redis transaction buffers).
|
||||
"""
|
||||
redis = FakeRedisLockStore()
|
||||
pod = PodLockManager(redis_cache=redis)
|
||||
lock_key = PodLockManager.get_redis_lock_key("db_spend_update_job")
|
||||
|
||||
acquired = await pod.acquire_lock(cronjob_id="db_spend_update_job")
|
||||
assert acquired is True
|
||||
assert lock_key in redis.store
|
||||
|
||||
await pod.release_lock(cronjob_id="db_spend_update_job")
|
||||
assert lock_key not in redis.store
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_release_lock_preserves_lock_held_by_other_pod():
|
||||
"""
|
||||
A pod must not release a lock currently held by a different pod, even with
|
||||
the encoding fix in place.
|
||||
"""
|
||||
redis = FakeRedisLockStore()
|
||||
holder = PodLockManager(redis_cache=redis)
|
||||
other = PodLockManager(redis_cache=redis)
|
||||
lock_key = PodLockManager.get_redis_lock_key("db_spend_update_job")
|
||||
|
||||
assert await holder.acquire_lock(cronjob_id="db_spend_update_job") is True
|
||||
|
||||
await other.release_lock(cronjob_id="db_spend_update_job")
|
||||
assert redis.store.get(lock_key) == json.dumps(holder.pod_id)
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_release_lock_falls_back_to_get_del_when_lua_execution_fails(
|
||||
pod_lock_manager, mock_redis
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue