Coverage for .venv/lib/python3.13/site-packages/litellm/proxy/db/db_transaction_queue/base_update_queue.py: 74%

25 statements  

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

1""" 

2Base class for in memory buffer for database transactions 

3""" 

4 

5import asyncio 

6 

7from litellm._logging import verbose_proxy_logger 

8from litellm._service_logger import ServiceLogging 

9 

10service_logger_obj = ServiceLogging() # used for tracking metrics for In memory buffer, redis buffer, pod lock manager 

11from typing import Final 

12 

13from litellm.constants import ( 

14 LITELLM_ASYNCIO_QUEUE_MAXSIZE, 

15 MAX_IN_MEMORY_QUEUE_FLUSH_COUNT, 

16 MAX_SIZE_IN_MEMORY_QUEUE, 

17) 

18 

19 

20class BaseUpdateQueue: 

21 """Base class for in memory buffer for database transactions""" 

22 

23 def __init__(self): 

24 self.update_queue = asyncio.Queue(maxsize=LITELLM_ASYNCIO_QUEUE_MAXSIZE) 

25 self.MAX_SIZE_IN_MEMORY_QUEUE = MAX_SIZE_IN_MEMORY_QUEUE 

26 if MAX_SIZE_IN_MEMORY_QUEUE >= LITELLM_ASYNCIO_QUEUE_MAXSIZE: 26 ↛ 27line 26 didn't jump to line 27 because the condition on line 26 was never true

27 verbose_proxy_logger.warning( 

28 "Misconfigured queue thresholds: MAX_SIZE_IN_MEMORY_QUEUE (%d) >= LITELLM_ASYNCIO_QUEUE_MAXSIZE (%d). " 

29 "The spend aggregation check will never trigger because the asyncio.Queue blocks at %d items. " 

30 "Set MAX_SIZE_IN_MEMORY_QUEUE to a value less than LITELLM_ASYNCIO_QUEUE_MAXSIZE (recommended: 80%% of it).", 

31 MAX_SIZE_IN_MEMORY_QUEUE, 

32 LITELLM_ASYNCIO_QUEUE_MAXSIZE, 

33 LITELLM_ASYNCIO_QUEUE_MAXSIZE, 

34 ) 

35 

36 async def add_update(self, update): 

37 """Enqueue an update.""" 

38 verbose_proxy_logger.debug("Adding update to queue: %s", update) 

39 await self.update_queue.put(update) 

40 await self._emit_new_item_added_to_queue_event(queue_size=self.update_queue.qsize()) 

41 

42 async def flush_all_updates_from_in_memory_queue(self): 

43 """Get all updates from the queue.""" 

44 updates: Final = [] 

45 while not self.update_queue.empty(): 

46 # Circuit breaker to ensure we're not stuck dequeuing updates. Protect CPU utilization 

47 if len(updates) >= MAX_IN_MEMORY_QUEUE_FLUSH_COUNT: 47 ↛ 48line 47 didn't jump to line 48 because the condition on line 47 was never true

48 verbose_proxy_logger.debug("Max in memory queue flush count reached, stopping flush") 

49 break 

50 updates.append(await self.update_queue.get()) 

51 return updates 

52 

53 async def _emit_new_item_added_to_queue_event( 

54 self, 

55 queue_size: int | None = None, 

56 ): 

57 """placeholder, emit event when a new item is added to the queue"""