Coverage for .venv/lib/python3.13/site-packages/litellm/proxy/db/db_spend_update_writer.py: 52%

1046 statements  

« prev     ^ index     » next       coverage.py v7.15.2, created at 2026-10-10 12:01 +0000

1""" 

2Module responsible for 

3 

41. Writing spend increments to either in memory list of transactions or to redis 

52. Reading increments from redis or in memory list of transactions and committing them to db 

6""" 

7 

8import asyncio 

9import copy 

10import json 

11import os 

12import random 

13import time 

14import traceback 

15from collections.abc import Callable, Coroutine, Mapping, Sequence 

16from contextvars import ContextVar 

17from datetime import datetime, timedelta, timezone 

18from types import MappingProxyType 

19from typing import TYPE_CHECKING, Any, Final, Literal, Protocol, TypeAlias, TypeVar, cast, overload 

20from urllib.parse import quote, unquote 

21 

22from pydantic import TypeAdapter 

23from typing_extensions import LiteralString, ReadOnly, TypedDict 

24 

25import litellm 

26from litellm._logging import verbose_proxy_logger 

27from litellm.caching import RedisCache 

28from litellm.constants import ( 

29 DB_DAILY_TAG_SPEND_UPDATE_JOB_NAME, 

30 DB_SPEND_UPDATE_JOB_NAME, 

31 INTERNAL_CALL_ORIGIN_METADATA_KEY, 

32) 

33from litellm.litellm_core_utils.litellm_logging import coerce_model_access_groups 

34from litellm.litellm_core_utils.safe_json_loads import safe_json_loads 

35from litellm.proxy._types import ( 

36 DB_CONNECTION_ERROR_TYPES, 

37 DB_RETRY_SAFE_ERROR_TYPES, 

38 BaseDailySpendTransaction, 

39 DailyAgentSpendTransaction, 

40 DailyEndUserSpendTransaction, 

41 DailyOrganizationSpendTransaction, 

42 DailyTagSpendTransaction, 

43 DailyTeamSpendTransaction, 

44 DailyUserSpendTransaction, 

45 DBSpendUpdateTransactions, 

46 Litellm_EntityType, 

47 SpendLogsMetadata, 

48 SpendLogsPayload, 

49 SpendUpdateQueueItem, 

50 ToolDiscoveryQueueItem, 

51) 

52from litellm.proxy.common_utils.user_api_key_cache import project_cache_key 

53from litellm.proxy.db.daily_spend_bulk_upsert import ( 

54 DAILY_SPEND_TABLES, 

55 build_bulk_upsert, 

56 daily_spend_entity_ids, 

57 merge_by_conflict_key, 

58) 

59from litellm.proxy.db.db_transaction_queue.daily_spend_update_queue import ( 

60 DailySpendUpdateQueue, 

61) 

62from litellm.proxy.db.db_transaction_queue.pod_lock_manager import PodLockManager 

63from litellm.proxy.db.db_transaction_queue.redis_update_buffer import RedisUpdateBuffer 

64from litellm.proxy.db.db_transaction_queue.spend_update_queue import SpendUpdateQueue 

65from litellm.proxy.db.db_transaction_queue.tool_discovery_queue import ( 

66 ToolDiscoveryQueue, 

67) 

68from litellm.proxy.db.db_transaction_queue.window_spend_update_queue import ( 

69 WindowSpendTransaction, 

70 WindowSpendUpdateQueue, 

71) 

72from litellm.proxy.db.exception_handler import PrismaDBExceptionHandler 

73from litellm.proxy.route_llm_request import ROUTE_ENDPOINT_MAPPING 

74from litellm.proxy.spend_tracking.compression_savings import ( 

75 extract_compression_saved_tokens, 

76) 

77from litellm.proxy.spend_tracking.savings import ( 

78 compute_savings_spend, 

79 extract_cache_creation_tokens, 

80 extract_cache_read_tokens, 

81 marks_gateway_injection, 

82) 

83from litellm.proxy.spend_tracking.spend_log_error_logger import spend_log_error 

84from litellm.repositories.prisma_protocols import BatchTable 

85from litellm.types.utils import CallTypes 

86 

87if TYPE_CHECKING: 87 ↛ 88line 87 didn't jump to line 88 because the condition on line 87 was never true

88 from litellm.proxy.db.autorouter_session_rollup import AutoRouterTurnTransaction 

89 from litellm.proxy.db.baseline_accounting import DailyBaselineAttribution 

90 from litellm.proxy.utils import PrismaClient, ProxyLogging 

91else: 

92 PrismaClient = Any 

93 ProxyLogging = Any 

94 

95 

96RESPONSES_SESSION_CALL_TYPES: Final = frozenset({CallTypes.responses.value, CallTypes.aresponses.value}) 

97_SPEND_METADATA_ADAPTER: Final = TypeAdapter(Mapping[str, object]) 

98 

99 

100def _org_member_transaction_key(org_id: str, user_id: str) -> str: 

101 return f"organization_id::{quote(org_id, safe='')}::user_id::{quote(user_id, safe='')}" 

102 

103 

104def _is_batch_cost_row(payload: SpendLogsPayload) -> bool: 

105 return payload.get("call_type") == CallTypes.aretrieve_batch.value and payload.get("status") == "success" 

106 

107 

108_BATCH_COST_CLAIM_FIELDS: Final = frozenset({"request_id", "call_type", "spend", "startTime", "endTime", "status"}) 

109 

110 

111def _batch_cost_row_to_write(payload: SpendLogsPayload, disable_spend_logs: bool) -> Mapping[str, object]: 

112 """Reduce a batch's cost row to what tells the retrieves apart when logging is off. 

113 

114 A proxy run with spend logs disabled still needs one row per batch to charge it once, 

115 so the row is written either way, but it carries no request of its own: no metadata, 

116 no requester IP, no key, model, or token counts (LIT-7048). 

117 """ 

118 if disable_spend_logs is False: 

119 return payload 

120 return MappingProxyType({field: value for field, value in payload.items() if field in _BATCH_COST_CLAIM_FIELDS}) 

121 

122 

123class _SpendIncrement(TypedDict): 

124 increment: ReadOnly[float] 

125 

126 

127class _MemberSpendRow(TypedDict): 

128 user_id: ReadOnly[str] 

129 team_id: ReadOnly[str] 

130 cost: ReadOnly[float] 

131 

132 

133class _SpendBatch(Protocol): 

134 litellm_usertable: BatchTable 

135 litellm_verificationtoken: BatchTable 

136 litellm_teamtable: BatchTable 

137 litellm_teammembership: BatchTable 

138 litellm_organizationtable: BatchTable 

139 litellm_organizationmembership: BatchTable 

140 litellm_projecttable: BatchTable 

141 litellm_tagtable: BatchTable 

142 litellm_agentstable: BatchTable 

143 litellm_modelaccessgroupbudgettable: BatchTable 

144 

145 

146_EntitySpendTable: TypeAlias = Literal[ 

147 "litellm_tagtable", "litellm_agentstable", "litellm_modelaccessgroupbudgettable", "litellm_projecttable" 

148] 

149 

150 

151_ENTITY_SPEND_TABLES: Final[Mapping[_EntitySpendTable, Callable[[_SpendBatch], BatchTable]]] = MappingProxyType( 

152 { 

153 "litellm_tagtable": lambda batcher: batcher.litellm_tagtable, 

154 "litellm_agentstable": lambda batcher: batcher.litellm_agentstable, 

155 "litellm_modelaccessgroupbudgettable": lambda batcher: batcher.litellm_modelaccessgroupbudgettable, 

156 "litellm_projecttable": lambda batcher: batcher.litellm_projecttable, 

157 } 

158) 

159 

160 

161def _entity_spend_table(batcher: _SpendBatch, table_accessor: _EntitySpendTable) -> BatchTable: 

162 return _ENTITY_SPEND_TABLES[table_accessor](batcher) 

163 

164 

165class _SpendBatchManager(Protocol): 

166 async def __aenter__(self) -> _SpendBatch: ... 166 ↛ exitline 166 didn't return from function '__aenter__' because

167 

168 async def __aexit__(self, exc_type: object, exc_value: object, traceback: object) -> bool | None: ... 168 ↛ exitline 168 didn't return from function '__aexit__' because

169 

170 

171class _SpendTransaction(Protocol): 

172 def batch_(self) -> _SpendBatchManager: ... 172 ↛ exitline 172 didn't return from function 'batch_' because

173 

174 async def execute_raw(self, query: LiteralString, *args: object) -> int: ... 174 ↛ exitline 174 didn't return from function 'execute_raw' because

175 

176 

177class _SpendTransactionManager(Protocol): 

178 async def __aenter__(self) -> _SpendTransaction: ... 178 ↛ exitline 178 didn't return from function '__aenter__' because

179 

180 async def __aexit__(self, exc_type: object, exc_value: object, traceback: object) -> bool | None: ... 180 ↛ exitline 180 didn't return from function '__aexit__' because

181 

182 

183_DailySpendTransactionT = TypeVar("_DailySpendTransactionT", bound=BaseDailySpendTransaction) 

184 

185 

186class _DailySpendCommit(Protocol[_DailySpendTransactionT]): 

187 async def __call__( 187 ↛ exitline 187 didn't return from function '__call__' because

188 self, 

189 *, 

190 n_retry_times: int, 

191 prisma_client: PrismaClient, 

192 proxy_logging_obj: ProxyLogging, 

193 daily_spend_transactions: dict[str, _DailySpendTransactionT], 

194 ) -> None: ... 

195 

196 

197_DATA_REJECTED_SQLSTATE_CLASSES: Final = frozenset({"22", "23"}) 

198 

199 

200def _spend_commit_failure_is_requeue_safe(e: Exception) -> bool: 

201 if isinstance(e, DB_CONNECTION_ERROR_TYPES): 

202 return isinstance(e, DB_RETRY_SAFE_ERROR_TYPES) 

203 sqlstate: Final = PrismaDBExceptionHandler.postgres_sqlstate(e) 

204 return sqlstate is None or sqlstate[:2] not in _DATA_REJECTED_SQLSTATE_CLASSES 

205 

206 

207_SpendTableName = Literal[ 

208 "user_list_transactions", 

209 "end_user_list_transactions", 

210 "key_list_transactions", 

211 "team_list_transactions", 

212 "team_member_list_transactions", 

213 "org_list_transactions", 

214 "org_member_list_transactions", 

215 "project_list_transactions", 

216 "tag_list_transactions", 

217 "model_access_group_list_transactions", 

218 "agent_list_transactions", 

219] 

220_SPEND_TABLE_COMMIT_ORDER: Final[tuple[_SpendTableName, ...]] = ( 

221 "user_list_transactions", 

222 "end_user_list_transactions", 

223 "key_list_transactions", 

224 "team_list_transactions", 

225 "team_member_list_transactions", 

226 "org_list_transactions", 

227 "org_member_list_transactions", 

228 "project_list_transactions", 

229 "tag_list_transactions", 

230 "model_access_group_list_transactions", 

231 "agent_list_transactions", 

232) 

233 

234 

235def _spend_tables_left_to_send( 

236 transactions: DBSpendUpdateTransactions, 

237 committed: Sequence[_SpendTableName], 

238 failure: Exception, 

239) -> DBSpendUpdateTransactions | None: 

240 in_flight: Final[_SpendTableName | None] = ( 

241 _SPEND_TABLE_COMMIT_ORDER[len(committed)] if len(committed) < len(_SPEND_TABLE_COMMIT_ORDER) else None 

242 ) 

243 dropped: Final[frozenset[_SpendTableName]] = ( 

244 frozenset() if in_flight is None or _spend_commit_failure_is_requeue_safe(failure) else frozenset({in_flight}) 

245 ) 

246 if dropped and in_flight is not None: 

247 spend_log_error( 

248 "Spend tracking - dropped %d %s increments: the failed statement may have applied or the " 

249 "database refused the data, so re-sending it is not safe. Error: %s", 

250 len(cast(dict[str, dict[str, float] | None], transactions).get(in_flight) or ()), 

251 in_flight, 

252 str(failure), 

253 exc=failure, 

254 ) 

255 remaining: Final = { 

256 name: (None if name in committed or name in dropped else txns) for name, txns in transactions.items() 

257 } 

258 return cast(DBSpendUpdateTransactions, remaining) if any(remaining.values()) else None 

259 

260 

261def _timed_request_duration_ms( 

262 payload: dict | SpendLogsPayload, 

263 request_status: Literal["success", "failure"], 

264 is_internal_call: bool, 

265) -> int | None: 

266 if is_internal_call or request_status != "success": 

267 return None 

268 duration_ms: Final = payload.get("request_duration_ms") 

269 if not isinstance(duration_ms, int) or duration_ms < 0: 269 ↛ 270line 269 didn't jump to line 270 because the condition on line 269 was never true

270 return None 

271 return duration_ms 

272 

273 

274def _spend_update_tx(prisma_client: PrismaClient) -> _SpendTransactionManager: 

275 tx: Final[_SpendTransactionManager] = prisma_client.db.tx(timeout=timedelta(seconds=60)) 

276 return tx 

277 

278 

279_daily_spend_commit_started: Final[ContextVar[asyncio.Event | None]] = ContextVar( 

280 "_daily_spend_commit_started", default=None 

281) 

282 

283 

284def _mark_daily_spend_commit_started() -> None: 

285 started: Final = _daily_spend_commit_started.get() 

286 if started is not None: 286 ↛ exitline 286 didn't return from function '_mark_daily_spend_commit_started' because the condition on line 286 was always true

287 started.set() 

288 

289 

290def _mark_daily_spend_commit_finished() -> None: 

291 started: Final = _daily_spend_commit_started.get() 

292 if started is not None: 292 ↛ exitline 292 didn't return from function '_mark_daily_spend_commit_finished' because the condition on line 292 was always true

293 started.clear() 

294 

295 

296def _start_daily_spend_commit( 

297 commit_started: asyncio.Event, commit: Callable[[], Coroutine[object, object, None]] 

298) -> "asyncio.Task[None]": 

299 token: Final = _daily_spend_commit_started.set(commit_started) 

300 try: 

301 return asyncio.ensure_future(commit()) 

302 finally: 

303 _daily_spend_commit_started.reset(token) 

304 

305 

306def _track_interrupted_commit(commits: set[asyncio.Task[None]], settle: Coroutine[object, object, None]) -> None: 

307 task: Final = asyncio.ensure_future(settle) 

308 commits.add(task) 

309 task.add_done_callback(commits.discard) 

310 

311 

312async def _settle_interrupted_commits(commits: set[asyncio.Task[None]]) -> None: 

313 while commits: 

314 await asyncio.wait(tuple(commits)) 

315 

316 

317async def _restore_tag_spend_the_commit_left_behind( 

318 commit_task: "asyncio.Task[None]", 

319 redis_update_buffer: RedisUpdateBuffer, 

320 transactions: dict[str, DailyTagSpendTransaction], 

321) -> None: 

322 await asyncio.wait({commit_task}) 

323 if commit_task.cancelled() or commit_task.exception() is None: 

324 return 

325 await redis_update_buffer.restore_transactions_to_redis( 

326 daily_tag_spend_update_transactions=transactions, 

327 ) 

328 

329 

330async def _requeue_daily_spend_the_commit_left_behind( 

331 commit_task: "asyncio.Task[None]", 

332 queue: DailySpendUpdateQueue, 

333 entity_type: str, 

334 transactions: dict[str, BaseDailySpendTransaction], 

335) -> None: 

336 await asyncio.wait({commit_task}) 

337 if commit_task.cancelled() or not transactions: 

338 return 

339 failure: Final = commit_task.exception() 

340 if failure is None: 

341 return 

342 spend_log_error( 

343 "Spend tracking - daily %s spend commit interrupted by shutdown failed. Re-queued %d rows for the " 

344 "shutdown flush. Error: %s", 

345 entity_type, 

346 len(transactions), 

347 str(failure), 

348 exc=failure, 

349 ) 

350 await queue.add_update(transactions) 

351 

352 

353# The per-team advisory lock the team endpoints hold while changing a roster (TEAM_ADVISORY_LOCK_SQL), 

354# so the roster check below cannot interleave with their writes. A row lock would deadlock with the 

355# access-group endpoints, which lock a team row after an access-group lock. 

356_TEAM_ADVISORY_LOCK_SQL: Final = "SELECT pg_advisory_xact_lock(hashtext($1)) IS NULL AS locked" 

357 

358# One statement adds every member's cost to their membership row. A missing row is created only 

359# while the user is still on the team's roster, so a spend flush landing after a removal never 

360# recreates the member. The rows travel as one JSON document, not as a numeric array: Prisma 

361# types a raw array parameter from the first batch a connection sees, so after an all-$0 batch 

362# (integers) every later fractional batch on that connection failed with "improper binary format". 

363_TEAM_MEMBER_SPEND_SQL: Final = """ 

364INSERT INTO "LiteLLM_TeamMembership" (user_id, team_id, spend, total_spend) 

365SELECT member.user_id, member.team_id, member.cost, member.cost 

366FROM jsonb_to_recordset($1::jsonb) AS member(user_id text, team_id text, cost float8) 

367WHERE EXISTS ( 

368 SELECT 1 FROM "LiteLLM_TeamTable" t 

369 WHERE t.team_id = member.team_id 

370 AND t.members_with_roles @> jsonb_build_array(jsonb_build_object('user_id', member.user_id)) 

371) 

372 OR EXISTS ( 

373 SELECT 1 FROM "LiteLLM_TeamMembership" m 

374 WHERE m.user_id = member.user_id AND m.team_id = member.team_id 

375) 

376ON CONFLICT (user_id, team_id) DO UPDATE 

377SET spend = "LiteLLM_TeamMembership".spend + EXCLUDED.spend, 

378 total_spend = "LiteLLM_TeamMembership".total_spend + EXCLUDED.total_spend 

379""" 

380 

381 

382async def _write_team_member_spend(transaction: _SpendTransaction, spend_by_member_key: Mapping[str, float]) -> None: 

383 # key is "team_id::<value>::user_id::<value>"; locks are taken in sorted team_id order like the team endpoints 

384 rows: Final = sorted((key.split("::")[1], key.split("::")[3], cost) for key, cost in spend_by_member_key.items()) 

385 for team_id in dict.fromkeys(team_id for team_id, _user_id, _cost in rows): 

386 _ = await transaction.execute_raw(_TEAM_ADVISORY_LOCK_SQL, team_id) 

387 members: Final = tuple( 

388 _MemberSpendRow(user_id=user_id, team_id=team_id, cost=cost) for team_id, user_id, cost in rows 

389 ) 

390 _ = await transaction.execute_raw(_TEAM_MEMBER_SPEND_SQL, json.dumps(members)) 

391 

392 

393def get_llm_router(): 

394 """The proxy's router, or None outside a running proxy. 

395 

396 Injected rather than imported where it is used, so the savings computation stays 

397 a pure function of its arguments and the caller owns where the router comes from. 

398 """ 

399 try: 

400 from litellm.proxy.proxy_server import llm_router 

401 

402 return llm_router 

403 except Exception: # noqa: BLE001 # no proxy in scope; savings degrade to zero 

404 return None 

405 

406 

407class _DeploymentLookup(Protocol): 

408 def get_model_info(self, id: str) -> Mapping[str, object] | None: ... 408 ↛ exitline 408 didn't return from function 'get_model_info' because

409 

410 

411def _served_model_access_groups( 

412 router: _DeploymentLookup | None, 

413 served_model_id: str | None, 

414) -> frozenset[str] | None: 

415 """Access groups declared by the deployment that actually served the request. 

416 

417 None when the served deployment cannot be identified, in which case the set 

418 attributed at auth time stands unchanged. 

419 """ 

420 if router is None or not served_model_id: 

421 return None 

422 deployment: Final = router.get_model_info(id=served_model_id) 

423 if deployment is None: 

424 return None 

425 model_info: Final = deployment.get("model_info") 

426 if not isinstance(model_info, Mapping): 

427 return None 

428 declared: Final = model_info.get("access_groups") 

429 if not isinstance(declared, (list, tuple)): 

430 return frozenset() 

431 return frozenset(group for group in declared if isinstance(group, str)) 

432 

433 

434def debitable_model_access_groups( 

435 attributed: Sequence[str] | None, 

436 served_model_id: str | None, 

437 router: _DeploymentLookup | None, 

438) -> tuple[str, ...]: 

439 """Groups to debit: the set attributed at auth time, narrowed to those the served model belongs to. 

440 

441 The router may fall back to a model outside the pool auth reserved against, so the 

442 attributed set is the hard upper bound: a group absent from it is never debited. 

443 """ 

444 ordered: Final = coerce_model_access_groups(attributed) 

445 if not ordered: 445 ↛ 447line 445 didn't jump to line 447 because the condition on line 445 was always true

446 return () 

447 served: Final = _served_model_access_groups(router=router, served_model_id=served_model_id) 

448 if served is None: 

449 return ordered 

450 return tuple(group for group in ordered if group in served) 

451 

452 

453class DBSpendUpdateWriter: 

454 """ 

455 Module responsible for 

456 

457 1. Writing spend increments to either in memory list of transactions or to redis 

458 2. Reading increments from redis or in memory list of transactions and committing them to db 

459 """ 

460 

461 def __init__( 

462 self, 

463 redis_cache: RedisCache | None = None, 

464 ): 

465 self.redis_cache = redis_cache 

466 self.redis_update_buffer = RedisUpdateBuffer(redis_cache=self.redis_cache) 

467 self.pod_lock_manager = PodLockManager() 

468 self.spend_update_queue = SpendUpdateQueue() 

469 self.tool_discovery_queue = ToolDiscoveryQueue() 

470 self.daily_spend_update_queue = DailySpendUpdateQueue() 

471 self.daily_team_spend_update_queue = DailySpendUpdateQueue() 

472 self.daily_end_user_spend_update_queue = DailySpendUpdateQueue() 

473 self.daily_agent_spend_update_queue = DailySpendUpdateQueue() 

474 self.daily_org_spend_update_queue = DailySpendUpdateQueue() 

475 self.daily_tag_spend_update_queue = DailySpendUpdateQueue() 

476 self.window_spend_update_queue = WindowSpendUpdateQueue() 

477 self.interrupted_tag_commits: set[asyncio.Task[None]] = ( 

478 set() 

479 ) # mutable-ok: same registry as DailySpendUpdateQueue.interrupted_commits 

480 

481 async def update_database( 

482 # LiteLLM management object fields 

483 self, 

484 token: str | None, 

485 user_id: str | None, 

486 end_user_id: str | None, 

487 team_id: str | None, 

488 org_id: str | None, 

489 # Completion object fields 

490 kwargs: dict | None, 

491 completion_response: object, 

492 start_time: datetime, 

493 end_time: datetime, 

494 response_cost: float | None, 

495 project_id: str | None = None, 

496 ) -> bool: 

497 """Record the request's spend, answering whether its cost still needs charging. 

498 

499 False only for a batch retrieve whose cost row another retrieve already wrote, 

500 so the caller leaves the key, team, and user counters alone (LIT-7048). 

501 """ 

502 from litellm.proxy.proxy_server import ( 

503 disable_spend_logs, 

504 litellm_proxy_budget_name, 

505 prisma_client, 

506 ) 

507 from litellm.proxy.utils import ProxyUpdateSpend, hash_token 

508 

509 try: 

510 verbose_proxy_logger.debug( 

511 "Enters prisma db call, response_cost: %s, token: %s; user_id: %s; team_id: %s", 

512 response_cost, 

513 token, 

514 user_id, 

515 team_id, 

516 ) 

517 if ProxyUpdateSpend.disable_spend_updates() is True: 517 ↛ 518line 517 didn't jump to line 518 because the condition on line 517 was never true

518 return True 

519 if token is not None and isinstance(token, str) and token.startswith("sk-"): 519 ↛ 520line 519 didn't jump to line 520 because the condition on line 519 was never true

520 hashed_token = hash_token(token=token) 

521 else: 

522 hashed_token = token 

523 

524 ## CREATE SPEND LOG PAYLOAD ## 

525 from litellm.proxy.spend_tracking.spend_tracking_utils import ( 

526 get_logging_payload, 

527 get_request_model_access_groups, 

528 ) 

529 

530 payload: Final = get_logging_payload( 

531 kwargs=kwargs, 

532 response_obj=completion_response, 

533 start_time=start_time, 

534 end_time=end_time, 

535 llm_router=get_llm_router(), 

536 ) 

537 payload["spend"] = response_cost or 0.0 

538 if isinstance(payload["startTime"], datetime): 538 ↛ 540line 538 didn't jump to line 540 because the condition on line 538 was always true

539 payload["startTime"] = payload["startTime"].isoformat() 

540 if isinstance(payload["endTime"], datetime): 540 ↛ 543line 540 didn't jump to line 543 because the condition on line 540 was always true

541 payload["endTime"] = payload["endTime"].isoformat() 

542 

543 if org_id is not None and org_id != "": 543 ↛ 544line 543 didn't jump to line 544 because the condition on line 543 was never true

544 payload["organization_id"] = org_id 

545 

546 if team_id is not None and team_id != "": 546 ↛ 547line 546 didn't jump to line 547 because the condition on line 546 was never true

547 payload["team_id"] = team_id 

548 

549 if not await self._record_spend_log( 549 ↛ 552line 549 didn't jump to line 552 because the condition on line 549 was never true

550 payload=payload, prisma_client=prisma_client, disable_spend_logs=disable_spend_logs 

551 ): 

552 return False 

553 

554 if disable_spend_logs is False: 554 ↛ 566line 554 didn't jump to line 566 because the condition on line 554 was always true

555 await self._enqueue_tool_usage_transaction( 

556 payload=payload, 

557 completion_response=completion_response, 

558 prisma_client=prisma_client, 

559 kwargs=kwargs, 

560 ) 

561 await self._enqueue_autorouter_turn_transaction( 

562 payload=payload, 

563 prisma_client=prisma_client, 

564 ) 

565 else: 

566 verbose_proxy_logger.debug( 

567 "disable_spend_logs=True. Skipping writing spend logs to db. Other spend updates - Key/User/Team table will still occur." 

568 ) 

569 

570 # Single task replaces 11 create_task() calls 

571 asyncio.create_task( 

572 self._batch_database_updates( 

573 response_cost=response_cost, 

574 user_id=user_id, 

575 hashed_token=hashed_token, 

576 team_id=team_id, 

577 org_id=org_id, 

578 project_id=project_id, 

579 end_user_id=end_user_id, 

580 prisma_client=prisma_client, 

581 litellm_proxy_budget_name=litellm_proxy_budget_name, 

582 payload=payload, 

583 request_model_access_groups=get_request_model_access_groups(kwargs), 

584 ) 

585 ) 

586 

587 self._enqueue_tool_registry_upsert( 

588 kwargs=kwargs, 

589 completion_response=completion_response, 

590 hashed_token=hashed_token, 

591 team_id=team_id, 

592 ) 

593 

594 verbose_proxy_logger.debug("Runs spend update on all tables") 

595 return True 

596 except Exception: 

597 spend_log_error( 

598 "Spend tracking - update_database failed. Spend log insertion or daily transaction enqueue " 

599 "may not have completed for this request. " 

600 "response_cost=%s, token=%s, user_id=%s, team_id=%s, org_id=%s, end_user_id=%s", 

601 response_cost, 

602 token, 

603 user_id, 

604 team_id, 

605 org_id, 

606 end_user_id, 

607 ) 

608 return True 

609 

610 async def _record_spend_log( 

611 self, payload: SpendLogsPayload, prisma_client: "PrismaClient | None", disable_spend_logs: bool 

612 ) -> bool: 

613 if prisma_client is not None and _is_batch_cost_row(payload): 613 ↛ 614line 613 didn't jump to line 614 because the condition on line 613 was never true

614 return await self._claim_batch_cost_spend_log( 

615 payload=payload, prisma_client=prisma_client, disable_spend_logs=disable_spend_logs 

616 ) 

617 if disable_spend_logs is False: 617 ↛ 619line 617 didn't jump to line 619 because the condition on line 617 was always true

618 await self._insert_spend_log_to_db(payload=payload, prisma_client=prisma_client) 

619 return True 

620 

621 async def _claim_batch_cost_spend_log( 

622 self, payload: SpendLogsPayload, prisma_client: "PrismaClient", disable_spend_logs: bool 

623 ) -> bool: 

624 """Write the batch's cost row now, or learn that another retrieve already did. 

625 

626 Every retrieve of one batch shares this row, so the insert that lands first owns 

627 the charge and every later one finds the row and charges nothing (LIT-7048). Only 

628 a row that recorded a charge counts: a failed retrieve, a request whose client 

629 picked the batch id as its call id, and the $0 row an older proxy left behind 

630 while the batch was still running all leave the charge to be made. 

631 """ 

632 from litellm.repositories.table_repositories import SpendLogsRepository 

633 

634 request_id: Final = payload["request_id"] 

635 row: Final = _batch_cost_row_to_write(payload, disable_spend_logs) 

636 spend_logs: Final = SpendLogsRepository(prisma_client).table 

637 try: 

638 claimed: Final = await spend_logs.create_many( 

639 data=[prisma_client.jsonify_object(row)], # mutable-ok: prisma create_many takes a list 

640 skip_duplicates=True, 

641 ) 

642 if claimed == 1: 

643 return True 

644 existing: Final = await spend_logs.find_unique( 

645 where={"request_id": request_id} # mutable-ok: prisma where clause 

646 ) 

647 except Exception as e: # noqa: BLE001 # prisma raises its own hierarchy; an unreachable DB queues the row like any other spend log 

648 verbose_proxy_logger.warning( 

649 "Could not claim spend row %s for a batch's cost, queueing it: %s", request_id, e 

650 ) 

651 await self._insert_spend_log_to_db(payload=prisma_client.jsonify_object(row), prisma_client=prisma_client) 

652 return True 

653 if existing is None or existing.call_type != CallTypes.aretrieve_batch.value or existing.status != "success": 

654 verbose_proxy_logger.warning( 

655 "Spend row %s belongs to a %s request, so this batch's cost is charged without a row of its own", 

656 request_id, 

657 getattr(existing, "call_type", None), 

658 ) 

659 return True 

660 if existing.spend > 0: 

661 verbose_proxy_logger.debug("Cost tracking skipped: spend row %s already charged this batch", request_id) 

662 return False 

663 return await self._take_over_uncharged_batch_cost_row(payload=payload, prisma_client=prisma_client, row=row) 

664 

665 async def _take_over_uncharged_batch_cost_row( 

666 self, payload: SpendLogsPayload, prisma_client: "PrismaClient", row: Mapping[str, object] 

667 ) -> bool: 

668 """Take the batch's cost row over from the poll that left it charging nothing. 

669 

670 A pre-upgrade proxy wrote that row every time it polled the batch while it was still 

671 running, so the charge is still to be made and the row still has to end up carrying 

672 it. The row stops matching the moment it carries a charge, so it is one retrieve that 

673 takes it over and charges, and every later one reads the charge and charges nothing. 

674 """ 

675 from litellm.repositories.table_repositories import SpendLogsRepository 

676 

677 request_id: Final = payload["request_id"] 

678 if payload["spend"] <= 0: 

679 verbose_proxy_logger.debug( 

680 "Cost tracking skipped: this batch costs nothing and spend row %s says so", request_id 

681 ) 

682 return False 

683 try: 

684 taken_over: Final = await SpendLogsRepository(prisma_client).table.update_many( 

685 data=prisma_client.jsonify_object( 

686 MappingProxyType({field: value for field, value in row.items() if field != "request_id"}) 

687 ), 

688 where={ # mutable-ok: prisma where clause 

689 "request_id": request_id, 

690 "call_type": CallTypes.aretrieve_batch.value, 

691 "status": "success", 

692 "spend": 0.0, 

693 }, 

694 ) 

695 except Exception as e: # noqa: BLE001 # prisma raises its own hierarchy; the next retrieve takes the row over 

696 verbose_proxy_logger.warning( 

697 "Could not take over spend row %s, leaving this batch's cost to the next retrieve: %s", request_id, e 

698 ) 

699 return False 

700 if taken_over == 0: 

701 verbose_proxy_logger.debug("Cost tracking skipped: spend row %s already charged this batch", request_id) 

702 return False 

703 return True 

704 

705 async def _enqueue_tool_usage_transaction( 

706 self, 

707 payload: SpendLogsPayload, 

708 completion_response: object, 

709 prisma_client: "PrismaClient | None", 

710 kwargs: "dict | None" = None, 

711 ) -> None: 

712 try: 

713 if prisma_client is None: 713 ↛ 714line 713 didn't jump to line 714 because the condition on line 713 was never true

714 return 

715 from litellm.proxy.db.spend_log_tool_index import ( 

716 build_tool_usage_transaction, 

717 ) 

718 

719 transaction: Final = build_tool_usage_transaction( 

720 request_id=payload["request_id"], 

721 start_time_iso=str(payload["startTime"]), 

722 mcp_namespaced_tool_name=payload.get("mcp_namespaced_tool_name"), 

723 spend=payload["spend"], 

724 total_tokens=payload["total_tokens"], 

725 completion_response=completion_response, 

726 realtime_tool_calls=(kwargs or {}).get("realtime_tool_calls"), 

727 ) 

728 if transaction is None: 728 ↛ 730line 728 didn't jump to line 730 because the condition on line 728 was always true

729 return 

730 async with prisma_client._tool_usage_transactions_lock: 

731 prisma_client.tool_usage_transactions.append(transaction) 

732 except Exception as e: 

733 verbose_proxy_logger.debug("_enqueue_tool_usage_transaction error (non-blocking): %s", e) 

734 

735 async def _enqueue_autorouter_turn_transaction( 

736 self, 

737 payload: SpendLogsPayload, 

738 prisma_client: "PrismaClient | None", 

739 ) -> None: 

740 try: 

741 if prisma_client is None: 741 ↛ 742line 741 didn't jump to line 742 because the condition on line 741 was never true

742 return 

743 metadata_raw: Final = payload.get("metadata") 

744 if not metadata_raw: 744 ↛ 745line 744 didn't jump to line 745 because the condition on line 744 was never true

745 return 

746 metadata: Final = _SPEND_METADATA_ADAPTER.validate_json(metadata_raw) 

747 routing_decision: Final = metadata.get("routing_decision") 

748 if not isinstance(routing_decision, Mapping) or not routing_decision: 748 ↛ 750line 748 didn't jump to line 750 because the condition on line 748 was always true

