mirror of
https://github.com/BerriAI/litellm.git
synced 2026-10-03 02:22:24 +00:00
* fix(ci): stop five stale or flaky CI reds and retry CyberArk policy-load conflicts The Langfuse redaction unit test exports to a local OTLP capture instead of polling Langfuse Cloud through a recorded lookup. The passthrough worker-kill test only requires spend rows for requests the surviving worker served. The spend-routes sweep treats the intentional /spend/capture_rate 503 as expected. CyberArk retries a 409 policy load in Python, Rust and the e2e Conjur helper instead of reading it as "variable exists". The integration egress guard now matches the script's own cgroup, so it no longer blocks the CircleCI agent, which runs as the same user. * fix(ci): keep the policy-load backoff typed as float * fix(ci): retry CyberArk policy loads without blocking the event loop and tighten the worker-kill and Langfuse tests * fix(secrets): load CyberArk policy one request at a time per manager * test(secrets): pin that non-conflict CyberArk policy failures are not retried * test(unit): run tests/unit with only an allowlisted host environment CircleCI's unit job inherits every project env var, so real provider keys, REDIS_HOST, DATABASE_URL and AWS or Azure credentials reached tests that assume none are set. Locally, litellm's import-time load_dotenv did the same from any .env up the tree. The unit conftest now drops every variable outside a small allowlist and disables dotenv before litellm is imported. * test(e2e): name a failed search and the stuck batch status instead of misattributing them The websearch session test read an empty web_search_tool_result_error block as a successful search, so a failing search tool surfaced as a session billing bug. The batch cancellation timeout now reports the last status the proxy returned. * fix(ci): scrub the host environment per unit test instead of for the whole pytest process GHA shards run tests/unit next to other suites in one process, so the import-time scrub deleted MCP_TEST_PEER_PYTHON before tests/mcp_tests read it and the MCP upstream fell back to the SDK2 interpreter. The two websearch tests that called OpenAI and Perplexity live are removed: tests/unit no longer sees their keys. * fix(ci): scrub only the host variables present before litellm is imported The per-test scrub also deleted TIKTOKEN_CACHE_DIR, which litellm sets at import to its bundled encodings, so tokenizer paths tried to download them and hit the socket guard. The prisma setup test now passes its own database URL instead of reading one another test leaked into the process environment. * fix(ci): stop the order-dependent unit reds and settle logging tasks on their own queue LoggingWorker marked a task done on whichever queue was current when the callback finished, so a callback that outlived an event-loop change raised "task_done() called too many times" or undercounted the new loop's queue. It now settles the queue the task came from. The rest are test isolation fixes for failures that only appeared when another file ran first on the same xdist worker: a replaced user_api_key_cache, breaker metrics unregistered by prometheus tests, semantic_router's health-check filter on uvicorn.access, logging tasks carried over from bedrock tests, a Router-written model_cost entry, and a stray post captured by the langflow test. The token counter check now asserts bounded chunking instead of wall-clock time. * test(e2e/ui): wait for the logout redirect before visiting a protected page Logout revokes the session server-side before clearing cookies and navigating, so an immediate page.goto either ran with the cookie still set or was aborted by the logout redirect (net::ERR_ABORTED). * test(unit): restore the prometheus metrics config per test and settle logs carried from earlier tests in the a2a cost tests * test(router): pin the router clock in the usage counter tests so a minute rollover cannot empty the read * test(e2e/ui): wait for logout to clear the token cookie instead of for a login redirect * test(integration/mcp): answer the model-info probe another test's proxy sends to the model double
148 lines
6.2 KiB
Python
148 lines
6.2 KiB
Python
from builtins import ExceptionGroup
|
|
from collections.abc import Callable
|
|
from itertools import count
|
|
from time import monotonic, sleep
|
|
from typing import Final, Protocol
|
|
|
|
from batch_client import BatchObject, FileDeleteResponse
|
|
from capabilities import is_cloud_storage_id, is_managed_id
|
|
from e2e_http import NetworkError, RateLimitedError, Result, Success, UnknownApiError
|
|
from pydantic import BaseModel
|
|
|
|
CLEANUP_DELAYS: Final = (1.0, 2.0, 4.0)
|
|
BATCH_TERMINAL_STATUSES: Final = frozenset({"completed", "failed", "expired", "cancelled"})
|
|
BATCH_PENDING_STATUSES: Final = frozenset({"validating", "in_progress", "finalizing", "cancelling"})
|
|
BATCH_CANCEL_TIMEOUT_SECONDS: Final = 660.0
|
|
BATCH_CANCEL_POLL_SECONDS: Final = 10.0
|
|
|
|
|
|
class BatchCleanupClient(Protocol):
|
|
def delete_file(self, file_id: str, *, key: str, provider: str | None = None) -> Result[FileDeleteResponse]: ...
|
|
|
|
def delete_file_as_admin(self, file_id: str, *, provider: str | None = None) -> Result[FileDeleteResponse]: ...
|
|
|
|
def retrieve_batch(self, batch_id: str, *, key: str, provider: str | None = None) -> Result[BatchObject]: ...
|
|
|
|
def cancel_batch(self, batch_id: str, *, key: str, provider: str | None = None) -> Result[BatchObject]: ...
|
|
|
|
|
|
def cleanup_result[R: BaseModel](
|
|
action: Callable[[], Result[R]], *, wait: Callable[[float], None] = sleep
|
|
) -> Result[R]:
|
|
for delay, result in ((delay, action()) for delay in CLEANUP_DELAYS):
|
|
match result:
|
|
case NetworkError() | RateLimitedError():
|
|
wait(delay)
|
|
case UnknownApiError(status_code=code) if code in {408, 429, 500, 502, 503, 504}:
|
|
wait(delay)
|
|
case _:
|
|
return result
|
|
return action()
|
|
|
|
|
|
def _require_cleanup_success[R: BaseModel](result: Result[R], operation: str) -> R:
|
|
match result:
|
|
case Success(data=data):
|
|
return data
|
|
case UnknownApiError(status_code=code):
|
|
raise AssertionError(f"{operation} failed: HTTP {code}")
|
|
case _:
|
|
raise AssertionError(f"{operation} failed: {result.kind}")
|
|
|
|
|
|
def cleanup_file(client: BatchCleanupClient, file_id: str, *, key: str, provider: str | None = None) -> None:
|
|
delete: Final[Callable[[], Result[FileDeleteResponse]]] = (
|
|
(lambda: client.delete_file_as_admin(file_id, provider=provider))
|
|
if is_cloud_storage_id(file_id)
|
|
else (lambda: client.delete_file(file_id, key=key, provider=provider))
|
|
)
|
|
result: Final = cleanup_result(delete)
|
|
if isinstance(result, UnknownApiError) and result.status_code == 404:
|
|
return
|
|
deleted: Final = _require_cleanup_success(result, f"Delete file {file_id}")
|
|
assert deleted.deleted is True or (
|
|
deleted.deleted is None and is_managed_id(file_id) and deleted.id == file_id and deleted.object == "file"
|
|
), f"Delete file {file_id} did not confirm deletion"
|
|
|
|
|
|
def cleanup_batch(
|
|
client: BatchCleanupClient,
|
|
batch_id: str,
|
|
*,
|
|
key: str,
|
|
provider: str | None = None,
|
|
delete_output_files: bool = False,
|
|
wait: Callable[[float], None] = sleep,
|
|
clock: Callable[[], float] = monotonic,
|
|
) -> None:
|
|
needs_terminal_state: Final = is_managed_id(batch_id)
|
|
fetched: Final = _require_cleanup_success(
|
|
cleanup_result(lambda: client.retrieve_batch(batch_id, key=key, provider=provider)),
|
|
f"Retrieve batch {batch_id} for cleanup",
|
|
)
|
|
if fetched.status in BATCH_TERMINAL_STATUSES:
|
|
if delete_output_files:
|
|
_cleanup_batch_outputs(client, fetched, key=key, provider=provider)
|
|
return
|
|
if fetched.status == "cancelling" and not needs_terminal_state:
|
|
return
|
|
result: Final = (
|
|
Success(status_code=200, data=fetched)
|
|
if fetched.status == "cancelling"
|
|
else cleanup_result(lambda: client.cancel_batch(batch_id, key=key, provider=provider))
|
|
)
|
|
conflicted: Final = isinstance(result, UnknownApiError) and result.status_code in {400, 409}
|
|
if not conflicted:
|
|
cancelled: Final = _require_cleanup_success(result, f"Cancel batch {batch_id}")
|
|
assert cancelled.status in BATCH_TERMINAL_STATUSES | BATCH_PENDING_STATUSES, (
|
|
f"Cancel batch {batch_id} left status {cancelled.status}"
|
|
)
|
|
if cancelled.status in BATCH_TERMINAL_STATUSES:
|
|
if delete_output_files:
|
|
_cleanup_batch_outputs(client, cancelled, key=key, provider=provider)
|
|
return
|
|
if cancelled.status == "cancelling" and not needs_terminal_state:
|
|
return
|
|
deadline: Final = clock() + BATCH_CANCEL_TIMEOUT_SECONDS
|
|
for current in (
|
|
_require_cleanup_success(
|
|
cleanup_result(lambda: client.retrieve_batch(batch_id, key=key, provider=provider)),
|
|
f"Retrieve batch {batch_id} after cancellation",
|
|
)
|
|
for _ in count()
|
|
):
|
|
if current.status in BATCH_TERMINAL_STATUSES:
|
|
if delete_output_files:
|
|
_cleanup_batch_outputs(client, current, key=key, provider=provider)
|
|
return
|
|
assert current.status in ({"cancelling"} if conflicted else BATCH_PENDING_STATUSES), (
|
|
f"Cancel batch {batch_id} left status {current.status}"
|
|
)
|
|
if current.status == "cancelling" and not needs_terminal_state:
|
|
return
|
|
assert clock() < deadline, (
|
|
f"Batch {batch_id} cancellation did not finish within {BATCH_CANCEL_TIMEOUT_SECONDS}s, "
|
|
f"last status {current.status}"
|
|
)
|
|
wait(BATCH_CANCEL_POLL_SECONDS)
|
|
|
|
|
|
def _cleanup_batch_outputs(client: BatchCleanupClient, batch: BatchObject, *, key: str, provider: str | None) -> None:
|
|
errors: Final = tuple(
|
|
error
|
|
for file_id in dict.fromkeys((batch.output_file_id, batch.error_file_id))
|
|
if file_id is not None and file_id != batch.input_file_id
|
|
if (error := _output_cleanup_error(client, file_id, key=key, provider=provider)) is not None
|
|
)
|
|
if errors:
|
|
raise ExceptionGroup(f"Batch {batch.id} output cleanup failed", errors)
|
|
|
|
|
|
def _output_cleanup_error(
|
|
client: BatchCleanupClient, file_id: str, *, key: str, provider: str | None
|
|
) -> Exception | None:
|
|
try:
|
|
cleanup_file(client, file_id, key=key, provider=provider)
|
|
except Exception as error:
|
|
return error
|
|
return None
|