From 187b205b34a691a67325aafa42b601374c0fa32b Mon Sep 17 00:00:00 2001 From: Yassin Kortam Date: Wed, 17 Jun 2026 17:01:04 -0700 Subject: [PATCH] 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 "". 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 --- .../db_transaction_queue/pod_lock_manager.py | 6 +- .../test_pod_lock_manager.py | 74 ++++++++++++++++++- 2 files changed, 78 insertions(+), 2 deletions(-) diff --git a/litellm/proxy/db/db_transaction_queue/pod_lock_manager.py b/litellm/proxy/db/db_transaction_queue/pod_lock_manager.py index 5e0ddef9eaa..2cbc0646567 100644 --- a/litellm/proxy/db/db_transaction_queue/pod_lock_manager.py +++ b/litellm/proxy/db/db_transaction_queue/pod_lock_manager.py @@ -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: diff --git a/tests/test_litellm/proxy/db/db_transaction_queue/test_pod_lock_manager.py b/tests/test_litellm/proxy/db/db_transaction_queue/test_pod_lock_manager.py index 27fe9202276..f2745052faa 100644 --- a/tests/test_litellm/proxy/db/db_transaction_queue/test_pod_lock_manager.py +++ b/tests/test_litellm/proxy/db/db_transaction_queue/test_pod_lock_manager.py @@ -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