749 return 

750 from litellm.proxy.db.autorouter_session_rollup import ( 

751 build_autorouter_turn_transaction, 

752 ) 

753 

754 usage_object_raw: Final = metadata.get("usage_object") 

755 cost_breakdown: Final = metadata.get("cost_breakdown") 

756 savings_estimate: Final = metadata.get("autorouter_savings_estimate") 

757 savings_spend: Final = compute_savings_spend( 

758 model=payload.get("model"), 

759 custom_llm_provider=payload.get("custom_llm_provider"), 

760 compression_saved_tokens=0, 

761 gateway_injected_cache=marks_gateway_injection(metadata, payload.get("model_id")), 

762 routing_decision=routing_decision, 

763 usage_object=usage_object_raw if isinstance(usage_object_raw, dict) else None, 

764 model_id=payload.get("model_id"), 

765 llm_router=get_llm_router, 

766 cost_breakdown=cost_breakdown if isinstance(cost_breakdown, Mapping) else None, 

767 recorded_autorouter_savings=metadata.get("autorouter_savings"), 

768 recorded_autorouter_savings_estimate=( 

769 savings_estimate if isinstance(savings_estimate, Mapping) else None 

770 ), 

771 billed_at=payload.get("endTime"), 

772 ) 

773 transaction: Final = build_autorouter_turn_transaction( 

774 payload=payload, 

775 metadata=metadata, 

776 saved_spend=savings_spend.autorouter, 

777 ) 

778 try: 

779 if await self._enqueue_baseline_accounting(payload, metadata, transaction, prisma_client): 

780 return 

781 except Exception: # noqa: BLE001 # optional baseline capture must preserve the original actual-spend rollup 

782 verbose_proxy_logger.warning("Auto-router baseline observation was unavailable; actual turn retained") 

783 if transaction is None: 

784 return 

785 async with prisma_client._autorouter_turn_transactions_lock: 

786 prisma_client.autorouter_turn_transactions.append(transaction) 

787 except Exception as e: # noqa: BLE001 # a metrics enqueue must never fail the spend write 

788 verbose_proxy_logger.debug("_enqueue_autorouter_turn_transaction error (non-blocking): %s", e) 

789 

790 async def _enqueue_baseline_accounting( 

791 self, 

792 payload: SpendLogsPayload, 

793 metadata: Mapping[str, object], 

794 turn: "AutoRouterTurnTransaction | None", 

795 prisma_client: "PrismaClient", 

796 ) -> bool: 

797 from litellm.proxy.db.baseline_accounting import ( 

798 BaselineAccountingRecord, 

799 ) 

800 from litellm.proxy.hooks.autorouter_baseline_cache import CapturedBaselineObservation 

801 from litellm.proxy.spend_tracking.savings import baseline_cost_snapshot 

802 

803 serialized: Final = metadata.get("autorouter_baseline_observation") 

804 if not isinstance(serialized, str): 

805 return False 

806 captured: Final = CapturedBaselineObservation.model_validate_json(serialized) 

807 if captured.api_key != payload["api_key"] or captured.session_id != payload["session_id"]: 

808 return False 

809 decision: Final = _SPEND_METADATA_ADAPTER.validate_python( 

810 metadata.get("routing_decision") or MappingProxyType({}) 

811 ) 

812 breakdown: Final = _SPEND_METADATA_ADAPTER.validate_python( 

813 metadata.get("cost_breakdown") or MappingProxyType({}) 

814 ) 

815 daily: Final = await self._baseline_daily_attribution(payload, prisma_client) 

816 record: Final = BaselineAccountingRecord( 

817 scope=captured.scope, 

818 api_key=captured.api_key, 

819 session_id=captured.session_id, 

820 router_name=captured.router_name, 

821 baseline_model=captured.baseline_model, 

822 observation=captured.observation.model_copy(update=MappingProxyType({"request_id": payload["request_id"]})), 

823 pricing=baseline_cost_snapshot(captured.model, captured.prices, payload["spend"], breakdown, decision), 

824 turn=turn, 

825 daily=daily, 

826 ) 

827 async with prisma_client.baseline_accounting_lock: 

828 if len(prisma_client.baseline_accounting_transactions) >= 10000: 

829 verbose_proxy_logger.warning("Auto-router baseline observation queue is full") 

830 return False 

831 prisma_client.baseline_accounting_transactions.append(record) 

832 from litellm.proxy.utils import request_spend_log_flush 

833 

834 request_spend_log_flush(prisma_client) 

835 return True 

836 

837 async def _baseline_daily_attribution( 

838 self, 

839 payload: SpendLogsPayload, 

840 prisma_client: "PrismaClient", 

841 ) -> "DailyBaselineAttribution | None": 

842 from litellm.proxy.db.baseline_accounting import DailyBaselineAttribution, DailyBaselineTarget 

843 

844 normalized: Final = cast(SpendLogsPayload, MappingProxyType({**payload, "end_user_id": payload["end_user"]})) 

845 bases: Final = tuple( 

846 zip( 

847 DAILY_SPEND_TABLES, 

848 await asyncio.gather( 

849 *( 

850 self._common_add_spend_log_transaction_to_daily_transaction( # pyright: ignore[reportUnknownMemberType] # legacy payload union; this caller supplies a validated spend payload 

851 normalized, 

852 prisma_client, 

853 "request_tags" if entity == "tag" else entity, 

854 ) 

855 for entity in DAILY_SPEND_TABLES 

856 ) 

857 ), 

858 ) 

859 ) 

860 base: Final = next((base for _, base in bases if base is not None), None) 

861 if base is None: 

862 return None 

863 return DailyBaselineAttribution( 

864 date=base["date"], 

865 api_key=base["api_key"], 

866 model=base.get("model"), 

867 custom_llm_provider=base.get("custom_llm_provider"), 

868 model_group=base.get("model_group"), 

869 endpoint=base.get("endpoint"), 

870 mcp_namespaced_tool_name=base.get("mcp_namespaced_tool_name"), 

871 targets=tuple( 

872 DailyBaselineTarget(entity=entity, entity_id=identity) 

873 for entity, values in bases 

874 if values is not None 

875 for identity in daily_spend_entity_ids(payload, entity) 

876 ), 

877 ) 

878 

879 def _enqueue_tool_registry_upsert( 

880 self, 

881 kwargs: dict | None, 

882 completion_response: object, 

883 hashed_token: str | None = None, 

884 team_id: str | None = None, 

885 ) -> None: 

886 """ 

887 Extract tool names from the LLM request and response and enqueue them 

888 for upsert into LiteLLM_ToolTable via ToolDiscoveryQueue. 

889 

890 Handles four sources: 

891 - MCP tools: standard_logging_object.mcp_tool_call_metadata.namespaced_tool_name 

892 - Response tool_calls (OpenAI / Anthropic pass-through converted to OpenAI format): 

893 completion_response.choices[].message.tool_calls[].function.name 

894 - Request tools array (OpenAI format): kwargs["tools"][].function.name 

895 - Request tools array (Anthropic /messages format): kwargs["passthrough_logging_payload"] 

896 ["request_body"]["tools"][].name 

897 """ 

898 try: 

899 if kwargs is None: 899 ↛ 900line 899 didn't jump to line 900 because the condition on line 899 was never true

900 return 

901 

902 # Extract key_alias from kwargs metadata if available 

903 key_alias: str | None = None 

904 _litellm_params: Final = kwargs.get("litellm_params") or {} 

905 _metadata: Final = _litellm_params.get("metadata") or {} 

906 key_alias = _metadata.get("user_api_key_alias") or None 

907 user_agent: Final = _metadata.get("user_agent") or None 

908 

909 def _enqueue(tool_name: str, origin: str = "user_defined") -> None: 

910 self.tool_discovery_queue.add_update( 

911 ToolDiscoveryQueueItem( 

912 tool_name=tool_name, 

913 origin=origin, 

914 key_hash=hashed_token, 

915 team_id=team_id, 

916 key_alias=key_alias, 

917 user_agent=user_agent, 

918 ) 

919 ) 

920 

921 # --- MCP tool calls --- 

922 sl_object: Final = kwargs.get("standard_logging_object") 

923 if sl_object is not None: 

924 mcp_metadata: Final = (sl_object.get("metadata", {}) or {}).get("mcp_tool_call_metadata") 

925 if mcp_metadata and isinstance(mcp_metadata, dict): 925 ↛ 926line 925 didn't jump to line 926 because the condition on line 925 was never true

926 tool_name = mcp_metadata.get("namespaced_tool_name") or mcp_metadata.get("name") 

927 mcp_server_name: Final = mcp_metadata.get("mcp_server_name") 

928 if tool_name: 

929 _enqueue(tool_name, origin=mcp_server_name or "user_defined") 

930 

931 # --- Tools from request body (OpenAI format: tools[].function.name) --- 

932 request_tools: Final = kwargs.get("tools") or [] 

933 for tool_def in request_tools: 

934 if not isinstance(tool_def, dict): 

935 continue 

936 fn = tool_def.get("function") or {} 

937 name = fn.get("name") if isinstance(fn, dict) else None 

938 if name: 938 ↛ 939line 938 didn't jump to line 939 because the condition on line 938 was never true

939 _enqueue(name) 

940 

941 # --- Tools from Anthropic /messages pass-through request body 

942 # (Anthropic format: tools[].name, no "function" wrapper) --- 

943 passthrough_payload: Final = kwargs.get("passthrough_logging_payload") or {} 

944 request_body: Final = ( 

945 passthrough_payload.get("request_body") if isinstance(passthrough_payload, dict) else None 

946 ) or {} 

947 for tool_def in request_body.get("tools") or []: 947 ↛ 948line 947 didn't jump to line 948 because the loop on line 947 never started

948 if not isinstance(tool_def, dict): 

949 continue 

950 name = tool_def.get("name") 

951 if name: 

952 _enqueue(name) 

953 

954 # --- Response tool_calls (OpenAI format; Anthropic pass-through converts tool_use here) --- 

955 from litellm.proxy.db.spend_log_tool_index import response_tool_call_names 

956 

957 for tool_name in response_tool_call_names(completion_response): 957 ↛ 958line 957 didn't jump to line 958 because the loop on line 957 never started

958 _enqueue(tool_name) 

959 except Exception as e: 

960 verbose_proxy_logger.debug("_enqueue_tool_registry_upsert error (non-blocking): %s", e) 

961 

962 async def _batch_database_updates( 

963 self, 

964 *, 

965 response_cost: float | None, 

966 user_id: str | None, 

967 hashed_token: str | None, 

968 team_id: str | None, 

969 org_id: str | None, 

970 end_user_id: str | None, 

971 prisma_client: PrismaClient | None, 

972 litellm_proxy_budget_name: str | None, 

973 payload: SpendLogsPayload, 

974 request_model_access_groups: Sequence[str] = (), 

975 project_id: str | None = None, 

976 ): 

977 """ 

978 Runs all 13 spend-update helpers sequentially inside a single asyncio task. 

979 

980 Each helper is wrapped in try/except so one failure doesn't prevent the others. 

981 

982 The deepcopy runs here, off the awaited request path, so the daily spend 

983 helpers get a payload isolated from the spend-log queue entry and the caller. 

984 """ 

985 payload_copy: Final = copy.deepcopy(payload) 

986 request_tags: Final = payload_copy.get("request_tags") 

987 try: 

988 await self._update_user_db( 

989 response_cost=response_cost, 

990 user_id=user_id, 

991 prisma_client=prisma_client, 

992 litellm_proxy_budget_name=litellm_proxy_budget_name, 

993 end_user_id=end_user_id, 

994 ) 

995 except Exception: 

996 verbose_proxy_logger.debug( 

997 "_batch_database_updates: _update_user_db failed: %s", 

998 traceback.format_exc(), 

999 ) 

1000 

1001 try: 

1002 await self._update_key_db( 

1003 response_cost=response_cost, 

1004 hashed_token=hashed_token, 

1005 prisma_client=prisma_client, 

1006 ) 

1007 except Exception: 

1008 verbose_proxy_logger.debug( 

1009 "_batch_database_updates: _update_key_db failed: %s", 

1010 traceback.format_exc(), 

1011 ) 

1012 

1013 try: 

1014 await self._update_team_db( 

1015 response_cost=response_cost, 

1016 team_id=team_id, 

1017 user_id=user_id, 

1018 prisma_client=prisma_client, 

1019 ) 

1020 except Exception: 

1021 verbose_proxy_logger.debug( 

1022 "_batch_database_updates: _update_team_db failed: %s", 

1023 traceback.format_exc(), 

1024 ) 

1025 

1026 try: 

1027 await self._update_org_db( 

1028 response_cost=response_cost, 

1029 org_id=org_id, 

1030 user_id=user_id, 

1031 prisma_client=prisma_client, 

1032 ) 

1033 except Exception: 

1034 verbose_proxy_logger.debug( 

1035 "_batch_database_updates: _update_org_db failed: %s", 

1036 traceback.format_exc(), 

1037 ) 

1038 

1039 try: 

1040 await self._update_project_db( 

1041 response_cost=response_cost, 

1042 project_id=project_id, 

1043 prisma_client=prisma_client, 

1044 ) 

1045 except Exception: # noqa: BLE001 # a project enqueue failure must not skip the sibling spend writes 

1046 verbose_proxy_logger.debug( 

1047 "_batch_database_updates: _update_project_db failed: %s", 

1048 traceback.format_exc(), 

1049 ) 

1050 

1051 try: 

1052 await self._update_tag_db( 

1053 response_cost=response_cost, 

1054 request_tags=request_tags, 

1055 prisma_client=prisma_client, 

1056 ) 

1057 except Exception: 

1058 verbose_proxy_logger.debug( 

1059 "_batch_database_updates: _update_tag_db failed: %s", 

1060 traceback.format_exc(), 

1061 ) 

1062 

1063 await self._update_model_access_group_db( 

1064 response_cost=response_cost, 

1065 request_model_access_groups=request_model_access_groups, 

1066 served_model_id=payload_copy.get("model_id"), 

1067 prisma_client=prisma_client, 

1068 router=get_llm_router(), 

1069 ) 

1070 

1071 _agent_id_for_spend: Final = payload_copy.get("agent_id") 

1072 try: 

1073 await self._update_agent_db( 

1074 response_cost=response_cost, 

1075 agent_id=_agent_id_for_spend, 

1076 prisma_client=prisma_client, 

1077 ) 

1078 except Exception: 

1079 verbose_proxy_logger.debug( 

1080 "_batch_database_updates: _update_agent_db failed: %s", 

1081 traceback.format_exc(), 

1082 ) 

1083 

1084 try: 

1085 await self.add_spend_log_transaction_to_daily_user_transaction( 

1086 payload=payload_copy, 

1087 prisma_client=prisma_client, 

1088 ) 

1089 except Exception: 

1090 verbose_proxy_logger.debug( 

1091 "_batch_database_updates: add_spend_log_transaction_to_daily_user_transaction failed: %s", 

1092 traceback.format_exc(), 

1093 ) 

1094 

1095 try: 

1096 await self.add_spend_log_transaction_to_daily_end_user_transaction( 

1097 payload=payload_copy, 

1098 prisma_client=prisma_client, 

1099 ) 

1100 except Exception: 

1101 verbose_proxy_logger.debug( 

1102 "_batch_database_updates: add_spend_log_transaction_to_daily_end_user_transaction failed: %s", 

1103 traceback.format_exc(), 

1104 ) 

1105 

1106 try: 

1107 await self.add_spend_log_transaction_to_daily_agent_transaction( 

1108 payload=payload_copy, 

1109 prisma_client=prisma_client, 

1110 ) 

1111 except Exception: 

1112 verbose_proxy_logger.debug( 

1113 "_batch_database_updates: add_spend_log_transaction_to_daily_agent_transaction failed: %s", 

1114 traceback.format_exc(), 

1115 ) 

1116 

1117 try: 

1118 await self.add_spend_log_transaction_to_daily_team_transaction( 

1119 payload=payload_copy, 

1120 prisma_client=prisma_client, 

1121 ) 

1122 except Exception: 

1123 verbose_proxy_logger.debug( 

1124 "_batch_database_updates: add_spend_log_transaction_to_daily_team_transaction failed: %s", 

1125 traceback.format_exc(), 

1126 ) 

1127 

1128 try: 

1129 await self.add_spend_log_transaction_to_daily_org_transaction( 

1130 payload=payload_copy, 

1131 org_id=org_id, 

1132 prisma_client=prisma_client, 

1133 ) 

1134 except Exception: 

1135 verbose_proxy_logger.debug( 

1136 "_batch_database_updates: add_spend_log_transaction_to_daily_org_transaction failed: %s", 

1137 traceback.format_exc(), 

1138 ) 

1139 

1140 try: 

1141 await self.add_spend_log_transaction_to_daily_tag_transaction( 

1142 payload=payload_copy, 

1143 prisma_client=prisma_client, 

1144 ) 

1145 except Exception: 

1146 verbose_proxy_logger.debug( 

1147 "_batch_database_updates: add_spend_log_transaction_to_daily_tag_transaction failed: %s", 

1148 traceback.format_exc(), 

1149 ) 

1150 

1151 async def _update_key_db( 

1152 self, 

1153 response_cost: float | None, 

1154 hashed_token: str | None, 

1155 prisma_client: PrismaClient | None, 

1156 ): 

