mirror of
https://github.com/BerriAI/litellm.git
synced 2026-09-21 00:21:49 +00:00
fix(proxy): drop daily spend batches that cannot be re-sent safely instead of requeueing them
This commit is contained in:
parent
3edbf60e9c
commit
b3cf45e9f2
4 changed files with 112 additions and 3 deletions
|
|
@ -31,6 +31,7 @@ from litellm.constants import (
|
|||
from litellm.litellm_core_utils.litellm_logging import coerce_model_access_groups
|
||||
from litellm.litellm_core_utils.safe_json_loads import safe_json_loads
|
||||
from litellm.proxy._types import (
|
||||
DB_CONNECTION_ERROR_TYPES,
|
||||
DB_RETRY_SAFE_ERROR_TYPES,
|
||||
BaseDailySpendTransaction,
|
||||
DailyAgentSpendTransaction,
|
||||
|
|
@ -64,6 +65,7 @@ from litellm.proxy.db.db_transaction_queue.window_spend_update_queue import (
|
|||
WindowSpendTransaction,
|
||||
WindowSpendUpdateQueue,
|
||||
)
|
||||
from litellm.proxy.db.exception_handler import PrismaDBExceptionHandler
|
||||
from litellm.proxy.route_llm_request import ROUTE_ENDPOINT_MAPPING
|
||||
from litellm.proxy.spend_tracking.compression_savings import (
|
||||
extract_compression_saved_tokens,
|
||||
|
|
@ -157,6 +159,16 @@ class _DailySpendCommit(Protocol[_DailySpendTransactionT]):
|
|||
) -> None: ...
|
||||
|
||||
|
||||
_DATA_REJECTED_SQLSTATE_CLASSES: Final = frozenset({"22", "23"})
|
||||
|
||||
|
||||
def _daily_spend_commit_failure_is_requeue_safe(e: Exception) -> bool:
|
||||
if isinstance(e, DB_CONNECTION_ERROR_TYPES):
|
||||
return isinstance(e, DB_RETRY_SAFE_ERROR_TYPES)
|
||||
sqlstate: Final = PrismaDBExceptionHandler.postgres_sqlstate(e)
|
||||
return sqlstate is None or sqlstate[:2] not in _DATA_REJECTED_SQLSTATE_CLASSES
|
||||
|
||||
|
||||
def _timed_request_duration_ms(
|
||||
payload: dict | SpendLogsPayload,
|
||||
request_status: Literal["success", "failure"],
|
||||
|
|
@ -1319,7 +1331,17 @@ class DBSpendUpdateWriter:
|
|||
proxy_logging_obj=proxy_logging_obj,
|
||||
daily_spend_transactions=cast(dict[str, _DailySpendTransactionT], transactions),
|
||||
)
|
||||
except Exception as e: # noqa: BLE001 # the uncommitted rows go back on the queue; the other tables must still flush
|
||||
except Exception as e: # noqa: BLE001 # whatever failed here, the other tables must still flush
|
||||
if not _daily_spend_commit_failure_is_requeue_safe(e):
|
||||
spend_log_error(
|
||||
"Spend tracking - dropped %d daily %s spend rows: the failed commit may have applied "
|
||||
"or the database refused the data, so re-sending it is not safe. Error: %s",
|
||||
len(transactions),
|
||||
entity_type,
|
||||
str(e),
|
||||
exc=e,
|
||||
)
|
||||
return
|
||||
spend_log_error(
|
||||
"Spend tracking - failed to commit daily %s spend updates. "
|
||||
"Re-queued %d rows for retry on next tick. Error: %s",
|
||||
|
|
|
|||
|
|
@ -1,6 +1,8 @@
|
|||
from collections.abc import Awaitable, Callable, Iterator
|
||||
from typing import Any, Final, TypeVar
|
||||
|
||||
from pydantic import TypeAdapter, ValidationError
|
||||
|
||||
from litellm._logging import verbose_proxy_logger
|
||||
from litellm.proxy._types import (
|
||||
DB_CONNECTION_ERROR_TYPES,
|
||||
|
|
@ -17,6 +19,8 @@ _TRANSIENT_DB_UNAVAILABLE_MESSAGE: Final = (
|
|||
"Service Unavailable, the authentication database is temporarily unreachable. Please retry shortly."
|
||||
)
|
||||
|
||||
_DATABASE_ERROR_META: Final = TypeAdapter(dict[str, object])
|
||||
|
||||
|
||||
def _exception_chain(e: BaseException) -> Iterator[BaseException]:
|
||||
current = e # rebind-ok: advances one link per iteration of the bounded walk
|
||||
|
|
@ -221,6 +225,20 @@ class PrismaDBExceptionHandler:
|
|||
or "write conflict or a deadlock" in error_message
|
||||
)
|
||||
|
||||
@staticmethod
|
||||
def postgres_sqlstate(e: Exception) -> str | None:
|
||||
"""The SQLSTATE Postgres attached to a failed statement, as prisma surfaces it, or None."""
|
||||
import prisma
|
||||
|
||||
if not isinstance(e, _exception_types(prisma.errors.DataError)):
|
||||
return None
|
||||
try:
|
||||
meta: Final = _DATABASE_ERROR_META.validate_python(getattr(e, "meta", None))
|
||||
except ValidationError:
|
||||
return None
|
||||
code: Final = meta.get("code")
|
||||
return code if isinstance(code, str) else None
|
||||
|
||||
@staticmethod
|
||||
def is_read_only_transaction_error(e: Exception) -> bool:
|
||||
"""True iff ``e`` is Postgres SQLSTATE 25006 surfaced through prisma: the
|
||||
|
|
|
|||
|
|
@ -11,7 +11,9 @@ from types import SimpleNamespace
|
|||
from typing import Final
|
||||
from unittest.mock import AsyncMock, MagicMock, call, patch
|
||||
|
||||
import httpx
|
||||
import pytest
|
||||
from prisma.errors import RawQueryError
|
||||
from redis.exceptions import DataError
|
||||
|
||||
import litellm
|
||||
|
|
@ -2812,14 +2814,15 @@ async def test_failed_window_spend_commit_requeues_the_increments_and_continues_
|
|||
class _DailySpendFakeDB(_WindowSpendFakeDB):
|
||||
"""Records the daily rollup upserts it is handed and fails the ones aimed at one table."""
|
||||
|
||||
def __init__(self, failing_table: str | None) -> None:
|
||||
def __init__(self, failing_table: str | None, failure: Exception | None = None) -> None:
|
||||
super().__init__()
|
||||
self.failing_table = failing_table
|
||||
self.failure = failure
|
||||
self.execute_raw_calls: list[Statement] = []
|
||||
|
||||
async def execute_raw(self, query: str, *args: object) -> int:
|
||||
if self.failing_table is not None and self.failing_table in query:
|
||||
raise Exception("connection reset")
|
||||
raise self.failure if self.failure is not None else Exception("connection reset")
|
||||
self.execute_raw_calls.append((query, args))
|
||||
return len(args)
|
||||
|
||||
|
|
@ -2828,6 +2831,50 @@ def _daily_upserts(db: _DailySpendFakeDB, table: str) -> list[Statement]:
|
|||
return [statement for statement in db.execute_raw_calls if table in statement[0]]
|
||||
|
||||
|
||||
def _postgres_rejection(sqlstate: str) -> RawQueryError:
|
||||
return RawQueryError(
|
||||
data={"user_facing_error": {"error_code": "P2010", "meta": {"code": sqlstate, "message": "db error"}}}
|
||||
)
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
("failure", "lands_on_the_next_tick"),
|
||||
[
|
||||
pytest.param(httpx.ReadTimeout("no reply"), False, id="reply lost after the statement was sent"),
|
||||
pytest.param(httpx.ConnectError("refused"), True, id="statement never reached the database"),
|
||||
pytest.param(_postgres_rejection("22021"), False, id="postgres refused the data itself"),
|
||||
pytest.param(_postgres_rejection("23502"), False, id="postgres refused a constraint violation"),
|
||||
pytest.param(_postgres_rejection("42P01"), True, id="table missing"),
|
||||
pytest.param(_postgres_rejection("57014"), True, id="statement cancelled"),
|
||||
],
|
||||
)
|
||||
@pytest.mark.asyncio
|
||||
async def test_failed_daily_spend_commit_is_requeued_only_when_the_rows_are_provably_uncommitted(
|
||||
failure: Exception, lands_on_the_next_tick: bool
|
||||
):
|
||||
"""A lost reply means the statement may already have applied, and re-sending it stacks a
|
||||
second increment into the same transaction (LIT-4823); a row Postgres refuses would fail
|
||||
every tick forever. Both are dropped loudly. Every other failure left nothing committed,
|
||||
so its rows go back on the queue and land on the next tick."""
|
||||
db_writer = DBSpendUpdateWriter()
|
||||
await db_writer.daily_spend_update_queue.add_update({"user-key": _daily_txn(user_id="user-1")})
|
||||
db = _DailySpendFakeDB(failing_table="LiteLLM_DailyUserSpend", failure=failure)
|
||||
db_writer._flush_tool_discovery_queue = AsyncMock()
|
||||
proxy_logging_obj = MagicMock()
|
||||
proxy_logging_obj.failure_handler = AsyncMock()
|
||||
|
||||
await db_writer._commit_spend_updates_to_db_without_redis_buffer(
|
||||
prisma_client=_WindowSpendFakePrisma(db), n_retry_times=0, proxy_logging_obj=proxy_logging_obj
|
||||
)
|
||||
db.failing_table = None
|
||||
await db_writer._commit_spend_updates_to_db_without_redis_buffer(
|
||||
prisma_client=_WindowSpendFakePrisma(db), n_retry_times=0, proxy_logging_obj=proxy_logging_obj
|
||||
)
|
||||
|
||||
assert len(_daily_upserts(db, "LiteLLM_DailyUserSpend")) == (1 if lands_on_the_next_tick else 0)
|
||||
assert db_writer.daily_spend_update_queue.update_queue.empty()
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_failed_daily_spend_commit_requeues_the_rows_and_flushes_the_other_tables():
|
||||
"""With the Redis buffer off, a daily batch that failed to commit was discarded along
|
||||
|
|
|
|||
|
|
@ -665,6 +665,28 @@ def test_is_deadlock_error_excludes_non_deadlocks(error):
|
|||
assert PrismaDBExceptionHandler.is_deadlock_error(error) is False
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
("error", "sqlstate"),
|
||||
[
|
||||
(
|
||||
RawQueryError(
|
||||
data={"user_facing_error": {"error_code": "P2010", "meta": {"code": "22021", "message": "m"}}}
|
||||
),
|
||||
"22021",
|
||||
),
|
||||
(RawQueryError(data={"user_facing_error": {"error_code": "P2010", "meta": {"message": "m"}}}), None),
|
||||
(RawQueryError(data={"user_facing_error": {"error_code": "P2010", "meta": {"code": 42, "message": "m"}}}), None),
|
||||
(prisma_errors.DataError(data={"user_facing_error": {"meta": None}}), None),
|
||||
(PrismaError("db error"), None),
|
||||
(httpx.ReadTimeout("no reply"), None),
|
||||
],
|
||||
)
|
||||
def test_postgres_sqlstate_reads_the_code_prisma_attached_to_the_failed_statement(error, sqlstate):
|
||||
"""Only a prisma data error carrying Postgres's own error code yields a SQLSTATE; a
|
||||
codeless or malformed payload, an engine-level error, and a transport error yield None."""
|
||||
assert PrismaDBExceptionHandler.postgres_sqlstate(error) == sqlstate
|
||||
|
||||
|
||||
READ_ONLY_CONNECTOR_ERROR: Final = (
|
||||
"Error occurred during query execution:\nConnectorError(ConnectorError { user_facing_error: None, "
|
||||
'kind: QueryError(PostgresError { code: "25006", message: "cannot execute UPDATE in a read-only transaction", '
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue