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

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 

8 

9from fastapi import HTTPException, Request 

10from pydantic import BaseModel 

11 

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 

42 

43 

44class _PrismaRecord(Protocol): 

45 """Row surface the management helpers read back from Prisma.""" 

46 

47 def model_dump(self) -> Mapping[str, object]: ... 47 ↛ exitline 47 didn't return from function 'model_dump' because

48 

49 

50class _PrismaUserRecord(Protocol): 

51 """User row surface the management helpers read back from Prisma.""" 

52 

53 user_id: str 

54 

55 def model_dump(self) -> Mapping[str, object]: ... 55 ↛ exitline 55 didn't return from function 'model_dump' because

56 

57 

58class _PrismaBudgetRecord(Protocol): 

59 """Budget row surface the management helpers read back from Prisma.""" 

60 

61 budget_id: str 

62 

63 def model_dump(self) -> Mapping[str, object]: ... 63 ↛ exitline 63 didn't return from function 'model_dump' because

64 

65 

66class _PrismaBudgetTable(Protocol): 

67 """Budget table actions the management helpers issue.""" 

68 

69 async def create(self, *, data: Mapping[str, object]) -> _PrismaBudgetRecord: ... 69 ↛ exitline 69 didn't return from function 'create' because

70 

71 async def find_unique(self, *, where: Mapping[str, object]) -> _PrismaBudgetRecord | None: ... 71 ↛ exitline 71 didn't return from function 'find_unique' because

72 

73 

74class _PrismaUserTable(Protocol): 

75 """User table actions the management helpers issue.""" 

76 

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

78 

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: ... 

82 

83 async def find_many(self, *, where: Mapping[str, object]) -> Sequence[_PrismaUserRecord]: ... 83 ↛ exitline 83 didn't return from function 'find_many' because

84 

85 

86class _PrismaTeamMembershipTable(Protocol): 

87 """Team membership table actions the management helpers issue.""" 

88 

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: ... 

92 

93 

94class MemberWriteTx(Protocol): 

95 """Transaction surface `add_new_member` writes through when the caller owns one. 

96 

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 """ 

101 

102 @property 

103 def litellm_usertable(self) -> _PrismaUserTable: ... 103 ↛ exitline 103 didn't return from function 'litellm_usertable' because

104 

105 @property 

106 def litellm_budgettable(self) -> _PrismaBudgetTable: ... 106 ↛ exitline 106 didn't return from function 'litellm_budgettable' because

107 

108 @property 

109 def litellm_teammembership(self) -> _PrismaTeamMembershipTable: ... 109 ↛ exitline 109 didn't return from function 'litellm_teammembership' because

110 

111 

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 

114 

115 

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 

118 

119 

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 

122 

123 

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

135 

136 

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. 

141 

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 ) 

151 

152 

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

159 

160 

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 {} 

163 

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 } 

172 

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 

178 

179 

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

190 

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) 

195 

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 

202 

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 ) 

209 

210 # Get all budget field names 

211 budget_params: Final = LiteLLM_BudgetTable.model_fields.keys() 

212 

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 ) 

223 

224 # Check if budget_id is explicitly provided in the data 

225 data_budget_id: Final[str | None] = getattr(data, "budget_id", None) 

226 

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

236 

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 ) 

244 

245 return _budget.budget_id 

246 else: 

247 # No budget fields provided, no budget to create 

248 return None 

249 

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 ) 

258 

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 

262 

263 # Otherwise, keep using the existing budget_id 

264 return existing_budget_id 

265 

266 

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) 

281 

282 

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. 

295 

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. 

299 

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 

308 

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 

323 

324 if budget_duration_override is not None: 

325 cloned_data["budget_duration"] = budget_duration_override 

326 

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"]) 

331 

332 new_budget: Final[_PrismaBudgetRecord] = await budget_table.create(data=cloned_data) 

333 return new_budget.budget_id 

334 

335 

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. 

348 

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 

357 

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 

363 

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 ) 

373 

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 

376 

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 

391 

392 

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. 

397 

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 ) 

409 

410 

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 

425 

426 - add team id to user table 

427 - add team member to team member table, linked to a budget when one resolves 

428 

429 Returns created/existing user + team membership (``budget_id`` is ``None`` when no budget applies) 

430 

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) 

459 

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 ) 

471 

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 ) 

482 

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 ) 

494 

495 returned_team_membership = LiteLLM_TeamMembership.model_validate(_returned_team_membership.model_dump()) 

496 

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!") 

499 

500 return returned_user, returned_team_membership 

501 

502 

503def _delete_user_id_from_cache(kwargs): 

504 from litellm.proxy.proxy_server import user_api_key_cache 

505 

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) 

510 

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) 

515 

516 

517def _delete_api_key_from_cache(kwargs): 

518 from litellm.proxy.proxy_server import user_api_key_cache 

519 

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) 

524 

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) 

529 

530 

531def _delete_team_id_from_cache(kwargs): 

532 from litellm.proxy.proxy_server import user_api_key_cache 

533 

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) 

538 

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) 

543 

544 

545def _delete_customer_id_from_cache(kwargs): 

546 from litellm.proxy.proxy_server import user_api_key_cache 

547 

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) 

552 

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) 

557 

558 

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 

572 

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 } 

586 

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] 

592 

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 ) 

599 

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 ) 

607 

608 

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 

612 

613 

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 

617 

618 

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 } 

627 

628 

629def _redact_record_env_vars(record: object) -> object: 

630 """Return ``record`` with its ``env_vars[].value`` blanked. 

631 

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 

648 

649 

650def _redact_env_var_values(response: MutableMapping[str, object]) -> None: 

651 """Blank ``env_vars[].value`` in a management response before telemetry. 

652 

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] 

663 

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] 

667 

668 

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. 

679 

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 

688 

689 if open_telemetry_logger is None: 

690 return 

691 

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 

697 

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 ) 

704 

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 = {} 

710 

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 ) 

724 

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 

733 

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 ) 

742 

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 ) 

753 

754 

755def management_endpoint_wrapper(func): 

756 """ 

757 This wrapper does the following: 

758 

759 1. Log I/O, Exceptions to OTEL 

760 2. Create an Audit log for success calls 

761 """ 

762 

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

771 

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 ) 

787 

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

796 

797 return result 

798 except Exception as e: 

799 end_time = datetime.now() 

800 

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 ) 

820 

821 raise e 

822 

823 return wrapper