Coverage for .venv/lib/python3.13/site-packages/litellm/proxy/db/db_lookup_gate.py: 68%

61 statements  

« prev     ^ index     » next       coverage.py v7.15.2, created at 2026-10-10 12:01 +0000

1import asyncio 

2import time 

3from collections.abc import Awaitable, Callable 

4from typing import Final, TypeVar 

5 

6from litellm.constants import ( 

7 PROXY_DB_LOOKUP_DEADLINE_SECONDS, 

8 PROXY_DB_LOOKUP_MAX_CONCURRENCY, 

9) 

10 

11LookupT = TypeVar("LookupT") 

12 

13 

14class LoopBoundSemaphore: 

15 __slots__ = ("_loop", "_semaphore", "_value") 

16 

17 def __init__(self, value: int) -> None: 

18 self._value: Final = value 

19 self._loop: asyncio.AbstractEventLoop | None = None 

20 self._semaphore: asyncio.Semaphore | None = None 

21 

22 def current(self) -> asyncio.Semaphore: 

23 loop: Final = asyncio.get_running_loop() 

24 if self._semaphore is None or self._loop is not loop: 

25 self._semaphore = asyncio.Semaphore(self._value) 

26 self._loop = loop 

27 return self._semaphore 

28 

29 

30class DBLookupDeadlineExceeded(asyncio.TimeoutError): 

31 def __init__(self, lookup: str, deadline_seconds: float) -> None: 

32 super().__init__(f"{lookup} lookup did not answer within {deadline_seconds:g}s") 

33 self.lookup: Final = lookup 

34 self.deadline_seconds: Final = deadline_seconds 

35 

36 

37class DBLookupStallTracker: 

38 __slots__ = ("_clock", "_last_hit") 

39 

40 def __init__(self, clock: Callable[[], float] = time.monotonic) -> None: 

41 self._clock: Final = clock 

42 self._last_hit: float | None = None 

43 

44 def record_hit(self) -> None: 

45 self._last_hit = self._clock() 

46 

47 def clear(self) -> None: 

48 self._last_hit = None 

49 

50 def stalled_within(self, window_seconds: float) -> bool: 

51 if self._last_hit is None: 51 ↛ 53line 51 didn't jump to line 53 because the condition on line 51 was always true

52 return False 

53 return self._clock() - self._last_hit < window_seconds 

54 

55 

56db_lookup_gate: Final = LoopBoundSemaphore(PROXY_DB_LOOKUP_MAX_CONCURRENCY) 

57db_lookup_stall_tracker: Final = DBLookupStallTracker() 

58 

59 

60def _consume_abandoned_lookup(task: asyncio.Future[LookupT]) -> None: 

61 if not task.cancelled(): 

62 task.exception() 

63 

64 

65async def bounded_db_lookup( 

66 lookup: Awaitable[LookupT], 

67 *, 

68 name: str, 

69 deadline_seconds: float | None = None, 

70 tracker: DBLookupStallTracker = db_lookup_stall_tracker, 

71) -> LookupT: 

72 timeout: Final = PROXY_DB_LOOKUP_DEADLINE_SECONDS if deadline_seconds is None else deadline_seconds 

73 task: Final = asyncio.ensure_future(lookup) 

74 try: 

75 done, _ = await asyncio.wait({task}, timeout=timeout) 

76 except asyncio.CancelledError: 

77 task.cancel() 

78 raise 

79 if task not in done: 79 ↛ 80line 79 didn't jump to line 80 because the condition on line 79 was never true

80 task.cancel() 

81 task.add_done_callback(_consume_abandoned_lookup) 

82 tracker.record_hit() 

83 raise DBLookupDeadlineExceeded(name, timeout) 

84 try: 

85 return task.result() 

86 except DBLookupDeadlineExceeded: 

87 raise 

88 except asyncio.TimeoutError as e: 

89 tracker.record_hit() 

90 raise DBLookupDeadlineExceeded(name, timeout) from e