From 661295634e8908b629221107c874382cf97f64ab Mon Sep 17 00:00:00 2001 From: yucheng Date: Sun, 13 Sep 2026 02:21:43 +0000 Subject: [PATCH] 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> --- litellm/integrations/gcs_bucket/gcs_bucket.py | 13 ++++++++----- .../integrations/gcs_bucket/test_gcs_bucket.py | 16 ++++++++++++++++ 2 files changed, 24 insertions(+), 5 deletions(-) diff --git a/litellm/integrations/gcs_bucket/gcs_bucket.py b/litellm/integrations/gcs_bucket/gcs_bucket.py index 373726da29a..06397c2b541 100644 --- a/litellm/integrations/gcs_bucket/gcs_bucket.py +++ b/litellm/integrations/gcs_bucket/gcs_bucket.py @@ -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), diff --git a/tests/test_litellm/integrations/gcs_bucket/test_gcs_bucket.py b/tests/test_litellm/integrations/gcs_bucket/test_gcs_bucket.py index b6e43fff6b7..0d2344ef9d6 100644 --- a/tests/test_litellm/integrations/gcs_bucket/test_gcs_bucket.py +++ b/tests/test_litellm/integrations/gcs_bucket/test_gcs_bucket.py @@ -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()