Coverage for .venv/lib/python3.13/site-packages/litellm/proxy/management_helpers/utils.py: 53%
292 statements
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-10 12:01 +0000
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-10 12:01 +0000
1# What is this?
2## Helper utils for the management endpoints (keys/users/teams)
3from collections.abc import Callable, Mapping, MutableMapping, Sequence
4from datetime import datetime
5from functools import wraps
6from types import MappingProxyType
7from typing import Any, Final, Protocol
9from fastapi import HTTPException, Request
10from pydantic import BaseModel
12import litellm
13from litellm._logging import verbose_logger
14from litellm._uuid import uuid
15from litellm.integrations.otel.model.config import is_otel_v2_enabled
16from litellm.proxy._types import ( # key request types; user request types; team request types; customer request types
17 BudgetNewRequest,
18 DeleteCustomerRequest,
19 DeleteTeamRequest,
20 DeleteUserRequest,
21 KeyRequest,
22 LiteLLM_BudgetTable,
23 LiteLLM_TeamMembership,
24 LiteLLM_UserTable,
25 ManagementEndpointLoggingPayload,
26 Member,
27 Span,
28 SSOUserDefinedValues,
29 UpdateCustomerRequest,
30 UpdateKeyRequest,
31 UpdateTeamRequest,
32 UpdateUserRequest,
33 UserAPIKeyAuth,
34 VirtualKeyEvent,
35)
36from litellm.proxy.common_utils.http_parsing_utils import _read_request_body
37from litellm.proxy.common_utils.timezone_utils import get_budget_reset_time
38from litellm.proxy.utils import PrismaClient, jsonify_object
39from litellm.repositories.budget_repository import BudgetRepository
40from litellm.repositories.table_repositories import TeamMembershipRepository
41from litellm.repositories.user_repository import UserRepository
44class _PrismaRecord(Protocol):
45 """Row surface the management helpers read back from Prisma."""
47 def model_dump(self) -> Mapping[str, object]: ... 47 ↛ exitline 47 didn't return from function 'model_dump' because
50class _PrismaUserRecord(Protocol):
51 """User row surface the management helpers read back from Prisma."""
53 user_id: str
55 def model_dump(self) -> Mapping[str, object]: ... 55 ↛ exitline 55 didn't return from function 'model_dump' because
58class _PrismaBudgetRecord(Protocol):
59 """Budget row surface the management helpers read back from Prisma."""
61 budget_id: str
63 def model_dump(self) -> Mapping[str, object]: ... 63 ↛ exitline 63 didn't return from function 'model_dump' because
66class _PrismaBudgetTable(Protocol):
67 """Budget table actions the management helpers issue."""
69 async def create(self, *, data: Mapping[str, object]) -> _PrismaBudgetRecord: ... 69 ↛ exitline 69 didn't return from function 'create' because
71 async def find_unique(self, *, where: Mapping[str, object]) -> _PrismaBudgetRecord | None: ... 71 ↛ exitline 71 didn't return from function 'find_unique' because
74class _PrismaUserTable(Protocol):
75 """User table actions the management helpers issue."""
77 async def update_many(self, *, where: Mapping[str, object], data: Mapping[str, object]) -> int: ... 77 ↛ exitline 77 didn't return from function 'update_many' because
79 async def upsert( 79 ↛ exitline 79 didn't return from function 'upsert' because
80 self, *, where: Mapping[str, object], data: Mapping[str, Mapping[str, object]]
81 ) -> _PrismaUserRecord | None: ...
83 async def find_many(self, *, where: Mapping[str, object]) -> Sequence[_PrismaUserRecord]: ... 83 ↛ exitline 83 didn't return from function 'find_many' because
86class _PrismaTeamMembershipTable(Protocol):
87 """Team membership table actions the management helpers issue."""
89 async def upsert( 89 ↛ exitline 89 didn't return from function 'upsert' because
90 self, *, where: Mapping[str, object], data: Mapping[str, Mapping[str, object]], include: Mapping[str, bool]
91 ) -> _PrismaRecord: ...
94class MemberWriteTx(Protocol):
95 """Transaction surface `add_new_member` writes through when the caller owns one.
97 A caller already holding a transaction, and with it a pooled connection plus that
98 transaction's locks, passes it here so these writes reuse that connection rather than
99 checking out another one that lock waiters may already have drained from the pool.
100 """
102 @property
103 def litellm_usertable(self) -> _PrismaUserTable: ... 103 ↛ exitline 103 didn't return from function 'litellm_usertable' because
105 @property
106 def litellm_budgettable(self) -> _PrismaBudgetTable: ... 106 ↛ exitline 106 didn't return from function 'litellm_budgettable' because
108 @property
109 def litellm_teammembership(self) -> _PrismaTeamMembershipTable: ... 109 ↛ exitline 109 didn't return from function 'litellm_teammembership' because
112def _user_table(prisma_client: PrismaClient, tx: MemberWriteTx | None) -> _PrismaUserTable:
113 return tx.litellm_usertable if tx is not None else UserRepository(prisma_client).table
116def _budget_table(prisma_client: PrismaClient, tx: MemberWriteTx | None) -> _PrismaBudgetTable:
117 return tx.litellm_budgettable if tx is not None else BudgetRepository(prisma_client).table
120def _team_membership_table(prisma_client: PrismaClient, tx: MemberWriteTx | None) -> _PrismaTeamMembershipTable:
121 return tx.litellm_teammembership if tx is not None else TeamMembershipRepository(prisma_client).table
124async def _find_users_by_email(
125 prisma_client: PrismaClient, tx: MemberWriteTx | None, user_email: str
126) -> Sequence[_PrismaUserRecord]:
127 if tx is not None:
128 return await tx.litellm_usertable.find_many(where={"user_email": user_email})
129 rows: Final[Sequence[_PrismaUserRecord] | None] = await prisma_client.get_data(
130 key_val={"user_email": user_email},
131 table_name="user",
132 query_type="find_all",
133 )
134 return rows if rows is not None else ()
137async def _upsert_user_row(
138 user_table: _PrismaUserTable, user_id: str, create_data: Mapping[str, object]
139) -> _PrismaUserRecord | None:
140 """Insert the user row if it is absent, leaving an existing row as it is.
142 Upserting keeps concurrent provisioning of the same new user from racing on create.
143 The update branch re-states user_id rather than being empty because Prisma only
144 compiles an upsert down to INSERT ... ON CONFLICT when the update is non-empty, and
145 otherwise falls back to a racy SELECT-then-INSERT.
146 """
147 return await user_table.upsert(
148 where={"user_id": user_id},
149 data={"create": create_data, "update": {"user_id": user_id}},
150 )
153async def _create_user_row(
154 prisma_client: PrismaClient, tx: MemberWriteTx | None, user_data: dict[str, object]
155) -> _PrismaUserRecord | None:
156 if tx is not None:
157 return await _upsert_user_row(tx.litellm_usertable, str(user_data["user_id"]), jsonify_object(user_data))
158 return await prisma_client.insert_data(data=user_data, table_name="user")
161def get_new_internal_user_defaults(user_id: str, user_email: str | None = None) -> dict[str, object]:
162 user_info: Final = litellm.default_internal_user_params or {}
164 returned_dict: Final[SSOUserDefinedValues] = {
165 "models": user_info.get("models") or [],
166 "max_budget": user_info.get("max_budget", litellm.max_internal_user_budget),
167 "budget_duration": user_info.get("budget_duration", litellm.internal_user_budget_duration),
168 "user_email": user_email or user_info.get("user_email", None),
169 "user_id": user_id,
170 "user_role": "internal_user",
171 }
173 non_null_dict: Final = {}
174 for k, v in returned_dict.items():
175 if v is not None:
176 non_null_dict[k] = v
177 return non_null_dict
180async def handle_budget_for_entity(
181 data,
182 existing_budget_id: str | None,
183 user_api_key_dict: UserAPIKeyAuth,
184 prisma_client: PrismaClient,
185 litellm_proxy_admin_name: str,
186 budget_duration_cleared: bool = False,
187) -> str | None:
188 """
189 Common helper to handle budget creation/updates for entities (organizations, tags, etc).
191 This function:
192 1. Creates a new budget if budget_id is None but budget fields are provided
193 2. Updates an existing budget if budget fields are provided and budget_id exists
194 3. Returns the budget_id to use (existing or newly created)
196 Args:
197 data: The request object (e.g., TagNewRequest, NewOrganizationRequest, etc.) containing budget fields
198 existing_budget_id: The existing budget_id if updating an entity, None if creating new
199 user_api_key_dict: User authentication info
200 prisma_client: Database client
201 litellm_proxy_admin_name: Admin name for audit trail
203 Returns:
204 Optional[str]: The budget_id to use, or None if no budget was created/updated
205 """
206 from litellm.proxy.management_endpoints.budget_management_endpoints import (
207 update_budget,
208 )
210 # Get all budget field names
211 budget_params: Final = LiteLLM_BudgetTable.model_fields.keys()
213 # Extract budget fields from data
214 _json_data: Final = data.model_dump(exclude_none=True) if hasattr(data, "model_dump") else data
215 _budget_data: Final = MappingProxyType(
216 {
217 k: _json_data.get(k)
218 for k in budget_params
219 if k in _json_data
220 or (k == "budget_duration" and existing_budget_id is not None and budget_duration_cleared)
221 }
222 )
224 # Check if budget_id is explicitly provided in the data
225 data_budget_id: Final[str | None] = getattr(data, "budget_id", None)
227 # Case 1: Creating new entity - no existing budget_id
228 if existing_budget_id is None:
229 if data_budget_id is not None:
230 # Use the provided budget_id
231 return data_budget_id
232 elif _budget_data: 232 ↛ 234line 232 didn't jump to line 234 because the condition on line 232 was never true
233 # Create a new budget with the provided fields
234 budget_row: Final = LiteLLM_BudgetTable(**_budget_data)
235 new_budget_data: Final = prisma_client.jsonify_object(budget_row.model_dump(exclude_none=True))
237 _budget: Final[_PrismaBudgetRecord] = await BudgetRepository(prisma_client).table.create(
238 data={
239 **new_budget_data,
240 "created_by": user_api_key_dict.user_id or litellm_proxy_admin_name,
241 "updated_by": user_api_key_dict.user_id or litellm_proxy_admin_name,
242 }
243 )
245 return _budget.budget_id
246 else:
247 # No budget fields provided, no budget to create
248 return None
250 # Case 2: Updating existing entity - has existing budget_id
251 else:
252 # If budget fields are provided, update the existing budget
253 if _budget_data:
254 await update_budget(
255 budget_obj=BudgetNewRequest(budget_id=existing_budget_id, **_budget_data),
256 user_api_key_dict=user_api_key_dict,
257 )
259 # If a different budget_id is explicitly provided, use that instead
260 if data_budget_id is not None and data_budget_id != existing_budget_id: 260 ↛ 261line 260 didn't jump to line 261 because the condition on line 260 was never true
261 return data_budget_id
263 # Otherwise, keep using the existing budget_id
264 return existing_budget_id
267# Fields on LiteLLM_BudgetTable that represent the budget's *configuration*
268# (i.e. the values an admin sets). We copy these when cloning a team's
269# default member-budget into an individual member-budget so that the new
270# row starts with the same limits as the default.
271_CLONABLE_BUDGET_FIELDS: Final[tuple[str, ...]] = (
272 "max_budget",
273 "soft_budget",
274 "max_parallel_requests",
275 "tpm_limit",
276 "rpm_limit",
277 "model_max_budget",
278 "budget_duration",
279 "allowed_models",
280)
283async def _clone_team_default_budget_for_member(
284 prisma_client: PrismaClient,
285 default_team_budget_id: str,
286 user_api_key_dict: UserAPIKeyAuth,
287 litellm_proxy_admin_name: str,
288 budget_duration_override: str | None = None,
289 tx: MemberWriteTx | None = None,
290) -> str | None:
291 """
292 Create a new budget row that copies the values from the team's default
293 member budget. Returns the new budget_id, or None if the default budget
294 no longer exists in the DB.
296 Used when adding a new team member with a per-member ``budget_duration``
297 but no other per-member limit, so the member keeps the team default's
298 values in their own private budget row while the reset window differs.
300 ``budget_duration_override`` replaces the default's reset window for this
301 member while keeping the default's other limits, so an admin can set a
302 member's reset cadence without discarding the team default's max_budget.
303 """
304 budget_table: Final[_PrismaBudgetTable] = _budget_table(prisma_client, tx)
305 default_budget: Final = await budget_table.find_unique(where={"budget_id": default_team_budget_id})
306 if default_budget is None:
307 return None
309 default_budget_dict: Final = default_budget.model_dump()
310 cloned_data: Final[dict] = {
311 "created_by": user_api_key_dict.user_id or litellm_proxy_admin_name,
312 "updated_by": user_api_key_dict.user_id or litellm_proxy_admin_name,
313 }
314 for field in _CLONABLE_BUDGET_FIELDS:
315 value = default_budget_dict.get(field)
316 if value is None:
317 continue
318 # Skip empty list defaults (e.g. allowed_models = []) so the cloned
319 # row matches the "no value set" shape rather than carrying a default.
320 if isinstance(value, list) and len(value) == 0:
321 continue
322 cloned_data[field] = value
324 if budget_duration_override is not None:
325 cloned_data["budget_duration"] = budget_duration_override
327 # Start the member's budget window at clone time, not the pool's reset
328 # timestamp — otherwise a member joining mid-cycle inherits a stale reset.
329 if cloned_data.get("budget_duration"):
330 cloned_data["budget_reset_at"] = get_budget_reset_time(cloned_data["budget_duration"])
332 new_budget: Final[_PrismaBudgetRecord] = await budget_table.create(data=cloned_data)
333 return new_budget.budget_id
336async def _resolve_member_budget_id(
337 prisma_client: PrismaClient,
338 user_api_key_dict: UserAPIKeyAuth,
339 litellm_proxy_admin_name: str,
340 max_budget_in_team: float | None,
341 allowed_models: list[str] | None,
342 budget_duration: str | None,
343 default_team_budget_id: str | None,
344 tx: MemberWriteTx | None = None,
345) -> str | None:
346 """
347 Resolve the budget a new team member should be linked to.
349 Explicit per-member limits create a fresh budget. Otherwise the member is
350 linked to the team's shared default member budget, so later ``/team/update``
351 changes reach them; ``/team/member_update`` clones that row on first write.
352 A lone ``budget_duration`` clones the default with the reset window
353 overridden, or creates a window-only budget when there is no team default.
354 With nothing set the member gets no budget, though ``add_new_member`` still writes its membership row.
355 """
356 has_explicit_limit: Final = max_budget_in_team is not None or allowed_models is not None
358 if not has_explicit_limit and default_team_budget_id is not None and budget_duration is None:
359 default_budget: Final = await _budget_table(prisma_client, tx).find_unique(
360 where={"budget_id": default_team_budget_id}
361 )
362 return default_team_budget_id if default_budget is not None else None
364 if not has_explicit_limit and default_team_budget_id is not None: 364 ↛ 365line 364 didn't jump to line 365 because the condition on line 364 was never true
365 return await _clone_team_default_budget_for_member(
366 prisma_client=prisma_client,
367 default_team_budget_id=default_team_budget_id,
368 user_api_key_dict=user_api_key_dict,
369 litellm_proxy_admin_name=litellm_proxy_admin_name,
370 budget_duration_override=budget_duration,
371 tx=tx,
372 )
374 if not has_explicit_limit and budget_duration is None: 374 ↛ 377line 374 didn't jump to line 377 because the condition on line 374 was always true
375 return None
377 budget_data: Final[dict[str, object]] = {
378 "created_by": user_api_key_dict.user_id or litellm_proxy_admin_name,
379 "updated_by": user_api_key_dict.user_id or litellm_proxy_admin_name,
380 }
381 if max_budget_in_team is not None:
382 budget_data["max_budget"] = max_budget_in_team
383 if allowed_models is not None:
384 budget_data["allowed_models"] = allowed_models
385 if budget_duration is not None:
386 budget_data["budget_duration"] = budget_duration
387 budget_data["budget_reset_at"] = get_budget_reset_time(budget_duration=budget_duration)
388 budget_table: Final[_PrismaBudgetTable] = _budget_table(prisma_client, tx)
389 response: Final = await budget_table.create(data=budget_data)
390 return response.budget_id
393async def _append_team_id_if_absent(
394 prisma_client: PrismaClient, user_id: str, team_id: str, tx: MemberWriteTx | None = None
395) -> None:
396 """Append team_id to a user's teams array, only if it is not already present.
398 The row-level filter makes the append a no-op once the team is present, so
399 repeated or concurrent adds of the same team cannot accumulate duplicate
400 team ids in user.teams (a duplicate also breaks auth logic that keys off the
401 number of teams a user belongs to). Teams added concurrently for a different
402 team id are unaffected, since each update filters on its own team id.
403 """
404 user_table: Final[_PrismaUserTable] = _user_table(prisma_client, tx)
405 await user_table.update_many(
406 where={"user_id": user_id, "NOT": {"teams": {"has": team_id}}},
407 data={"teams": {"push": [team_id]}},
408 )
411async def add_new_member(
412 new_member: Member,
413 max_budget_in_team: float | None,
414 prisma_client: PrismaClient,
415 team_id: str,
416 user_api_key_dict: UserAPIKeyAuth,
417 litellm_proxy_admin_name: str,
418 default_team_budget_id: str | None = None,
419 allowed_models: list[str] | None = None,
420 budget_duration: str | None = None,
421 tx: MemberWriteTx | None = None,
422) -> tuple[LiteLLM_UserTable, LiteLLM_TeamMembership | None]:
423 """
424 Add a new member to a team
426 - add team id to user table
427 - add team member to team member table, linked to a budget when one resolves
429 Returns created/existing user + team membership (``budget_id`` is ``None`` when no budget applies)
431 Callers already inside a transaction pass it as ``tx`` so every write here runs on that
432 connection instead of borrowing more from the pool while the caller's locks are held.
433 """
434 returned_user: LiteLLM_UserTable | None = None
435 returned_team_membership: LiteLLM_TeamMembership | None = None
436 ## ADD TEAM ID, to USER TABLE IF NEW ##
437 if new_member.user_id is not None: 437 ↛ 449line 437 didn't jump to line 449 because the condition on line 437 was always true
438 new_user_defaults = get_new_internal_user_defaults(user_id=new_member.user_id)
439 # The teams append lives in the filtered update below rather than the upsert's
440 # update branch so an already-existing user does not get a duplicate team id.
441 _returned_user: _PrismaUserRecord | None = await _upsert_user_row(
442 _user_table(prisma_client, tx),
443 new_member.user_id,
444 {"teams": [team_id], **new_user_defaults},
445 )
446 await _append_team_id_if_absent(prisma_client, new_member.user_id, team_id, tx)
447 if _returned_user is not None: 447 ↛ 472line 447 didn't jump to line 472 because the condition on line 447 was always true
448 returned_user = LiteLLM_UserTable.model_validate(_returned_user.model_dump())
449 elif new_member.user_email is not None:
450 new_user_defaults = get_new_internal_user_defaults(user_id=str(uuid.uuid4()), user_email=new_member.user_email)
451 ## user email is not unique acc. to prisma schema -> future improvement
452 ### for now: check if it exists in db, if not - insert it
453 existing_user_row: Final[Sequence[_PrismaUserRecord]] = await _find_users_by_email(
454 prisma_client, tx, new_member.user_email
455 )
456 if len(existing_user_row) == 0:
457 new_user_defaults["teams"] = [team_id]
458 _returned_user = await _create_user_row(prisma_client, tx, new_user_defaults)
460 if _returned_user is not None:
461 returned_user = LiteLLM_UserTable.model_validate(_returned_user.model_dump())
462 elif len(existing_user_row) == 1:
463 user_info: Final = existing_user_row[0]
464 await _append_team_id_if_absent(prisma_client, user_info.user_id, team_id, tx)
465 returned_user = LiteLLM_UserTable.model_validate(user_info.model_dump())
466 elif len(existing_user_row) > 1:
467 raise HTTPException(
468 status_code=400,
469 detail={"error": "Multiple users with this email found in db. Please use 'user_id' instead."},
470 )
472 _budget_id: Final = await _resolve_member_budget_id(
473 prisma_client=prisma_client,
474 user_api_key_dict=user_api_key_dict,
475 litellm_proxy_admin_name=litellm_proxy_admin_name,
476 max_budget_in_team=max_budget_in_team,
477 allowed_models=allowed_models,
478 budget_duration=budget_duration,
479 default_team_budget_id=default_team_budget_id,
480 tx=tx,
481 )
483 if returned_user is not None and returned_user.user_id is not None: 483 ↛ 497line 483 didn't jump to line 497 because the condition on line 483 was always true
484 membership_table: Final[_PrismaTeamMembershipTable] = _team_membership_table(prisma_client, tx)
485 membership_key: Final[Mapping[str, object]] = {"user_id": returned_user.user_id, "team_id": team_id}
486 budget_link: Final[Mapping[str, str]] = (
487 MappingProxyType({"budget_id": _budget_id}) if _budget_id is not None else MappingProxyType({})
488 )
489 _returned_team_membership: Final = await membership_table.upsert(
490 where={"user_id_team_id": membership_key},
491 data={"create": {**membership_key, **budget_link}, "update": {}},
492 include={"litellm_budget_table": True},
493 )
495 returned_team_membership = LiteLLM_TeamMembership.model_validate(_returned_team_membership.model_dump())
497 if returned_user is None: 497 ↛ 498line 497 didn't jump to line 498 because the condition on line 497 was never true
498 raise Exception("Unable to update user table with membership information!")
500 return returned_user, returned_team_membership
503def _delete_user_id_from_cache(kwargs):
504 from litellm.proxy.proxy_server import user_api_key_cache
506 if kwargs.get("data") is not None:
507 update_user_request: Final = kwargs.get("data")
508 if isinstance(update_user_request, UpdateUserRequest): 508 ↛ 509line 508 didn't jump to line 509 because the condition on line 508 was never true
509 user_api_key_cache.delete_cache(key=update_user_request.user_id)
511 # delete user request
512 if isinstance(update_user_request, DeleteUserRequest):
513 for user_id in update_user_request.user_ids: 513 ↛ 514line 513 didn't jump to line 514 because the loop on line 513 never started
514 user_api_key_cache.delete_cache(key=user_id)
517def _delete_api_key_from_cache(kwargs):
518 from litellm.proxy.proxy_server import user_api_key_cache
520 if kwargs.get("data") is not None:
521 update_request: Final = kwargs.get("data")
522 if isinstance(update_request, UpdateKeyRequest): 522 ↛ 523line 522 didn't jump to line 523 because the condition on line 522 was never true
523 user_api_key_cache.delete_cache(key=update_request.key)
525 # delete key request
526 if isinstance(update_request, KeyRequest) and update_request.keys: 526 ↛ 527line 526 didn't jump to line 527 because the condition on line 526 was never true
527 for key in update_request.keys:
528 user_api_key_cache.delete_cache(key=key)
531def _delete_team_id_from_cache(kwargs):
532 from litellm.proxy.proxy_server import user_api_key_cache
534 if kwargs.get("data") is not None:
535 update_request: Final = kwargs.get("data")
536 if isinstance(update_request, UpdateTeamRequest):
537 user_api_key_cache.delete_cache(key=update_request.team_id)
539 # delete team request
540 if isinstance(update_request, DeleteTeamRequest):
541 for team_id in update_request.team_ids: 541 ↛ 542line 541 didn't jump to line 542 because the loop on line 541 never started
542 user_api_key_cache.delete_cache(key=team_id)
545def _delete_customer_id_from_cache(kwargs):
546 from litellm.proxy.proxy_server import user_api_key_cache
548 if kwargs.get("data") is not None:
549 update_request: Final = kwargs.get("data")
550 if isinstance(update_request, UpdateCustomerRequest): 550 ↛ 551line 550 didn't jump to line 551 because the condition on line 550 was never true
551 user_api_key_cache.delete_cache(key=update_request.user_id)
553 # delete customer request
554 if isinstance(update_request, DeleteCustomerRequest): 554 ↛ 555line 554 didn't jump to line 555 because the condition on line 554 was never true
555 for user_id in update_request.user_ids:
556 user_api_key_cache.delete_cache(key=user_id)
559async def send_management_endpoint_alert(
560 request_kwargs: dict,
561 user_api_key_dict: UserAPIKeyAuth,
562 function_name: str,
563):
564 """
565 Sends a slack alert when:
566 - A virtual key is created, updated, or deleted
567 - An internal user is created, updated, or deleted
568 - A team is created, updated, or deleted
569 """
570 from litellm.proxy.proxy_server import proxy_logging_obj
571 from litellm.types.integrations.slack_alerting import AlertType
573 management_function_to_event_name: Final = {
574 "generate_key_fn": AlertType.new_virtual_key_created,
575 "update_key_fn": AlertType.virtual_key_updated,
576 "delete_key_fn": AlertType.virtual_key_deleted,
577 # Team events
578 "new_team": AlertType.new_team_created,
579 "update_team": AlertType.team_updated,
580 "delete_team": AlertType.team_deleted,
581 # Internal User events
582 "new_user": AlertType.new_internal_user_created,
583 "user_update": AlertType.internal_user_updated,
584 "delete_user": AlertType.internal_user_deleted,
585 }
587 # Check if alerting is enabled
588 if proxy_logging_obj is not None and proxy_logging_obj.slack_alerting_instance is not None: 588 ↛ exitline 588 didn't return from function 'send_management_endpoint_alert' because the condition on line 588 was always true
589 # Virtual Key Events
590 if function_name in management_function_to_event_name:
591 _event_name: Final[AlertType] = management_function_to_event_name[function_name]
593 key_event: Final = VirtualKeyEvent(
594 created_by_user_id=user_api_key_dict.user_id or "Unknown",
595 created_by_user_role=user_api_key_dict.user_role or "Unknown",
596 created_by_key_alias=user_api_key_dict.key_alias,
597 request_kwargs=request_kwargs,
598 )
600 # replace all "_" with " " and capitalize
601 event_name: Final = _event_name.replace("_", " ").title()
602 await proxy_logging_obj.slack_alerting_instance.send_virtual_key_event_slack(
603 key_event=key_event,
604 event_name=event_name,
605 alert_type=_event_name,
606 )
609def _object_mapping(value: object) -> Mapping[str, object] | None:
610 """Return ``value`` as an opaque mapping when it is a dict."""
611 return value if isinstance(value, dict) else None
614def _object_list(value: object) -> Sequence[object] | None:
615 """Return ``value`` as an opaque sequence when it is a list."""
616 return value if isinstance(value, list) else None
619def _redacted_env_var(entry: object) -> dict[str, object]:
620 get: Final[Callable[[str], object]] = entry.get if isinstance(entry, dict) else lambda k: getattr(entry, k, None)
621 return {
622 "name": get("name"),
623 "scope": get("scope"),
624 "description": get("description"),
625 "value": "",
626 }
629def _redact_record_env_vars(record: object) -> object:
630 """Return ``record`` with its ``env_vars[].value`` blanked.
632 Copies rather than mutating, because the record aliases the live response
633 object that is also returned to the caller. Records without an ``env_vars``
634 list are returned unchanged.
635 """
636 record_map: Final = _object_mapping(record)
637 env_vars: Final = _object_list(
638 record_map.get("env_vars") if record_map is not None else getattr(record, "env_vars", None)
639 )
640 if env_vars is None:
641 return record
642 redacted: Final = [_redacted_env_var(entry) for entry in env_vars]
643 if record_map is not None:
644 return {**record_map, "env_vars": redacted}
645 if isinstance(record, BaseModel):
646 return record.model_copy(update={"env_vars": redacted})
647 return record
650def _redact_env_var_values(response: MutableMapping[str, object]) -> None:
651 """Blank ``env_vars[].value`` in a management response before telemetry.
653 MCP endpoints return decrypted ``scope="global"`` env var values so the admin
654 UI can pre-fill the edit form; those values are upstream credentials and must
655 not be serialized verbatim into OTEL spans, where an observability user could
656 read them. The values surface both at the top level (single-server
657 create/update) and nested under ``items`` (the submissions queue), so both are
658 scrubbed. Names, scopes, and descriptions are kept so traces stay useful.
659 """
660 env_vars: Final = _object_list(response.get("env_vars"))
661 if env_vars is not None:
662 response["env_vars"] = [_redacted_env_var(entry) for entry in env_vars]
664 items: Final = _object_list(response.get("items"))
665 if items is not None:
666 response["items"] = [_redact_record_env_vars(item) for item in items]
669async def _emit_management_endpoint_otel_span(
670 func: Callable,
671 kwargs: dict,
672 parent_otel_span: Span | None,
673 start_time: datetime,
674 end_time: datetime,
675 result: Any = None,
676 exception: Exception | None = None,
677) -> None:
678 """Stamp + end the parent OTEL SERVER span for a management endpoint.
680 Routes the request/response (or exception) through the OTEL success/failure
681 hook. Falls back to ``func.__name__`` for the route when the handler has no
682 ``http_request`` param — endpoints like ``/key/generate`` never receive one,
683 and gating the hook on it leaked their SERVER span (created in auth, never
684 ended → never exported). Always emitting keeps both success and failure
685 paths consistent.
686 """
687 from litellm.proxy.proxy_server import open_telemetry_logger
689 if open_telemetry_logger is None:
690 return
692 # Under V2 OTel, management endpoints are ordinary FastAPI routes already
693 # spanned by the mounted instrumentor — there is no management hook to fire, so
694 # skip the payload build entirely. The legacy logger still needs the hook.
695 if is_otel_v2_enabled():
696 return
698 http_request: Final[Request | None] = kwargs.get("http_request")
699 if http_request is not None:
700 # Inline import — auth_utils participates in a proxy import cycle.
701 from litellm.proxy.auth.auth_utils import ( # noqa: PLC0415
702 get_request_route,
703 )
705 route = get_request_route(http_request)
706 request_body: dict = await _read_request_body(request=http_request)
707 else:
708 route = func.__name__
709 request_body = {}
711 _CREDENTIAL_FIELDS: Final = frozenset(
712 {
713 "key",
714 "token",
715 "api_key",
716 "secret",
717 "password",
718 "access_token",
719 "refresh_token",
720 "private_key",
721 "service_account_key",
722 }
723 )
725 _response: dict[str, object] | None = None
726 if exception is None and result is not None:
727 try:
728 raw: Final[Mapping[str, object]] = dict(result)
729 _response = {k: v for k, v in raw.items() if k not in _CREDENTIAL_FIELDS}
730 _redact_env_var_values(_response)
731 except Exception:
732 _response = None
734 logging_payload: Final = ManagementEndpointLoggingPayload(
735 route=route,
736 request_data=request_body,
737 response=_response,
738 start_time=start_time,
739 end_time=end_time,
740 exception=exception,
741 )
743 if exception is None:
744 await open_telemetry_logger.async_management_endpoint_success_hook(
745 logging_payload=logging_payload,
746 parent_otel_span=parent_otel_span,
747 )
748 else:
749 await open_telemetry_logger.async_management_endpoint_failure_hook(
750 logging_payload=logging_payload,
751 parent_otel_span=parent_otel_span,
752 )
755def management_endpoint_wrapper(func):
756 """
757 This wrapper does the following:
759 1. Log I/O, Exceptions to OTEL
760 2. Create an Audit log for success calls
761 """
763 @wraps(func)
764 async def wrapper(*args, **kwargs):
765 start_time: Final = datetime.now()
766 try:
767 result: Final = await func(*args, **kwargs)
768 end_time = datetime.now()
769 try:
770 user_api_key_dict: UserAPIKeyAuth = kwargs.get("user_api_key_dict") or UserAPIKeyAuth()
772 await send_management_endpoint_alert(
773 request_kwargs=kwargs,
774 user_api_key_dict=user_api_key_dict,
775 function_name=func.__name__,
776 )
777 parent_otel_span: Span | None = getattr(user_api_key_dict, "parent_otel_span", None)
778 if parent_otel_span is not None: 778 ↛ 779line 778 didn't jump to line 779 because the condition on line 778 was never true
779 await _emit_management_endpoint_otel_span(
780 func=func,
781 kwargs=kwargs,
782 parent_otel_span=parent_otel_span,
783 start_time=start_time,
784 end_time=end_time,
785 result=result,
786 )
788 # Delete updated/deleted info from cache
789 _delete_api_key_from_cache(kwargs=kwargs)
790 _delete_user_id_from_cache(kwargs=kwargs)
791 _delete_team_id_from_cache(kwargs=kwargs)
792 _delete_customer_id_from_cache(kwargs=kwargs)
793 except Exception as e:
794 # Non-Blocking Exception
795 verbose_logger.debug("Error in management endpoint wrapper: %s", str(e))
797 return result
798 except Exception as e:
799 end_time = datetime.now()
801 user_api_key_dict: UserAPIKeyAuth = kwargs.get("user_api_key_dict") or UserAPIKeyAuth()
802 parent_otel_span = getattr(user_api_key_dict, "parent_otel_span", None)
803 if parent_otel_span is not None: 803 ↛ 804line 803 didn't jump to line 804 because the condition on line 803 was never true
804 try:
805 await _emit_management_endpoint_otel_span(
806 func=func,
807 kwargs=kwargs,
808 parent_otel_span=parent_otel_span,
809 start_time=start_time,
810 end_time=end_time,
811 exception=e,
812 )
813 except Exception as otel_exc:
814 # Non-Blocking Exception - never let OTEL failures swallow
815 # the original management-endpoint exception.
816 verbose_logger.debug(
817 "Error emitting OTEL span in management endpoint wrapper failure path: %s",
818 str(otel_exc),
819 )
821 raise e
823 return wrapper