litellm/tests/e2e/management/management_client.py
devin-ai-integration[bot] f285229b51
fix(proxy): delete large teams without per-member transaction fan-out (#42998)
* fix(proxy): delete large teams without per-member transaction fan-out

Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* fix(proxy): evict email-only member caches and reset team members metric on delete

Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* fix(proxy): keep new delete-team literals within the LIT002 ceiling

Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* fix(proxy): resolve deleted-team member ids before the locked delete

Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* fix(proxy): resolve email-only deleted-team members with one case-insensitive lookup

`_deleted_team_member_user_ids` looked each email-only roster entry up with its own
`find_users_by_email` call inside an unbounded `asyncio.gather`: one exact-match query
per email, so a large roster fanned out against the pool again and a roster email that
differed in case from its user row was missed. Add `UserRepository.find_by_emails`, a
single case-insensitive `in` query, and call it once before the locked delete.
`management_helpers/utils.py` goes back to its main-branch shape since the single-email
helper no longer needs exporting.

* fix(repositories): slice find_by_emails into bounded IN statements

The unbounded-IN lint flagged the case-insensitive email lookup added for
/team/delete cache eviction. chunked_in.find_many_in cannot carry Prisma's
insensitive mode, so the repository slices the deduplicated list into
IN_LIST_CHUNK_SIZE statements itself and concatenates the pages. Empty input
still returns () without a query.

* fix(proxy): delete a team once when /team/delete repeats its id

The audit sent {"team_ids": [T, T]}: main answered 400 "User not found in
team" after deleting the keys and memberships and writing two tombstones,
leaving the team row behind; this branch answered 200 but still wrote the
tombstone, audit row and eviction twice. DeleteTeamRequest now collapses
repeated ids in order, so every later step sees each team once and the
response lists each deleted team once.

* test(integration): audit cells for /team/delete on large, legacy and concurrent teams

Thirty-eight deterministic cells in tests/integration/management/ (the CircleCI
integration-management group) covering the /team/delete happy, sad, edge and chaos rows:
250 members against a pool limit of five on two workers, the advisory-lock wait, email-only
legacy roster entries in every casing, member and team cache eviction on both proxies for
every client and endpoint, the Prometheus gauge, audit rows, malformed and duplicate input,
the route gate, and a worker kill, a Redis outage and a proxy restart mid-burst.

Every cell runs against the real proxy, Postgres and Redis with the scripted upstream; no
component is mocked. On the merge base the rows this fix changes are red (P2028 on the
250-member team, two lock waiters, case-mismatched email lookups, duplicate ids, orphaned
LiteLLM_UserTable.teams references under a concurrent burst); on the tip every cell is green
twice with identical selections.

Two pre-existing behaviours are pinned as observed rather than fixed here: a roster entry with
neither user_id nor user_email answers 500, and the LiteLLM_DeletedTeamTable row is committed
before the locked transaction, so a delete that dies in between leaves a tombstone for a live
team and the retry adds a second.

* test(integration): pin each chaos outage to a live /team/delete

The three chaos cells applied the outage once three deletes had answered, which on a fast
run let the whole burst finish before the worker kill, Redis stop or SIGTERM landed, so the
cells passed without exercising the failure. Each cell now holds the first team's advisory
lock from a test-owned transaction, waits until that team's delete is queued behind it in
Postgres with its request unanswered, applies the outage, and only then releases the lock,
so an in-flight delete meets the failure on every run and both legs. The pinned team's
outcome and the number of deletes answered before the outage are recorded as junit
properties (pinned_delete, answered_before_outage).

---------

Co-authored-by: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
Co-authored-by: ryan-crabbe-berri <ryan@berri.ai>
2026-09-30 13:49:18 -07:00

670 lines
24 KiB
Python

"""Client for the management-routes e2e suite: the shared ProxyClient 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
import warnings
from dataclasses import dataclass, field, replace
import jwt
from e2e_config import MASTER_KEY
from e2e_http import (
AuthHeaders,
NetworkError,
NoBody,
ProbeResult,
Result,
StreamingResponse,
Success,
UnknownApiError,
retry_attempts,
unwrap,
)
from models import (
AuditLogPage,
AuditLogParams,
ChatBody,
ChatMessage,
ConnectionTestBody,
ConnectionTestResponse,
CustomerDeleteBody,
CustomerInfoParams,
CustomerNewBody,
CustomerResponse,
KeyBlockBody,
KeyDeleteBody,
KeyDeleteByAliasBody,
KeyGenerateBody,
KeyGenerateResponse,
KeyInfoParams,
KeyInfoResponse,
KeyListParams,
KeyListResponse,
KeyRegenerateBody,
KeyResetSpendBody,
KeyResetSpendResponse,
KeyUpdateBody,
McpServerCreateBody,
McpServerRow,
McpServerUpdateBody,
ModelDeleteBody,
OrgDeleteBody,
OrgInfoParams,
OrgInfoResponse,
OrgNewBody,
OrgNewResponse,
OrgUpdateBody,
TagDeleteBody,
TagListEntry,
TagListResponse,
TagNewBody,
TeamData,
TeamDeleteBody,
TeamInfoParams,
TeamInfoResponse,
TeamListResponse,
TeamMemberAddBody,
TeamMemberDeleteBody,
TeamMemberEntry,
TeamNewBody,
TeamNewResponse,
TeamUpdateBody,
UiLoginBody,
UiLoginResponse,
UiSessionClaims,
UserDeleteBody,
UserDeleteResponse,
UserInfoParams,
UserInfoResponse,
UserListParams,
UserListResponse,
UserNewBody,
UserNewResponse,
UserUpdateBody,
)
from proxy_client import Caller, ProxyClient
MODEL_ACCESS_DENIED_MARKER = "key_model_access_denied"
ROUTE_NOT_ALLOWED_MARKER = "not allowed to call this route"
DASHBOARD_SESSION_TEAM_ID = "litellm-dashboard"
_TEAM_READY_ATTEMPTS = 15
_TEAM_READY_SLEEP_SECONDS = 0.4
_KEY_WRITE_ATTEMPTS = 5
_TRANSIENT_BACKEND_MARKERS = ("connecting to redis", "name resolution")
@dataclass(frozen=True, slots=True)
class DashboardSession:
"""What a dashboard sign-in hands the Admin UI: the session key it sends as
its bearer on every subsequent call, the claims it renders the signed-in user
from, and where it lands the browser."""
session_key: str = field(repr=False)
claims: UiSessionClaims
redirect_url: str
@dataclass(frozen=True, slots=True)
class ManagementClient:
proxy: ProxyClient
master_key: str = field(repr=False)
def with_caller(self, caller: Caller) -> ManagementClient:
return replace(self, proxy=self.proxy.with_caller(caller))
def llm_only_key(self) -> str:
return self.proxy.generate_key(KeyGenerateBody(models=[], allowed_routes=["llm_api_routes"]))
def generate_key(self, body: KeyGenerateBody, *, caller_key: str | None = None) -> Result[KeyGenerateResponse]:
"""POST /key/generate. `caller_key` is who is creating the key: the master
key by default, or a virtual key (an admin filling in Create New Key on the
dashboard creates it under the session key their sign-in minted). Returns
the outcome rather than unwrapping it, so a caller can poll a route that is
only transiently refusing."""
headers = self.proxy.management_headers(caller_key)
return self.proxy.transport.post(
"/key/generate",
headers=headers,
json=body,
response_type=KeyGenerateResponse,
)
def update_key(self, body: KeyUpdateBody, *, caller_key: str | None = None) -> Result[NoBody]:
"""POST /key/update. `caller_key` is who is editing: the master key by
default, or a virtual key (the dashboard edits under the session key its
sign-in minted, never the master key). Returns the outcome rather than
unwrapping it, so a caller can poll a route that is only transiently
refusing; `update_key_models` is the unwrapping shorthand."""
headers = self.proxy.management_headers(caller_key)
last: Result[NoBody] = NetworkError(message="/key/update was never attempted")
for attempt in range(retry_attempts(_KEY_WRITE_ATTEMPTS)):
last = self.proxy.transport.post(
"/key/update",
headers=headers,
json=body,
response_type=NoBody,
)
match last:
case UnknownApiError(body=error_body) if any(
marker in error_body.lower() for marker in _TRANSIENT_BACKEND_MARKERS
):
warnings.warn(f"Transient backend response on attempt {attempt + 1}", RuntimeWarning, stacklevel=2)
time.sleep(0.5 * (attempt + 1))
continue
case _:
break
return last
def update_key_models(self, key: str, models: list[str]) -> None:
_ = unwrap(self.update_key(KeyUpdateBody(key=key, models=models)))
def delete_key_by_alias(self, key_alias: str) -> None:
_ = unwrap(
self.proxy.transport.post(
"/key/delete",
headers=self.proxy.management_headers(),
json=KeyDeleteByAliasBody(key_aliases=[key_alias]),
response_type=NoBody,
)
)
def key_deleted_audit_logs(self, token_hash: str) -> AuditLogPage:
return unwrap(
self.proxy.transport.get(
"/audit",
headers=self.proxy.management_headers(),
params=AuditLogParams(
object_id=token_hash,
action="deleted",
table_name="LiteLLM_VerificationToken",
page_size=100,
),
response_type=AuditLogPage,
)
)
def key_info_as(self, key: str, *, caller_key: str | None = None) -> Result[KeyInfoResponse]:
return self.proxy.transport.get(
"/key/info",
headers=self.proxy.management_headers(caller_key),
params=KeyInfoParams(key=key),
response_type=KeyInfoResponse,
)
def delete_key_strict(self, key: str, *, caller_key: str | None = None, missing_ok: bool = False) -> None:
"""Strict delete for the act phase of a test: a failed delete is a hard
failure, unlike the warn-only ProxyClient.delete_key used at teardown."""
result = self.proxy.transport.post(
"/key/delete",
headers=self.proxy.management_headers(caller_key),
json=KeyDeleteBody(keys=[key]),
response_type=NoBody,
)
if missing_ok and isinstance(result, UnknownApiError) and result.status_code == 404:
return
_ = unwrap(result)
def delete_model_strict(self, model_id: str) -> None:
"""Strict delete for the act phase of a test: a failed delete is a hard
failure, unlike the warn-only ProxyClient.delete_model used at teardown."""
_ = unwrap(
self.proxy.transport.post(
"/model/delete",
headers=self.proxy.management_headers(),
json=ModelDeleteBody(id=model_id),
response_type=NoBody,
)
)
def connection_test(self, body: ConnectionTestBody) -> Result[ConnectionTestResponse]:
"""POST /health/test_connection, the call behind the Admin UI's Test
Connection button, probing the live provider with the supplied params."""
return self.proxy.transport.post(
"/health/test_connection",
headers=self.proxy.management_headers(),
json=body,
response_type=ConnectionTestResponse,
timeout=120.0,
)
def block_key(self, key: str) -> None:
_ = unwrap(
self.proxy.transport.post(
"/key/block",
headers=self.proxy.management_headers(),
json=KeyBlockBody(key=key),
response_type=NoBody,
)
)
def regenerate_key(self, key: str, *, grace_period: str | None = None) -> str:
return unwrap(
self.proxy.transport.post(
"/key/regenerate",
headers=self.proxy.management_headers(),
json=KeyRegenerateBody(key=key, grace_period=grace_period),
response_type=KeyGenerateResponse,
)
).key
def reset_key_spend(self, key: str, reset_to: float) -> KeyResetSpendResponse:
return unwrap(
self.proxy.transport.post(
f"/key/{key}/reset_spend",
headers=self.proxy.management_headers(),
json=KeyResetSpendBody(reset_to=reset_to),
response_type=KeyResetSpendResponse,
)
)
def key_list(self, key_alias: str, *, caller_key: str | None = None) -> Result[KeyListResponse]:
"""GET /key/list, the Virtual Keys page's own inventory call. `caller_key` is
who is asking: the master key by default, or a virtual key."""
headers = self.proxy.management_headers(caller_key)
return self.proxy.transport.get(
"/key/list",
headers=headers,
params=KeyListParams(key_alias=key_alias),
response_type=KeyListResponse,
)
def key_alias_count(self, key_alias: str) -> int:
return unwrap(self.key_list(key_alias)).total_count
def dashboard_login(self, username: str, password: str) -> DashboardSession:
"""POST /v2/login, the call the Admin UI's sign-in form makes.
The proxy authenticates the credentials, mints a UI session key for the
signed-in user, and hands it back inside a JWT signed with the master key.
Decoding that JWT is the only way to reach the session key, and it is what
the dashboard itself does before it can call a single management route."""
response = unwrap(
self.proxy.transport.post(
"/v2/login",
headers=AuthHeaders(),
json=UiLoginBody(username=username, password=password),
response_type=UiLoginResponse,
)
)
decoded: object = jwt.decode(response.token, self.master_key, algorithms=["HS256"])
claims = UiSessionClaims.model_validate(decoded)
return DashboardSession(
session_key=claims.key,
claims=claims,
redirect_url=response.redirect_url,
)
def create_team(self, body: TeamNewBody) -> str:
team_id = unwrap(
self.proxy.transport.post(
"/team/new",
headers=self.proxy.management_headers(),
json=body,
response_type=TeamNewResponse,
)
).team_id
self._wait_for_team(team_id)
return team_id
def update_team(self, body: TeamUpdateBody) -> None:
last: Result[NoBody] | None = None
for attempt in range(retry_attempts(5)):
last = self.proxy.transport.post(
"/team/update",
headers=self.proxy.management_headers(),
json=body,
response_type=NoBody,
)
match last:
case Success():
return
case UnknownApiError(body=body_text) if (
"connecting to redis" in body_text.lower() or "name resolution" in body_text.lower()
):
warnings.warn(f"Transient backend response on attempt {attempt + 1}", RuntimeWarning, stacklevel=2)
time.sleep(0.5 * (attempt + 1))
continue
case _:
break
assert last is not None
raise AssertionError(last)
def delete_team(self, team_id: str) -> None:
_ = self.proxy.transport.post(
"/team/delete",
headers=self.proxy.management_headers(),
json=TeamDeleteBody(team_ids=[team_id]),
response_type=NoBody,
)
def team_info(self, team_id: str) -> TeamData:
return unwrap(
self.proxy.transport.get(
"/team/info",
headers=self.proxy.management_headers(),
params=TeamInfoParams(team_id=team_id),
response_type=TeamInfoResponse,
)
).team_info
def team_list_ids(self) -> tuple[str, ...]:
return tuple(
entry.team_id
for entry in unwrap(
self.proxy.transport.get(
"/team/list",
headers=self.proxy.management_headers(),
params=NoBody(),
response_type=TeamListResponse,
)
).root
)
def team_info_status(self, team_id: str) -> ProbeResult:
return self.proxy.transport.probe(
"/team/info", params=TeamInfoParams(team_id=team_id), headers=self.proxy.management_headers()
)
def _wait_for_team(self, team_id: str) -> None:
last: Result[TeamInfoResponse] | None = None
for _ in range(retry_attempts(_TEAM_READY_ATTEMPTS)):
last = self.proxy.transport.get(
"/team/info",
headers=self.proxy.management_headers(),
params=TeamInfoParams(team_id=team_id),
response_type=TeamInfoResponse,
)
match last:
case Success():
return
case _:
warnings.warn("Repeating team read while the team becomes available", RuntimeWarning, stacklevel=2)
time.sleep(_TEAM_READY_SLEEP_SECONDS)
assert last is not None
raise AssertionError(last)
def add_team_member(self, team_id: str, user_id: str) -> None:
last: Result[NoBody] | None = None
for attempt in range(retry_attempts(_TEAM_READY_ATTEMPTS)):
last = self.proxy.transport.post(
"/team/member_add",
headers=self.proxy.management_headers(),
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 < retry_attempts(
_TEAM_READY_ATTEMPTS
):
warnings.warn(
"Retrying team membership while the team becomes available", RuntimeWarning, stacklevel=2
)
time.sleep(_TEAM_READY_SLEEP_SECONDS)
continue
case _:
break
assert last is not None
raise AssertionError(last)
def add_team_members(self, team_id: str, members: list[TeamMemberEntry]) -> None:
"""Bulk form of /team/member_add: `member` accepts a list, so one call
seeds a whole roster the way an admin import does."""
_ = unwrap(
self.proxy.transport.post(
"/team/member_add",
headers=self.proxy.management_headers(),
json=TeamMemberAddBody(team_id=team_id, member=members),
response_type=NoBody,
)
)
def delete_team_status(self, team_id: str) -> StreamingResponse:
"""POST /team/delete judged by HTTP outcome: the raw status and body, so a
test can assert on what a caller actually sees when the delete fails."""
return self.proxy.transport.send(
"/team/delete",
headers=self.proxy.management_headers(),
json=TeamDeleteBody(team_ids=[team_id]),
)
def delete_team_member(self, team_id: str, user_id: str) -> None:
_ = unwrap(
self.proxy.transport.post(
"/team/member_delete",
headers=self.proxy.management_headers(),
json=TeamMemberDeleteBody(team_id=team_id, user_id=user_id),
response_type=NoBody,
)
)
def create_user(self, body: UserNewBody) -> str:
return unwrap(
self.proxy.transport.post(
"/user/new",
headers=self.proxy.management_headers(),
json=body,
response_type=UserNewResponse,
)
).user_id
def create_customer(self, user_id: str) -> str:
_ = unwrap(
self.proxy.transport.post(
"/customer/new",
headers=self.proxy.management_headers(),
json=CustomerNewBody(user_id=user_id),
response_type=CustomerResponse,
)
)
return user_id
def customer_info(self, end_user_id: str) -> CustomerResponse:
return unwrap(
self.proxy.transport.get(
"/customer/info",
headers=self.proxy.management_headers(),
params=CustomerInfoParams(end_user_id=end_user_id),
response_type=CustomerResponse,
)
)
def delete_customer(self, user_id: str) -> None:
_ = self.proxy.transport.post(
"/customer/delete",
headers=self.proxy.management_headers(),
json=CustomerDeleteBody(user_ids=[user_id]),
response_type=NoBody,
)
def update_user(self, body: UserUpdateBody) -> None:
_ = unwrap(
self.proxy.transport.post(
"/user/update",
headers=self.proxy.management_headers(),
json=body,
response_type=NoBody,
)
)
def delete_user(self, user_id: str) -> None:
_ = self.proxy.transport.post(
"/user/delete",
headers=self.proxy.management_headers(),
json=UserDeleteBody(user_ids=[user_id]),
response_type=NoBody,
)
def delete_user_strict(self, user_id: str) -> None:
"""Strict delete for the act phase of a test: a failed delete is a hard
failure, unlike the warn-only delete_user used at teardown."""
_ = unwrap(
self.proxy.transport.post(
"/user/delete",
headers=self.proxy.management_headers(),
json=UserDeleteBody(user_ids=[user_id]),
response_type=UserDeleteResponse,
)
)
def user_info(self, user_id: str | None = None) -> UserInfoResponse:
return unwrap(
self.proxy.transport.get(
"/user/info",
headers=self.proxy.management_headers(),
params=UserInfoParams(user_id=user_id),
response_type=UserInfoResponse,
)
)
def user_count(self, user_id: str) -> int:
return unwrap(
self.proxy.transport.get(
"/user/list",
headers=self.proxy.management_headers(),
params=UserListParams(user_ids=user_id),
response_type=UserListResponse,
)
).total
def user_list_ids(self, user_id: str) -> tuple[str, ...]:
listing = unwrap(
self.proxy.transport.get(
"/user/list",
headers=self.proxy.management_headers(),
params=UserListParams(user_ids=user_id),
response_type=UserListResponse,
)
)
return tuple(row.user_id for row in listing.users)
def create_org(self, body: OrgNewBody) -> str:
return unwrap(
self.proxy.transport.post(
"/organization/new",
headers=self.proxy.management_headers(),
json=body,
response_type=OrgNewResponse,
)
).organization_id
def update_org(self, body: OrgUpdateBody) -> None:
_ = unwrap(
self.proxy.transport.patch(
"/organization/update",
headers=self.proxy.management_headers(),
json=body,
response_type=NoBody,
)
)
def delete_org(self, organization_id: str) -> None:
_ = self.proxy.transport.delete(
"/organization/delete",
headers=self.proxy.management_headers(),
json=OrgDeleteBody(organization_ids=[organization_id]),
response_type=NoBody,
)
def org_info(self, organization_id: str) -> OrgInfoResponse:
return unwrap(
self.proxy.transport.get(
"/organization/info",
headers=self.proxy.management_headers(),
params=OrgInfoParams(organization_id=organization_id),
response_type=OrgInfoResponse,
)
)
def org_info_status(self, organization_id: str) -> ProbeResult:
return self.proxy.transport.probe(
"/organization/info",
params=OrgInfoParams(organization_id=organization_id),
headers=self.proxy.management_headers(),
)
def create_tag(self, body: TagNewBody) -> None:
_ = unwrap(
self.proxy.transport.post(
"/tag/new",
headers=self.proxy.management_headers(),
json=body,
response_type=NoBody,
)
)
def delete_tag(self, name: str) -> None:
_ = self.proxy.transport.post(
"/tag/delete",
headers=self.proxy.management_headers(),
json=TagDeleteBody(name=name),
response_type=NoBody,
)
def tag_list(self) -> tuple[TagListEntry, ...]:
return tuple(
unwrap(
self.proxy.transport.get(
"/tag/list",
headers=self.proxy.management_headers(),
params=NoBody(),
response_type=TagListResponse,
)
).root
)
def create_mcp_server(self, body: McpServerCreateBody) -> McpServerRow:
return unwrap(
self.proxy.transport.post(
"/v1/mcp/server",
headers=self.proxy.management_headers(),
json=body,
response_type=McpServerRow,
)
)
def update_mcp_server(self, body: McpServerUpdateBody) -> McpServerRow:
"""PUT /v1/mcp/server, the call behind the dashboard's Save Changes: a partial
update where a field left unset keeps its stored value and None clears it."""
return unwrap(
self.proxy.transport.put(
"/v1/mcp/server",
headers=self.proxy.management_headers(),
json=body,
response_type=McpServerRow,
)
)
def delete_mcp_server(self, server_id: str) -> Result[NoBody]:
"""DELETE /v1/mcp/server/{server_id}. Returns the outcome so the act phase can
unwrap it while a deferred teardown can ignore an already-deleted server."""
return self.proxy.transport.delete(
f"/v1/mcp/server/{server_id}",
headers=self.proxy.management_headers(),
json=NoBody(),
response_type=NoBody,
)
def chat_status(self, key: str, model: str, content: str) -> StreamingResponse:
return self.proxy.transport.send(
"/chat/completions",
headers=self.proxy.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.proxy.transport.send("/key/generate", headers=self.proxy.transport.bearer(key), json=body)
def team_new_status(self, key: str, body: TeamNewBody) -> StreamingResponse:
return self.proxy.transport.send("/team/new", headers=self.proxy.transport.bearer(key), json=body)
def user_new_status(self, key: str, body: UserNewBody) -> StreamingResponse:
return self.proxy.transport.send("/user/new", headers=self.proxy.transport.bearer(key), json=body)
def build_client(proxy: ProxyClient) -> ManagementClient:
return ManagementClient(proxy=proxy, master_key=MASTER_KEY)