From deec7f26d152887df989a0d1ed68b609e769682b Mon Sep 17 00:00:00 2001 From: atul naik Date: Sat, 3 Oct 2026 22:25:48 +0530 Subject: [PATCH] fix(usage): retain new throughput after upgrades Co-authored-by: Cursor --- litellm/proxy/db/daily_spend_bulk_upsert.py | 24 ++++++++++++++----- .../daily_spend_update_queue.py | 13 ++++++++-- litellm/repositories/daily_activity_sql.py | 6 +++-- .../test_daily_spend_update_queue.py | 20 +++++++++++++++- .../proxy/db/test_daily_spend_bulk_upsert.py | 21 ++++++++++++++-- 5 files changed, 71 insertions(+), 13 deletions(-) diff --git a/litellm/proxy/db/daily_spend_bulk_upsert.py b/litellm/proxy/db/daily_spend_bulk_upsert.py index b2bfb68813a..bfc49de3dd6 100644 --- a/litellm/proxy/db/daily_spend_bulk_upsert.py +++ b/litellm/proxy/db/daily_spend_bulk_upsert.py @@ -127,14 +127,17 @@ def _as_float(value: object) -> float: def _counter_total(column: str, group: Sequence[SpendRow]) -> int | None: values: Final = tuple(row.get(column) for row in group) - if column == "timed_completion_tokens" and any(value is None for value in values): + has_unknown_timed_tokens: Final = column == "timed_completion_tokens" and any( + row.get("timed_completion_tokens") is None and _as_int(row.get("timed_requests")) > 0 for row in group + ) + if has_unknown_timed_tokens: return None return sum(_as_int(value) for value in values) def _counter_value(column: str, transaction: SpendRow) -> int | None: value: Final = transaction.get(column) - if column == "timed_completion_tokens" and value is None: + if column == "timed_completion_tokens" and value is None and _as_int(transaction.get("timed_requests")) > 0: return None return _as_int(value) @@ -214,9 +217,18 @@ def build_bulk_upsert( + ", (NOW() AT TIME ZONE 'UTC'))" for row_index in range(len(batch)) ) - increments: Final = ", ".join( - f'"{column}" = {quoted_table}."{column}" + EXCLUDED."{column}"' - for column in (*_COUNTER_COLUMNS, *_SPEND_COLUMNS) + regular_increment_columns: Final = tuple( + column for column in (*_COUNTER_COLUMNS, *_SPEND_COLUMNS) if column != "timed_completion_tokens" + ) + regular_increments: Final = ", ".join( + f'"{column}" = {quoted_table}."{column}" + EXCLUDED."{column}"' for column in regular_increment_columns + ) + timed_completion_tokens_increment: Final = ( + f'"timed_completion_tokens" = CASE ' + f'WHEN ({quoted_table}."timed_requests" > 0 AND {quoted_table}."timed_completion_tokens" IS NULL) ' + 'OR (EXCLUDED."timed_requests" > 0 AND EXCLUDED."timed_completion_tokens" IS NULL) THEN NULL ' + f'ELSE COALESCE({quoted_table}."timed_completion_tokens", 0) ' + '+ COALESCE(EXCLUDED."timed_completion_tokens", 0) END' ) # request_id names one arbitrary contributing request, so an entry carrying none must # not blank out the one already recorded. @@ -229,7 +241,7 @@ def build_bulk_upsert( f'INSERT INTO {quoted_table} ({_quoted(columns)}, "updated_at")\n' f"VALUES {rows}\n" f"ON CONFLICT ({_quoted((table.entity_id_column, *_KEY_COLUMNS))}) DO UPDATE SET\n" - f" {increments}{request_id_update},\n" + f" {regular_increments}, {timed_completion_tokens_increment}{request_id_update},\n" f" \"updated_at\" = (NOW() AT TIME ZONE 'UTC')" ) return sql, tuple(value for key, transaction in batch for value in _row_params(table, key, transaction)) diff --git a/litellm/proxy/db/db_transaction_queue/daily_spend_update_queue.py b/litellm/proxy/db/db_transaction_queue/daily_spend_update_queue.py index 6bfaab45ad1..7dc3d42b818 100644 --- a/litellm/proxy/db/db_transaction_queue/daily_spend_update_queue.py +++ b/litellm/proxy/db/db_transaction_queue/daily_spend_update_queue.py @@ -22,8 +22,15 @@ def _daily_spend_updates( yield from update.items() -def _sum_timed_completion_tokens(left: int | None, right: int | None) -> int | None: - return left + right if left is not None and right is not None else None +def _sum_timed_completion_tokens( + left_tokens: int | None, + left_requests: int, + right_tokens: int | None, + right_requests: int, +) -> int | None: + if (left_tokens is None and left_requests > 0) or (right_tokens is None and right_requests > 0): + return None + return (left_tokens or 0) + (right_tokens or 0) def _merge_daily_spend_transactions( @@ -57,7 +64,9 @@ def _merge_daily_spend_transactions( "timed_requests": (existing.get("timed_requests", 0) or 0) + (payload.get("timed_requests", 0) or 0), "timed_completion_tokens": _sum_timed_completion_tokens( existing.get("timed_completion_tokens"), + existing.get("timed_requests", 0) or 0, payload.get("timed_completion_tokens"), + payload.get("timed_requests", 0) or 0, ), } diff --git a/litellm/repositories/daily_activity_sql.py b/litellm/repositories/daily_activity_sql.py index 3a253c5de3a..005372a598e 100644 --- a/litellm/repositories/daily_activity_sql.py +++ b/litellm/repositories/daily_activity_sql.py @@ -131,8 +131,10 @@ def _rollup_metric_select(table: DailyActivityTable) -> str: SUM(total_response_time_ms)::bigint AS total_response_time_ms, SUM(timed_requests)::bigint AS timed_requests, CASE - WHEN COUNT(timed_completion_tokens) = COUNT(*) THEN SUM(timed_completion_tokens)::bigint - ELSE NULL::bigint + WHEN COUNT(*) FILTER ( + WHERE timed_requests > 0 AND timed_completion_tokens IS NULL + ) > 0 THEN NULL::bigint + ELSE COALESCE(SUM(timed_completion_tokens), 0)::bigint END AS timed_completion_tokens""" diff --git a/tests/unit/proxy/db/db_transaction_queue/test_daily_spend_update_queue.py b/tests/unit/proxy/db/db_transaction_queue/test_daily_spend_update_queue.py index 5212c29a6c7..0b203b9703d 100644 --- a/tests/unit/proxy/db/db_transaction_queue/test_daily_spend_update_queue.py +++ b/tests/unit/proxy/db/db_transaction_queue/test_daily_spend_update_queue.py @@ -591,4 +591,22 @@ async def test_optional_metric_missing_from_an_older_payload_still_aggregates( assert updates[0][test_key]["autorouter_savings_spend"] == pytest.approx(0.25) assert updates[0][test_key]["total_response_time_ms"] == 900 assert updates[0][test_key]["timed_requests"] == 1 - assert updates[0][test_key]["timed_completion_tokens"] is None + assert updates[0][test_key]["timed_completion_tokens"] == 5 + + +def test_legacy_timed_request_keeps_timed_tokens_unknown(): + key = "user1_2023-01-01_key123_gpt-4o_openai" + base = { + "spend": 1.0, + "prompt_tokens": 10, + "completion_tokens": 5, + "api_requests": 1, + "successful_requests": 1, + "failed_requests": 0, + "timed_requests": 1, + } + updates = [{key: base}, {key: {**base, "timed_completion_tokens": 5}}] + + aggregated = DailySpendUpdateQueue.get_aggregated_daily_spend_update_transactions(updates) + + assert aggregated[key]["timed_completion_tokens"] is None diff --git a/tests/unit/proxy/db/test_daily_spend_bulk_upsert.py b/tests/unit/proxy/db/test_daily_spend_bulk_upsert.py index f63ef1acb64..81fb28ffe58 100644 --- a/tests/unit/proxy/db/test_daily_spend_bulk_upsert.py +++ b/tests/unit/proxy/db/test_daily_spend_bulk_upsert.py @@ -72,12 +72,21 @@ def test_null_and_empty_provider_merge_into_one_row(order): assert folded["api_requests"] == 4 -def test_unknown_timed_tokens_remain_unknown_when_rows_merge(): +def test_untimed_legacy_rows_do_not_hide_known_timed_tokens_when_rows_merge(): legacy = tag_txn() del legacy["timed_completion_tokens"] merged = merge_by_conflict_key(TAG_TABLE, (legacy, tag_txn(timed_completion_tokens=7))) + assert merged[0][1]["timed_completion_tokens"] == 7 + + +def test_legacy_timed_requests_keep_timed_tokens_unknown_when_rows_merge(): + legacy = tag_txn(timed_requests=1) + del legacy["timed_completion_tokens"] + + merged = merge_by_conflict_key(TAG_TABLE, (legacy, tag_txn(timed_completion_tokens=7))) + assert merged[0][1]["timed_completion_tokens"] is None @@ -125,7 +134,6 @@ def test_conflict_target_is_the_full_unique_constraint(): "failed_requests", "total_response_time_ms", "timed_requests", - "timed_completion_tokens", ], ) def test_counters_increment_rather_than_overwrite(column): @@ -135,6 +143,15 @@ def test_counters_increment_rather_than_overwrite(column): assert f'"{column}" = "LiteLLM_DailyTagSpend"."{column}" + EXCLUDED."{column}"' in sql +def test_timed_token_upsert_preserves_unknown_legacy_measurements(): + sql, _ = build_bulk_upsert(TAG_TABLE, merge_by_conflict_key(TAG_TABLE, (tag_txn(),))) + + assert '"timed_requests" > 0 AND "LiteLLM_DailyTagSpend"."timed_completion_tokens" IS NULL' in sql + assert 'EXCLUDED."timed_requests" > 0 AND EXCLUDED."timed_completion_tokens" IS NULL' in sql + assert 'COALESCE("LiteLLM_DailyTagSpend"."timed_completion_tokens", 0)' in sql + assert '+ COALESCE(EXCLUDED."timed_completion_tokens", 0)' in sql + + def test_request_id_is_preserved_when_a_later_batch_carries_none(): sql, params = build_bulk_upsert(TAG_TABLE, merge_by_conflict_key(TAG_TABLE, (tag_txn(request_id=None),)))