Coverage for .venv/lib/python3.13/site-packages/litellm/proxy/_experimental/mcp_server/outbound_credentials/redis_distributed_lock.py: 50%
42 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
1"""Concrete ``DistributedLock`` over a Redis client: ``SET NX PX`` / owner-only renew / delete.
3The cross-replica lock the ``RedisRefreshCoordinator`` elects refreshers with. ``acquire`` is an
4atomic ``SET key token NX PX ttl`` (only the first caller wins; the entry self-expires so a crashed
5holder can't wedge refresh). ``extend`` renews the lease only when the token still matches, and
6``release`` deletes the key only when it still holds this caller's token, so a holder whose lock already
7PX-expired and was re-acquired by another worker cannot delete the new holder's lock. ``is_held`` is
8``EXISTS``. Every key is run through the injected ``namespace_key`` before it reaches Redis, so lock
9keys carry the same namespace as cache keys and cannot collide with another deployment sharing Redis.
11The Redis client is injected (in production the async client from LiteLLM's ``RedisCache``), so the
12lock is unit-testable with a fake. A transport error on ``acquire`` returns ``LockAcquisition.ERROR`` -
13distinct from ``HELD`` - so the coordinator refreshes anyway instead of mistaking a dead backend for a
14busy holder; a Redis blip degrades to an extra refresh, never a stale bearer.
15"""
17from __future__ import annotations
19from collections.abc import Callable
20from dataclasses import KW_ONLY, dataclass
21from typing import Final, Protocol
23from litellm._logging import verbose_logger
24from litellm.proxy._experimental.mcp_server.outbound_credentials.redis_refresh_coordinator import (
25 LockAcquisition,
26)
28# Delete the key only if it still holds this caller's token, so a holder whose lock already expired
29# (PX) and was re-acquired by another worker cannot delete the new holder's lock.
30_RELEASE_IF_OWNER = "if redis.call('get', KEYS[1]) == ARGV[1] then return redis.call('del', KEYS[1]) else return 0 end"
31_EXTEND_IF_OWNER: Final = (
32 "if redis.call('get', KEYS[1]) == ARGV[1] then return redis.call('pexpire', KEYS[1], ARGV[2]) else return 0 end"
33)
36class RedisCommands(Protocol):
37 """The slice of the async Redis client this lock needs."""
39 async def set(self, name: str, value: str, *, nx: bool = False, px: int | None = None) -> object | None: ... 39 ↛ exitline 39 didn't return from function 'set' because
41 async def eval(self, script: str, numkeys: int, *keys_and_args: str) -> object: ... 41 ↛ exitline 41 didn't return from function 'eval' because
43 async def exists(self, *names: str) -> int: ... 43 ↛ exitline 43 didn't return from function 'exists' because
46@dataclass(frozen=True, slots=True)
47class RedisDistributedLock:
48 client: RedisCommands
49 _: KW_ONLY
50 namespace_key: Callable[[str], str] = lambda key: key
52 async def acquire(self, key: str, token: str, ttl_seconds: float) -> LockAcquisition:
53 try:
54 result: Final = await self.client.set(self.namespace_key(key), token, nx=True, px=int(ttl_seconds * 1000))
55 # Degrade on any Redis client error: redis.exceptions narrows only via an import that
56 # is Unknown under basedpyright, and the lock must never crash the resolve path.
57 except Exception as exc: # noqa: BLE001
58 verbose_logger.warning("RedisDistributedLock.acquire failed: %s", exc)
59 return LockAcquisition.ERROR
60 return LockAcquisition.ACQUIRED if result is not None else LockAcquisition.HELD
62 async def extend(self, key: str, token: str, ttl_seconds: float) -> bool:
63 try:
64 result: Final = await self.client.eval(
65 _EXTEND_IF_OWNER,
66 1,
67 self.namespace_key(key),
68 token,
69 str(int(ttl_seconds * 1000)),
70 )
71 # Degrade on any Redis client error: redis.exceptions narrows only via an import that
72 # is Unknown under basedpyright, and the lock must never crash the resolve path.
73 except Exception as exc: # noqa: BLE001
74 verbose_logger.warning("RedisDistributedLock.extend failed: %s", exc)
75 return False
76 return result == 1
78 async def release(self, key: str, token: str) -> None:
79 try:
80 await self.client.eval(_RELEASE_IF_OWNER, 1, self.namespace_key(key), token)
81 # Degrade on any Redis client error: redis.exceptions narrows only via an import that
82 # is Unknown under basedpyright, and the lock must never crash the resolve path.
83 except Exception as exc: # noqa: BLE001
84 verbose_logger.warning("RedisDistributedLock.release failed: %s", exc)
86 async def is_held(self, key: str) -> bool:
87 try:
88 return await self.client.exists(self.namespace_key(key)) > 0
89 # Degrade on any Redis client error: redis.exceptions narrows only via an import that
90 # is Unknown under basedpyright, and the lock must never crash the resolve path.
91 except Exception as exc: # noqa: BLE001
92 # On error, report "not held" so a waiter stops waiting and re-reads rather than blocking.
93 verbose_logger.warning("RedisDistributedLock.is_held failed: %s", exc)
94 return False