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
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-10 12:01 +0000
1"""
2Module responsible for
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"""
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
22from pydantic import TypeAdapter
23from typing_extensions import LiteralString, ReadOnly, TypedDict
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
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
96RESPONSES_SESSION_CALL_TYPES: Final = frozenset({CallTypes.responses.value, CallTypes.aresponses.value})
97_SPEND_METADATA_ADAPTER: Final = TypeAdapter(Mapping[str, object])
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='')}"
104def _is_batch_cost_row(payload: SpendLogsPayload) -> bool:
105 return payload.get("call_type") == CallTypes.aretrieve_batch.value and payload.get("status") == "success"
108_BATCH_COST_CLAIM_FIELDS: Final = frozenset({"request_id", "call_type", "spend", "startTime", "endTime", "status"})
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.
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})
123class _SpendIncrement(TypedDict):
124 increment: ReadOnly[float]
127class _MemberSpendRow(TypedDict):
128 user_id: ReadOnly[str]
129 team_id: ReadOnly[str]
130 cost: ReadOnly[float]
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
146_EntitySpendTable: TypeAlias = Literal[
147 "litellm_tagtable", "litellm_agentstable", "litellm_modelaccessgroupbudgettable", "litellm_projecttable"
148]
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)
161def _entity_spend_table(batcher: _SpendBatch, table_accessor: _EntitySpendTable) -> BatchTable:
162 return _ENTITY_SPEND_TABLES[table_accessor](batcher)
165class _SpendBatchManager(Protocol):
166 async def __aenter__(self) -> _SpendBatch: ... 166 ↛ exitline 166 didn't return from function '__aenter__' because
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
171class _SpendTransaction(Protocol):
172 def batch_(self) -> _SpendBatchManager: ... 172 ↛ exitline 172 didn't return from function 'batch_' because
174 async def execute_raw(self, query: LiteralString, *args: object) -> int: ... 174 ↛ exitline 174 didn't return from function 'execute_raw' because
177class _SpendTransactionManager(Protocol):
178 async def __aenter__(self) -> _SpendTransaction: ... 178 ↛ exitline 178 didn't return from function '__aenter__' because
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
183_DailySpendTransactionT = TypeVar("_DailySpendTransactionT", bound=BaseDailySpendTransaction)
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: ...
197_DATA_REJECTED_SQLSTATE_CLASSES: Final = frozenset({"22", "23"})
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
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)
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
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
274def _spend_update_tx(prisma_client: PrismaClient) -> _SpendTransactionManager:
275 tx: Final[_SpendTransactionManager] = prisma_client.db.tx(timeout=timedelta(seconds=60))
276 return tx
279_daily_spend_commit_started: Final[ContextVar[asyncio.Event | None]] = ContextVar(
280 "_daily_spend_commit_started", default=None
281)
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()
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()
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)
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)
312async def _settle_interrupted_commits(commits: set[asyncio.Task[None]]) -> None:
313 while commits:
314 await asyncio.wait(tuple(commits))
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 )
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)
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"
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"""
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))
393def get_llm_router():
394 """The proxy's router, or None outside a running proxy.
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
402 return llm_router
403 except Exception: # noqa: BLE001 # no proxy in scope; savings degrade to zero
404 return None
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
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.
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))
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.
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)
453class DBSpendUpdateWriter:
454 """
455 Module responsible for
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 """
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
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.
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
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
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 )
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()
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
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
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
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 )
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 )
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 )
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
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
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.
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
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)
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.
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
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
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 )
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)
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 )
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)
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
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
834 request_spend_log_flush(prisma_client)
835 return True
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
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 )
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.
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
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
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 )
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")
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)
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)
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
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)
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.
980 Each helper is wrapped in try/except so one failure doesn't prevent the others.
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 )
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 )
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 )
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 )
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 )
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 )
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 )
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 )
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 )
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 )
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 )
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 )
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 )
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 )
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
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
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)
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 )
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 )
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
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 )
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
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
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 )
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
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
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
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
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.
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
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
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
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.
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
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 )
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
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.")
1480 return prisma_client
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
1491 `UPDATES` can lead to deadlocks, hence we handle them separately
1493 Args:
1494 prisma_client: PrismaClient object
1495 n_retry_times: int, number of retry times
1496 proxy_logging_obj: ProxyLogging object
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 )
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 )
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
1528 This is a v2 scalable approach to first commit spend updates to redis, then commit to db
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 )
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")
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
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()
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 }
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)
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)
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)
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)
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)
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 )
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)
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
1739 This is the regular flow of committing to db without using a redis buffer
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 """
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 )
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 )
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 )
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 )
1789 # NOTE: Daily tag spend is committed by a separate scheduler job.
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 )
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 )
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 )
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 )
1843 ################## Tool Registry Upserts ##################
1844 await self._flush_tool_discovery_queue(prisma_client=prisma_client)
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 )
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.
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 )
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 )
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.
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 )
1920 await commit_window_spend_updates(
1921 prisma_client=prisma_client,
1922 transactions=window_spend_transactions,
1923 )
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.
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
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
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
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)
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
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)))
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
2021 """
2022 from litellm.proxy.utils import ProxyUpdateSpend
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")
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")
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")
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))
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")
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")
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")
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")
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 )
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")
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")
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")
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))
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).
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 )
2347 # fmt: off
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 ...
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 ...
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 ...
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 ...
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 ...
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
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
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()
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 )
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
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
2518 verbose_proxy_logger.debug(
2519 f"Processed {len(transactions_to_process)} daily {entity_type} transactions in {time.time() - start_time:.2f}s"
2520 )
2522 # Remove processed transactions
2523 for key in transactions_to_process:
2524 daily_spend_transactions.pop(key, None)
2526 break
2528 except Exception as e:
2529 from litellm.proxy.db.exception_handler import (
2530 PrismaDBExceptionHandler,
2531 )
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 )
2553 except Exception as e:
2554 _raise_failed_update_spend_exception(e=e, start_time=start_time, proxy_logging_obj=proxy_logging_obj)
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 )
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 )
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 )
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 )
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 )
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 )
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
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
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)
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)
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
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`
2791 Key = @@unique([user_id, date, api_key, model, custom_llm_provider]) )
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
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
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})
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
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
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})
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
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
2847 payload_with_org: Final = cast(
2848 SpendLogsPayload,
2849 {
2850 **payload,
2851 "organization_id": org_id,
2852 },
2853 )
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
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})
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
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
2880 payload_with_end_user_id: Final = cast(
2881 SpendLogsPayload,
2882 {
2883 **payload,
2884 "end_user_id": end_user_id,
2885 },
2886 )
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
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})
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})
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
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
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 )
2954 await self.daily_tag_spend_update_queue.add_update(update={daily_transaction_key: daily_transaction})