mirror of
https://github.com/BerriAI/litellm.git
synced 2026-10-06 02:48:13 +00:00
fix(gcs_bucket): drain the queue before uploading so a health flush never retries its own requeued batch
Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
This commit is contained in:
parent
caa0a66760
commit
661295634e
2 changed files with 24 additions and 5 deletions
|
|
@ -313,14 +313,15 @@ class GCSBucketLogger(GCSBucketBase, AdditionalLoggingUtils):
|
|||
await self._send_queued_events()
|
||||
|
||||
async def _send_queued_events(self) -> GCSFlushResult:
|
||||
items_to_process: Final = self._drain_queue_batch()
|
||||
return await self._send_items(self._drain_queue_batch())
|
||||
|
||||
if not items_to_process:
|
||||
async def _send_items(self, items: list[GCSLogQueueItem]) -> GCSFlushResult:
|
||||
if not items:
|
||||
return GCSFlushResult(sent_ids=(), failed_ids=())
|
||||
|
||||
if self.use_batched_logging:
|
||||
return await self._send_grouped_batches(items_to_process)
|
||||
return await self._send_individual_logs(items_to_process)
|
||||
return await self._send_grouped_batches(items)
|
||||
return await self._send_individual_logs(items)
|
||||
|
||||
def _get_object_name(self, kwargs: dict, logging_payload: StandardLoggingPayload, response_obj: Any) -> str:
|
||||
"""
|
||||
|
|
@ -411,10 +412,12 @@ class GCSBucketLogger(GCSBucketBase, AdditionalLoggingUtils):
|
|||
async def flush_queue_and_report(self) -> GCSFlushResult:
|
||||
"""
|
||||
Flush everything queued at call time, waiting for any in-flight periodic flush first, and report every event id.
|
||||
Events are drained before any upload starts, so a batch that fails and is requeued is not retried in this call.
|
||||
"""
|
||||
async with self.flush_lock:
|
||||
batch_count: Final = math.ceil(self.log_queue.qsize() / self.batch_size)
|
||||
results: Final = tuple([await self._send_queued_events() for _ in range(batch_count)])
|
||||
batches: Final = tuple(self._drain_queue_batch() for _ in range(batch_count))
|
||||
results: Final = tuple([await self._send_items(batch) for batch in batches])
|
||||
self.last_flush_time = time.time()
|
||||
return GCSFlushResult(
|
||||
sent_ids=tuple(event_id for result in results for event_id in result.sent_ids),
|
||||
|
|
|
|||
|
|
@ -73,6 +73,22 @@ async def test_failed_batch_stays_queued_and_is_retried_on_the_next_flush():
|
|||
assert logger.uploaded == [["req-1", "req-2"]]
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_multi_batch_flush_tries_every_queued_event_once_and_leaves_requeued_failures_for_the_next_flush():
|
||||
logger = _FakeUploadGCSLogger()
|
||||
logger.batch_size = 2
|
||||
logger.failing_ids = frozenset({"req-1"})
|
||||
await logger.enqueue("req-1")
|
||||
await logger.enqueue("req-2")
|
||||
await logger.enqueue("req-3")
|
||||
|
||||
result = await logger.flush_queue_and_report()
|
||||
|
||||
assert result == GCSFlushResult(sent_ids=("req-3",), failed_ids=("req-1", "req-2"))
|
||||
assert logger.uploaded == [["req-3"]]
|
||||
assert logger.queued_ids() == ["req-1", "req-2"]
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_individual_mode_requeues_only_the_failed_items():
|
||||
logger = _FakeUploadGCSLogger()
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue