mirror of
https://github.com/BerriAI/litellm.git
synced 2026-10-04 02:31:27 +00:00
* fix(e2e): wire batch provider secrets for docker and k8s
Point batch deployments at the credential field names and os.environ refs
the gateway actually resolves from process env (compose .env or EKS secret
mounts). Missing secrets skip instead of failing red so a red run means a
product bug. Mirror S3 bucket env aliases in docker-compose for provider_fallback
* fix(e2e): drop batch provider_env unit tests
The batches suite is live e2e only; no monkeypatch or unit-level tests
* fix: batch credentials, provider list, and team db lookup
Keep object-storage fields through CredentialLiteLLMParams and resolve
os.environ/ refs when reading deployment credentials so Vertex/Bedrock
batch file uploads see bucket and AWS keys from K8s/docker env
Skip managed batch list when the request is provider-scoped so
/{provider}/v1/batches list works instead of 500
Force DB on check_db_only team lookups and stop masking non-404 errors
as "team doesn't exist"
Drop e2e runner-side skip helpers; hard-fail on missing gateway secrets
* fix: tag reseed, team window spend, and remaining e2e flakes
Reseed spend:tag counters from LiteLLM_TagTable so cold redis still
enforces after the spend writer flushes
When applying post-call cost to team multi-window counters, load the
team from the DB if it is missing from the management cache so window
spend is not dropped on cache misses
Harden cold-counter reseed e2e (namespace-aware keys, burst success,
poll). Give tag budget more headroom. Retry /key/update on redis DNS
blips. Ensure NLTK punkt_tab is present for pipecat realtime audio
* revert: drop product code changes; e2e-only scope
Reverts all litellm/ and unit-test product edits. This branch is limited
to tests/e2e per contributor instruction
* fix(e2e): harden batch list and team member setup races
provider_fallback list falls back when managed batches reject provider
filtering. Team create waits for /team/info and member_add retries on
transient team-not-found so split control-plane lag does not red the suite
* fix(e2e): remove .env.example
Leave local .env and docker-compose env wiring as the secret source
* fix(e2e): wire files_settings and faster budget rescheduler for compose
OpenAI/Azure batch file uploads need files_settings; budget reset e2e needs a
short rescheduler window. Drop unsupported bedrock-encoded create_batch cells,
tolerate bedrock file.bytes=0, and surface team-info wait failures instead of
hanging silently
* chore(e2e): strip verbose comments from batch capabilities
* fix(e2e): assert managed list fallback before provider_fallback skip
When provider-scoped list is rejected, still fetch the unfiltered list and
check the envelope. Only skip membership when the id is a raw
provider_fallback batch that managed list cannot index
268 lines
9.2 KiB
Python
268 lines
9.2 KiB
Python
"""Client for the management-routes e2e suite: the shared Gateway plus the
|
|
key/team/user/organization writes, the info/list read-backs the tests assert,
|
|
and the raw-status calls judged by HTTP outcome (chat under a scoped key, an
|
|
llm-only key hitting a management route).
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import time
|
|
from dataclasses import dataclass
|
|
|
|
from e2e_gateway import Gateway, build_gateway
|
|
from e2e_http import NoBody, ProbeResult, Result, StreamingResponse, Success, UnknownApiError, unwrap
|
|
from models import (
|
|
ChatBody,
|
|
ChatMessage,
|
|
KeyDeleteBody,
|
|
KeyGenerateBody,
|
|
KeyListParams,
|
|
KeyListResponse,
|
|
KeyUpdateBody,
|
|
OrgDeleteBody,
|
|
OrgInfoParams,
|
|
OrgInfoResponse,
|
|
OrgNewBody,
|
|
OrgNewResponse,
|
|
TeamData,
|
|
TeamDeleteBody,
|
|
TeamInfoParams,
|
|
TeamInfoResponse,
|
|
TeamMemberAddBody,
|
|
TeamMemberDeleteBody,
|
|
TeamMemberEntry,
|
|
TeamNewBody,
|
|
TeamNewResponse,
|
|
UserDeleteBody,
|
|
UserInfoParams,
|
|
UserInfoResponse,
|
|
UserListParams,
|
|
UserListResponse,
|
|
UserNewBody,
|
|
UserNewResponse,
|
|
)
|
|
|
|
MODEL_ACCESS_DENIED_MARKER = "key_model_access_denied"
|
|
ROUTE_NOT_ALLOWED_MARKER = "not allowed to call this route"
|
|
_TEAM_READY_ATTEMPTS = 15
|
|
_TEAM_READY_SLEEP_SECONDS = 0.4
|
|
|
|
|
|
@dataclass(frozen=True, slots=True)
|
|
class ManagementClient:
|
|
gateway: Gateway
|
|
|
|
def llm_only_key(self) -> str:
|
|
return self.gateway.generate_key(KeyGenerateBody(models=[], allowed_routes=["llm_api_routes"]))
|
|
|
|
def update_key_models(self, key: str, models: list[str]) -> None:
|
|
last: Result[NoBody] | None = None
|
|
for attempt in range(5):
|
|
last = self.gateway.transport.post(
|
|
"/key/update",
|
|
headers=self.gateway.transport.master,
|
|
json=KeyUpdateBody(key=key, models=models),
|
|
response_type=NoBody,
|
|
)
|
|
match last:
|
|
case Success():
|
|
return
|
|
case UnknownApiError(body=body) if (
|
|
"connecting to redis" in body.lower() or "name resolution" in body.lower()
|
|
):
|
|
time.sleep(0.5 * (attempt + 1))
|
|
continue
|
|
case _:
|
|
break
|
|
assert last is not None
|
|
_ = unwrap(last)
|
|
|
|
def delete_key_strict(self, key: str) -> None:
|
|
"""Strict delete for the act phase of a test: a failed delete is a hard
|
|
failure, unlike the warn-only Gateway.delete_key used at teardown."""
|
|
_ = unwrap(
|
|
self.gateway.transport.post(
|
|
"/key/delete",
|
|
headers=self.gateway.transport.master,
|
|
json=KeyDeleteBody(keys=[key]),
|
|
response_type=NoBody,
|
|
)
|
|
)
|
|
|
|
def key_alias_count(self, key_alias: str) -> int:
|
|
return unwrap(
|
|
self.gateway.transport.get(
|
|
"/key/list",
|
|
headers=self.gateway.transport.master,
|
|
params=KeyListParams(key_alias=key_alias),
|
|
response_type=KeyListResponse,
|
|
)
|
|
).total_count
|
|
|
|
def create_team(self, body: TeamNewBody) -> str:
|
|
team_id = unwrap(
|
|
self.gateway.transport.post(
|
|
"/team/new",
|
|
headers=self.gateway.transport.master,
|
|
json=body,
|
|
response_type=TeamNewResponse,
|
|
)
|
|
).team_id
|
|
self._wait_for_team(team_id)
|
|
return team_id
|
|
|
|
def delete_team(self, team_id: str) -> None:
|
|
_ = self.gateway.transport.post(
|
|
"/team/delete",
|
|
headers=self.gateway.transport.master,
|
|
json=TeamDeleteBody(team_ids=[team_id]),
|
|
response_type=NoBody,
|
|
)
|
|
|
|
def team_info(self, team_id: str) -> TeamData:
|
|
return unwrap(
|
|
self.gateway.transport.get(
|
|
"/team/info",
|
|
headers=self.gateway.transport.master,
|
|
params=TeamInfoParams(team_id=team_id),
|
|
response_type=TeamInfoResponse,
|
|
)
|
|
).team_info
|
|
|
|
def team_info_status(self, team_id: str) -> ProbeResult:
|
|
return self.gateway.transport.probe("/team/info", params=TeamInfoParams(team_id=team_id))
|
|
|
|
def _wait_for_team(self, team_id: str) -> None:
|
|
last: Result[TeamInfoResponse] | None = None
|
|
for _ in range(_TEAM_READY_ATTEMPTS):
|
|
last = self.gateway.transport.get(
|
|
"/team/info",
|
|
headers=self.gateway.transport.master,
|
|
params=TeamInfoParams(team_id=team_id),
|
|
response_type=TeamInfoResponse,
|
|
)
|
|
match last:
|
|
case Success():
|
|
return
|
|
case _:
|
|
time.sleep(_TEAM_READY_SLEEP_SECONDS)
|
|
assert last is not None
|
|
_ = unwrap(last)
|
|
|
|
def add_team_member(self, team_id: str, user_id: str) -> None:
|
|
last: Result[NoBody] | None = None
|
|
for attempt in range(_TEAM_READY_ATTEMPTS):
|
|
last = self.gateway.transport.post(
|
|
"/team/member_add",
|
|
headers=self.gateway.transport.master,
|
|
json=TeamMemberAddBody(team_id=team_id, member=TeamMemberEntry(role="user", user_id=user_id)),
|
|
response_type=NoBody,
|
|
)
|
|
match last:
|
|
case Success():
|
|
return
|
|
case UnknownApiError(body=body) if (
|
|
"doesn't exist" in body and attempt + 1 < _TEAM_READY_ATTEMPTS
|
|
):
|
|
time.sleep(_TEAM_READY_SLEEP_SECONDS)
|
|
continue
|
|
case _:
|
|
break
|
|
assert last is not None
|
|
_ = unwrap(last)
|
|
|
|
def delete_team_member(self, team_id: str, user_id: str) -> None:
|
|
_ = unwrap(
|
|
self.gateway.transport.post(
|
|
"/team/member_delete",
|
|
headers=self.gateway.transport.master,
|
|
json=TeamMemberDeleteBody(team_id=team_id, user_id=user_id),
|
|
response_type=NoBody,
|
|
)
|
|
)
|
|
|
|
def create_user(self, body: UserNewBody) -> str:
|
|
return unwrap(
|
|
self.gateway.transport.post(
|
|
"/user/new",
|
|
headers=self.gateway.transport.master,
|
|
json=body,
|
|
response_type=UserNewResponse,
|
|
)
|
|
).user_id
|
|
|
|
def delete_user(self, user_id: str) -> None:
|
|
_ = self.gateway.transport.post(
|
|
"/user/delete",
|
|
headers=self.gateway.transport.master,
|
|
json=UserDeleteBody(user_ids=[user_id]),
|
|
response_type=NoBody,
|
|
)
|
|
|
|
def user_info(self, user_id: str) -> UserInfoResponse:
|
|
return unwrap(
|
|
self.gateway.transport.get(
|
|
"/user/info",
|
|
headers=self.gateway.transport.master,
|
|
params=UserInfoParams(user_id=user_id),
|
|
response_type=UserInfoResponse,
|
|
)
|
|
)
|
|
|
|
def user_count(self, user_id: str) -> int:
|
|
return unwrap(
|
|
self.gateway.transport.get(
|
|
"/user/list",
|
|
headers=self.gateway.transport.master,
|
|
params=UserListParams(user_ids=user_id),
|
|
response_type=UserListResponse,
|
|
)
|
|
).total
|
|
|
|
def create_org(self, body: OrgNewBody) -> str:
|
|
return unwrap(
|
|
self.gateway.transport.post(
|
|
"/organization/new",
|
|
headers=self.gateway.transport.master,
|
|
json=body,
|
|
response_type=OrgNewResponse,
|
|
)
|
|
).organization_id
|
|
|
|
def delete_org(self, organization_id: str) -> None:
|
|
_ = self.gateway.transport.delete(
|
|
"/organization/delete",
|
|
headers=self.gateway.transport.master,
|
|
json=OrgDeleteBody(organization_ids=[organization_id]),
|
|
response_type=NoBody,
|
|
)
|
|
|
|
def org_info(self, organization_id: str) -> OrgInfoResponse:
|
|
return unwrap(
|
|
self.gateway.transport.get(
|
|
"/organization/info",
|
|
headers=self.gateway.transport.master,
|
|
params=OrgInfoParams(organization_id=organization_id),
|
|
response_type=OrgInfoResponse,
|
|
)
|
|
)
|
|
|
|
def chat_status(self, key: str, model: str, content: str) -> StreamingResponse:
|
|
return self.gateway.transport.send(
|
|
"/chat/completions",
|
|
headers=self.gateway.transport.bearer(key),
|
|
json=ChatBody(model=model, messages=[ChatMessage(role="user", content=content)], max_tokens=16),
|
|
)
|
|
|
|
def key_generate_status(self, key: str, body: KeyGenerateBody) -> StreamingResponse:
|
|
return self.gateway.transport.send("/key/generate", headers=self.gateway.transport.bearer(key), json=body)
|
|
|
|
def team_new_status(self, key: str, body: TeamNewBody) -> StreamingResponse:
|
|
return self.gateway.transport.send("/team/new", headers=self.gateway.transport.bearer(key), json=body)
|
|
|
|
def user_new_status(self, key: str, body: UserNewBody) -> StreamingResponse:
|
|
return self.gateway.transport.send("/user/new", headers=self.gateway.transport.bearer(key), json=body)
|
|
|
|
|
|
def build_client() -> ManagementClient:
|
|
return ManagementClient(gateway=build_gateway())
|