fix(usage): retain new throughput after upgrades

Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
atul naik 2026-10-03 22:25:48 +05:30
parent ed862773d1
commit deec7f26d1
5 changed files with 71 additions and 13 deletions

View file

@ -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))

View file

@ -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,
),
}

View file

@ -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"""

View file

@ -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

View file

@ -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),)))