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
« 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
6from litellm.constants import (
7 PROXY_DB_LOOKUP_DEADLINE_SECONDS,
8 PROXY_DB_LOOKUP_MAX_CONCURRENCY,
9)
11LookupT = TypeVar("LookupT")
14class LoopBoundSemaphore:
15 __slots__ = ("_loop", "_semaphore", "_value")
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
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
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
37class DBLookupStallTracker:
38 __slots__ = ("_clock", "_last_hit")
40 def __init__(self, clock: Callable[[], float] = time.monotonic) -> None:
41 self._clock: Final = clock
42 self._last_hit: float | None = None
44 def record_hit(self) -> None:
45 self._last_hit = self._clock()
47 def clear(self) -> None:
48 self._last_hit = None
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
56db_lookup_gate: Final = LoopBoundSemaphore(PROXY_DB_LOOKUP_MAX_CONCURRENCY)
57db_lookup_stall_tracker: Final = DBLookupStallTracker()
60def _consume_abandoned_lookup(task: asyncio.Future[LookupT]) -> None:
61 if not task.cancelled():
62 task.exception()
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