1157 try: 

1158 if hashed_token is None or prisma_client is None: 

1159 return 

1160 

1161 await self.spend_update_queue.add_update( 

1162 update=SpendUpdateQueueItem( 

1163 entity_type=Litellm_EntityType.KEY, 

1164 entity_id=hashed_token, 

1165 response_cost=response_cost, 

1166 ) 

1167 ) 

1168 except Exception as e: 

1169 spend_log_error("Update Key DB Call failed to execute - %s", str(e), exc=e) 

1170 raise e 

1171 

1172 async def _update_user_db( 

1173 self, 

1174 response_cost: float | None, 

1175 user_id: str | None, 

1176 prisma_client: PrismaClient | None, 

1177 litellm_proxy_budget_name: str | None, 

1178 end_user_id: str | None = None, 

1179 ): 

1180 """ 

1181 - Update that user's row 

1182 - Update litellm-proxy-budget row (global proxy spend) 

1183 """ 

1184 try: 

1185 if prisma_client is not None: # update 1185 ↛ exitline 1185 didn't return from function '_update_user_db' because the condition on line 1185 was always true

1186 user_ids: Final = [user_id] 

1187 if litellm.max_budget > 0: # track global proxy budget, if user set max budget 1187 ↛ 1188line 1187 didn't jump to line 1188 because the condition on line 1187 was never true

1188 user_ids.append(litellm_proxy_budget_name) 

1189 

1190 for _id in user_ids: 

1191 if _id is not None: 

1192 await self.spend_update_queue.add_update( 

1193 update=SpendUpdateQueueItem( 

1194 entity_type=Litellm_EntityType.USER, 

1195 entity_id=_id, 

1196 response_cost=response_cost, 

1197 ) 

1198 ) 

1199 

1200 if end_user_id is not None: 

1201 await self.spend_update_queue.add_update( 

1202 update=SpendUpdateQueueItem( 

1203 entity_type=Litellm_EntityType.END_USER, 

1204 entity_id=end_user_id, 

1205 response_cost=response_cost, 

1206 ) 

1207 ) 

1208 except Exception as e: 

1209 spend_log_error( 

1210 "Spend tracking - failed to enqueue user spend update. " 

1211 "user_id=%s, end_user_id=%s, response_cost=%s - %s", 

1212 user_id, 

1213 end_user_id, 

1214 response_cost, 

1215 str(e), 

1216 exc=e, 

1217 ) 

1218 

1219 async def _update_team_db( 

1220 self, 

1221 response_cost: float | None, 

1222 team_id: str | None, 

1223 user_id: str | None, 

1224 prisma_client: PrismaClient | None, 

1225 ): 

1226 try: 

1227 if team_id is None or prisma_client is None: 1227 ↛ 1233line 1227 didn't jump to line 1233 because the condition on line 1227 was always true

1228 verbose_proxy_logger.debug( 

1229 "track_cost_callback: team_id is None or prisma_client is None. Not tracking spend for team" 

1230 ) 

1231 return 

1232 

1233 await self.spend_update_queue.add_update( 

1234 update=SpendUpdateQueueItem( 

1235 entity_type=Litellm_EntityType.TEAM, 

1236 entity_id=team_id, 

1237 response_cost=response_cost, 

1238 ) 

1239 ) 

1240 

1241 try: 

1242 # Track spend of the team member within this team 

1243 if user_id is not None: 

1244 # key is "team_id::<value>::user_id::<value>" 

1245 team_member_key: Final = f"team_id::{team_id}::user_id::{user_id}" 

1246 await self.spend_update_queue.add_update( 

1247 update=SpendUpdateQueueItem( 

1248 entity_type=Litellm_EntityType.TEAM_MEMBER, 

1249 entity_id=team_member_key, 

1250 response_cost=response_cost, 

1251 ) 

1252 ) 

1253 except Exception as e: 

1254 spend_log_error( 

1255 "Spend tracking - failed to enqueue team member spend update. " 

1256 "team_id=%s, user_id=%s, response_cost=%s - %s", 

1257 team_id, 

1258 user_id, 

1259 response_cost, 

1260 str(e), 

1261 exc=e, 

1262 ) 

1263 except Exception as e: 

1264 spend_log_error( 

1265 "Spend tracking - failed to enqueue team spend update. team_id=%s, response_cost=%s - %s", 

1266 team_id, 

1267 response_cost, 

1268 str(e), 

1269 exc=e, 

1270 ) 

1271 raise e 

1272 

1273 async def _update_org_db( 

1274 self, 

1275 response_cost: float | None, 

1276 org_id: str | None, 

1277 user_id: str | None, 

1278 prisma_client: PrismaClient | None, 

1279 ): 

1280 try: 

1281 if org_id is None or prisma_client is None: 1281 ↛ 1287line 1281 didn't jump to line 1287 because the condition on line 1281 was always true

1282 verbose_proxy_logger.debug( 

1283 "track_cost_callback: org_id is None or prisma_client is None. Not tracking spend for org" 

1284 ) 

1285 return 

1286 

1287 await self.spend_update_queue.add_update( 

1288 update=SpendUpdateQueueItem( 

1289 entity_type=Litellm_EntityType.ORGANIZATION, 

1290 entity_id=org_id, 

1291 response_cost=response_cost, 

1292 ) 

1293 ) 

1294 

1295 if user_id is not None: 

1296 await self.spend_update_queue.add_update( 

1297 update=SpendUpdateQueueItem( 

1298 entity_type=Litellm_EntityType.ORGANIZATION_MEMBER, 

1299 entity_id=_org_member_transaction_key(org_id, user_id), 

1300 response_cost=response_cost, 

1301 ) 

1302 ) 

1303 except Exception as e: 

1304 spend_log_error( 

1305 "Spend tracking - failed to enqueue org spend update. org_id=%s, response_cost=%s - %s", 

1306 org_id, 

1307 response_cost, 

1308 str(e), 

1309 exc=e, 

1310 ) 

1311 raise e 

1312 

1313 async def _update_project_db( 

1314 self, 

1315 response_cost: float | None, 

1316 project_id: str | None, 

1317 prisma_client: PrismaClient | None, 

1318 ) -> None: 

1319 if project_id is None or prisma_client is None: 1319 ↛ 1321line 1319 didn't jump to line 1321 because the condition on line 1319 was always true

1320 return 

1321 try: 

1322 await self.spend_update_queue.add_update( 

1323 update=SpendUpdateQueueItem( 

1324 entity_type=Litellm_EntityType.PROJECT, 

1325 entity_id=project_id, 

1326 response_cost=response_cost, 

1327 ) 

1328 ) 

1329 except Exception as e: 

1330 spend_log_error( 

1331 "Spend tracking - failed to enqueue project spend update. project_id=%s, response_cost=%s - %s", 

1332 project_id, 

1333 response_cost, 

1334 str(e), 

1335 exc=e, 

1336 ) 

1337 raise e 

1338 

1339 async def _update_agent_db( 

1340 self, 

1341 response_cost: float | None, 

1342 agent_id: str | None, 

1343 prisma_client: PrismaClient | None, 

1344 ): 

1345 try: 

1346 if agent_id is None or prisma_client is None: 

1347 return 

1348 

1349 await self.spend_update_queue.add_update( 

1350 update=SpendUpdateQueueItem( 

1351 entity_type=Litellm_EntityType.AGENT, 

1352 entity_id=agent_id, 

1353 response_cost=response_cost, 

1354 ) 

1355 ) 

1356 except Exception as e: 

1357 spend_log_error( 

1358 "Spend tracking - failed to enqueue agent spend update. agent_id=%s, response_cost=%s - %s", 

1359 agent_id, 

1360 response_cost, 

1361 str(e), 

1362 exc=e, 

1363 ) 

1364 raise e 

1365 

1366 async def _update_tag_db( 

1367 self, 

1368 response_cost: float | None, 

1369 request_tags: str | None, 

1370 prisma_client: PrismaClient | None, 

1371 ): 

1372 """ 

1373 Update spend for all tags in the request. 

1374 

1375 Args: 

1376 response_cost: Cost of the request 

1377 request_tags: JSON string of tags list e.g. '["prod-tag", "test-tag"]' 

1378 prisma_client: Prisma client instance 

1379 """ 

1380 try: 

1381 if request_tags is None or prisma_client is None: 1381 ↛ 1382line 1381 didn't jump to line 1382 because the condition on line 1381 was never true

1382 return 

1383 

1384 # Parse tags from JSON string 

1385 tags: Sequence[object] = [] 

1386 if isinstance(request_tags, str): 1386 ↛ 1391line 1386 didn't jump to line 1391 because the condition on line 1386 was always true

1387 tags = safe_json_loads(request_tags, default=[]) 

1388 if not tags: 

1389 verbose_proxy_logger.debug("Failed to parse request_tags JSON: %s", request_tags) 

1390 return 

1391 elif isinstance(request_tags, list): 

1392 tags = request_tags 

1393 else: 

1394 return 

1395 

1396 # Update spend for each tag 

1397 for tag_name in tags: 

1398 if tag_name and isinstance(tag_name, str): 1398 ↛ 1397line 1398 didn't jump to line 1397 because the condition on line 1398 was always true

1399 await self.spend_update_queue.add_update( 

1400 update=SpendUpdateQueueItem( 

1401 entity_type=Litellm_EntityType.TAG, 

1402 entity_id=tag_name, 

1403 response_cost=response_cost, 

1404 ) 

1405 ) 

1406 except Exception as e: 

1407 spend_log_error( 

1408 "Spend tracking - failed to enqueue tag spend update. request_tags=%s, response_cost=%s - %s", 

1409 request_tags, 

1410 response_cost, 

1411 str(e), 

1412 exc=e, 

1413 ) 

1414 raise e 

1415 

1416 async def _update_model_access_group_db( 

1417 self, 

1418 response_cost: float | None, 

1419 request_model_access_groups: Sequence[str] | None, 

1420 served_model_id: str | None, 

1421 prisma_client: PrismaClient | None, 

1422 router: _DeploymentLookup | None = None, 

1423 ) -> None: 

1424 """ 

1425 Update spend for every model access group this request is billed against. 

1426 

1427 Args: 

1428 response_cost: Cost of the request, charged in full to each group 

1429 request_model_access_groups: Groups attributed at auth time, the upper bound on what may be debited 

1430 served_model_id: Deployment id actually served, used to narrow the attributed set 

1431 prisma_client: Prisma client instance 

1432 router: Deployment lookup used to re-resolve groups after a fallback 

1433 """ 

1434 try: 

1435 if prisma_client is None: 1435 ↛ 1436line 1435 didn't jump to line 1436 because the condition on line 1435 was never true

1436 return 

1437 

1438 for model_access_group in debitable_model_access_groups( 1438 ↛ 1443line 1438 didn't jump to line 1443 because the loop on line 1438 never started

1439 attributed=request_model_access_groups, 

1440 served_model_id=served_model_id, 

1441 router=router, 

1442 ): 

1443 await self.spend_update_queue.add_update( 

1444 update=SpendUpdateQueueItem( 

1445 entity_type=Litellm_EntityType.MODEL_ACCESS_GROUP, 

1446 entity_id=model_access_group, 

1447 response_cost=response_cost, 

1448 ) 

1449 ) 

1450 except Exception as e: # noqa: BLE001 # isolation: a helper failure must not stop the batch 

1451 spend_log_error( 

1452 "Spend tracking - failed to enqueue model access group spend update. " 

1453 "model_access_groups=%s, response_cost=%s - %s", 

1454 request_model_access_groups, 

1455 response_cost, 

1456 str(e), 

1457 exc=e, 

1458 ) 

1459 

1460 async def _insert_spend_log_to_db( 

1461 self, 

1462 payload: dict | SpendLogsPayload, 

1463 prisma_client: PrismaClient | None = None, 

1464 spend_logs_url: str | None = os.getenv("SPEND_LOGS_URL"), 

1465 ) -> PrismaClient | None: 

1466 verbose_proxy_logger.debug( 

1467 "Writing spend log to db - request_id: {}, spend: {}".format( 

1468 payload.get("request_id"), payload.get("spend") 

1469 ) 

1470 ) 

1471 if prisma_client is not None and spend_logs_url is not None or prisma_client is not None: 1471 ↛ 1478line 1471 didn't jump to line 1478 because the condition on line 1471 was always true

1472 from litellm.proxy.utils import enqueue_spend_logs, request_spend_log_flush 

1473 

1474 await enqueue_spend_logs(prisma_client, (payload,)) 

1475 if payload.get("call_type") in RESPONSES_SESSION_CALL_TYPES: 

1476 request_spend_log_flush(prisma_client) 

1477 else: 

1478 verbose_proxy_logger.debug("prisma_client is None. Skipping writing spend logs to db.") 

1479 

1480 return prisma_client 

1481 

1482 async def db_update_spend_transaction_handler( 

1483 self, 

1484 prisma_client: PrismaClient, 

1485 n_retry_times: int, 

1486 proxy_logging_obj: ProxyLogging, 

1487 ): 

1488 """ 

1489 Handles commiting update spend transactions to db 

1490 

1491 `UPDATES` can lead to deadlocks, hence we handle them separately 

1492 

1493 Args: 

1494 prisma_client: PrismaClient object 

1495 n_retry_times: int, number of retry times 

1496 proxy_logging_obj: ProxyLogging object 

1497 

1498 How this works: 

1499 - Check `general_settings.use_redis_transaction_buffer` 

1500 - If enabled, write in-memory transactions to Redis 

1501 - Check if this Pod should read from the DB 

1502 else: 

1503 - Regular flow of this method 

1504 """ 

1505 if RedisUpdateBuffer._should_commit_spend_updates_to_redis(): 1505 ↛ 1506line 1505 didn't jump to line 1506 because the condition on line 1505 was never true

1506 await self._commit_spend_updates_to_db_with_redis( 

1507 prisma_client=prisma_client, 

1508 n_retry_times=n_retry_times, 

1509 proxy_logging_obj=proxy_logging_obj, 

1510 ) 

1511 

1512 else: 

1513 await self._commit_spend_updates_to_db_without_redis_buffer( 

1514 prisma_client=prisma_client, 

1515 n_retry_times=n_retry_times, 

1516 proxy_logging_obj=proxy_logging_obj, 

1517 ) 

1518 

1519 async def _commit_spend_updates_to_db_with_redis( 

1520 self, 

1521 prisma_client: PrismaClient, 

1522 n_retry_times: int, 

1523 proxy_logging_obj: ProxyLogging, 

1524 ): 

1525 """ 

1526 Handler to commit spend updates to Redis and attempt to acquire lock to commit to db 

1527 

1528 This is a v2 scalable approach to first commit spend updates to redis, then commit to db 

1529 

1530 This minimizes DB Deadlocks since 

1531 - All pods only need to write their spend updates to redis 

1532 - Only 1 pod will commit to db at a time (based on if it can acquire the lock over writing to DB) 

1533 """ 

1534 await self.redis_update_buffer.store_in_memory_spend_updates_in_redis( 

1535 spend_update_queue=self.spend_update_queue, 

1536 daily_spend_update_queue=self.daily_spend_update_queue, 

1537 daily_team_spend_update_queue=self.daily_team_spend_update_queue, 

1538 daily_org_spend_update_queue=self.daily_org_spend_update_queue, 

1539 daily_end_user_spend_update_queue=self.daily_end_user_spend_update_queue, 

1540 daily_agent_spend_update_queue=self.daily_agent_spend_update_queue, 

1541 window_spend_update_queue=self.window_spend_update_queue, 

1542 ) 

1543 

1544 # Only commit from redis to db if this pod is the leader 

1545 if await self.pod_lock_manager.acquire_lock( 

1546 cronjob_id=DB_SPEND_UPDATE_JOB_NAME, 

1547 ): 

1548 verbose_proxy_logger.debug("acquired lock for spend updates") 

1549 

1550 uncommitted: dict[str, Any] = {} # mutable-ok: tracks popped categories still needing commit 

1551 committed_spend_tables: Final[list[_SpendTableName]] = [] # mutable-ok: filled as each table lands 

1552 

1553 try: 

1554 ( 

1555 db_spend_update_transactions, 

1556 daily_spend_update_transactions, 

1557 daily_team_spend_update_transactions, 

1558 daily_org_spend_update_transactions, 

1559 daily_end_user_spend_update_transactions, 

1560 daily_agent_spend_update_transactions, 

1561 window_spend_update_transactions, 

1562 ) = await self.redis_update_buffer.get_all_transactions_from_redis_buffer_pipeline() 

1563 

1564 uncommitted = { # mutable-ok: drives which popped categories still need re-queuing 

1565 "db_spend_update_transactions": db_spend_update_transactions, 

1566 "daily_spend_update_transactions": daily_spend_update_transactions, 

1567 "daily_team_spend_update_transactions": daily_team_spend_update_transactions, 

1568 "daily_org_spend_update_transactions": daily_org_spend_update_transactions, 

1569 "daily_end_user_spend_update_transactions": daily_end_user_spend_update_transactions, 

1570 "daily_agent_spend_update_transactions": daily_agent_spend_update_transactions, 

1571 "window_spend_update_transactions": window_spend_update_transactions, 

1572 } 

1573 

1574 if db_spend_update_transactions is not None: 

