From cffa6405b4c9074ac95edf29831ea83e483ba302 Mon Sep 17 00:00:00 2001 From: yucheng Date: Sun, 27 Sep 2026 00:52:17 +0000 Subject: [PATCH] fix(batches): release line-item claim when fan-out emitted nothing Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- litellm/batches/batch_line_item_logging.py | 13 +++- .../batches/test_batch_line_item_logging.py | 72 +++++++++++++++++++ 2 files changed, 82 insertions(+), 3 deletions(-) diff --git a/litellm/batches/batch_line_item_logging.py b/litellm/batches/batch_line_item_logging.py index f24ae9fb432..21dd23bdb59 100644 --- a/litellm/batches/batch_line_item_logging.py +++ b/litellm/batches/batch_line_item_logging.py @@ -397,9 +397,8 @@ async def log_batch_line_items( batch.id, ) return 0 - claim: Final = await claim_cache.async_increment_cache( - f"batch_line_items_emitted:{batch.id}", 1, ttl=_LINE_ITEM_CLAIM_TTL_SECONDS - ) + claim_key: Final = f"batch_line_items_emitted:{batch.id}" + claim: Final = await claim_cache.async_increment_cache(claim_key, 1, ttl=_LINE_ITEM_CLAIM_TTL_SECONDS) if claim is not None and claim > 1: verbose_logger.debug("batch line items already emitted for batch_id=%s, skipping", batch.id) return 0 @@ -447,4 +446,12 @@ async def log_batch_line_items( "batch line item logging failed for batch_id=%s; aggregate logging unaffected", batch.id, ) + if emitted == 0: + try: + await claim_cache.async_delete_cache(claim_key) + except Exception: # noqa: BLE001 # the claim release must never raise; worst case the batch stays claimed + verbose_logger.debug( + "batch line item claim release failed for batch_id=%s, claim persists until ttl", + batch.id, + ) return emitted diff --git a/tests/unit/batches/test_batch_line_item_logging.py b/tests/unit/batches/test_batch_line_item_logging.py index 9a855b602e5..f5514c5814c 100644 --- a/tests/unit/batches/test_batch_line_item_logging.py +++ b/tests/unit/batches/test_batch_line_item_logging.py @@ -892,3 +892,75 @@ async def test_line_items_emit_when_the_claim_backend_returns_nothing(recorder): assert emitted == 2 assert len(recorder.success_events) == 1 assert len(recorder.failure_events) == 1 + + +_BROKEN_ERROR_FILE: Final = {**_FILE_BYTES, "error-file-1": "not bytes"} + + +def _broken_error_file_content(file_id: str, **_kwargs): + return SimpleNamespace(content=_BROKEN_ERROR_FILE[file_id]) + + +@pytest.mark.asyncio +async def test_line_items_retry_after_a_failed_fanout(recorder): + claim_cache: Final = DualCache() + batch: Final = _batch() + parent: Final = _parent_logging() + with patch("litellm.files.main.afile_content", new_callable=AsyncMock, side_effect=ValueError("boom")): # test-quality-ok: afile_content is the provider boundary; no injection seam for managed file fetch + first: Final = await log_batch_line_items( + batch=batch, + custom_llm_provider="openai", + parent=parent, + model_name="gpt-4o", + litellm_params=None, + model_info=None, + claim_cache=claim_cache, + ) + file_mock: Final = AsyncMock(side_effect=_file_content) + with patch("litellm.files.main.afile_content", file_mock): # test-quality-ok: afile_content is the provider boundary; no injection seam for managed file fetch + second: Final = await log_batch_line_items( + batch=batch, + custom_llm_provider="openai", + parent=parent, + model_name="gpt-4o", + litellm_params=None, + model_info=None, + claim_cache=claim_cache, + ) + + assert first == 0 + assert second == 2 + assert len(recorder.success_events) == 1 + assert len(recorder.failure_events) == 1 + + +@pytest.mark.asyncio +async def test_line_items_claim_kept_after_partial_emission(recorder): + claim_cache: Final = DualCache() + batch: Final = _batch() + parent: Final = _parent_logging() + file_mock: Final = AsyncMock(side_effect=_broken_error_file_content) + with patch("litellm.files.main.afile_content", file_mock): # test-quality-ok: afile_content is the provider boundary; no injection seam for managed file fetch + first: Final = await log_batch_line_items( + batch=batch, + custom_llm_provider="openai", + parent=parent, + model_name="gpt-4o", + litellm_params=None, + model_info=None, + claim_cache=claim_cache, + ) + second: Final = await log_batch_line_items( + batch=batch, + custom_llm_provider="openai", + parent=parent, + model_name="gpt-4o", + litellm_params=None, + model_info=None, + claim_cache=claim_cache, + ) + + assert first == 1 + assert second == 0 + assert len(recorder.success_events) == 1 + assert len(recorder.failure_events) == 0