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
« 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"""
5import asyncio
7from litellm._logging import verbose_proxy_logger
8from litellm._service_logger import ServiceLogging
10service_logger_obj = ServiceLogging() # used for tracking metrics for In memory buffer, redis buffer, pod lock manager
11from typing import Final
13from litellm.constants import (
14 LITELLM_ASYNCIO_QUEUE_MAXSIZE,
15 MAX_IN_MEMORY_QUEUE_FLUSH_COUNT,
16 MAX_SIZE_IN_MEMORY_QUEUE,
17)
20class BaseUpdateQueue:
21 """Base class for in memory buffer for database transactions"""
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 )
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())
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
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"""