1575 verbose_proxy_logger.info( 

1576 "Spend tracking - committing spend updates from Redis to DB: " 

1577 "keys=%d, users=%d, teams=%d, orgs=%d, end_users=%d, team_members=%d, org_members=%d, " 

1578 "projects=%d, tags=%d, agents=%d, model_access_groups=%d", 

1579 len(db_spend_update_transactions.get("key_list_transactions") or ()), 

1580 len(db_spend_update_transactions.get("user_list_transactions") or ()), 

1581 len(db_spend_update_transactions.get("team_list_transactions") or ()), 

1582 len(db_spend_update_transactions.get("org_list_transactions") or ()), 

1583 len(db_spend_update_transactions.get("end_user_list_transactions") or ()), 

1584 len(db_spend_update_transactions.get("team_member_list_transactions") or ()), 

1585 len(db_spend_update_transactions.get("org_member_list_transactions") or ()), 

1586 len(db_spend_update_transactions.get("project_list_transactions") or ()), 

1587 len(db_spend_update_transactions.get("tag_list_transactions") or ()), 

1588 len(db_spend_update_transactions.get("agent_list_transactions") or ()), 

1589 len(db_spend_update_transactions.get("model_access_group_list_transactions") or ()), 

1590 ) 

1591 try: 

1592 await self._commit_spend_updates_to_db( 

1593 prisma_client=prisma_client, 

1594 n_retry_times=n_retry_times, 

1595 proxy_logging_obj=proxy_logging_obj, 

1596 db_spend_update_transactions=db_spend_update_transactions, 

1597 on_table_committed=committed_spend_tables.append, 

1598 ) 

1599 except Exception as e: 

1600 uncommitted["db_spend_update_transactions"] = _spend_tables_left_to_send( 

1601 db_spend_update_transactions, committed_spend_tables, e 

1602 ) 

1603 raise 

1604 uncommitted.pop("db_spend_update_transactions", None) 

1605 

1606 if daily_spend_update_transactions is not None: 

1607 await DBSpendUpdateWriter.update_daily_user_spend( 

1608 n_retry_times=n_retry_times, 

1609 prisma_client=prisma_client, 

1610 proxy_logging_obj=proxy_logging_obj, 

1611 daily_spend_transactions=daily_spend_update_transactions, 

1612 ) 

1613 uncommitted.pop("daily_spend_update_transactions", None) 

1614 

1615 if daily_team_spend_update_transactions is not None: 

1616 await DBSpendUpdateWriter.update_daily_team_spend( 

1617 n_retry_times=n_retry_times, 

1618 prisma_client=prisma_client, 

1619 proxy_logging_obj=proxy_logging_obj, 

1620 daily_spend_transactions=daily_team_spend_update_transactions, 

1621 ) 

1622 uncommitted.pop("daily_team_spend_update_transactions", None) 

1623 

1624 if daily_org_spend_update_transactions is not None: 

1625 await DBSpendUpdateWriter.update_daily_org_spend( 

1626 n_retry_times=n_retry_times, 

1627 prisma_client=prisma_client, 

1628 proxy_logging_obj=proxy_logging_obj, 

1629 daily_spend_transactions=daily_org_spend_update_transactions, 

1630 ) 

1631 uncommitted.pop("daily_org_spend_update_transactions", None) 

1632 

1633 if daily_end_user_spend_update_transactions is not None: 

1634 await DBSpendUpdateWriter.update_daily_end_user_spend( 

1635 n_retry_times=n_retry_times, 

1636 prisma_client=prisma_client, 

1637 proxy_logging_obj=proxy_logging_obj, 

1638 daily_spend_transactions=daily_end_user_spend_update_transactions, 

1639 ) 

1640 uncommitted.pop("daily_end_user_spend_update_transactions", None) 

1641 

1642 if daily_agent_spend_update_transactions is not None: 

1643 await DBSpendUpdateWriter.update_daily_agent_spend( 

1644 n_retry_times=n_retry_times, 

1645 prisma_client=prisma_client, 

1646 proxy_logging_obj=proxy_logging_obj, 

1647 daily_spend_transactions=daily_agent_spend_update_transactions, 

1648 ) 

1649 uncommitted.pop("daily_agent_spend_update_transactions", None) 

1650 if window_spend_update_transactions is not None: 

1651 try: 

1652 await DBSpendUpdateWriter._commit_window_spend_updates( 

1653 prisma_client=prisma_client, 

1654 window_spend_transactions=window_spend_update_transactions, 

1655 ) 

1656 except Exception as e: 

1657 if not _spend_commit_failure_is_requeue_safe(e): 

1658 uncommitted.pop("window_spend_update_transactions", None) 

1659 spend_log_error( 

1660 "Spend tracking - dropped %d budget window increments: the failed statement may have " 

1661 "applied or the database refused the data, so re-sending it is not safe. Error: %s", 

1662 len(window_spend_update_transactions), 

1663 str(e), 

1664 exc=e, 

1665 ) 

1666 raise 

1667 uncommitted.pop("window_spend_update_transactions", None) 

1668 except Exception as e: 

1669 spend_log_error( 

1670 "Spend tracking - failed to commit spend updates from Redis to DB. " 

1671 "Re-queuing uncommitted transactions to Redis for retry on next tick. Error: %s", 

1672 str(e), 

1673 exc=e, 

1674 ) 

1675 finally: 

1676 to_restore = { # mutable-ok: transient kwargs payload consumed immediately below 

1677 name: txns for name, txns in uncommitted.items() if txns is not None 

1678 } 

1679 if to_restore: 

1680 await self.redis_update_buffer.restore_transactions_to_redis(**to_restore) 

1681 await self.pod_lock_manager.release_lock( 

1682 cronjob_id=DB_SPEND_UPDATE_JOB_NAME, 

1683 ) 

1684 

1685 async def _flush_daily_spend_queue( 

1686 self, 

1687 queue: DailySpendUpdateQueue, 

1688 entity_type: Literal["user", "team", "org", "tag", "end_user", "agent"], 

1689 commit: _DailySpendCommit[_DailySpendTransactionT], 

1690 n_retry_times: int, 

1691 prisma_client: PrismaClient, 

1692 proxy_logging_obj: ProxyLogging, 

1693 ) -> None: 

1694 transactions: Final = await queue.flush_and_get_aggregated_daily_spend_update_transactions() 

1695 commit_started: Final = asyncio.Event() 

1696 commit_task: Final = _start_daily_spend_commit( 

1697 commit_started, 

1698 lambda: commit( 

1699 n_retry_times=n_retry_times, 

1700 prisma_client=prisma_client, 

1701 proxy_logging_obj=proxy_logging_obj, 

1702 daily_spend_transactions=cast(dict[str, _DailySpendTransactionT], transactions), 

1703 ), 

1704 ) 

1705 try: 

1706 await asyncio.shield(commit_task) 

1707 except asyncio.CancelledError: 

1708 if commit_started.is_set(): 

1709 queue.track_interrupted_commit( 

1710 _requeue_daily_spend_the_commit_left_behind(commit_task, queue, entity_type, transactions) 

1711 ) 

1712 raise 

1713 commit_task.cancel() 

1714 if transactions: 

1715 await queue.add_update(transactions) 

1716 raise 

1717 except Exception as e: # noqa: BLE001 # whatever failed here, the other tables must still flush 

1718 if not transactions: 1718 ↛ 1719line 1718 didn't jump to line 1719 because the condition on line 1718 was never true

1719 return 

1720 spend_log_error( 

1721 "Spend tracking - failed to commit daily %s spend updates. " 

1722 "Re-queued %d rows for retry on next tick. Error: %s", 

1723 entity_type, 

1724 len(transactions), 

1725 str(e), 

1726 exc=e, 

1727 ) 

1728 await queue.add_update(transactions) 

1729 

1730 async def _commit_spend_updates_to_db_without_redis_buffer( 

1731 self, 

1732 prisma_client: PrismaClient, 

1733 n_retry_times: int, 

1734 proxy_logging_obj: ProxyLogging, 

1735 ): 

1736 """ 

1737 Commits all the spend `UPDATE` transactions to the Database 

1738 

1739 This is the regular flow of committing to db without using a redis buffer 

1740 

1741 Note: This flow causes Deadlocks in production (1K RPS+). Use self._commit_spend_updates_to_db_with_redis() instead if you expect 1K+ RPS. 

1742 """ 

1743 

1744 # Aggregate all in memory spend updates (key, user, end_user, team, team_member, org) and commit to db 

1745 ################## Spend Update Transactions ################## 

1746 db_spend_update_transactions: Final = ( 

1747 await self.spend_update_queue.flush_and_get_aggregated_db_spend_update_transactions() 

1748 ) 

1749 await self._commit_spend_updates_to_db( 

1750 prisma_client=prisma_client, 

1751 n_retry_times=n_retry_times, 

1752 proxy_logging_obj=proxy_logging_obj, 

1753 db_spend_update_transactions=db_spend_update_transactions, 

1754 ) 

1755 

1756 ################## Daily Spend Update Transactions ################## 

1757 # Aggregate all in memory daily spend transactions and commit to db 

1758 await self._flush_daily_spend_queue( 

1759 queue=self.daily_spend_update_queue, 

1760 entity_type="user", 

1761 commit=DBSpendUpdateWriter.update_daily_user_spend, 

1762 n_retry_times=n_retry_times, 

1763 prisma_client=prisma_client, 

1764 proxy_logging_obj=proxy_logging_obj, 

1765 ) 

1766 

1767 ################## Daily Team Spend Update Transactions ################## 

1768 # Aggregate all in memory daily team spend transactions and commit to db 

1769 await self._flush_daily_spend_queue( 

1770 queue=self.daily_team_spend_update_queue, 

1771 entity_type="team", 

1772 commit=DBSpendUpdateWriter.update_daily_team_spend, 

1773 n_retry_times=n_retry_times, 

1774 prisma_client=prisma_client, 

1775 proxy_logging_obj=proxy_logging_obj, 

1776 ) 

1777 

1778 ################## Daily Organization Spend Update Transactions ################## 

1779 # Aggregate all in memory daily org spend transactions and commit to db 

1780 await self._flush_daily_spend_queue( 

1781 queue=self.daily_org_spend_update_queue, 

1782 entity_type="org", 

1783 commit=DBSpendUpdateWriter.update_daily_org_spend, 

1784 n_retry_times=n_retry_times, 

1785 prisma_client=prisma_client, 

1786 proxy_logging_obj=proxy_logging_obj, 

1787 ) 

1788 

1789 # NOTE: Daily tag spend is committed by a separate scheduler job. 

1790 

1791 ################## Daily End-User Spend Update Transactions ################## 

1792 # Aggregate all in memory daily end-user spend transactions and commit to db 

1793 await self._flush_daily_spend_queue( 

1794 queue=self.daily_end_user_spend_update_queue, 

1795 entity_type="end_user", 

1796 commit=DBSpendUpdateWriter.update_daily_end_user_spend, 

1797 n_retry_times=n_retry_times, 

1798 prisma_client=prisma_client, 

1799 proxy_logging_obj=proxy_logging_obj, 

1800 ) 

1801 

1802 ################## Daily Agent Spend Update Transactions ################## 

1803 # Aggregate all in memory daily agent spend transactions and commit to db 

1804 await self._flush_daily_spend_queue( 

1805 queue=self.daily_agent_spend_update_queue, 

1806 entity_type="agent", 

1807 commit=DBSpendUpdateWriter.update_daily_agent_spend, 

1808 n_retry_times=n_retry_times, 

1809 prisma_client=prisma_client, 

1810 proxy_logging_obj=proxy_logging_obj, 

1811 ) 

1812 

1813 ################## Budget Window Spend Update Transactions ################## 

1814 # Aggregate all in memory budget window spend transactions and commit to db 

1815 window_spend_update_transactions: Final = ( 

1816 await self.window_spend_update_queue.flush_and_get_aggregated_window_spend_transactions() 

1817 ) 

1818 

1819 try: 

1820 await DBSpendUpdateWriter._commit_window_spend_updates( 

1821 prisma_client=prisma_client, 

1822 window_spend_transactions=window_spend_update_transactions, 

1823 ) 

1824 except Exception as e: # noqa: BLE001 # the increments go back on the queue; the rest of the flush must run 

1825 if _spend_commit_failure_is_requeue_safe(e): 

1826 spend_log_error( 

1827 "Spend tracking - failed to commit budget window spend updates. " 

1828 "Re-queued %d window increments for retry on next tick. Error: %s", 

1829 len(window_spend_update_transactions), 

1830 str(e), 

1831 exc=e, 

1832 ) 

1833 await self.window_spend_update_queue.update_queue.put(window_spend_update_transactions) 

1834 else: 

1835 spend_log_error( 

1836 "Spend tracking - dropped %d budget window increments: the failed statement may have " 

1837 "applied or the database refused the data, so re-sending it is not safe. Error: %s", 

1838 len(window_spend_update_transactions), 

1839 str(e), 

1840 exc=e, 

1841 ) 

1842 

1843 ################## Tool Registry Upserts ################## 

1844 await self._flush_tool_discovery_queue(prisma_client=prisma_client) 

1845 

1846 async def _commit_daily_tag_spend_to_db( 

1847 self, 

1848 prisma_client: PrismaClient, 

1849 n_retry_times: int, 

1850 proxy_logging_obj: ProxyLogging, 

1851 ): 

1852 """ 

1853 Commit only tag spend updates to database. 

1854 This is called by a separate scheduler job at a longer interval. 

1855 """ 

1856 await self._flush_daily_spend_queue( 

1857 queue=self.daily_tag_spend_update_queue, 

1858 entity_type="tag", 

1859 commit=DBSpendUpdateWriter.update_daily_tag_spend, 

1860 n_retry_times=n_retry_times, 

1861 prisma_client=prisma_client, 

1862 proxy_logging_obj=proxy_logging_obj, 

1863 ) 

1864 

1865 async def _commit_daily_tag_spend_to_db_with_redis( 

1866 self, 

1867 prisma_client: PrismaClient, 

1868 n_retry_times: int, 

1869 proxy_logging_obj: ProxyLogging, 

1870 ): 

1871 """ 

1872 Commit daily tag spend updates using Redis buffering. 

1873 

1874 This lets the dedicated daily tag scheduler drain both in-memory and 

1875 Redis-backed tag transactions. 

1876 """ 

1877 await self.redis_update_buffer.store_in_memory_daily_tag_spend_updates_in_redis( 

1878 daily_tag_spend_update_queue=self.daily_tag_spend_update_queue, 

1879 ) 

1880 

1881 if await self.pod_lock_manager.acquire_lock( 

1882 cronjob_id=DB_DAILY_TAG_SPEND_UPDATE_JOB_NAME, 

1883 ): 

1884 verbose_proxy_logger.debug("acquired lock for daily tag spend updates") 

1885 try: 

1886 await self._drain_and_commit_daily_tag_spend_from_redis( 

1887 prisma_client=prisma_client, 

1888 n_retry_times=n_retry_times, 

1889 proxy_logging_obj=proxy_logging_obj, 

1890 ) 

1891 except Exception as e: 

1892 spend_log_error( 

1893 "Spend tracking - failed to commit daily tag spend updates from Redis to DB. " 

1894 "Re-queuing to Redis for retry on next tick. Error: %s", 

1895 str(e), 

1896 exc=e, 

1897 ) 

1898 finally: 

1899 await self.pod_lock_manager.release_lock( 

1900 cronjob_id=DB_DAILY_TAG_SPEND_UPDATE_JOB_NAME, 

1901 ) 

1902 

1903 @staticmethod 

1904 async def _commit_window_spend_updates( 

1905 prisma_client: PrismaClient, 

1906 window_spend_transactions: Sequence[WindowSpendTransaction], 

1907 ) -> None: 

1908 """ 

1909 Commit per-budget-window spend increments to LiteLLM_BudgetWindowSpend. 

1910 

1911 Raises on failure so the caller re-queues the increments: budget 

1912 enforcement trusts a current row without reconciling it against 

1913 LiteLLM_SpendLogs, so a dropped increment would let the entity spend 

1914 past its window limit after the next counter reseed. 

1915 """ 

1916 from litellm.proxy.db.budget_window_spend_writer import ( 

1917 commit_window_spend_updates, 

1918 ) 

1919 

1920 await commit_window_spend_updates( 

1921 prisma_client=prisma_client, 

1922 transactions=window_spend_transactions, 

1923 ) 

1924 

1925 async def _drain_and_commit_daily_tag_spend_from_redis( 

1926 self, 

1927 prisma_client: PrismaClient, 

1928 n_retry_times: int, 

1929 proxy_logging_obj: ProxyLogging, 

1930 ) -> None: 

1931 """ 

1932 Drain the Redis tag spend buffer and commit it, restoring the drained transactions if the commit fails. 

1933 

1934 The drain is destructive, so a failed commit must push the transactions back for the next tick 

1935 or their spend is lost permanently. 

1936 """ 

1937 await _settle_interrupted_commits(self.interrupted_tag_commits) 

1938 daily_tag_spend_update_transactions: Final = ( 

1939 await self.redis_update_buffer.get_all_daily_tag_spend_update_transactions_from_redis_buffer() 

1940 ) 

1941 if not daily_tag_spend_update_transactions: 

1942 return 

1943 

1944 commit_started: Final = asyncio.Event() 

