From bdfd0e46915295c19e215443038eb79756ff0ca6 Mon Sep 17 00:00:00 2001 From: todayim <809634488@qq.com> Date: Wed, 5 Aug 2026 16:50:25 +0800 Subject: [PATCH 01/12] feat(proxy): add POST /spend/usage to ingest externally measured usage into spend tracking --- litellm/proxy/_types.py | 1 + litellm/proxy/proxy_server.py | 4 + .../usage_ingestion_endpoints.py | 204 ++++++++++++++++++ .../test_usage_ingestion_endpoints.py | 183 ++++++++++++++++ 4 files changed, 392 insertions(+) create mode 100644 litellm/proxy/spend_tracking/usage_ingestion_endpoints.py create mode 100644 tests/test_litellm/proxy/spend_tracking/test_usage_ingestion_endpoints.py diff --git a/litellm/proxy/_types.py b/litellm/proxy/_types.py index 7d6829aca70..dfe6bc672bf 100644 --- a/litellm/proxy/_types.py +++ b/litellm/proxy/_types.py @@ -632,6 +632,7 @@ class LiteLLMRoutes(enum.Enum): "/jwt/key/mapping/delete", "/jwt/key/mapping/list", "/jwt/key/mapping/info", + "/spend/usage", ] + key_management_routes + mcp_management_routes diff --git a/litellm/proxy/proxy_server.py b/litellm/proxy/proxy_server.py index e343d46f872..56169bed74c 100644 --- a/litellm/proxy/proxy_server.py +++ b/litellm/proxy/proxy_server.py @@ -549,6 +549,9 @@ from litellm.proxy.spend_tracking.spend_management_endpoints import ( router as spend_management_router, ) from litellm.proxy.spend_tracking.spend_tracking_utils import get_logging_payload +from litellm.proxy.spend_tracking.usage_ingestion_endpoints import ( + router as usage_ingestion_router, +) from litellm.proxy.types_utils.utils import get_instance_fn from litellm.proxy.ui_crud_endpoints.proxy_setting_endpoints import ( router as ui_crud_endpoints_router, @@ -16449,6 +16452,7 @@ app.include_router(organization_router) app.include_router(customer_router) app.include_router(management_v1_router) app.include_router(spend_management_router) +app.include_router(usage_ingestion_router) app.include_router(caching_router) app.include_router(analytics_router) app.include_router(callback_management_endpoints_router) diff --git a/litellm/proxy/spend_tracking/usage_ingestion_endpoints.py b/litellm/proxy/spend_tracking/usage_ingestion_endpoints.py new file mode 100644 index 00000000000..8e3e3dd448d --- /dev/null +++ b/litellm/proxy/spend_tracking/usage_ingestion_endpoints.py @@ -0,0 +1,204 @@ +import uuid +from collections.abc import Awaitable, Callable +from dataclasses import dataclass +from datetime import datetime +from typing import Final, Literal, NamedTuple + +from fastapi import APIRouter, Depends, status +from pydantic import BaseModel, Field, model_validator + +import litellm +from litellm._logging import verbose_proxy_logger +from litellm.proxy._types import ProxyException +from litellm.proxy.auth.user_api_key_auth import user_api_key_auth +from litellm.proxy.utils import hash_token + +router: Final = APIRouter() + +MAX_RECORDS_PER_REQUEST: Final = 1000 + + +class ExternalUsageRecord(BaseModel): + api_key: str = Field(min_length=1, description="Raw virtual key (sk-...) to attribute usage to. Never logged.") + model: str = Field(min_length=1) + prompt_tokens: int = Field(ge=0) + completion_tokens: int = Field(ge=0) + start_time: datetime + end_time: datetime | None = None + cost: float | None = Field(default=None, ge=0, description="Explicit cost in USD. Computed from litellm pricing when omitted.") + idempotency_key: str | None = Field(default=None, max_length=255, description="Becomes the spend-log request_id for dedup on retries.") + tags: list[str] | None = None + end_user_id: str | None = None + + @model_validator(mode="after") + def end_time_not_before_start_time(self) -> "ExternalUsageRecord": + if self.end_time is not None and self.end_time < self.start_time: + raise ValueError("end_time must not be before start_time") + return self + + +class UsageIngestRequest(BaseModel): + records: list[ExternalUsageRecord] = Field(min_length=1, max_length=MAX_RECORDS_PER_REQUEST) + + +class UsageIngestRecordResult(BaseModel): + request_id: str + status: Literal["recorded", "duplicate", "error"] + spend: float | None = None + error: str | None = None + + +class UsageIngestResponse(BaseModel): + results: list[UsageIngestRecordResult] + + +class KeyAttribution(NamedTuple): + user_id: str | None + team_id: str | None + organization_id: str | None + + +RecordSpendFn = Callable[..., Awaitable[None]] + + +@dataclass(frozen=True, slots=True) +class UsageIngestionDeps: + lookup_key: Callable[[str], Awaitable[KeyAttribution | None]] + spend_log_exists: Callable[[str], Awaitable[bool]] + record_spend: RecordSpendFn + compute_cost: Callable[[litellm.ModelResponse, str], float] + generate_request_id: Callable[[], str] + + +def _build_usage_kwargs(record: ExternalUsageRecord, hashed_token: str) -> dict[str, object]: + metadata = { + "user_api_key": hashed_token, + "user_api_key_end_user_id": record.end_user_id, + "tags": list(record.tags) if record.tags else [], + } + return { + "model": record.model, + "call_type": "ingest_external_usage", + "litellm_params": {"model": record.model, "metadata": metadata}, + } + + +def _build_completion_response(record: ExternalUsageRecord, request_id: str) -> litellm.ModelResponse: + total_tokens: Final = record.prompt_tokens + record.completion_tokens + usage: Final = litellm.Usage( + prompt_tokens=record.prompt_tokens, + completion_tokens=record.completion_tokens, + total_tokens=total_tokens, + ) + return litellm.ModelResponse( + id=request_id, + model=record.model, + created=int(record.start_time.timestamp()), + usage=usage, + ) + + +def _resolve_cost(deps: UsageIngestionDeps, record: ExternalUsageRecord, response: litellm.ModelResponse) -> float: + if record.cost is not None: + return record.cost + return deps.compute_cost(response, record.model) + + +async def process_external_usage_record(record: ExternalUsageRecord, deps: UsageIngestionDeps) -> UsageIngestRecordResult: + request_id: Final = record.idempotency_key or deps.generate_request_id() + hashed_token: Final = hash_token(record.api_key) + + key: Final = await deps.lookup_key(hashed_token) + if key is None: + return UsageIngestRecordResult(request_id=request_id, status="error", error="api key not found") + + if record.idempotency_key is not None and await deps.spend_log_exists(request_id): + return UsageIngestRecordResult(request_id=request_id, status="duplicate") + + response: Final = _build_completion_response(record, request_id) + + try: + cost: Final = _resolve_cost(deps, record, response) + except Exception as e: # noqa: BLE001 # pricing lookup raises arbitrary provider-specific errors; any failure means the record is unpriceable and must carry an explicit cost + verbose_proxy_logger.info("ingest usage: cost computation failed for model %s: %s", record.model, e) + return UsageIngestRecordResult( + request_id=request_id, + status="error", + error="could not compute cost for this model, pass an explicit cost", + ) + + await deps.record_spend( + token=hashed_token, + user_id=key.user_id, + end_user_id=record.end_user_id, + team_id=key.team_id, + kwargs=_build_usage_kwargs(record, hashed_token), + completion_response=response, + start_time=record.start_time, + end_time=record.end_time or record.start_time, + response_cost=cost, + org_id=key.organization_id, + ) + return UsageIngestRecordResult(request_id=request_id, status="recorded", spend=cost) + + +def _attribution_of(key_row: object) -> KeyAttribution: + return KeyAttribution( + user_id=getattr(key_row, "user_id", None), + team_id=getattr(key_row, "team_id", None), + organization_id=getattr(key_row, "organization_id", None), + ) + + +def default_ingestion_deps() -> UsageIngestionDeps: + from litellm.proxy.proxy_server import prisma_client, proxy_logging_obj + + if prisma_client is None: + raise ProxyException( + message="Prisma Client is not initialized", + type="internal_error", + param="None", + code=status.HTTP_500_INTERNAL_SERVER_ERROR, + ) + + async def lookup_key(hashed_token: str) -> KeyAttribution | None: + row: Final = await prisma_client.db.litellm_verificationtoken.find_unique(where={"token": hashed_token}) + if row is None: + return None + return _attribution_of(row) + + async def spend_log_exists(request_id: str) -> bool: + row: Final = await prisma_client.db.litellm_spendlogs.find_unique(where={"request_id": request_id}) + return row is not None + + return UsageIngestionDeps( + lookup_key=lookup_key, + spend_log_exists=spend_log_exists, + record_spend=proxy_logging_obj.db_spend_update_writer.update_database, + compute_cost=lambda resp, model: litellm.completion_cost(completion_response=resp, model=model), + generate_request_id=lambda: str(uuid.uuid4()), + ) + + +@router.post( + "/spend/usage", + tags=["Budget & Spend Tracking"], + dependencies=[Depends(user_api_key_auth)], + response_model=UsageIngestResponse, +) +async def ingest_external_usage(request: UsageIngestRequest) -> UsageIngestResponse: + """ + PROXY_ADMIN ONLY: record externally measured usage into the same spend pipeline as proxy-routed traffic. + + For inference traffic that legitimately bypasses the proxy (for example async batch processors + dispatching directly to model gateways), so budgets and spend stay coherent in litellm as the + single metering system. + + Attribution (user/team/org) is derived from the given virtual key. Records accept an optional + idempotency_key, stored as the spend-log request_id, so retries are deduplicated. When cost is + omitted it is computed from litellm pricing; records whose model cannot be priced are rejected + with an error instead of being booked as zero spend. + """ + deps: Final = default_ingestion_deps() + results: Final = [await process_external_usage_record(record, deps) for record in request.records] + return UsageIngestResponse(results=results) diff --git a/tests/test_litellm/proxy/spend_tracking/test_usage_ingestion_endpoints.py b/tests/test_litellm/proxy/spend_tracking/test_usage_ingestion_endpoints.py new file mode 100644 index 00000000000..360f2e1af6e --- /dev/null +++ b/tests/test_litellm/proxy/spend_tracking/test_usage_ingestion_endpoints.py @@ -0,0 +1,183 @@ +import asyncio +import os +import sys +from datetime import datetime, timezone +from typing import Any + +import pytest +from pydantic import ValidationError + +sys.path.insert(0, os.path.abspath("../../..")) + +from litellm.proxy.spend_tracking.usage_ingestion_endpoints import ( + ExternalUsageRecord, + KeyAttribution, + UsageIngestionDeps, + process_external_usage_record, +) +from litellm.proxy.utils import hash_token + +RAW_KEY = "sk-test-batch-dispatch-key" +GENERATED_ID = "generated-uuid-1" +DEFAULT_KEY = KeyAttribution(user_id="u-1", team_id="t-1", organization_id="o-1") + + +class RecordingDeps: + def __init__( + self, + key: KeyAttribution | None = DEFAULT_KEY, + existing_ids: frozenset[str] = frozenset(), + compute_cost_result: float = 0.05, + compute_cost_error: Exception | None = None, + ): + self._key = key + self._existing_ids = existing_ids + self._compute_cost_result = compute_cost_result + self._compute_cost_error = compute_cost_error + self.spend_calls: list[dict[str, Any]] = [] + self.compute_cost_calls: list[tuple[Any, str]] = [] + self.exists_calls: list[str] = [] + + def as_deps(self) -> UsageIngestionDeps: + async def lookup_key(hashed: str) -> KeyAttribution | None: + self.looked_up_hashed = hashed + return self._key + + async def spend_log_exists(request_id: str) -> bool: + self.exists_calls.append(request_id) + return request_id in self._existing_ids + + async def record_spend(**kwargs: Any) -> None: + self.spend_calls.append(kwargs) + + def compute_cost(response: Any, model: str) -> float: + self.compute_cost_calls.append((response, model)) + if self._compute_cost_error is not None: + raise self._compute_cost_error + return self._compute_cost_result + + return UsageIngestionDeps( + lookup_key=lookup_key, + spend_log_exists=spend_log_exists, + record_spend=record_spend, + compute_cost=compute_cost, + generate_request_id=lambda: GENERATED_ID, + ) + + +def make_record(**overrides: Any) -> ExternalUsageRecord: + base: dict[str, Any] = { + "api_key": RAW_KEY, + "model": "gpt-4o-mini", + "prompt_tokens": 100, + "completion_tokens": 50, + "start_time": datetime(2026, 8, 5, 12, 0, 0, tzinfo=timezone.utc), + } + base.update(overrides) + return ExternalUsageRecord(**base) + + +def run(coro: Any) -> Any: + return asyncio.run(coro) + + +def test_records_spend_with_explicit_cost_without_calling_pricing(): + deps = RecordingDeps(compute_cost_error=RuntimeError("pricing must not be consulted")) + result = run(process_external_usage_record(make_record(cost=0.123, idempotency_key="batch-1-line-1"), deps.as_deps())) + assert result.status == "recorded" + assert result.spend == 0.123 + assert result.request_id == "batch-1-line-1" + assert len(deps.spend_calls) == 1 + assert deps.compute_cost_calls == [] + + call = deps.spend_calls[0] + assert call["token"] == hash_token(RAW_KEY) + assert call["token"] != RAW_KEY + assert call["user_id"] == "u-1" + assert call["team_id"] == "t-1" + assert call["org_id"] == "o-1" + assert call["response_cost"] == 0.123 + + response = call["completion_response"] + assert response.id == "batch-1-line-1" + assert response.usage.prompt_tokens == 100 + assert response.usage.completion_tokens == 50 + assert response.usage.total_tokens == 150 + + kwargs = call["kwargs"] + assert kwargs["call_type"] == "ingest_external_usage" + metadata = kwargs["litellm_params"]["metadata"] + assert metadata["user_api_key"] == hash_token(RAW_KEY) + assert metadata["user_api_key"] != RAW_KEY + + +def test_computed_cost_used_when_no_explicit_cost(): + deps = RecordingDeps(compute_cost_result=0.07) + result = run(process_external_usage_record(make_record(idempotency_key="k-2"), deps.as_deps())) + assert result.status == "recorded" + assert result.spend == 0.07 + assert len(deps.compute_cost_calls) == 1 + assert deps.compute_cost_calls[0][1] == "gpt-4o-mini" + assert deps.spend_calls[0]["response_cost"] == 0.07 + + +def test_unpriceable_model_without_explicit_cost_is_error_not_zero_spend(): + deps = RecordingDeps(compute_cost_error=ValueError("unknown model")) + result = run(process_external_usage_record(make_record(idempotency_key="k-3"), deps.as_deps())) + assert result.status == "error" + assert result.spend is None + assert "explicit cost" in (result.error or "") + assert len(deps.spend_calls) == 0 + + +def test_unknown_key_is_rejected_and_never_books_spend(): + deps = RecordingDeps(key=None) + result = run(process_external_usage_record(make_record(idempotency_key="k-4"), deps.as_deps())) + assert result.status == "error" + assert result.error == "api key not found" + assert len(deps.spend_calls) == 0 + + +def test_duplicate_idempotency_key_is_skipped_and_never_rebooks(): + deps = RecordingDeps(existing_ids=frozenset({"k-5"})) + result = run(process_external_usage_record(make_record(idempotency_key="k-5"), deps.as_deps())) + assert result.status == "duplicate" + assert result.request_id == "k-5" + assert len(deps.spend_calls) == 0 + + +def test_missing_idempotency_key_generates_request_id_and_skips_dedup_probe(): + deps = RecordingDeps() + result = run(process_external_usage_record(make_record(cost=0.01), deps.as_deps())) + assert result.status == "recorded" + assert result.request_id == GENERATED_ID + assert deps.spend_calls[0]["completion_response"].id == GENERATED_ID + assert deps.exists_calls == [] + + +def test_end_time_defaults_to_start_time_and_end_before_start_rejected(): + deps = RecordingDeps() + start = datetime(2026, 8, 5, 12, 0, 0, tzinfo=timezone.utc) + run(process_external_usage_record(make_record(cost=0.01, start_time=start), deps.as_deps())) + assert deps.spend_calls[0]["start_time"] == start + assert deps.spend_calls[0]["end_time"] == start + + earlier = datetime(2026, 8, 5, 11, 0, 0, tzinfo=timezone.utc) + with pytest.raises(ValidationError): + make_record(start_time=start, end_time=earlier) + + +def test_tags_and_end_user_flow_into_payload(): + deps = RecordingDeps() + result = run( + process_external_usage_record( + make_record(cost=0.01, idempotency_key="k-8", tags=["batch:job-42"], end_user_id="tenant-a"), + deps.as_deps(), + ) + ) + assert result.status == "recorded" + call = deps.spend_calls[0] + metadata = call["kwargs"]["litellm_params"]["metadata"] + assert metadata["tags"] == ["batch:job-42"] + assert metadata["user_api_key_end_user_id"] == "tenant-a" + assert call["end_user_id"] == "tenant-a" From b9ce4980098d9547ef84a40129d6c70fd25901c9 Mon Sep 17 00:00:00 2001 From: todayim <809634488@qq.com> Date: Wed, 5 Aug 2026 17:28:21 +0800 Subject: [PATCH 02/12] fix(proxy): enforce proxy-admin role, atomic idempotency, key-helper lookup for /spend/usage --- .../usage_ingestion_endpoints.py | 165 +++++++++++++++--- .../test_usage_ingestion_endpoints.py | 137 ++++++++++----- ui/litellm-dashboard/src/lib/http/schema.d.ts | 136 +++++++++++++++ 3 files changed, 368 insertions(+), 70 deletions(-) diff --git a/litellm/proxy/spend_tracking/usage_ingestion_endpoints.py b/litellm/proxy/spend_tracking/usage_ingestion_endpoints.py index 8e3e3dd448d..ff1b3678e2d 100644 --- a/litellm/proxy/spend_tracking/usage_ingestion_endpoints.py +++ b/litellm/proxy/spend_tracking/usage_ingestion_endpoints.py @@ -1,16 +1,18 @@ +import asyncio import uuid from collections.abc import Awaitable, Callable from dataclasses import dataclass from datetime import datetime -from typing import Final, Literal, NamedTuple +from typing import Annotated, Final, Literal, NamedTuple -from fastapi import APIRouter, Depends, status +from fastapi import APIRouter, Depends, HTTPException, status from pydantic import BaseModel, Field, model_validator import litellm from litellm._logging import verbose_proxy_logger -from litellm.proxy._types import ProxyException +from litellm.proxy._types import LitellmUserRoles, ProxyException, SpendLogsPayload, UserAPIKeyAuth from litellm.proxy.auth.user_api_key_auth import user_api_key_auth +from litellm.proxy.spend_tracking.spend_tracking_utils import get_logging_payload from litellm.proxy.utils import hash_token router: Final = APIRouter() @@ -25,8 +27,12 @@ class ExternalUsageRecord(BaseModel): completion_tokens: int = Field(ge=0) start_time: datetime end_time: datetime | None = None - cost: float | None = Field(default=None, ge=0, description="Explicit cost in USD. Computed from litellm pricing when omitted.") - idempotency_key: str | None = Field(default=None, max_length=255, description="Becomes the spend-log request_id for dedup on retries.") + cost: float | None = Field( + default=None, ge=0, description="Explicit cost in USD. Computed from litellm pricing when omitted." + ) + idempotency_key: str | None = Field( + default=None, max_length=255, description="Becomes the spend-log request_id for dedup on retries." + ) tags: list[str] | None = None end_user_id: str | None = None @@ -58,14 +64,14 @@ class KeyAttribution(NamedTuple): organization_id: str | None -RecordSpendFn = Callable[..., Awaitable[None]] +ReserveSpendFn = Callable[[ExternalUsageRecord, str, str, KeyAttribution, float], Awaitable[bool]] @dataclass(frozen=True, slots=True) class UsageIngestionDeps: lookup_key: Callable[[str], Awaitable[KeyAttribution | None]] - spend_log_exists: Callable[[str], Awaitable[bool]] - record_spend: RecordSpendFn + reserve_spend_log: ReserveSpendFn + record_spend: Callable[..., Awaitable[None]] compute_cost: Callable[[litellm.ModelResponse, str], float] generate_request_id: Callable[[], str] @@ -98,13 +104,40 @@ def _build_completion_response(record: ExternalUsageRecord, request_id: str) -> ) +def build_spend_log_payload( + record: ExternalUsageRecord, + request_id: str, + hashed_token: str, + key: KeyAttribution, + cost: float, +) -> SpendLogsPayload: + payload: SpendLogsPayload = get_logging_payload( + kwargs=_build_usage_kwargs(record, hashed_token), + response_obj=_build_completion_response(record, request_id), + start_time=record.start_time, + end_time=record.end_time or record.start_time, + ) + payload["spend"] = cost + if isinstance(payload["startTime"], datetime): + payload["startTime"] = payload["startTime"].isoformat() + if isinstance(payload["endTime"], datetime): + payload["endTime"] = payload["endTime"].isoformat() + if key.organization_id is not None and key.organization_id != "": + payload["organization_id"] = key.organization_id + if key.team_id is not None and key.team_id != "": + payload["team_id"] = key.team_id + return payload + + def _resolve_cost(deps: UsageIngestionDeps, record: ExternalUsageRecord, response: litellm.ModelResponse) -> float: if record.cost is not None: return record.cost return deps.compute_cost(response, record.model) -async def process_external_usage_record(record: ExternalUsageRecord, deps: UsageIngestionDeps) -> UsageIngestRecordResult: +async def process_external_usage_record( + record: ExternalUsageRecord, deps: UsageIngestionDeps +) -> UsageIngestRecordResult: request_id: Final = record.idempotency_key or deps.generate_request_id() hashed_token: Final = hash_token(record.api_key) @@ -112,14 +145,11 @@ async def process_external_usage_record(record: ExternalUsageRecord, deps: Usage if key is None: return UsageIngestRecordResult(request_id=request_id, status="error", error="api key not found") - if record.idempotency_key is not None and await deps.spend_log_exists(request_id): - return UsageIngestRecordResult(request_id=request_id, status="duplicate") - response: Final = _build_completion_response(record, request_id) try: cost: Final = _resolve_cost(deps, record, response) - except Exception as e: # noqa: BLE001 # pricing lookup raises arbitrary provider-specific errors; any failure means the record is unpriceable and must carry an explicit cost + except Exception as e: # noqa: BLE001 verbose_proxy_logger.info("ingest usage: cost computation failed for model %s: %s", record.model, e) return UsageIngestRecordResult( request_id=request_id, @@ -127,6 +157,12 @@ async def process_external_usage_record(record: ExternalUsageRecord, deps: Usage error="could not compute cost for this model, pass an explicit cost", ) + if record.idempotency_key is not None: + reserved: Final = await deps.reserve_spend_log(record, request_id, hashed_token, key, cost) + if reserved is False: + return UsageIngestRecordResult(request_id=request_id, status="duplicate") + return UsageIngestRecordResult(request_id=request_id, status="recorded", spend=cost) + await deps.record_spend( token=hashed_token, user_id=key.user_id, @@ -150,8 +186,72 @@ def _attribution_of(key_row: object) -> KeyAttribution: ) +async def reserve_spend_log_atomic( + record: ExternalUsageRecord, + request_id: str, + hashed_token: str, + key: KeyAttribution, + cost: float, +) -> bool: + from litellm.proxy.proxy_server import ( + disable_spend_logs, + litellm_proxy_budget_name, + prisma_client, + proxy_logging_obj, + ) + from litellm.proxy.utils import ProxyUpdateSpend + from litellm.repositories.table_repositories import SpendLogsRepository + + if ProxyUpdateSpend.disable_spend_updates() is True: + return True + + if disable_spend_logs is False: + payload: Final = prisma_client.jsonify_object( + {**build_spend_log_payload(record, request_id, hashed_token, key, cost)} + ) + from prisma.errors import UniqueViolationError + + try: + await SpendLogsRepository(prisma_client).table.create(data=payload) + except UniqueViolationError: + return False + + writer: Final = proxy_logging_obj.db_spend_update_writer + counter_calls: Final = ( + writer._update_key_db( + response_cost=cost, + hashed_token=hashed_token, + prisma_client=prisma_client, + ), + writer._update_user_db( + response_cost=cost, + user_id=key.user_id, + prisma_client=prisma_client, + litellm_proxy_budget_name=litellm_proxy_budget_name, + end_user_id=record.end_user_id, + ), + writer._update_team_db( + response_cost=cost, + team_id=key.team_id, + user_id=key.user_id, + prisma_client=prisma_client, + ), + writer._update_org_db( + response_cost=cost, + org_id=key.organization_id, + prisma_client=prisma_client, + ), + ) + + results: Final = await asyncio.gather(*counter_calls, return_exceptions=True) + for counter_result in results: + if isinstance(counter_result, Exception): + verbose_proxy_logger.debug("ingest usage: spend counter update failed: %s", counter_result) + return True + + def default_ingestion_deps() -> UsageIngestionDeps: - from litellm.proxy.proxy_server import prisma_client, proxy_logging_obj + from litellm.proxy.proxy_server import prisma_client, proxy_logging_obj, user_api_key_cache if prisma_client is None: raise ProxyException( @@ -162,18 +262,21 @@ def default_ingestion_deps() -> UsageIngestionDeps: ) async def lookup_key(hashed_token: str) -> KeyAttribution | None: - row: Final = await prisma_client.db.litellm_verificationtoken.find_unique(where={"token": hashed_token}) - if row is None: - return None - return _attribution_of(row) + from litellm.proxy.auth.auth_checks import get_key_object - async def spend_log_exists(request_id: str) -> bool: - row: Final = await prisma_client.db.litellm_spendlogs.find_unique(where={"request_id": request_id}) - return row is not None + try: + key_row: Final = await get_key_object( + hashed_token=hashed_token, + prisma_client=prisma_client, + user_api_key_cache=user_api_key_cache, + ) + except ProxyException: + return None + return _attribution_of(key_row) return UsageIngestionDeps( lookup_key=lookup_key, - spend_log_exists=spend_log_exists, + reserve_spend_log=reserve_spend_log_atomic, record_spend=proxy_logging_obj.db_spend_update_writer.update_database, compute_cost=lambda resp, model: litellm.completion_cost(completion_response=resp, model=model), generate_request_id=lambda: str(uuid.uuid4()), @@ -186,7 +289,10 @@ def default_ingestion_deps() -> UsageIngestionDeps: dependencies=[Depends(user_api_key_auth)], response_model=UsageIngestResponse, ) -async def ingest_external_usage(request: UsageIngestRequest) -> UsageIngestResponse: +async def ingest_external_usage( + request: UsageIngestRequest, + user_api_key_dict: Annotated[UserAPIKeyAuth, Depends(user_api_key_auth)], +) -> UsageIngestResponse: """ PROXY_ADMIN ONLY: record externally measured usage into the same spend pipeline as proxy-routed traffic. @@ -195,10 +301,17 @@ async def ingest_external_usage(request: UsageIngestRequest) -> UsageIngestRespo single metering system. Attribution (user/team/org) is derived from the given virtual key. Records accept an optional - idempotency_key, stored as the spend-log request_id, so retries are deduplicated. When cost is - omitted it is computed from litellm pricing; records whose model cannot be priced are rejected - with an error instead of being booked as zero spend. + idempotency_key, stored as the spend-log request_id: the reservation insert, counter updates and + dedup are checked atomically at the database primary key, so overlapping retries are safe. When + cost is omitted it is computed from litellm pricing; records whose model cannot be priced are + rejected with an error instead of being booked as zero spend. """ + if user_api_key_dict.user_role != LitellmUserRoles.PROXY_ADMIN: + raise HTTPException( + status_code=status.HTTP_403_FORBIDDEN, + detail="Only proxy admins ingest spend records here. Use a key with the proxy_admin role.", + ) + deps: Final = default_ingestion_deps() results: Final = [await process_external_usage_record(record, deps) for record in request.records] return UsageIngestResponse(results=results) diff --git a/tests/test_litellm/proxy/spend_tracking/test_usage_ingestion_endpoints.py b/tests/test_litellm/proxy/spend_tracking/test_usage_ingestion_endpoints.py index 360f2e1af6e..5384cd18a82 100644 --- a/tests/test_litellm/proxy/spend_tracking/test_usage_ingestion_endpoints.py +++ b/tests/test_litellm/proxy/spend_tracking/test_usage_ingestion_endpoints.py @@ -1,4 +1,5 @@ import asyncio +import json import os import sys from datetime import datetime, timezone @@ -13,6 +14,7 @@ from litellm.proxy.spend_tracking.usage_ingestion_endpoints import ( ExternalUsageRecord, KeyAttribution, UsageIngestionDeps, + build_spend_log_payload, process_external_usage_record, ) from litellm.proxy.utils import hash_token @@ -35,17 +37,31 @@ class RecordingDeps: self._compute_cost_result = compute_cost_result self._compute_cost_error = compute_cost_error self.spend_calls: list[dict[str, Any]] = [] + self.reserve_calls: list[dict[str, Any]] = [] self.compute_cost_calls: list[tuple[Any, str]] = [] - self.exists_calls: list[str] = [] def as_deps(self) -> UsageIngestionDeps: async def lookup_key(hashed: str) -> KeyAttribution | None: self.looked_up_hashed = hashed return self._key - async def spend_log_exists(request_id: str) -> bool: - self.exists_calls.append(request_id) - return request_id in self._existing_ids + async def reserve_spend_log( + record: ExternalUsageRecord, + request_id: str, + hashed_token: str, + key: KeyAttribution, + cost: float, + ) -> bool: + self.reserve_calls.append( + { + "record": record, + "request_id": request_id, + "hashed_token": hashed_token, + "key": key, + "cost": cost, + } + ) + return request_id not in self._existing_ids async def record_spend(**kwargs: Any) -> None: self.spend_calls.append(kwargs) @@ -58,7 +74,7 @@ class RecordingDeps: return UsageIngestionDeps( lookup_key=lookup_key, - spend_log_exists=spend_log_exists, + reserve_spend_log=reserve_spend_log, record_spend=record_spend, compute_cost=compute_cost, generate_request_id=lambda: GENERATED_ID, @@ -81,61 +97,82 @@ def run(coro: Any) -> Any: return asyncio.run(coro) -def test_records_spend_with_explicit_cost_without_calling_pricing(): +def test_spend_log_payload_matches_funnel_shape(): + record = make_record(idempotency_key="batch-1-line-1", tags=["batch:job-42"], end_user_id="tenant-a") + payload = build_spend_log_payload( + record=record, + request_id="batch-1-line-1", + hashed_token=hash_token(RAW_KEY), + key=DEFAULT_KEY, + cost=0.123, + ) + assert payload["request_id"] == "batch-1-line-1" + assert payload["spend"] == 0.123 + assert payload["total_tokens"] == 150 + assert payload["prompt_tokens"] == 100 + assert payload["completion_tokens"] == 50 + assert payload["api_key"] == hash_token(RAW_KEY) + assert payload["api_key"] != RAW_KEY + assert payload["team_id"] == "t-1" + assert payload["organization_id"] == "o-1" + assert payload["end_user"] == "tenant-a" + + metadata = json.loads(payload["metadata"]) if isinstance(payload["metadata"], str) else payload["metadata"] + assert metadata["user_api_key"] == hash_token(RAW_KEY) + + request_tags = ( + json.loads(payload["request_tags"]) if isinstance(payload["request_tags"], str) else payload["request_tags"] + ) + assert request_tags == ["batch:job-42"] + + +def test_idempotent_record_books_through_atomic_reserve_not_funnel(): deps = RecordingDeps(compute_cost_error=RuntimeError("pricing must not be consulted")) - result = run(process_external_usage_record(make_record(cost=0.123, idempotency_key="batch-1-line-1"), deps.as_deps())) + result = run( + process_external_usage_record(make_record(cost=0.123, idempotency_key="batch-1-line-1"), deps.as_deps()) + ) assert result.status == "recorded" assert result.spend == 0.123 assert result.request_id == "batch-1-line-1" - assert len(deps.spend_calls) == 1 + + assert len(deps.reserve_calls) == 1 + reserve_call = deps.reserve_calls[0] + assert reserve_call["request_id"] == "batch-1-line-1" + assert reserve_call["hashed_token"] == hash_token(RAW_KEY) + assert reserve_call["cost"] == 0.123 + assert reserve_call["key"] == DEFAULT_KEY + + assert deps.spend_calls == [] assert deps.compute_cost_calls == [] - call = deps.spend_calls[0] - assert call["token"] == hash_token(RAW_KEY) - assert call["token"] != RAW_KEY - assert call["user_id"] == "u-1" - assert call["team_id"] == "t-1" - assert call["org_id"] == "o-1" - assert call["response_cost"] == 0.123 - response = call["completion_response"] - assert response.id == "batch-1-line-1" - assert response.usage.prompt_tokens == 100 - assert response.usage.completion_tokens == 50 - assert response.usage.total_tokens == 150 - - kwargs = call["kwargs"] - assert kwargs["call_type"] == "ingest_external_usage" - metadata = kwargs["litellm_params"]["metadata"] - assert metadata["user_api_key"] == hash_token(RAW_KEY) - assert metadata["user_api_key"] != RAW_KEY - - -def test_computed_cost_used_when_no_explicit_cost(): +def test_computed_cost_is_resolved_before_reserving(): deps = RecordingDeps(compute_cost_result=0.07) result = run(process_external_usage_record(make_record(idempotency_key="k-2"), deps.as_deps())) assert result.status == "recorded" assert result.spend == 0.07 assert len(deps.compute_cost_calls) == 1 assert deps.compute_cost_calls[0][1] == "gpt-4o-mini" - assert deps.spend_calls[0]["response_cost"] == 0.07 + assert deps.reserve_calls[0]["cost"] == 0.07 -def test_unpriceable_model_without_explicit_cost_is_error_not_zero_spend(): +def test_unpriceable_model_without_explicit_cost_is_error_and_books_nothing(): deps = RecordingDeps(compute_cost_error=ValueError("unknown model")) result = run(process_external_usage_record(make_record(idempotency_key="k-3"), deps.as_deps())) assert result.status == "error" assert result.spend is None assert "explicit cost" in (result.error or "") - assert len(deps.spend_calls) == 0 + assert deps.reserve_calls == [] + assert deps.spend_calls == [] -def test_unknown_key_is_rejected_and_never_books_spend(): +def test_unknown_key_is_rejected_and_books_nothing(): deps = RecordingDeps(key=None) result = run(process_external_usage_record(make_record(idempotency_key="k-4"), deps.as_deps())) assert result.status == "error" assert result.error == "api key not found" - assert len(deps.spend_calls) == 0 + assert deps.reserve_calls == [] + assert deps.spend_calls == [] def test_duplicate_idempotency_key_is_skipped_and_never_rebooks(): @@ -143,16 +180,30 @@ def test_duplicate_idempotency_key_is_skipped_and_never_rebooks(): result = run(process_external_usage_record(make_record(idempotency_key="k-5"), deps.as_deps())) assert result.status == "duplicate" assert result.request_id == "k-5" - assert len(deps.spend_calls) == 0 + assert len(deps.reserve_calls) == 1 + assert deps.spend_calls == [] -def test_missing_idempotency_key_generates_request_id_and_skips_dedup_probe(): +def test_missing_idempotency_key_uses_funnel_with_generated_request_id(): deps = RecordingDeps() result = run(process_external_usage_record(make_record(cost=0.01), deps.as_deps())) assert result.status == "recorded" assert result.request_id == GENERATED_ID - assert deps.spend_calls[0]["completion_response"].id == GENERATED_ID - assert deps.exists_calls == [] + assert deps.reserve_calls == [] + assert len(deps.spend_calls) == 1 + + call = deps.spend_calls[0] + assert call["token"] == hash_token(RAW_KEY) + assert call["user_id"] == "u-1" + assert call["team_id"] == "t-1" + assert call["org_id"] == "o-1" + assert call["response_cost"] == 0.01 + response = call["completion_response"] + assert response.id == GENERATED_ID + assert response.usage.total_tokens == 150 + metadata = call["kwargs"]["litellm_params"]["metadata"] + assert metadata["user_api_key"] == hash_token(RAW_KEY) + assert metadata["user_api_key"] != RAW_KEY def test_end_time_defaults_to_start_time_and_end_before_start_rejected(): @@ -167,17 +218,15 @@ def test_end_time_defaults_to_start_time_and_end_before_start_rejected(): make_record(start_time=start, end_time=earlier) -def test_tags_and_end_user_flow_into_payload(): +def test_record_without_idempotency_key_still_flows_tags_to_funnel_kwargs(): deps = RecordingDeps() result = run( process_external_usage_record( - make_record(cost=0.01, idempotency_key="k-8", tags=["batch:job-42"], end_user_id="tenant-a"), - deps.as_deps(), + make_record(cost=0.01, tags=["batch:job-42"], end_user_id="tenant-a"), deps.as_deps() ) ) assert result.status == "recorded" - call = deps.spend_calls[0] - metadata = call["kwargs"]["litellm_params"]["metadata"] + metadata = deps.spend_calls[0]["kwargs"]["litellm_params"]["metadata"] assert metadata["tags"] == ["batch:job-42"] assert metadata["user_api_key_end_user_id"] == "tenant-a" - assert call["end_user_id"] == "tenant-a" + assert deps.spend_calls[0]["end_user_id"] == "tenant-a" diff --git a/ui/litellm-dashboard/src/lib/http/schema.d.ts b/ui/litellm-dashboard/src/lib/http/schema.d.ts index df85decc676..d86f3bbb5bd 100644 --- a/ui/litellm-dashboard/src/lib/http/schema.d.ts +++ b/ui/litellm-dashboard/src/lib/http/schema.d.ts @@ -12859,6 +12859,36 @@ export interface paths { patch?: never; trace?: never; }; + "/spend/usage": { + parameters: { + query?: never; + header?: never; + path?: never; + cookie?: never; + }; + get?: never; + put?: never; + /** + * Ingest External Usage + * @description PROXY_ADMIN ONLY: record externally measured usage into the same spend pipeline as proxy-routed traffic. + * + * For inference traffic that legitimately bypasses the proxy (for example async batch processors + * dispatching directly to model gateways), so budgets and spend stay coherent in litellm as the + * single metering system. + * + * Attribution (user/team/org) is derived from the given virtual key. Records accept an optional + * idempotency_key, stored as the spend-log request_id: the reservation insert, counter updates and + * dedup are checked atomically at the database primary key, so overlapping retries are safe. When + * cost is omitted it is computed from litellm pricing; records whose model cannot be priced are + * rejected with an error instead of being booked as zero spend. + */ + post: operations["ingest_external_usage_spend_usage_post"]; + delete?: never; + options?: never; + head?: never; + patch?: never; + trace?: never; + }; "/spend/users": { parameters: { query?: never; @@ -22482,6 +22512,7 @@ export interface components { /** ChatCompletionAudioObject */ ChatCompletionAudioObject: { input_audio: components["schemas"]["InputAudio"]; + prompt_cache_breakpoint?: components["schemas"]["PromptCacheBreakpoint"]; /** * Type * @constant @@ -24401,6 +24432,41 @@ export interface components { /** Updated At */ updated_at?: number | null; }; + /** ExternalUsageRecord */ + ExternalUsageRecord: { + /** + * Api Key + * @description Raw virtual key (sk-...) to attribute usage to. Never logged. + */ + api_key: string; + /** Completion Tokens */ + completion_tokens: number; + /** + * Cost + * @description Explicit cost in USD. Computed from litellm pricing when omitted. + */ + cost?: number | null; + /** End Time */ + end_time?: string | null; + /** End User Id */ + end_user_id?: string | null; + /** + * Idempotency Key + * @description Becomes the spend-log request_id for dedup on retries. + */ + idempotency_key?: string | null; + /** Model */ + model: string; + /** Prompt Tokens */ + prompt_tokens: number; + /** + * Start Time + * Format: date-time + */ + start_time: string; + /** Tags */ + tags?: string[] | null; + }; /** * FacetListResponse * @description The distinct values one column takes over a filtered query. `data` holds bare values, not entity rows. @@ -30582,6 +30648,19 @@ export interface components { prompt_id: string; prompt_info?: components["schemas"]["PromptInfo"] | null; }; + /** + * PromptCacheBreakpoint + * @description Marks the exact end of a reusable prompt prefix. + * + * The breakpoint inherits its TTL from the request's `prompt_cache_options.ttl`; the boundary is not rounded to a token block. + */ + PromptCacheBreakpoint: { + /** + * Mode + * @constant + */ + mode: "explicit"; + }; /** PromptInfo */ PromptInfo: { /** @@ -34118,6 +34197,30 @@ export interface components { /** Type */ type: string; }; + /** UsageIngestRecordResult */ + UsageIngestRecordResult: { + /** Error */ + error?: string | null; + /** Request Id */ + request_id: string; + /** Spend */ + spend?: number | null; + /** + * Status + * @enum {string} + */ + status: "recorded" | "duplicate" | "error"; + }; + /** UsageIngestRequest */ + UsageIngestRequest: { + /** Records */ + records: components["schemas"]["ExternalUsageRecord"][]; + }; + /** UsageIngestResponse */ + UsageIngestResponse: { + /** Results */ + results: components["schemas"]["UsageIngestRecordResult"][]; + }; /** UsageLogEntry */ UsageLogEntry: { /** Action */ @@ -50983,6 +51086,39 @@ export interface operations { }; }; }; + ingest_external_usage_spend_usage_post: { + parameters: { + query?: never; + header?: never; + path?: never; + cookie?: never; + }; + requestBody: { + content: { + "application/json": components["schemas"]["UsageIngestRequest"]; + }; + }; + responses: { + /** @description Successful Response */ + 200: { + headers: { + [name: string]: unknown; + }; + content: { + "application/json": components["schemas"]["UsageIngestResponse"]; + }; + }; + /** @description Validation Error */ + 422: { + headers: { + [name: string]: unknown; + }; + content: { + "application/json": components["schemas"]["HTTPValidationError"]; + }; + }; + }; + }; spend_user_fn_spend_users_get: { parameters: { query?: { From 7a4f6a09156fdc5d946a5e975fbacd97be49fee5 Mon Sep 17 00:00:00 2001 From: todayim <809634488@qq.com> Date: Wed, 5 Aug 2026 21:35:01 +0800 Subject: [PATCH 03/12] refactor(proxy): satisfy type-discipline budget gates for /spend/usage --- .../usage_ingestion_endpoints.py | 55 ++++++++++--------- 1 file changed, 30 insertions(+), 25 deletions(-) diff --git a/litellm/proxy/spend_tracking/usage_ingestion_endpoints.py b/litellm/proxy/spend_tracking/usage_ingestion_endpoints.py index ff1b3678e2d..62b07260b93 100644 --- a/litellm/proxy/spend_tracking/usage_ingestion_endpoints.py +++ b/litellm/proxy/spend_tracking/usage_ingestion_endpoints.py @@ -1,8 +1,9 @@ import asyncio import uuid -from collections.abc import Awaitable, Callable +from collections.abc import Awaitable, Callable, Mapping from dataclasses import dataclass from datetime import datetime +from types import MappingProxyType from typing import Annotated, Final, Literal, NamedTuple from fastapi import APIRouter, Depends, HTTPException, status @@ -33,7 +34,7 @@ class ExternalUsageRecord(BaseModel): idempotency_key: str | None = Field( default=None, max_length=255, description="Becomes the spend-log request_id for dedup on retries." ) - tags: list[str] | None = None + tags: list[str] | None = None # mutable-ok: serialized as a JSON array by spend logs end_user_id: str | None = None @model_validator(mode="after") @@ -44,7 +45,7 @@ class ExternalUsageRecord(BaseModel): class UsageIngestRequest(BaseModel): - records: list[ExternalUsageRecord] = Field(min_length=1, max_length=MAX_RECORDS_PER_REQUEST) + records: tuple[ExternalUsageRecord, ...] = Field(min_length=1, max_length=MAX_RECORDS_PER_REQUEST) class UsageIngestRecordResult(BaseModel): @@ -55,7 +56,7 @@ class UsageIngestRecordResult(BaseModel): class UsageIngestResponse(BaseModel): - results: list[UsageIngestRecordResult] + results: tuple[UsageIngestRecordResult, ...] class KeyAttribution(NamedTuple): @@ -64,29 +65,31 @@ class KeyAttribution(NamedTuple): organization_id: str | None -ReserveSpendFn = Callable[[ExternalUsageRecord, str, str, KeyAttribution, float], Awaitable[bool]] - - @dataclass(frozen=True, slots=True) class UsageIngestionDeps: lookup_key: Callable[[str], Awaitable[KeyAttribution | None]] - reserve_spend_log: ReserveSpendFn + reserve_spend_log: Callable[[ExternalUsageRecord, str, str, KeyAttribution, float], Awaitable[bool]] record_spend: Callable[..., Awaitable[None]] compute_cost: Callable[[litellm.ModelResponse, str], float] generate_request_id: Callable[[], str] -def _build_usage_kwargs(record: ExternalUsageRecord, hashed_token: str) -> dict[str, object]: - metadata = { - "user_api_key": hashed_token, - "user_api_key_end_user_id": record.end_user_id, - "tags": list(record.tags) if record.tags else [], - } - return { - "model": record.model, - "call_type": "ingest_external_usage", - "litellm_params": {"model": record.model, "metadata": metadata}, - } +def _build_usage_kwargs(record: ExternalUsageRecord, hashed_token: str) -> Mapping[str, object]: + tags: Final = list(record.tags) if record.tags else [] # mutable-ok: real list required by json serializer + metadata: Final = MappingProxyType( + { + "user_api_key": hashed_token, + "user_api_key_end_user_id": record.end_user_id, + "tags": tags, + } + ) + return MappingProxyType( + { + "model": record.model, + "call_type": "ingest_external_usage", + "litellm_params": MappingProxyType({"model": record.model, "metadata": metadata}), + } + ) def _build_completion_response(record: ExternalUsageRecord, request_id: str) -> litellm.ModelResponse: @@ -111,7 +114,7 @@ def build_spend_log_payload( key: KeyAttribution, cost: float, ) -> SpendLogsPayload: - payload: SpendLogsPayload = get_logging_payload( + payload: Final[SpendLogsPayload] = get_logging_payload( kwargs=_build_usage_kwargs(record, hashed_token), response_obj=_build_completion_response(record, request_id), start_time=record.start_time, @@ -149,7 +152,7 @@ async def process_external_usage_record( try: cost: Final = _resolve_cost(deps, record, response) - except Exception as e: # noqa: BLE001 + except Exception as e: # noqa: BLE001 # pricing lookup raises arbitrary provider-specific errors; any failure means the model is unpriceable and the record must carry an explicit cost verbose_proxy_logger.info("ingest usage: cost computation failed for model %s: %s", record.model, e) return UsageIngestRecordResult( request_id=request_id, @@ -207,7 +210,7 @@ async def reserve_spend_log_atomic( if disable_spend_logs is False: payload: Final = prisma_client.jsonify_object( - {**build_spend_log_payload(record, request_id, hashed_token, key, cost)} + build_spend_log_payload(record, request_id, hashed_token, key, cost) ) from prisma.errors import UniqueViolationError @@ -285,8 +288,8 @@ def default_ingestion_deps() -> UsageIngestionDeps: @router.post( "/spend/usage", - tags=["Budget & Spend Tracking"], - dependencies=[Depends(user_api_key_auth)], + tags=["Budget & Spend Tracking"], # mutable-ok: fastapi decorator contract takes a list + dependencies=[Depends(user_api_key_auth)], # mutable-ok: fastapi decorator contract takes a list response_model=UsageIngestResponse, ) async def ingest_external_usage( @@ -313,5 +316,7 @@ async def ingest_external_usage( ) deps: Final = default_ingestion_deps() - results: Final = [await process_external_usage_record(record, deps) for record in request.records] + results: Final = tuple( + await asyncio.gather(*(process_external_usage_record(record, deps) for record in request.records)) + ) return UsageIngestResponse(results=results) From bd6a177bc8068a3340bf4b1858f93eebabda2250 Mon Sep 17 00:00:00 2001 From: todayim <809634488@qq.com> Date: Wed, 5 Aug 2026 22:05:45 +0800 Subject: [PATCH 04/12] chore(ui): regenerate schema.d.ts for the new /spend/usage route --- ui/litellm-dashboard/src/lib/http/schema.d.ts | 14 -------------- 1 file changed, 14 deletions(-) diff --git a/ui/litellm-dashboard/src/lib/http/schema.d.ts b/ui/litellm-dashboard/src/lib/http/schema.d.ts index d86f3bbb5bd..11c02cbd4d7 100644 --- a/ui/litellm-dashboard/src/lib/http/schema.d.ts +++ b/ui/litellm-dashboard/src/lib/http/schema.d.ts @@ -22512,7 +22512,6 @@ export interface components { /** ChatCompletionAudioObject */ ChatCompletionAudioObject: { input_audio: components["schemas"]["InputAudio"]; - prompt_cache_breakpoint?: components["schemas"]["PromptCacheBreakpoint"]; /** * Type * @constant @@ -30648,19 +30647,6 @@ export interface components { prompt_id: string; prompt_info?: components["schemas"]["PromptInfo"] | null; }; - /** - * PromptCacheBreakpoint - * @description Marks the exact end of a reusable prompt prefix. - * - * The breakpoint inherits its TTL from the request's `prompt_cache_options.ttl`; the boundary is not rounded to a token block. - */ - PromptCacheBreakpoint: { - /** - * Mode - * @constant - */ - mode: "explicit"; - }; /** PromptInfo */ PromptInfo: { /** From bc619dd785da40768e75442ddea2097d3130bc43 Mon Sep 17 00:00:00 2001 From: todayim <809634488@qq.com> Date: Wed, 5 Aug 2026 22:32:57 +0800 Subject: [PATCH 05/12] fix(proxy): always write the idempotency reservation row for /spend/usage records --- .../usage_ingestion_endpoints.py | 30 ++++++++----------- 1 file changed, 12 insertions(+), 18 deletions(-) diff --git a/litellm/proxy/spend_tracking/usage_ingestion_endpoints.py b/litellm/proxy/spend_tracking/usage_ingestion_endpoints.py index 62b07260b93..2af7f97bb97 100644 --- a/litellm/proxy/spend_tracking/usage_ingestion_endpoints.py +++ b/litellm/proxy/spend_tracking/usage_ingestion_endpoints.py @@ -196,28 +196,20 @@ async def reserve_spend_log_atomic( key: KeyAttribution, cost: float, ) -> bool: - from litellm.proxy.proxy_server import ( - disable_spend_logs, - litellm_proxy_budget_name, - prisma_client, - proxy_logging_obj, - ) + from litellm.proxy.proxy_server import litellm_proxy_budget_name, prisma_client, proxy_logging_obj from litellm.proxy.utils import ProxyUpdateSpend from litellm.repositories.table_repositories import SpendLogsRepository if ProxyUpdateSpend.disable_spend_updates() is True: return True - if disable_spend_logs is False: - payload: Final = prisma_client.jsonify_object( - build_spend_log_payload(record, request_id, hashed_token, key, cost) - ) - from prisma.errors import UniqueViolationError + payload: Final = prisma_client.jsonify_object(build_spend_log_payload(record, request_id, hashed_token, key, cost)) + from prisma.errors import UniqueViolationError - try: - await SpendLogsRepository(prisma_client).table.create(data=payload) - except UniqueViolationError: - return False + try: + await SpendLogsRepository(prisma_client).table.create(data=payload) + except UniqueViolationError: + return False writer: Final = proxy_logging_obj.db_spend_update_writer counter_calls: Final = ( @@ -305,9 +297,11 @@ async def ingest_external_usage( Attribution (user/team/org) is derived from the given virtual key. Records accept an optional idempotency_key, stored as the spend-log request_id: the reservation insert, counter updates and - dedup are checked atomically at the database primary key, so overlapping retries are safe. When - cost is omitted it is computed from litellm pricing; records whose model cannot be priced are - rejected with an error instead of being booked as zero spend. + dedup are checked atomically at the database primary key, so overlapping retries are safe. The + reservation row is always written (even when disable_spend_logs is set), because it is both the + dedup anchor and the audit record for the booked usage. When cost is omitted it is computed from + litellm pricing; records whose model cannot be priced are rejected with an error instead of + being booked as zero spend. """ if user_api_key_dict.user_role != LitellmUserRoles.PROXY_ADMIN: raise HTTPException( From 38329c22174d7603f849d6a170e792d04f8c3fdc Mon Sep 17 00:00:00 2001 From: todayim <809634488@qq.com> Date: Wed, 5 Aug 2026 22:35:08 +0800 Subject: [PATCH 06/12] chore(ui): regenerate schema.d.ts after docstring update --- ui/litellm-dashboard/src/lib/http/schema.d.ts | 8 +++++--- 1 file changed, 5 insertions(+), 3 deletions(-) diff --git a/ui/litellm-dashboard/src/lib/http/schema.d.ts b/ui/litellm-dashboard/src/lib/http/schema.d.ts index 11c02cbd4d7..3647ef4c396 100644 --- a/ui/litellm-dashboard/src/lib/http/schema.d.ts +++ b/ui/litellm-dashboard/src/lib/http/schema.d.ts @@ -12878,9 +12878,11 @@ export interface paths { * * Attribution (user/team/org) is derived from the given virtual key. Records accept an optional * idempotency_key, stored as the spend-log request_id: the reservation insert, counter updates and - * dedup are checked atomically at the database primary key, so overlapping retries are safe. When - * cost is omitted it is computed from litellm pricing; records whose model cannot be priced are - * rejected with an error instead of being booked as zero spend. + * dedup are checked atomically at the database primary key, so overlapping retries are safe. The + * reservation row is always written (even when disable_spend_logs is set), because it is both the + * dedup anchor and the audit record for the booked usage. When cost is omitted it is computed from + * litellm pricing; records whose model cannot be priced are rejected with an error instead of + * being booked as zero spend. */ post: operations["ingest_external_usage_spend_usage_post"]; delete?: never; From 1fc0abe027a7aa3d0ff0f08c6455f73bb163c8a9 Mon Sep 17 00:00:00 2001 From: todayim <809634488@qq.com> Date: Wed, 5 Aug 2026 22:55:09 +0800 Subject: [PATCH 07/12] chore: retrigger CI after a cancelled linting run From 73f31beb0d6b43777b1a05d01e7c039f6d178e5d Mon Sep 17 00:00:00 2001 From: todayim <809634488@qq.com> Date: Thu, 6 Aug 2026 01:00:04 +0800 Subject: [PATCH 08/12] fix(proxy): book idempotent /spend/usage records in one transaction, report disabled spend updates --- .../usage_ingestion_endpoints.py | 102 ++++++++++-------- .../test_usage_ingestion_endpoints.py | 30 +++++- 2 files changed, 88 insertions(+), 44 deletions(-) diff --git a/litellm/proxy/spend_tracking/usage_ingestion_endpoints.py b/litellm/proxy/spend_tracking/usage_ingestion_endpoints.py index 2af7f97bb97..9a49268fdae 100644 --- a/litellm/proxy/spend_tracking/usage_ingestion_endpoints.py +++ b/litellm/proxy/spend_tracking/usage_ingestion_endpoints.py @@ -4,7 +4,15 @@ from collections.abc import Awaitable, Callable, Mapping from dataclasses import dataclass from datetime import datetime from types import MappingProxyType -from typing import Annotated, Final, Literal, NamedTuple +from typing import ( + TYPE_CHECKING, + Annotated, + Final, + Literal, + NamedTuple, + TypeAlias, + cast, # noqa: TID251 # untyped tx boundary needs cast for the shim +) from fastapi import APIRouter, Depends, HTTPException, status from pydantic import BaseModel, Field, model_validator @@ -65,10 +73,13 @@ class KeyAttribution(NamedTuple): organization_id: str | None +ReservationOutcome: TypeAlias = Literal["reserved", "duplicate", "disabled"] + + @dataclass(frozen=True, slots=True) class UsageIngestionDeps: lookup_key: Callable[[str], Awaitable[KeyAttribution | None]] - reserve_spend_log: Callable[[ExternalUsageRecord, str, str, KeyAttribution, float], Awaitable[bool]] + reserve_spend_log: Callable[[ExternalUsageRecord, str, str, KeyAttribution, float], Awaitable[ReservationOutcome]] record_spend: Callable[..., Awaitable[None]] compute_cost: Callable[[litellm.ModelResponse, str], float] generate_request_id: Callable[[], str] @@ -161,9 +172,23 @@ async def process_external_usage_record( ) if record.idempotency_key is not None: - reserved: Final = await deps.reserve_spend_log(record, request_id, hashed_token, key, cost) - if reserved is False: + try: + reservation: Final = await deps.reserve_spend_log(record, request_id, hashed_token, key, cost) + except Exception as e: # noqa: BLE001 # booking raises arbitrary persistence errors; an aborted transaction means nothing was booked, so telling the caller to retry is safe + verbose_proxy_logger.info("ingest usage: transactional booking failed for %s: %s", request_id, e) + return UsageIngestRecordResult( + request_id=request_id, + status="error", + error="booking failed transactionally, nothing was recorded, safe to retry", + ) + if reservation == "duplicate": return UsageIngestRecordResult(request_id=request_id, status="duplicate") + if reservation == "disabled": + return UsageIngestRecordResult( + request_id=request_id, + status="error", + error="spend updates are disabled on this proxy, nothing was recorded", + ) return UsageIngestRecordResult(request_id=request_id, status="recorded", spend=cost) await deps.record_spend( @@ -189,60 +214,53 @@ def _attribution_of(key_row: object) -> KeyAttribution: ) +if TYPE_CHECKING: + from prisma.client import TransactionManager + + +class _TransactionClientShim: + def __init__(self, tx: "TransactionManager") -> None: + self.db: Final = tx + + async def reserve_spend_log_atomic( record: ExternalUsageRecord, request_id: str, hashed_token: str, key: KeyAttribution, cost: float, -) -> bool: +) -> ReservationOutcome: from litellm.proxy.proxy_server import litellm_proxy_budget_name, prisma_client, proxy_logging_obj - from litellm.proxy.utils import ProxyUpdateSpend + from litellm.proxy.utils import PrismaClient, ProxyUpdateSpend from litellm.repositories.table_repositories import SpendLogsRepository if ProxyUpdateSpend.disable_spend_updates() is True: - return True + return "disabled" payload: Final = prisma_client.jsonify_object(build_spend_log_payload(record, request_id, hashed_token, key, cost)) from prisma.errors import UniqueViolationError - try: - await SpendLogsRepository(prisma_client).table.create(data=payload) - except UniqueViolationError: - return False - writer: Final = proxy_logging_obj.db_spend_update_writer - counter_calls: Final = ( - writer._update_key_db( - response_cost=cost, - hashed_token=hashed_token, - prisma_client=prisma_client, - ), - writer._update_user_db( - response_cost=cost, - user_id=key.user_id, - prisma_client=prisma_client, - litellm_proxy_budget_name=litellm_proxy_budget_name, - end_user_id=record.end_user_id, - ), - writer._update_team_db( - response_cost=cost, - team_id=key.team_id, - user_id=key.user_id, - prisma_client=prisma_client, - ), - writer._update_org_db( - response_cost=cost, - org_id=key.organization_id, - prisma_client=prisma_client, - ), - ) - results: Final = await asyncio.gather(*counter_calls, return_exceptions=True) - for counter_result in results: - if isinstance(counter_result, Exception): - verbose_proxy_logger.debug("ingest usage: spend counter update failed: %s", counter_result) - return True + try: + async with prisma_client.tx() as tx: + shim: Final = cast(PrismaClient, _TransactionClientShim(tx)) # cast-ok: helper uses only .db (untyped) + await SpendLogsRepository(shim).table.create(data=payload) + await writer._update_key_db(response_cost=cost, hashed_token=hashed_token, prisma_client=shim) + await writer._update_user_db( + response_cost=cost, + user_id=key.user_id, + prisma_client=shim, + litellm_proxy_budget_name=litellm_proxy_budget_name, + end_user_id=record.end_user_id, + ) + await writer._update_team_db( + response_cost=cost, team_id=key.team_id, user_id=key.user_id, prisma_client=shim + ) + await writer._update_org_db(response_cost=cost, org_id=key.organization_id, prisma_client=shim) + except UniqueViolationError: + return "duplicate" + return "reserved" def default_ingestion_deps() -> UsageIngestionDeps: diff --git a/tests/test_litellm/proxy/spend_tracking/test_usage_ingestion_endpoints.py b/tests/test_litellm/proxy/spend_tracking/test_usage_ingestion_endpoints.py index 5384cd18a82..b5ae287269a 100644 --- a/tests/test_litellm/proxy/spend_tracking/test_usage_ingestion_endpoints.py +++ b/tests/test_litellm/proxy/spend_tracking/test_usage_ingestion_endpoints.py @@ -31,11 +31,15 @@ class RecordingDeps: existing_ids: frozenset[str] = frozenset(), compute_cost_result: float = 0.05, compute_cost_error: Exception | None = None, + reserve_outcome: str = "reserved", + reserve_raises: Exception | None = None, ): self._key = key self._existing_ids = existing_ids self._compute_cost_result = compute_cost_result self._compute_cost_error = compute_cost_error + self._reserve_outcome = reserve_outcome + self._reserve_raises = reserve_raises self.spend_calls: list[dict[str, Any]] = [] self.reserve_calls: list[dict[str, Any]] = [] self.compute_cost_calls: list[tuple[Any, str]] = [] @@ -51,7 +55,7 @@ class RecordingDeps: hashed_token: str, key: KeyAttribution, cost: float, - ) -> bool: + ) -> Any: self.reserve_calls.append( { "record": record, @@ -61,7 +65,11 @@ class RecordingDeps: "cost": cost, } ) - return request_id not in self._existing_ids + if self._reserve_raises is not None: + raise self._reserve_raises + if self._reserve_outcome == "disabled": + return "disabled" + return "duplicate" if request_id in self._existing_ids else "reserved" async def record_spend(**kwargs: Any) -> None: self.spend_calls.append(kwargs) @@ -230,3 +238,21 @@ def test_record_without_idempotency_key_still_flows_tags_to_funnel_kwargs(): assert metadata["tags"] == ["batch:job-42"] assert metadata["user_api_key_end_user_id"] == "tenant-a" assert deps.spend_calls[0]["end_user_id"] == "tenant-a" + + +def test_disabled_spend_updates_reports_error_instead_of_fake_recorded(): + deps = RecordingDeps(reserve_outcome="disabled") + result = run(process_external_usage_record(make_record(cost=0.01, idempotency_key="k-9"), deps.as_deps())) + assert result.status == "error" + assert "disabled" in (result.error or "") + assert result.spend is None + assert deps.spend_calls == [] + + +def test_failed_booking_is_retry_safe_error_not_permanent_duplicate(): + deps = RecordingDeps(reserve_raises=RuntimeError("db gone mid-tx")) + result = run(process_external_usage_record(make_record(cost=0.01, idempotency_key="k-10"), deps.as_deps())) + assert result.status == "error" + assert "safe to retry" in (result.error or "") + assert result.spend is None + assert deps.spend_calls == [] From d4fb8e9a1b7ca9eb12db5a542763510e0e4ecc79 Mon Sep 17 00:00:00 2001 From: todayim <809634488@qq.com> Date: Thu, 6 Aug 2026 01:07:18 +0800 Subject: [PATCH 09/12] fix(proxy): include tag budgets in the idempotent booking transaction --- litellm/proxy/spend_tracking/usage_ingestion_endpoints.py | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/litellm/proxy/spend_tracking/usage_ingestion_endpoints.py b/litellm/proxy/spend_tracking/usage_ingestion_endpoints.py index 9a49268fdae..63e0f20a395 100644 --- a/litellm/proxy/spend_tracking/usage_ingestion_endpoints.py +++ b/litellm/proxy/spend_tracking/usage_ingestion_endpoints.py @@ -237,7 +237,9 @@ async def reserve_spend_log_atomic( if ProxyUpdateSpend.disable_spend_updates() is True: return "disabled" - payload: Final = prisma_client.jsonify_object(build_spend_log_payload(record, request_id, hashed_token, key, cost)) + spend_payload: Final = build_spend_log_payload(record, request_id, hashed_token, key, cost) + payload: Final = prisma_client.jsonify_object(spend_payload) + request_tags: Final = spend_payload.get("request_tags") from prisma.errors import UniqueViolationError writer: Final = proxy_logging_obj.db_spend_update_writer @@ -258,6 +260,7 @@ async def reserve_spend_log_atomic( response_cost=cost, team_id=key.team_id, user_id=key.user_id, prisma_client=shim ) await writer._update_org_db(response_cost=cost, org_id=key.organization_id, prisma_client=shim) + await writer._update_tag_db(response_cost=cost, request_tags=request_tags, prisma_client=shim) except UniqueViolationError: return "duplicate" return "reserved" From b32ce3d0b3b143262982b201c91633867e181703 Mon Sep 17 00:00:00 2001 From: todayim <809634488@qq.com> Date: Thu, 6 Aug 2026 02:16:30 +0800 Subject: [PATCH 10/12] chore(ci): retrigger after base eslint budget fix From b3e8e9f54d432a301f746321ab4fd0e7fb348b46 Mon Sep 17 00:00:00 2001 From: todayim <809634488@qq.com> Date: Thu, 6 Aug 2026 02:32:49 +0800 Subject: [PATCH 11/12] feat(spend): accept pre-hashed virtual keys in usage ingestion records --- .../usage_ingestion_endpoints.py | 20 ++++++++++++++++--- .../test_usage_ingestion_endpoints.py | 17 ++++++++++++++++ ui/litellm-dashboard/src/lib/http/schema.d.ts | 10 ++++++++-- 3 files changed, 42 insertions(+), 5 deletions(-) diff --git a/litellm/proxy/spend_tracking/usage_ingestion_endpoints.py b/litellm/proxy/spend_tracking/usage_ingestion_endpoints.py index 63e0f20a395..9c3e5029173 100644 --- a/litellm/proxy/spend_tracking/usage_ingestion_endpoints.py +++ b/litellm/proxy/spend_tracking/usage_ingestion_endpoints.py @@ -30,7 +30,14 @@ MAX_RECORDS_PER_REQUEST: Final = 1000 class ExternalUsageRecord(BaseModel): - api_key: str = Field(min_length=1, description="Raw virtual key (sk-...) to attribute usage to. Never logged.") + api_key: str | None = Field( + default=None, min_length=1, description="Raw virtual key (sk-...) to attribute usage to. Never logged." + ) + api_key_hash: str | None = Field( + default=None, + min_length=1, + description="SHA-256 hash of the virtual key. Use instead of api_key to avoid submitting raw keys.", + ) model: str = Field(min_length=1) prompt_tokens: int = Field(ge=0) completion_tokens: int = Field(ge=0) @@ -51,6 +58,12 @@ class ExternalUsageRecord(BaseModel): raise ValueError("end_time must not be before start_time") return self + @model_validator(mode="after") + def exactly_one_key_identifier(self) -> "ExternalUsageRecord": + if (self.api_key is None) == (self.api_key_hash is None): + raise ValueError("exactly one of api_key or api_key_hash is required") + return self + class UsageIngestRequest(BaseModel): records: tuple[ExternalUsageRecord, ...] = Field(min_length=1, max_length=MAX_RECORDS_PER_REQUEST) @@ -153,7 +166,7 @@ async def process_external_usage_record( record: ExternalUsageRecord, deps: UsageIngestionDeps ) -> UsageIngestRecordResult: request_id: Final = record.idempotency_key or deps.generate_request_id() - hashed_token: Final = hash_token(record.api_key) + hashed_token: Final = record.api_key_hash if record.api_key_hash is not None else hash_token(record.api_key or "") key: Final = await deps.lookup_key(hashed_token) if key is None: @@ -316,7 +329,8 @@ async def ingest_external_usage( dispatching directly to model gateways), so budgets and spend stay coherent in litellm as the single metering system. - Attribution (user/team/org) is derived from the given virtual key. Records accept an optional + Attribution (user/team/org) is derived from the given virtual key, submitted either raw + (api_key) or pre-hashed (api_key_hash) to keep raw keys out of request bodies. Records accept an optional idempotency_key, stored as the spend-log request_id: the reservation insert, counter updates and dedup are checked atomically at the database primary key, so overlapping retries are safe. The reservation row is always written (even when disable_spend_logs is set), because it is both the diff --git a/tests/test_litellm/proxy/spend_tracking/test_usage_ingestion_endpoints.py b/tests/test_litellm/proxy/spend_tracking/test_usage_ingestion_endpoints.py index b5ae287269a..8a477c0ea81 100644 --- a/tests/test_litellm/proxy/spend_tracking/test_usage_ingestion_endpoints.py +++ b/tests/test_litellm/proxy/spend_tracking/test_usage_ingestion_endpoints.py @@ -256,3 +256,20 @@ def test_failed_booking_is_retry_safe_error_not_permanent_duplicate(): assert "safe to retry" in (result.error or "") assert result.spend is None assert deps.spend_calls == [] + + +def test_key_hash_resolves_without_raw_key_in_body(): + deps = RecordingDeps() + record = make_record(cost=0.01, idempotency_key="k-11") + record = ExternalUsageRecord(**{**record.model_dump(), "api_key": None, "api_key_hash": hash_token(RAW_KEY)}) + result = run(process_external_usage_record(record, deps.as_deps())) + assert result.status == "recorded" + assert deps.looked_up_hashed == hash_token(RAW_KEY) + assert deps.reserve_calls[0]["hashed_token"] == hash_token(RAW_KEY) + + +def test_exactly_one_key_identifier_required(): + with pytest.raises(ValidationError): + make_record(api_key=None) + with pytest.raises(ValidationError): + make_record(api_key_hash=hash_token(RAW_KEY)) diff --git a/ui/litellm-dashboard/src/lib/http/schema.d.ts b/ui/litellm-dashboard/src/lib/http/schema.d.ts index 3647ef4c396..35bd7151383 100644 --- a/ui/litellm-dashboard/src/lib/http/schema.d.ts +++ b/ui/litellm-dashboard/src/lib/http/schema.d.ts @@ -12876,7 +12876,8 @@ export interface paths { * dispatching directly to model gateways), so budgets and spend stay coherent in litellm as the * single metering system. * - * Attribution (user/team/org) is derived from the given virtual key. Records accept an optional + * Attribution (user/team/org) is derived from the given virtual key, submitted either raw + * (api_key) or pre-hashed (api_key_hash) to keep raw keys out of request bodies. Records accept an optional * idempotency_key, stored as the spend-log request_id: the reservation insert, counter updates and * dedup are checked atomically at the database primary key, so overlapping retries are safe. The * reservation row is always written (even when disable_spend_logs is set), because it is both the @@ -24439,7 +24440,12 @@ export interface components { * Api Key * @description Raw virtual key (sk-...) to attribute usage to. Never logged. */ - api_key: string; + api_key?: string | null; + /** + * Api Key Hash + * @description SHA-256 hash of the virtual key. Use instead of api_key to avoid submitting raw keys. + */ + api_key_hash?: string | null; /** Completion Tokens */ completion_tokens: number; /** From 66269f08d2779395a8bdae5e8386566b23f71fca Mon Sep 17 00:00:00 2001 From: todayim <809634488@qq.com> Date: Thu, 6 Aug 2026 03:36:23 +0800 Subject: [PATCH 12/12] chore(ci): retrigger stale secret-scan check-run