litellm/tests/e2e/management/management_client.py
mubashir1osmani 54d404ef2c
fix(e2e): batch credentials wiring and compose harness for live proxy suite (#32744)
* 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
2026-07-10 11:31:40 -07:00

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())