fix(batches): release line-item claim when fan-out emitted nothing
Some checks failed
LiteLLM Rust / rust-lint (push) Has been cancelled
LiteLLM Rust / rust-test (push) Has been cancelled
LiteLLM Rust / rust-wheel (push) Has been cancelled

Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
This commit is contained in:
yucheng 2026-09-27 00:52:17 +00:00
parent f65b5b1281
commit cffa6405b4
2 changed files with 82 additions and 3 deletions

View file

@ -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

View file

@ -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