1945 commit_task: Final = _start_daily_spend_commit( 

1946 commit_started, 

1947 lambda: DBSpendUpdateWriter.update_daily_tag_spend( 

1948 n_retry_times=n_retry_times, 

1949 prisma_client=prisma_client, 

1950 proxy_logging_obj=proxy_logging_obj, 

1951 daily_spend_transactions=daily_tag_spend_update_transactions, 

1952 ), 

1953 ) 

1954 try: 

1955 await asyncio.shield(commit_task) 

1956 except BaseException: # noqa: BLE001 # a cancel must restore the drained rows before its rollback returns 

1957 if commit_started.is_set(): 

1958 _track_interrupted_commit( 

1959 self.interrupted_tag_commits, 

1960 _restore_tag_spend_the_commit_left_behind( 

1961 commit_task, 

1962 self.redis_update_buffer, 

1963 daily_tag_spend_update_transactions, 

1964 ), 

1965 ) 

1966 raise 

1967 commit_task.cancel() 

1968 await self.redis_update_buffer.restore_transactions_to_redis( 

1969 daily_tag_spend_update_transactions=daily_tag_spend_update_transactions, 

1970 ) 

1971 raise 

1972 

1973 async def _flush_tool_discovery_queue( 

1974 self, 

1975 prisma_client: PrismaClient, 

1976 ) -> None: 

1977 """Flush ToolDiscoveryQueue and batch-upsert new tools into LiteLLM_ToolTable.""" 

1978 from litellm.proxy.db.tool_registry_writer import batch_upsert_tools 

1979 

1980 try: 

1981 items: Final = self.tool_discovery_queue.flush() 

1982 if items: 1982 ↛ 1983line 1982 didn't jump to line 1983 because the condition on line 1982 was never true

1983 await batch_upsert_tools(prisma_client=prisma_client, items=items) 

1984 except Exception as e: 

1985 verbose_proxy_logger.debug("_flush_tool_discovery_queue error (non-blocking): %s", e) 

1986 

1987 @staticmethod 

1988 async def _handle_spend_update_failure( 

1989 e: Exception, 

1990 attempt: int, 

1991 n_retry_times: int, 

1992 start_time: float, 

1993 proxy_logging_obj: ProxyLogging, 

1994 ) -> None: 

1995 """Retry a failed spend-update transaction on connection errors or deadlocks, else re-raise.""" 

1996 from litellm.proxy.db.exception_handler import PrismaDBExceptionHandler 

1997 from litellm.proxy.utils import _raise_failed_update_spend_exception 

1998 

1999 is_retryable = isinstance(e, DB_RETRY_SAFE_ERROR_TYPES) or PrismaDBExceptionHandler.is_deadlock_error(e) 

2000 if not is_retryable or attempt >= n_retry_times: 

2001 _raise_failed_update_spend_exception(e=e, start_time=start_time, proxy_logging_obj=proxy_logging_obj) 

2002 verbose_proxy_logger.warning( 

2003 "Retrying spend update after retryable DB error (attempt %s/%s): %s", 

2004 attempt + 1, 

2005 n_retry_times, 

2006 e, 

2007 ) 

2008 await asyncio.sleep(random.uniform(2**attempt, 2 ** (attempt + 1))) 

2009 

2010 async def _commit_spend_updates_to_db( 

2011 self, 

2012 prisma_client: PrismaClient, 

2013 n_retry_times: int, 

2014 proxy_logging_obj: ProxyLogging, 

2015 db_spend_update_transactions: DBSpendUpdateTransactions, 

2016 on_table_committed: Callable[[_SpendTableName], None] | None = None, 

2017 ): 

2018 """ 

2019 Commits all the spend `UPDATE` transactions to the Database 

2020 

2021 """ 

2022 from litellm.proxy.utils import ProxyUpdateSpend 

2023 

2024 ### UPDATE USER TABLE ### 

2025 user_list_transactions: Final = db_spend_update_transactions["user_list_transactions"] 

2026 verbose_proxy_logger.debug("User Spend transactions: %s", user_list_transactions) 

2027 if user_list_transactions is not None and len(user_list_transactions.keys()) > 0: 

2028 for i in range(n_retry_times + 1): 2028 ↛ 2050line 2028 didn't jump to line 2050 because the loop on line 2028 didn't complete

2029 start_time = time.time() 

2030 try: 

2031 async with _spend_update_tx(prisma_client) as transaction: 

2032 async with transaction.batch_() as batcher: 

2033 # Sort by ID for consistent lock ordering across pods to prevent deadlocks. 

2034 # batch_() issues statements sequentially within the tx, so iteration 

2035 # order = lock acquisition order. 

2036 for user_id, response_cost in sorted(user_list_transactions.items()): 

2037 batcher.litellm_usertable.update_many( 

2038 where={"user_id": user_id}, 

2039 data={"spend": {"increment": response_cost}}, 

2040 ) 

2041 break 

2042 except Exception as e: 

2043 await self._handle_spend_update_failure( 

2044 e=e, 

2045 attempt=i, 

2046 n_retry_times=n_retry_times, 

2047 start_time=start_time, 

2048 proxy_logging_obj=proxy_logging_obj, 

2049 ) 

2050 if on_table_committed is not None: 2050 ↛ 2051line 2050 didn't jump to line 2051 because the condition on line 2050 was never true

2051 on_table_committed("user_list_transactions") 

2052 

2053 ### UPDATE END-USER TABLE ### 

2054 end_user_list_transactions: Final = db_spend_update_transactions["end_user_list_transactions"] 

2055 verbose_proxy_logger.debug("End-User Spend transactions: %s", end_user_list_transactions) 

2056 if end_user_list_transactions is not None and len(end_user_list_transactions.keys()) > 0: 

2057 await ProxyUpdateSpend.update_end_user_spend( 

2058 n_retry_times=n_retry_times, 

2059 prisma_client=prisma_client, 

2060 proxy_logging_obj=proxy_logging_obj, 

2061 end_user_list_transactions=end_user_list_transactions, 

2062 ) 

2063 if on_table_committed is not None: 2063 ↛ 2064line 2063 didn't jump to line 2064 because the condition on line 2063 was never true

2064 on_table_committed("end_user_list_transactions") 

2065 ### UPDATE KEY TABLE ### 

2066 key_list_transactions: Final = db_spend_update_transactions["key_list_transactions"] 

2067 verbose_proxy_logger.debug("KEY Spend transactions: %s", key_list_transactions) 

2068 if key_list_transactions is not None and len(key_list_transactions.keys()) > 0: 

2069 for i in range(n_retry_times + 1): 2069 ↛ 2094line 2069 didn't jump to line 2094 because the loop on line 2069 didn't complete

2070 start_time = time.time() 

2071 try: 

2072 async with _spend_update_tx(prisma_client) as transaction: 

2073 async with transaction.batch_() as batcher: 

2074 # Sort by token for consistent lock ordering across pods to prevent deadlocks. 

2075 for token, response_cost in sorted(key_list_transactions.items()): 

2076 spend_increment: _SpendIncrement = {"increment": response_cost} 

2077 batcher.litellm_verificationtoken.update_many( # 'update_many' prevents error from being raised if no row exists 

2078 where={"token": token}, 

2079 data={ 

2080 "spend": spend_increment, 

2081 "total_spend": spend_increment, 

2082 "last_active": datetime.now(timezone.utc), 

2083 }, 

2084 ) 

2085 break 

2086 except Exception as e: 

2087 await self._handle_spend_update_failure( 

2088 e=e, 

2089 attempt=i, 

2090 n_retry_times=n_retry_times, 

2091 start_time=start_time, 

2092 proxy_logging_obj=proxy_logging_obj, 

2093 ) 

2094 if on_table_committed is not None: 2094 ↛ 2095line 2094 didn't jump to line 2095 because the condition on line 2094 was never true

2095 on_table_committed("key_list_transactions") 

2096 

2097 ### UPDATE TEAM TABLE ### 

2098 team_list_transactions: Final = db_spend_update_transactions["team_list_transactions"] 

2099 verbose_proxy_logger.debug("Team Spend transactions: %s", team_list_transactions) 

2100 if team_list_transactions is not None and len(team_list_transactions.keys()) > 0: 2100 ↛ 2101line 2100 didn't jump to line 2101 because the condition on line 2100 was never true

2101 for i in range(n_retry_times + 1): 

2102 start_time = time.time() 

2103 try: 

2104 async with _spend_update_tx(prisma_client) as transaction: 

2105 async with transaction.batch_() as batcher: 

2106 # Sort by team_id for consistent lock ordering across pods to prevent deadlocks. 

2107 for team_id, response_cost in sorted(team_list_transactions.items()): 

2108 verbose_proxy_logger.debug( 

2109 "Updating spend for team id=%s by %s", team_id, response_cost 

2110 ) 

2111 batcher.litellm_teamtable.update_many( # 'update_many' prevents error from being raised if no row exists 

2112 where={"team_id": team_id}, 

2113 data={"spend": {"increment": response_cost}}, 

2114 ) 

2115 break 

2116 except Exception as e: 

2117 await self._handle_spend_update_failure( 

2118 e=e, 

2119 attempt=i, 

2120 n_retry_times=n_retry_times, 

2121 start_time=start_time, 

2122 proxy_logging_obj=proxy_logging_obj, 

2123 ) 

2124 if on_table_committed is not None: 2124 ↛ 2125line 2124 didn't jump to line 2125 because the condition on line 2124 was never true

2125 on_table_committed("team_list_transactions") 

2126 

2127 ### UPDATE TEAM Membership TABLE with spend ### 

2128 team_member_list_transactions: Final = db_spend_update_transactions["team_member_list_transactions"] 

2129 verbose_proxy_logger.debug("Team Membership Spend transactions: %s", team_member_list_transactions) 

2130 if team_member_list_transactions is not None and len(team_member_list_transactions.keys()) > 0: 2130 ↛ 2132line 2130 didn't jump to line 2132 because the condition on line 2130 was never true

2131 # Track which team memberships will be updated for cache invalidation 

2132 team_memberships_to_invalidate: Final[list[tuple[str, str]]] = [] 

2133 for key in team_member_list_transactions: 

2134 # key is "team_id::<value>::user_id::<value>" 

2135 team_id = key.split("::")[1] 

2136 user_id = key.split("::")[3] 

2137 team_memberships_to_invalidate.append((user_id, team_id)) 

2138 

2139 for i in range(n_retry_times + 1): 

2140 start_time = time.time() 

2141 try: 

2142 async with _spend_update_tx(prisma_client) as transaction: 

2143 await _write_team_member_spend(transaction, team_member_list_transactions) 

2144 # Transaction succeeded, break out of retry loop 

2145 break 

2146 except Exception as e: 

2147 await self._handle_spend_update_failure( 

2148 e=e, 

2149 attempt=i, 

2150 n_retry_times=n_retry_times, 

2151 start_time=start_time, 

2152 proxy_logging_obj=proxy_logging_obj, 

2153 ) 

2154 if on_table_committed is not None: 

2155 on_table_committed("team_member_list_transactions") 

2156 

2157 # Invalidate cache for updated team memberships 

2158 # This ensures budget checks read fresh spend data from the database 

2159 if team_memberships_to_invalidate and proxy_logging_obj is not None: 

2160 user_api_key_cache: Final = proxy_logging_obj.call_details.get("user_api_key_cache") 

2161 if user_api_key_cache is not None: 

2162 for user_id, team_id in team_memberships_to_invalidate: 

2163 cache_key = f"team_membership:{user_id}:{team_id}" 

2164 await user_api_key_cache.async_delete_cache(key=cache_key) 

2165 verbose_proxy_logger.debug( 

2166 "Invalidated team membership cache for user_id=%s, team_id=%s", user_id, team_id 

2167 ) 

2168 elif on_table_committed is not None: 2168 ↛ 2169line 2168 didn't jump to line 2169 because the condition on line 2168 was never true

2169 on_table_committed("team_member_list_transactions") 

2170 

2171 ### UPDATE ORG TABLE ### 

2172 org_list_transactions: Final = db_spend_update_transactions["org_list_transactions"] 

2173 verbose_proxy_logger.debug("Org Spend transactions: %s", org_list_transactions) 

2174 if org_list_transactions is not None and len(org_list_transactions.keys()) > 0: 2174 ↛ 2175line 2174 didn't jump to line 2175 because the condition on line 2174 was never true

2175 for i in range(n_retry_times + 1): 

2176 start_time = time.time() 

2177 try: 

2178 async with _spend_update_tx(prisma_client) as transaction: 

2179 async with transaction.batch_() as batcher: 

2180 # Sort by org_id for consistent lock ordering across pods to prevent deadlocks. 

2181 for org_id, response_cost in sorted(org_list_transactions.items()): 

2182 batcher.litellm_organizationtable.update_many( # 'update_many' prevents error from being raised if no row exists 

2183 where={"organization_id": org_id}, 

2184 data={"spend": {"increment": response_cost}}, 

2185 ) 

2186 break 

2187 except Exception as e: 

2188 await self._handle_spend_update_failure( 

2189 e=e, 

2190 attempt=i, 

2191 n_retry_times=n_retry_times, 

2192 start_time=start_time, 

2193 proxy_logging_obj=proxy_logging_obj, 

2194 ) 

2195 if on_table_committed is not None: 2195 ↛ 2196line 2195 didn't jump to line 2196 because the condition on line 2195 was never true

2196 on_table_committed("org_list_transactions") 

2197 

2198 org_member_list_transactions: Final = db_spend_update_transactions.get("org_member_list_transactions") 

2199 verbose_proxy_logger.debug("Org Membership Spend transactions: %s", org_member_list_transactions) 

2200 if org_member_list_transactions is not None and len(org_member_list_transactions.keys()) > 0: 2200 ↛ 2201line 2200 didn't jump to line 2201 because the condition on line 2200 was never true

2201 for i in range(n_retry_times + 1): 

2202 start_time = time.time() 

2203 try: 

2204 async with _spend_update_tx(prisma_client) as transaction, transaction.batch_() as batcher: 

2205 for key, response_cost in sorted(org_member_list_transactions.items()): 

2206 _, quoted_org_id, _, quoted_user_id = key.split("::") 

2207 batcher.litellm_organizationmembership.update_many( 

2208 where={"organization_id": unquote(quoted_org_id), "user_id": unquote(quoted_user_id)}, 

2209 data={"spend": {"increment": response_cost}}, 

2210 ) 

2211 break 

2212 except Exception as e: 

2213 await self._handle_spend_update_failure( 

2214 e=e, 

2215 attempt=i, 

2216 n_retry_times=n_retry_times, 

2217 start_time=start_time, 

2218 proxy_logging_obj=proxy_logging_obj, 

2219 ) 

2220 if on_table_committed is not None: 2220 ↛ 2221line 2220 didn't jump to line 2221 because the condition on line 2220 was never true

2221 on_table_committed("org_member_list_transactions") 

2222 

2223 ### UPDATE PROJECT TABLE ### 

2224 project_list_transactions: Final = db_spend_update_transactions.get("project_list_transactions") 

2225 await DBSpendUpdateWriter._update_entity_spend_in_db( 

2226 entity_name="Project", 

2227 transactions=project_list_transactions, 

2228 table_accessor="litellm_projecttable", 

2229 where_field="project_id", 

2230 n_retry_times=n_retry_times, 

2231 prisma_client=prisma_client, 

2232 proxy_logging_obj=proxy_logging_obj, 

2233 ) 

2234 if on_table_committed is not None: 2234 ↛ 2235line 2234 didn't jump to line 2235 because the condition on line 2234 was never true

2235 on_table_committed("project_list_transactions") 

2236 await DBSpendUpdateWriter._invalidate_project_caches( 

2237 project_ids=tuple(project_list_transactions or ()), 

2238 proxy_logging_obj=proxy_logging_obj, 

2239 ) 

2240 

2241 ### UPDATE TAG TABLE ### 

2242 tag_list_transactions: Final = db_spend_update_transactions["tag_list_transactions"] 

2243 await DBSpendUpdateWriter._update_entity_spend_in_db( 

2244 entity_name="Tag", 

2245 transactions=tag_list_transactions, 

2246 table_accessor="litellm_tagtable", 

2247 where_field="tag_name", 

2248 n_retry_times=n_retry_times, 

2249 prisma_client=prisma_client, 

2250 proxy_logging_obj=proxy_logging_obj, 

2251 ) 

2252 if on_table_committed is not None: 2252 ↛ 2253line 2252 didn't jump to line 2253 because the condition on line 2252 was never true

2253 on_table_committed("tag_list_transactions") 

2254 

2255 ### UPDATE MODEL ACCESS GROUP TABLE ### 

2256 model_access_group_list_transactions: Final = db_spend_update_transactions.get( 

2257 "model_access_group_list_transactions" 

2258 ) 

2259 await DBSpendUpdateWriter._update_entity_spend_in_db( 

2260 entity_name="Model access group", 

2261 transactions=model_access_group_list_transactions, 

2262 table_accessor="litellm_modelaccessgroupbudgettable", 

2263 where_field="access_group_name", 

2264 n_retry_times=n_retry_times, 

2265 prisma_client=prisma_client, 

2266 proxy_logging_obj=proxy_logging_obj, 

2267 ) 

2268 if on_table_committed is not None: 2268 ↛ 2269line 2268 didn't jump to line 2269 because the condition on line 2268 was never true

2269 on_table_committed("model_access_group_list_transactions") 

2270 

