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

1"""Concrete ``DistributedLock`` over a Redis client: ``SET NX PX`` / owner-only renew / delete. 

2 

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. 

10 

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""" 

16 

17from __future__ import annotations 

18 

19from collections.abc import Callable 

20from dataclasses import KW_ONLY, dataclass 

21from typing import Final, Protocol 

22 

23from litellm._logging import verbose_logger 

24from litellm.proxy._experimental.mcp_server.outbound_credentials.redis_refresh_coordinator import ( 

25 LockAcquisition, 

26) 

27 

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) 

34 

35 

36class RedisCommands(Protocol): 

37 """The slice of the async Redis client this lock needs.""" 

38 

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

40 

41 async def eval(self, script: str, numkeys: int, *keys_and_args: str) -> object: ... 41 ↛ exitline 41 didn't return from function 'eval' because

42 

43 async def exists(self, *names: str) -> int: ... 43 ↛ exitline 43 didn't return from function 'exists' because

44 

45 

46@dataclass(frozen=True, slots=True) 

47class RedisDistributedLock: 

48 client: RedisCommands 

49 _: KW_ONLY 

50 namespace_key: Callable[[str], str] = lambda key: key 

51 

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 

61 

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 

77 

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) 

85 

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