Coverage for .venv/lib/python3.13/site-packages/litellm/proxy/db/db_transaction_queue/pod_lock_manager.py: 26%
85 statements
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-10 12:01 +0000
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-10 12:01 +0000
1import asyncio
2import json
3import logging
4from typing import TYPE_CHECKING, Any, Final
6from litellm._logging import verbose_proxy_logger
7from litellm._uuid import uuid
8from litellm.caching.redis_cache import RedisCache, log_redis_failure
9from litellm.constants import DEFAULT_CRON_JOB_LOCK_TTL_SECONDS
10from litellm.proxy.db.db_transaction_queue.base_update_queue import service_logger_obj
11from litellm.types.services import ServiceTypes
13if TYPE_CHECKING: 13 ↛ 14line 13 didn't jump to line 14 because the condition on line 13 was never true
14 ProxyLogging = Any
15else:
16 ProxyLogging = Any
19class PodLockManager:
20 """
21 Manager for acquiring and releasing locks for cron jobs using Redis.
23 Ensures that only one pod can run a cron job at a time.
24 """
26 _COMPARE_AND_DELETE_LOCK_SCRIPT = """
27if redis.call("get", KEYS[1]) == ARGV[1] then
28 return redis.call("del", KEYS[1])
29else
30 return 0
31end
32"""
34 def __init__(self, redis_cache: RedisCache | None = None):
35 self.pod_id = str(uuid.uuid4())
36 self.redis_cache = redis_cache
37 self._release_lock_script: Any | None = None
39 @staticmethod
40 def get_redis_lock_key(cronjob_id: str) -> str:
41 return f"cronjob_lock:{cronjob_id}"
43 async def acquire_lock(
44 self,
45 cronjob_id: str,
46 ttl: int | None = None,
47 allow_reentrant: bool = True,
48 ) -> bool | None:
49 """
50 Attempt to acquire the lock for a specific cron job using Redis.
51 Uses the SET command with NX and EX options to ensure atomicity.
53 Args:
54 cronjob_id: The ID of the cron job to lock
55 ttl: Optional custom TTL in seconds. Defaults to DEFAULT_CRON_JOB_LOCK_TTL_SECONDS.
56 Use a longer TTL for jobs that may take longer than the default 60s
57 (e.g. key rotation with many keys).
58 allow_reentrant: With the default True, a pod that already holds the lock
59 acquires it again (leader election semantics). Pass False when the live
60 lock marks work as already done for this window, so not even the holder
61 may redo it before the TTL expires.
62 """
63 if self.redis_cache is None:
64 verbose_proxy_logger.debug("redis_cache is None, skipping acquire_lock")
65 return None
66 try:
67 lock_ttl: Final = ttl or DEFAULT_CRON_JOB_LOCK_TTL_SECONDS
68 verbose_proxy_logger.debug(
69 "Pod %s attempting to acquire Redis lock for cronjob_id=%s (ttl=%ds)",
70 self.pod_id,
71 cronjob_id,
72 lock_ttl,
73 )
74 # Try to set the lock key with the pod_id as its value, only if it doesn't exist (NX)
75 # and with an expiration (EX) to avoid deadlocks.
76 lock_key: Final = PodLockManager.get_redis_lock_key(cronjob_id)
77 acquired: Final = await self.redis_cache.async_set_cache(
78 lock_key,
79 self.pod_id,
80 nx=True,
81 ttl=lock_ttl,
82 )
83 if acquired:
84 verbose_proxy_logger.info(
85 "Pod %s successfully acquired Redis lock for cronjob_id=%s",
86 self.pod_id,
87 cronjob_id,
88 )
90 return True
91 else:
92 # Check if the current pod already holds the lock
93 current_value = await self.redis_cache.async_get_cache(lock_key)
94 if current_value is not None:
95 if isinstance(current_value, bytes):
96 current_value = current_value.decode("utf-8")
97 if current_value == self.pod_id and allow_reentrant:
98 verbose_proxy_logger.info(
99 "Pod %s already holds the Redis lock for cronjob_id=%s",
100 self.pod_id,
101 cronjob_id,
102 )
103 self._emit_acquired_lock_event(cronjob_id, self.pod_id)
104 return True
105 verbose_proxy_logger.info(
106 "Pod %s could not acquire lock for cronjob_id=%s, held by pod %s.",
107 self.pod_id,
108 cronjob_id,
109 current_value,
110 )
111 return False
112 except Exception as e:
113 log_redis_failure(verbose_proxy_logger, logging.ERROR, f"Error acquiring Redis lock for {cronjob_id}", e)
114 return False
116 async def release_lock(
117 self,
118 cronjob_id: str,
119 ):
120 """
121 Release the lock if the current pod holds it.
122 Uses an atomic Lua compare-and-delete to prevent TOCTOU races where a
123 stale owner could delete a newly reacquired lock.
124 Falls back to GET + DEL for cache implementations that don't support
125 script registration.
126 """
127 if self.redis_cache is None:
128 verbose_proxy_logger.debug("redis_cache is None, skipping release_lock")
129 return
130 try:
131 verbose_proxy_logger.debug(
132 "Pod %s attempting to release Redis lock for cronjob_id=%s",
133 self.pod_id,
134 cronjob_id,
135 )
136 lock_key: Final = PodLockManager.get_redis_lock_key(cronjob_id)
137 result: Final = await self._compare_and_delete_lock(lock_key=lock_key)
138 if result == 1:
139 verbose_proxy_logger.info(
140 "Pod %s successfully released Redis lock for cronjob_id=%s",
141 self.pod_id,
142 cronjob_id,
143 )
144 self._emit_released_lock_event(
145 cronjob_id=cronjob_id,
146 pod_id=self.pod_id,
147 )
148 else:
149 verbose_proxy_logger.debug(
150 "Pod %s failed to release Redis lock for cronjob_id=%s (lock missing or held by another pod)",
151 self.pod_id,
152 cronjob_id,
153 )
154 except Exception as e:
155 log_redis_failure(verbose_proxy_logger, logging.ERROR, f"Error releasing Redis lock for {cronjob_id}", e)
157 async def _compare_and_delete_lock(self, lock_key: str) -> int:
158 """
159 Atomically delete lock key only if current pod owns it.
161 Falls back to get/delete for non-RedisCache implementations that do not
162 expose Lua script registration.
163 """
164 script_register: Final = getattr(self.redis_cache, "async_register_script", None)
165 if callable(script_register):
166 try:
167 if self._release_lock_script is None:
168 self._release_lock_script = script_register(self._COMPARE_AND_DELETE_LOCK_SCRIPT)
169 # acquire_lock stores the pod_id via async_set_cache, which
170 # JSON-encodes the value; compare against the same encoding so
171 # the Lua equality check matches and the lock is released
172 result = await self._release_lock_script(keys=[lock_key], args=[json.dumps(self.pod_id)])
173 return int(result or 0)
174 except Exception:
175 # Lua execution failed (e.g. Redis restart cleared loaded scripts,
176 # or scripting is disabled). Reset cached script handle and fall
177 # through to the GET + DEL fallback so the lock is still released.
178 self._release_lock_script = None
179 verbose_proxy_logger.warning(
180 "Lua compare-and-delete failed for lock_key=%s, falling back to GET+DEL",
181 lock_key,
182 )
184 current_value = await self.redis_cache.async_get_cache(lock_key)
185 if isinstance(current_value, bytes):
186 current_value = current_value.decode("utf-8")
187 if current_value != self.pod_id:
188 return 0
189 result = await self.redis_cache.async_delete_cache(lock_key)
190 return int(result or 0)
192 @staticmethod
193 def _emit_acquired_lock_event(cronjob_id: str, pod_id: str):
194 asyncio.create_task(
195 service_logger_obj.async_service_success_hook(
196 service=ServiceTypes.POD_LOCK_MANAGER,
197 duration=DEFAULT_CRON_JOB_LOCK_TTL_SECONDS,
198 call_type="_emit_acquired_lock_event",
199 event_metadata={
200 "gauge_labels": f"{cronjob_id}:{pod_id}",
201 "gauge_value": 1,
202 },
203 )
204 )
206 @staticmethod
207 def _emit_released_lock_event(cronjob_id: str, pod_id: str):
208 asyncio.create_task(
209 service_logger_obj.async_service_success_hook(
210 service=ServiceTypes.POD_LOCK_MANAGER,
211 duration=DEFAULT_CRON_JOB_LOCK_TTL_SECONDS,
212 call_type="_emit_released_lock_event",
213 event_metadata={
214 "gauge_labels": f"{cronjob_id}:{pod_id}",
215 "gauge_value": 0,
216 },
217 )
218 )