2271 ### UPDATE AGENT TABLE ### 

2272 agent_list_transactions: Final = db_spend_update_transactions["agent_list_transactions"] 

2273 await DBSpendUpdateWriter._update_entity_spend_in_db( 

2274 entity_name="Agent", 

2275 transactions=agent_list_transactions, 

2276 table_accessor="litellm_agentstable", 

2277 where_field="agent_id", 

2278 n_retry_times=n_retry_times, 

2279 prisma_client=prisma_client, 

2280 proxy_logging_obj=proxy_logging_obj, 

2281 ) 

2282 if on_table_committed is not None: 2282 ↛ 2283line 2282 didn't jump to line 2283 because the condition on line 2282 was never true

2283 on_table_committed("agent_list_transactions") 

2284 

2285 @staticmethod 

2286 async def _invalidate_project_caches(project_ids: Sequence[str], proxy_logging_obj: ProxyLogging | None) -> None: 

2287 if not project_ids or proxy_logging_obj is None: 2287 ↛ 2289line 2287 didn't jump to line 2289 because the condition on line 2287 was always true

2288 return 

2289 user_api_key_cache: Final = proxy_logging_obj.call_details.get("user_api_key_cache") 

2290 if user_api_key_cache is None: 

2291 return 

2292 for project_id in project_ids: 

2293 await user_api_key_cache.async_delete_cache(key=project_cache_key(project_id)) 

2294 

2295 @staticmethod 

2296 async def _update_entity_spend_in_db( 

2297 entity_name: str, 

2298 transactions: dict[str, float] | None, 

2299 table_accessor: _EntitySpendTable, 

2300 where_field: str, 

2301 n_retry_times: int, 

2302 prisma_client: PrismaClient, 

2303 proxy_logging_obj: ProxyLogging, 

2304 ): 

2305 """ 

2306 Helper function to update spend for any entity type (team, org, tag, etc). 

2307 

2308 Args: 

2309 entity_name: Name of entity for logging (e.g., "Team", "Org", "Tag") 

2310 transactions: Dictionary of {entity_id: response_cost} 

2311 table_accessor: Prisma table accessor (e.g., prisma_client.db.litellm_teamtable) 

2312 where_field: Field name for where clause (e.g., "team_id", "organization_id", "tag_name") 

2313 n_retry_times: Number of retries on failure 

2314 prisma_client: Prisma client instance 

2315 proxy_logging_obj: Proxy logging object 

2316 """ 

2317 verbose_proxy_logger.debug("%s Spend transactions: %s", entity_name, transactions) 

2318 if transactions is not None and len(transactions.keys()) > 0: 

2319 for i in range(n_retry_times + 1): 2319 ↛ exitline 2319 didn't return from function '_update_entity_spend_in_db' because the loop on line 2319 didn't complete

2320 start_time = time.time() 

2321 try: 

2322 async with _spend_update_tx(prisma_client) as transaction: 

2323 async with transaction.batch_() as batcher: 

2324 # Sort by entity_id for consistent lock ordering across pods to prevent deadlocks. 

2325 for entity_id, response_cost in sorted(transactions.items()): 

2326 verbose_proxy_logger.debug( 

2327 "Updating spend for %s %s=%s by %s", 

2328 entity_name, 

2329 where_field, 

2330 entity_id, 

2331 response_cost, 

2332 ) 

2333 _entity_spend_table(batcher, table_accessor).update_many( 

2334 where={where_field: entity_id}, 

2335 data={"spend": {"increment": response_cost}}, 

2336 ) 

2337 break 

2338 except Exception as e: 

2339 await DBSpendUpdateWriter._handle_spend_update_failure( 

2340 e=e, 

2341 attempt=i, 

2342 n_retry_times=n_retry_times, 

2343 start_time=start_time, 

2344 proxy_logging_obj=proxy_logging_obj, 

2345 ) 

2346 

2347 # fmt: off 

2348 

2349 @overload 

2350 @staticmethod 

2351 async def _update_daily_spend( 

2352 n_retry_times: int, 

2353 prisma_client: PrismaClient, 

2354 proxy_logging_obj: ProxyLogging, 

2355 daily_spend_transactions: dict[str, DailyUserSpendTransaction], 

2356 entity_type: Literal["user"], 

2357 entity_id_field: str, 

2358 ) -> None: 

2359 ... 

2360 

2361 @overload 

2362 @staticmethod 

2363 async def _update_daily_spend( 

2364 n_retry_times: int, 

2365 prisma_client: PrismaClient, 

2366 proxy_logging_obj: ProxyLogging, 

2367 daily_spend_transactions: dict[str, DailyTeamSpendTransaction], 

2368 entity_type: Literal["team"], 

2369 entity_id_field: str, 

2370 ) -> None: 

2371 ... 

2372 

2373 @overload 

2374 @staticmethod 

2375 async def _update_daily_spend( 

2376 n_retry_times: int, 

2377 prisma_client: PrismaClient, 

2378 proxy_logging_obj: ProxyLogging, 

2379 daily_spend_transactions: dict[str, DailyOrganizationSpendTransaction], 

2380 entity_type: Literal["org"], 

2381 entity_id_field: str, 

2382 ) -> None: 

2383 ... 

2384 

2385 @overload 

2386 @staticmethod 

2387 async def _update_daily_spend( 

2388 n_retry_times: int, 

2389 prisma_client: PrismaClient, 

2390 proxy_logging_obj: ProxyLogging, 

2391 daily_spend_transactions: dict[str, DailyEndUserSpendTransaction], 

2392 entity_type: Literal["end_user"], 

2393 entity_id_field: str, 

2394 ) -> None: 

2395 ... 

2396 

2397 @overload 

2398 @staticmethod 

2399 async def _update_daily_spend( 

2400 n_retry_times: int, 

2401 prisma_client: PrismaClient, 

2402 proxy_logging_obj: ProxyLogging, 

2403 daily_spend_transactions: dict[str, DailyAgentSpendTransaction], 

2404 entity_type: Literal["agent"], 

2405 entity_id_field: str, 

2406 ) -> None: 

2407 ... 

2408 

2409 @overload 

2410 @staticmethod 

2411 async def _update_daily_spend( 

2412 n_retry_times: int, 

2413 prisma_client: PrismaClient, 

2414 proxy_logging_obj: ProxyLogging, 

2415 daily_spend_transactions: dict[str, DailyTagSpendTransaction], 

2416 entity_type: Literal["tag"], 

2417 entity_id_field: str, 

2418 ) -> None: 

2419 ... 

2420 # fmt: on 

2421 

2422 @staticmethod 

2423 async def _update_daily_spend( 

2424 n_retry_times: int, 

2425 prisma_client: PrismaClient, 

2426 proxy_logging_obj: ProxyLogging, 

2427 daily_spend_transactions: dict[str, DailyUserSpendTransaction] 

2428 | dict[str, DailyTeamSpendTransaction] 

2429 | dict[str, DailyTagSpendTransaction] 

2430 | dict[str, DailyOrganizationSpendTransaction] 

2431 | dict[str, DailyEndUserSpendTransaction] 

2432 | dict[str, DailyAgentSpendTransaction], 

2433 entity_type: Literal["user", "team", "org", "tag", "end_user", "agent"], 

2434 entity_id_field: str, 

2435 ) -> None: 

2436 """ 

2437 Generic function to update daily spend for any entity type (user, team, org, tag, end_user, agent) 

2438 """ 

2439 from litellm.proxy.utils import _raise_failed_update_spend_exception 

2440 

2441 verbose_proxy_logger.debug( 

2442 "Daily %s Spend transactions: %s", entity_type.capitalize(), len(daily_spend_transactions) 

2443 ) 

2444 BATCH_SIZE: Final = 100 

2445 start_time: Final = time.time() 

2446 

2447 try: 

2448 while daily_spend_transactions: 

2449 for i in range(n_retry_times + 1): 2449 ↛ 2448line 2449 didn't jump to line 2448 because the loop on line 2449 didn't complete

2450 try: 

2451 # Sort the transactions to minimize the probability of deadlocks by reducing the chance of concurrent 

2452 # trasactions locking the same rows/ranges in different orders. 

2453 transactions_to_process = dict( 

2454 sorted( 

2455 daily_spend_transactions.items(), 

2456 # Normally to avoid deadlocks we would sort by the index, but since we have sprinkled indexes 

2457 # on our schema like we're discount Salt Bae, we just sort by all fields that have an index, 

2458 # in an ad-hoc (but hopefully sensible) order of indexes. The actual ordering matters less than 

2459 # ensuring that all concurrent transactions sort in the same order. 

2460 # We could in theory use the dict key, as it contains basically the same fields, but this is more 

2461 # robust to future changes in the key format. 

2462 # If _update_daily_spend ever gets the ability to write to multiple tables at once, the sorting 

2463 # should sort by the table first. 

2464 key=lambda x: ( 

2465 x[1].get("date") or "", 

2466 x[1].get(entity_id_field) or "", 

2467 x[1].get("api_key") or "", 

2468 x[1].get("model") or "", 

2469 x[1].get("custom_llm_provider") or "", 

2470 ), 

2471 )[:BATCH_SIZE] 

2472 ) 

2473 

2474 if len(transactions_to_process) == 0: 2474 ↛ 2475line 2474 didn't jump to line 2475 because the condition on line 2474 was never true

2475 verbose_proxy_logger.debug( 

2476 "No new transactions to process for daily %s spend update", entity_type 

2477 ) 

2478 return 

2479 

2480 table = DAILY_SPEND_TABLES[entity_type] 

2481 try: 

2482 # One statement per batch rather than per key: the same rows are 

2483 # aggregated, but concurrent writers no longer hold a batch's worth 

2484 # of row locks across a hundred round trips. 

2485 merged_batch = merge_by_conflict_key( 

2486 table=table, transactions=tuple(transactions_to_process.values()) 

2487 ) 

2488 sql, params = build_bulk_upsert(table=table, batch=merged_batch) 

2489 async with _spend_update_tx(prisma_client) as transaction: 

2490 await transaction.execute_raw(sql, *params) 

2491 _mark_daily_spend_commit_started() 

2492 _mark_daily_spend_commit_finished() 

2493 except Exception as batch_error: 

2494 if _spend_commit_failure_is_requeue_safe(batch_error): 

2495 spend_log_error( 

2496 "Daily %s spend batch upsert failed. Table: %s, Rows: %d, Error: %s", 

2497 entity_type, 

2498 table.name, 

2499 len(transactions_to_process), 

2500 str(batch_error), 

2501 exc=batch_error, 

2502 ) 

2503 raise 

2504 for key in transactions_to_process: 

2505 daily_spend_transactions.pop(key, None) 

2506 spend_log_error( 

2507 "Spend tracking - dropped %d daily %s spend rows: the failed statement may have " 

2508 "applied or the database refused the data, so re-sending it is not safe. " 

2509 "Table: %s, Error: %s", 

2510 len(transactions_to_process), 

2511 entity_type, 

2512 table.name, 

2513 str(batch_error), 

2514 exc=batch_error, 

2515 ) 

2516 raise 

2517 

2518 verbose_proxy_logger.debug( 

2519 f"Processed {len(transactions_to_process)} daily {entity_type} transactions in {time.time() - start_time:.2f}s" 

2520 ) 

2521 

2522 # Remove processed transactions 

2523 for key in transactions_to_process: 

2524 daily_spend_transactions.pop(key, None) 

2525 

2526 break 

2527 

2528 except Exception as e: 

2529 from litellm.proxy.db.exception_handler import ( 

2530 PrismaDBExceptionHandler, 

2531 ) 

2532 

2533 is_retryable = isinstance( 

2534 e, DB_RETRY_SAFE_ERROR_TYPES 

2535 ) or PrismaDBExceptionHandler.is_deadlock_error(e) 

2536 if not is_retryable: 2536 ↛ 2538line 2536 didn't jump to line 2538 because the condition on line 2536 was always true

2537 raise 

2538 if i >= n_retry_times: 

2539 _raise_failed_update_spend_exception( 

2540 e=e, 

2541 start_time=start_time, 

2542 proxy_logging_obj=proxy_logging_obj, 

2543 ) 

2544 await asyncio.sleep( 

2545 # Sleep a random amount to avoid retrying and deadlocking again: when two transactions deadlock they are 

2546 # cancelled basically at the same time, so if they wait the same time they will also retry at the same time 

2547 # and thus they are more likely to deadlock again. 

2548 # Instead, we sleep a random amount so that they retry at slightly different times, lowering the chance of 

2549 # repeated deadlocks, and therefore of exceeding the retry limit. 

2550 random.uniform(2**i, 2 ** (i + 1)) 

2551 ) 

2552 

2553 except Exception as e: 

2554 _raise_failed_update_spend_exception(e=e, start_time=start_time, proxy_logging_obj=proxy_logging_obj) 

2555 

2556 @staticmethod 

2557 async def update_daily_user_spend( 

2558 n_retry_times: int, 

2559 prisma_client: PrismaClient, 

2560 proxy_logging_obj: ProxyLogging, 

2561 daily_spend_transactions: dict[str, DailyUserSpendTransaction], 

2562 ): 

2563 """ 

2564 Batch job to update LiteLLM_DailyUserSpend table using in-memory daily_spend_transactions 

2565 """ 

2566 await DBSpendUpdateWriter._update_daily_spend( 

2567 n_retry_times=n_retry_times, 

2568 prisma_client=prisma_client, 

2569 proxy_logging_obj=proxy_logging_obj, 

2570 daily_spend_transactions=daily_spend_transactions, 

2571 entity_type="user", 

2572 entity_id_field="user_id", 

2573 ) 

2574 

2575 @staticmethod 

2576 async def update_daily_team_spend( 

2577 n_retry_times: int, 

2578 prisma_client: PrismaClient, 

2579 proxy_logging_obj: ProxyLogging, 

2580 daily_spend_transactions: dict[str, DailyTeamSpendTransaction], 

2581 ): 

2582 """ 

2583 Batch job to update LiteLLM_DailyTeamSpend table using in-memory daily_spend_transactions 

2584 """ 

2585 await DBSpendUpdateWriter._update_daily_spend( 

2586 n_retry_times=n_retry_times, 

2587 prisma_client=prisma_client, 

2588 proxy_logging_obj=proxy_logging_obj, 

2589 daily_spend_transactions=daily_spend_transactions, 

2590 entity_type="team", 

2591 entity_id_field="team_id", 

2592 ) 

2593 

2594 @staticmethod 

2595 async def update_daily_org_spend( 

2596 n_retry_times: int, 

2597 prisma_client: PrismaClient, 

2598 proxy_logging_obj: ProxyLogging, 

2599 daily_spend_transactions: dict[str, DailyOrganizationSpendTransaction], 

2600 ): 

2601 """ 

2602 Batch job to update LiteLLM_DailyOrganizationSpend table using in-memory daily_spend_transactions 

2603 """ 

2604 await DBSpendUpdateWriter._update_daily_spend( 

2605 n_retry_times=n_retry_times, 

2606 prisma_client=prisma_client, 

2607 proxy_logging_obj=proxy_logging_obj, 

2608 daily_spend_transactions=daily_spend_transactions, 

2609 entity_type="org", 

2610 entity_id_field="organization_id", 

2611 ) 

2612 

2613 @staticmethod 

2614 async def update_daily_end_user_spend( 

2615 n_retry_times: int, 

2616 prisma_client: PrismaClient, 

2617 proxy_logging_obj: ProxyLogging, 

2618 daily_spend_transactions: dict[str, DailyEndUserSpendTransaction], 

2619 ): 

2620 """ 

2621 Batch job to update LiteLLM_DailyEndUserSpend table using in-memory daily_spend_transactions 

2622 """ 

2623 await DBSpendUpdateWriter._update_daily_spend( 

2624 n_retry_times=n_retry_times, 

2625 prisma_client=prisma_client, 

2626 proxy_logging_obj=proxy_logging_obj, 

2627 daily_spend_transactions=daily_spend_transactions, 

2628 entity_type="end_user", 

2629 entity_id_field="end_user_id", 

2630 ) 

2631 

2632 @staticmethod 

2633 async def update_daily_agent_spend( 

2634 n_retry_times: int, 

2635 prisma_client: PrismaClient, 

2636 proxy_logging_obj: ProxyLogging, 

2637 daily_spend_transactions: dict[str, DailyAgentSpendTransaction], 

2638 ): 

2639 """ 

2640 Batch job to update LiteLLM_DailyAgentSpend table using in-memory daily_spend_transactions 

2641 """ 

2642 await DBSpendUpdateWriter._update_daily_spend( 

2643 n_retry_times=n_retry_times, 

2644 prisma_client=prisma_client, 

2645 proxy_logging_obj=proxy_logging_obj, 

2646 daily_spend_transactions=daily_spend_transactions, 

2647 entity_type="agent", 

2648 entity_id_field="agent_id", 

2649 ) 

2650 

2651 @staticmethod 

2652 async def update_daily_tag_spend( 

2653 n_retry_times: int, 

2654 prisma_client: PrismaClient, 

2655 proxy_logging_obj: ProxyLogging, 

2656 daily_spend_transactions: dict[str, DailyTagSpendTransaction], 

2657 ): 

2658 """ 

2659 Batch job to update LiteLLM_DailyTagSpend table using in-memory daily_spend_transactions 

2660 """ 

2661 await DBSpendUpdateWriter._update_daily_spend( 

2662 n_retry_times=n_retry_times, 

2663 prisma_client=prisma_client, 

2664 proxy_logging_obj=proxy_logging_obj, 

2665 daily_spend_transactions=daily_spend_transactions, 

2666 entity_type="tag", 

2667 entity_id_field="tag", 

2668 ) 

2669 

2670 async def _common_add_spend_log_transaction_to_daily_transaction( 

2671 self, 

2672 payload: dict | SpendLogsPayload, 

2673 prisma_client: PrismaClient, 

2674 type: Literal["user", "team", "org", "request_tags", "end_user", "agent"] = "user", 

2675 ) -> BaseDailySpendTransaction | None: 

2676 entity: Final = "tag" if type == "request_tags" else type 

2677 identity_payload: Final = ( 

2678 MappingProxyType({**payload, "end_user": payload.get("end_user_id")}) if type == "end_user" else payload 

2679 ) 

2680 if not daily_spend_entity_ids(identity_payload, entity): 

2681 return None 

2682 expected_keys: Final = ("startTime", "api_key") 

2683 if not all(key in payload for key in expected_keys): 2683 ↛ 2684line 2683 didn't jump to line 2684 because the condition on line 2683 was never true

2684 verbose_proxy_logger.debug( 

2685 "Missing expected keys: %s, in payload, skipping from daily_user_spend_transactions", expected_keys 

2686 ) 

2687 return None 

2688 

2689 any_expected_keys: Final = ["model", "mcp_namespaced_tool_name"] 

2690 if not any(key in payload for key in any_expected_keys): 2690 ↛ 2691line 2690 didn't jump to line 2691 because the condition on line 2690 was never true

2691 verbose_proxy_logger.debug( 

2692 "Missing any expected keys: %s, in payload, skipping from daily_user_spend_transactions", 

2693 any_expected_keys, 

2694 ) 

2695 return None 

2696 elif "mcp_namespaced_tool_name" in payload: 2696 ↛ 2698line 2696 didn't jump to line 2698 because the condition on line 2696 was always true

2697 pass 

2698 elif "model" in payload and ("custom_llm_provider" not in payload or "model_group" not in payload): 

2699 verbose_proxy_logger.debug( 

2700 "Missing custom_llm_provider or model_group in payload, skipping from daily_user_spend_transactions" 

2701 ) 

2702 return None 

2703 

2704 # TODO: remove the successful_requests/failed_requests counters below once the 

2705 # admin UI has fully migrated to LiteLLM_DailyGatewayRequests, which is now the 

2706 # source of truth for SGR. This path derives the counts from spend-log metadata 

2707 # rather than from what the gateway answered, so the two intentionally disagree 

2708 # (see litellm/proxy/middleware/billable_request_metrics_middleware.py). The 

2709 # spend, token and per-entity columns written here stay either way. 

2710 request_status: Final = prisma_client.get_request_status(payload) 

2711 verbose_proxy_logger.debug("Logged request status: %s", request_status) 

2712 _metadata: Final[SpendLogsMetadata] = json.loads(payload["metadata"]) 

2713 usage_obj: Final = _metadata.get("usage_object", {}) or {} 

2714 if isinstance(payload["startTime"], datetime): 2714 ↛ 2715line 2714 didn't jump to line 2715 because the condition on line 2714 was never true

2715 start_time: Final = payload["startTime"].isoformat() 

2716 date = start_time.split("T")[0] 

2717 elif isinstance(payload["startTime"], str): 2717 ↛ 2720line 2717 didn't jump to line 2720 because the condition on line 2717 was always true

2718 date = payload["startTime"].split("T")[0] 

2719 else: 

2720 verbose_proxy_logger.debug( 

2721 "Invalid start time: %s, skipping from daily_user_spend_transactions", payload["startTime"] 

2722 ) 

2723 return None 

2724 try: 

2725 # Map call_type to endpoint using ROUTE_ENDPOINT_MAPPING 

2726 call_type: Final = payload.get("call_type", None) 

2727 endpoint = None 

2728 if call_type: 

2729 endpoint = ROUTE_ENDPOINT_MAPPING.get(call_type, None) 

2730 

2731 is_internal_call: Final = bool(_metadata.get(INTERNAL_CALL_ORIGIN_METADATA_KEY)) 

2732 cache_read_input_tokens: Final = extract_cache_read_tokens(usage_obj) 

2733 compression_saved_tokens: Final = extract_compression_saved_tokens(_metadata) 

2734 savings_spend: Final = compute_savings_spend( 

2735 model=payload.get("model", None), 

2736 custom_llm_provider=payload.get("custom_llm_provider", None), 

2737 compression_saved_tokens=compression_saved_tokens, 

2738 gateway_injected_cache=marks_gateway_injection(_metadata, payload.get("model_id")), 

2739 routing_decision=_metadata.get("routing_decision"), 

2740 model_id=payload.get("model_id"), 

2741 llm_router=get_llm_router, 

2742 usage_object=usage_obj, 

2743 cost_breakdown=_metadata.get("cost_breakdown"), 

2744 recorded_autorouter_savings=_metadata.get("autorouter_savings"), 

2745 recorded_autorouter_savings_estimate=_metadata.get("autorouter_savings_estimate"), 

2746 billed_at=payload.get("endTime"), 

2747 ) 

2748 timed_duration_ms: Final = _timed_request_duration_ms(payload, request_status, is_internal_call) 

2749 

2750 daily_transaction: Final = BaseDailySpendTransaction( 

2751 date=date, 

2752 api_key=payload["api_key"], 

2753 model=payload.get("model", None), 

2754 model_group=payload.get("model_group", None), 

2755 mcp_namespaced_tool_name=payload.get("mcp_namespaced_tool_name", None), 

2756 custom_llm_provider=payload.get("custom_llm_provider", None), 

2757 endpoint=endpoint, 

2758 prompt_tokens=payload["prompt_tokens"], 

2759 completion_tokens=payload["completion_tokens"], 

2760 spend=payload["spend"], 

2761 # Internal sub-calls (auto-router classifier, shadow eval's shadow and 

2762 # judge) bill real spend and tokens to the key, but they are not 

2763 # requests the caller made: counting them inflates request-volume 

2764 # readers, and an auto-router savings figure computed on a shadow 

2765 # duplicate credits savings for traffic no user sent. 

2766 api_requests=0 if is_internal_call else 1, 

2767 successful_requests=1 if not is_internal_call and request_status == "success" else 0, 

2768 failed_requests=1 if not is_internal_call and request_status != "success" else 0, 

2769 cache_read_input_tokens=cache_read_input_tokens, 

2770 cache_creation_input_tokens=extract_cache_creation_tokens(usage_obj), 

2771 compression_saved_tokens=compression_saved_tokens, 

2772 compression_savings_spend=savings_spend.compression, 

2773 prompt_caching_savings_spend=savings_spend.prompt_caching, 

2774 gateway_injected_caching_savings_spend=savings_spend.gateway_injected_caching, 

2775 autorouter_savings_spend=0.0 if is_internal_call else savings_spend.autorouter, 

2776 total_response_time_ms=timed_duration_ms or 0, 

2777 timed_requests=0 if timed_duration_ms is None else 1, 

2778 ) 

2779 return daily_transaction 

2780 except Exception as e: 

2781 raise e 

2782 

2783 async def add_spend_log_transaction_to_daily_user_transaction( 

2784 self, 

2785 payload: dict | SpendLogsPayload, 

2786 prisma_client: PrismaClient | None = None, 

2787 ): 

2788 """ 

2789 Add a spend log transaction to the `daily_spend_update_queue` 

2790 

2791 Key = @@unique([user_id, date, api_key, model, custom_llm_provider]) ) 

2792 

2793 If key exists, update the transaction with the new spend and usage 

2794 """ 

2795 if prisma_client is None: 2795 ↛ 2796line 2795 didn't jump to line 2796 because the condition on line 2795 was never true

2796 verbose_proxy_logger.debug("prisma_client is None. Skipping writing spend logs to db.") 

2797 return 

2798 

2799 base_daily_transaction: Final = await self._common_add_spend_log_transaction_to_daily_transaction( 

2800 payload, prisma_client, "user" 

2801 ) 

2802 if base_daily_transaction is None: 2802 ↛ 2803line 2802 didn't jump to line 2803 because the condition on line 2802 was never true

2803 return 

2804 

2805 endpoint_str: Final = base_daily_transaction.get("endpoint") or "" 

2806 daily_transaction_key = f"{payload['user']}_{base_daily_transaction['date']}_{payload['api_key']}_{payload['model']}_{payload['custom_llm_provider']}_{endpoint_str}" 

2807 daily_transaction: Final = DailyUserSpendTransaction(user_id=payload["user"], **base_daily_transaction) 

2808 await self.daily_spend_update_queue.add_update(update={daily_transaction_key: daily_transaction}) 

2809 

2810 async def add_spend_log_transaction_to_daily_team_transaction( 

2811 self, 

2812 payload: SpendLogsPayload, 

2813 prisma_client: PrismaClient | None = None, 

2814 ) -> None: 

2815 if prisma_client is None: 2815 ↛ 2816line 2815 didn't jump to line 2816 because the condition on line 2815 was never true

2816 verbose_proxy_logger.debug("prisma_client is None. Skipping writing spend logs to db.") 

2817 return 

2818 

2819 base_daily_transaction: Final = await self._common_add_spend_log_transaction_to_daily_transaction( 

2820 payload, prisma_client, "team" 

2821 ) 

2822 if base_daily_transaction is None: 2822 ↛ 2823line 2822 didn't jump to line 2823 because the condition on line 2822 was never true

2823 return 

2824 if payload["team_id"] is None: 2824 ↛ 2825line 2824 didn't jump to line 2825 because the condition on line 2824 was never true

2825 verbose_proxy_logger.debug("team_id is None for request. Skipping incrementing team spend.") 

2826 return 

2827 

2828 endpoint_str: Final = base_daily_transaction.get("endpoint") or "" 

2829 daily_transaction_key = f"{payload['team_id']}_{base_daily_transaction['date']}_{payload['api_key']}_{payload['model']}_{payload['custom_llm_provider']}_{endpoint_str}" 

2830 daily_transaction: Final = DailyTeamSpendTransaction(team_id=payload["team_id"], **base_daily_transaction) 

2831 await self.daily_team_spend_update_queue.add_update(update={daily_transaction_key: daily_transaction}) 

2832 

2833 async def add_spend_log_transaction_to_daily_org_transaction( 

2834 self, 

2835 payload: SpendLogsPayload, 

2836 prisma_client: PrismaClient | None = None, 

2837 org_id: str | None = None, 

2838 ) -> None: 

2839 if prisma_client is None: 2839 ↛ 2840line 2839 didn't jump to line 2840 because the condition on line 2839 was never true

2840 verbose_proxy_logger.debug("prisma_client is None. Skipping writing spend logs to db.") 

2841 return 

2842 

2843 if org_id is None: 2843 ↛ 2847line 2843 didn't jump to line 2847 because the condition on line 2843 was always true

2844 verbose_proxy_logger.debug("organization_id is None for request. Skipping incrementing organization spend.") 

2845 return 

2846 

2847 payload_with_org: Final = cast( 

2848 SpendLogsPayload, 

2849 { 

2850 **payload, 

2851 "organization_id": org_id, 

2852 }, 

2853 ) 

2854 

2855 base_daily_transaction: Final = await self._common_add_spend_log_transaction_to_daily_transaction( 

2856 payload_with_org, prisma_client, "org" 

2857 ) 

2858 if base_daily_transaction is None: 

2859 return 

2860 

2861 endpoint_str: Final = base_daily_transaction.get("endpoint") or "" 

2862 daily_transaction_key = f"{org_id}_{base_daily_transaction['date']}_{payload_with_org['api_key']}_{payload_with_org['model']}_{payload_with_org['custom_llm_provider']}_{endpoint_str}" 

2863 daily_transaction: Final = DailyOrganizationSpendTransaction(organization_id=org_id, **base_daily_transaction) 

2864 await self.daily_org_spend_update_queue.add_update(update={daily_transaction_key: daily_transaction}) 

2865 

2866 async def add_spend_log_transaction_to_daily_end_user_transaction( 

2867 self, 

2868 payload: SpendLogsPayload, 

2869 prisma_client: PrismaClient | None = None, 

2870 ) -> None: 

2871 if prisma_client is None: 2871 ↛ 2872line 2871 didn't jump to line 2872 because the condition on line 2871 was never true

2872 verbose_proxy_logger.debug("prisma_client is None. Skipping writing spend logs to db.") 

2873 return 

2874 

2875 end_user_id: Final = payload.get("end_user") 

2876 if end_user_id is None or end_user_id == "": 

2877 verbose_proxy_logger.debug("end_user is None or empty for request. Skipping incrementing end user spend.") 

2878 return 

2879 

2880 payload_with_end_user_id: Final = cast( 

2881 SpendLogsPayload, 

2882 { 

2883 **payload, 

2884 "end_user_id": end_user_id, 

2885 }, 

2886 ) 

2887 

2888 base_daily_transaction: Final = await self._common_add_spend_log_transaction_to_daily_transaction( 

2889 payload_with_end_user_id, prisma_client, "end_user" 

2890 ) 

2891 if base_daily_transaction is None: 2891 ↛ 2892line 2891 didn't jump to line 2892 because the condition on line 2891 was never true

2892 return 

2893 

2894 endpoint_str: Final = base_daily_transaction.get("endpoint") or "" 

2895 daily_transaction_key = f"{end_user_id}_{base_daily_transaction['date']}_{payload_with_end_user_id['api_key']}_{payload_with_end_user_id['model']}_{payload_with_end_user_id['custom_llm_provider']}_{endpoint_str}" 

2896 daily_transaction: Final = DailyEndUserSpendTransaction(end_user_id=end_user_id, **base_daily_transaction) 

2897 await self.daily_end_user_spend_update_queue.add_update(update={daily_transaction_key: daily_transaction}) 

2898 

2899 async def add_spend_log_transaction_to_daily_agent_transaction( 

2900 self, 

2901 payload: SpendLogsPayload, 

2902 prisma_client: PrismaClient | None = None, 

2903 ) -> None: 

2904 if prisma_client is None: 2904 ↛ 2905line 2904 didn't jump to line 2905 because the condition on line 2904 was never true

2905 verbose_proxy_logger.debug("prisma_client is None. Skipping writing spend logs to db.") 

2906 return 

2907 if payload["agent_id"] is None: 

2908 return 

2909 payload_with_agent_id: Final = cast( 

2910 SpendLogsPayload, 

2911 { 

2912 **payload, 

2913 "agent_id": payload["agent_id"], 

2914 }, 

2915 ) 

2916 base_daily_transaction: Final = await self._common_add_spend_log_transaction_to_daily_transaction( 

2917 payload_with_agent_id, prisma_client, "agent" 

2918 ) 

2919 if base_daily_transaction is None: 2919 ↛ 2920line 2919 didn't jump to line 2920 because the condition on line 2919 was never true

2920 return 

2921 endpoint_str: Final = base_daily_transaction.get("endpoint") or "" 

2922 daily_transaction_key = f"{payload['agent_id']}_{base_daily_transaction['date']}_{payload_with_agent_id['api_key']}_{payload_with_agent_id['model']}_{payload_with_agent_id['custom_llm_provider']}_{endpoint_str}" 

2923 daily_transaction: Final = DailyAgentSpendTransaction(agent_id=payload["agent_id"], **base_daily_transaction) 

2924 await self.daily_agent_spend_update_queue.add_update(update={daily_transaction_key: daily_transaction}) 

2925 

2926 async def add_spend_log_transaction_to_daily_tag_transaction( 

2927 self, 

2928 payload: SpendLogsPayload, 

2929 prisma_client: PrismaClient | None = None, 

2930 ) -> None: 

2931 if prisma_client is None: 2931 ↛ 2932line 2931 didn't jump to line 2932 because the condition on line 2931 was never true

2932 verbose_proxy_logger.debug("prisma_client is None. Skipping writing spend logs to db.") 

2933 return 

2934 

2935 base_daily_transaction: Final = await self._common_add_spend_log_transaction_to_daily_transaction( 

2936 payload, prisma_client, "request_tags" 

2937 ) 

2938 if base_daily_transaction is None: 

2939 return 

2940 if payload["request_tags"] is None: 2940 ↛ 2941line 2940 didn't jump to line 2941 because the condition on line 2940 was never true

2941 verbose_proxy_logger.debug("request_tags is None for request. Skipping incrementing tag spend.") 

2942 return 

2943 

2944 request_tags: Final = daily_spend_entity_ids(payload, "tag") 

2945 for tag in request_tags: 

2946 if tag is None: 2946 ↛ 2947line 2946 didn't jump to line 2947 because the condition on line 2946 was never true

2947 continue 

2948 endpoint_str = base_daily_transaction.get("endpoint") or "" 

2949 daily_transaction_key = f"{tag}_{base_daily_transaction['date']}_{payload['api_key']}_{payload['model']}_{payload['custom_llm_provider']}_{endpoint_str}" 

2950 daily_transaction = DailyTagSpendTransaction( 

2951 tag=tag, **base_daily_transaction, request_id=payload["request_id"] 

2952 ) 

2953 

2954 await self.daily_tag_spend_update_queue.add_update(update={daily_transaction_key: daily_transaction})