mirror of
https://github.com/BerriAI/litellm.git
synced 2026-09-28 01:32:17 +00:00
Merge 4864658081 into 04b1b077de
This commit is contained in:
commit
2f16544c1c
8 changed files with 406 additions and 13 deletions
|
|
@ -2950,6 +2950,14 @@ class ConfigGeneralSettings(LiteLLMPydanticObjectBase):
|
|||
"Set this well above health_check_interval because /health and the UI read the latest row per model."
|
||||
),
|
||||
)
|
||||
maximum_daily_tag_spend_retention_period: str | None = Field(
|
||||
None,
|
||||
description=(
|
||||
"Maximum retention period for per-day tag spend aggregate rows (e.g., '90d'). Rows whose day is older "
|
||||
"than this are deleted by the spend log cleanup job, on that job's schedule. Unset means rows are never "
|
||||
"deleted. Only historical tag usage analytics are affected; tag budgets read the lifetime counter."
|
||||
),
|
||||
)
|
||||
use_spend_logs_partitioning: bool | None = Field(
|
||||
None,
|
||||
description="If True and LiteLLM_SpendLogs has been converted to a range-partitioned table (db_scripts/partition_spend_logs.sql), retention cleanup drops expired partitions instead of deleting rows, and pre-creates upcoming partitions. Default is False.",
|
||||
|
|
|
|||
|
|
@ -32,6 +32,17 @@ from litellm.proxy.utils import PrismaClient
|
|||
|
||||
StopReason: TypeAlias = Literal["exhausted", "budget_exhausted", "batch_cap_reached", "aborted"]
|
||||
|
||||
Cutoff: TypeAlias = datetime | str
|
||||
"""Rows strictly older than this are expired: a timestamp, or an ISO calendar day for tables keyed by day"""
|
||||
|
||||
|
||||
def _cutoff_cast(cutoff: Cutoff) -> str:
|
||||
return "timestamptz" if isinstance(cutoff, datetime) else "text"
|
||||
|
||||
|
||||
def _cutoff_text(cutoff: Cutoff) -> str:
|
||||
return cutoff.isoformat() if isinstance(cutoff, datetime) else cutoff
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True)
|
||||
class TableCleanupResult:
|
||||
|
|
@ -277,7 +288,7 @@ class SpendLogCleanup:
|
|||
return remaining
|
||||
|
||||
async def _execute_delete_batch(
|
||||
self, prisma_client: PrismaClient, delete_sql: str, cutoff_date: datetime, deadline: float
|
||||
self, prisma_client: PrismaClient, delete_sql: str, cutoff_date: Cutoff, deadline: float
|
||||
) -> int | None:
|
||||
"""
|
||||
Run one delete batch under a Postgres statement and lock timeout.
|
||||
|
|
@ -300,7 +311,7 @@ class SpendLogCleanup:
|
|||
return deleted_result if isinstance(deleted_result, int) else None
|
||||
|
||||
async def _count_remaining(
|
||||
self, prisma_client: PrismaClient, cutoff_date: datetime, table_name: str, time_column: str, deadline: float
|
||||
self, prisma_client: PrismaClient, cutoff_date: Cutoff, table_name: str, time_column: str, deadline: float
|
||||
) -> int | None:
|
||||
"""
|
||||
Count expired rows still outstanding, stopping at a cap.
|
||||
|
|
@ -313,7 +324,7 @@ class SpendLogCleanup:
|
|||
count_sql: Final = f"""
|
||||
SELECT count(*)::int AS remaining FROM (
|
||||
SELECT 1 FROM "{table_name}"
|
||||
WHERE "{time_column}" < $1::timestamptz
|
||||
WHERE "{time_column}" < $1::{_cutoff_cast(cutoff_date)}
|
||||
LIMIT $2
|
||||
) capped
|
||||
"""
|
||||
|
|
@ -331,7 +342,7 @@ class SpendLogCleanup:
|
|||
async def _delete_old_rows_batched(
|
||||
self,
|
||||
prisma_client: PrismaClient,
|
||||
cutoff_date: datetime,
|
||||
cutoff_date: Cutoff,
|
||||
table_name: str,
|
||||
key_columns: tuple[str, ...],
|
||||
time_column: str,
|
||||
|
|
@ -349,7 +360,7 @@ class SpendLogCleanup:
|
|||
DELETE FROM "{table_name}"
|
||||
WHERE ({key_list}) IN (
|
||||
SELECT {key_list} FROM "{table_name}"
|
||||
WHERE "{time_column}" < $1::timestamptz
|
||||
WHERE "{time_column}" < $1::{_cutoff_cast(cutoff_date)}
|
||||
LIMIT $2
|
||||
)
|
||||
"""
|
||||
|
|
@ -405,7 +416,7 @@ class SpendLogCleanup:
|
|||
run_count,
|
||||
consecutive_failures,
|
||||
self.batch_size,
|
||||
cutoff_date.isoformat(),
|
||||
_cutoff_text(cutoff_date),
|
||||
total_deleted,
|
||||
type(batch_exc).__name__,
|
||||
batch_exc,
|
||||
|
|
@ -453,7 +464,7 @@ class SpendLogCleanup:
|
|||
async def _finish_table(
|
||||
self,
|
||||
prisma_client: PrismaClient,
|
||||
cutoff_date: datetime,
|
||||
cutoff_date: Cutoff,
|
||||
table_name: str,
|
||||
time_column: str,
|
||||
rows_deleted: int,
|
||||
|
|
@ -540,6 +551,18 @@ class SpendLogCleanup:
|
|||
deadline=deadline,
|
||||
)
|
||||
|
||||
async def _delete_old_daily_tag_spend_rows(
|
||||
self, prisma_client: PrismaClient, cutoff_day: str, deadline: float
|
||||
) -> TableCleanupResult:
|
||||
return await self._delete_old_rows_batched(
|
||||
prisma_client,
|
||||
cutoff_day,
|
||||
table_name="LiteLLM_DailyTagSpend",
|
||||
key_columns=("id",),
|
||||
time_column="date",
|
||||
deadline=deadline,
|
||||
)
|
||||
|
||||
async def _clean_spend_log_tables(
|
||||
self, prisma_client: PrismaClient, deadline: float
|
||||
) -> tuple[TableCleanupResult, ...]:
|
||||
|
|
@ -621,6 +644,18 @@ class SpendLogCleanup:
|
|||
)
|
||||
return (health_checks_result,)
|
||||
|
||||
async def _clean_daily_tag_spend(
|
||||
self, prisma_client: PrismaClient, retention_seconds: int, deadline: float
|
||||
) -> tuple[TableCleanupResult, ...]:
|
||||
"""
|
||||
Prune per-day tag spend rows whose ISO day sorts before the horizon day; the horizon day itself is kept.
|
||||
"""
|
||||
horizon: Final = datetime.now(timezone.utc) - timedelta(seconds=float(retention_seconds))
|
||||
cutoff_day: Final = horizon.date().isoformat()
|
||||
result: Final = await self._delete_old_daily_tag_spend_rows(prisma_client, cutoff_day, deadline)
|
||||
verbose_proxy_logger.info("Deleted %s expired daily tag spend rows", result.rows_deleted)
|
||||
return (result,)
|
||||
|
||||
@staticmethod
|
||||
def _run_outcome(results: tuple[TableCleanupResult, ...]) -> RunOutcome:
|
||||
"""
|
||||
|
|
@ -657,10 +692,14 @@ class SpendLogCleanup:
|
|||
"maximum_autorouter_session_retention_period"
|
||||
)
|
||||
health_check_retention_seconds: Final = self._retention_seconds_for("maximum_health_check_retention_period")
|
||||
daily_tag_spend_retention_seconds: Final = self._retention_seconds_for(
|
||||
"maximum_daily_tag_spend_retention_period"
|
||||
)
|
||||
if (
|
||||
not delete_spend_logs
|
||||
and autorouter_retention_seconds is None
|
||||
and health_check_retention_seconds is None
|
||||
and daily_tag_spend_retention_seconds is None
|
||||
):
|
||||
SpendLogCleanupMetrics.record_run("skipped_disabled")
|
||||
return
|
||||
|
|
@ -692,6 +731,7 @@ class SpendLogCleanup:
|
|||
int(delete_spend_logs and self.retention_seconds is not None)
|
||||
+ int(autorouter_retention_seconds is not None)
|
||||
+ int(health_check_retention_seconds is not None)
|
||||
+ int(daily_tag_spend_retention_seconds is not None)
|
||||
)
|
||||
|
||||
spend_log_results: Final = (
|
||||
|
|
@ -702,8 +742,13 @@ class SpendLogCleanup:
|
|||
if delete_spend_logs and self.retention_seconds is not None
|
||||
else ()
|
||||
)
|
||||
remaining_groups_after_spend_logs: Final = int(autorouter_retention_seconds is not None) + int(
|
||||
health_check_retention_seconds is not None
|
||||
remaining_groups_after_spend_logs: Final = (
|
||||
int(autorouter_retention_seconds is not None)
|
||||
+ int(health_check_retention_seconds is not None)
|
||||
+ int(daily_tag_spend_retention_seconds is not None)
|
||||
)
|
||||
remaining_groups_after_sessions: Final = int(health_check_retention_seconds is not None) + int(
|
||||
daily_tag_spend_retention_seconds is not None
|
||||
)
|
||||
session_results: Final = (
|
||||
await self._clean_session_rollup(
|
||||
|
|
@ -718,14 +763,19 @@ class SpendLogCleanup:
|
|||
await self._clean_health_checks(
|
||||
prisma_client,
|
||||
health_check_retention_seconds,
|
||||
deadline,
|
||||
self._group_deadline(deadline, remaining_groups_after_sessions),
|
||||
)
|
||||
if health_check_retention_seconds is not None
|
||||
else ()
|
||||
)
|
||||
daily_tag_spend_results: Final = (
|
||||
await self._clean_daily_tag_spend(prisma_client, daily_tag_spend_retention_seconds, deadline)
|
||||
if daily_tag_spend_retention_seconds is not None
|
||||
else ()
|
||||
)
|
||||
|
||||
SpendLogCleanupMetrics.record_run(
|
||||
self._run_outcome(spend_log_results + session_results + health_check_results)
|
||||
self._run_outcome(spend_log_results + session_results + health_check_results + daily_tag_spend_results)
|
||||
)
|
||||
|
||||
except asyncio.CancelledError:
|
||||
|
|
|
|||
|
|
@ -7432,7 +7432,13 @@ class ProxyConfig:
|
|||
retention_period: Final = general_settings.get("maximum_spend_logs_retention_period")
|
||||
autorouter_retention: Final = general_settings.get("maximum_autorouter_session_retention_period")
|
||||
health_check_retention: Final = general_settings.get("maximum_health_check_retention_period")
|
||||
if retention_period is not None or autorouter_retention is not None or health_check_retention is not None:
|
||||
daily_tag_spend_retention: Final = general_settings.get("maximum_daily_tag_spend_retention_period")
|
||||
if (
|
||||
retention_period is not None
|
||||
or autorouter_retention is not None
|
||||
or health_check_retention is not None
|
||||
or daily_tag_spend_retention is not None
|
||||
):
|
||||
from litellm.proxy.db.db_transaction_queue.spend_log_cleanup import (
|
||||
SpendLogCleanup,
|
||||
)
|
||||
|
|
@ -7510,6 +7516,7 @@ class ProxyConfig:
|
|||
"maximum_spend_logs_retention_period",
|
||||
"maximum_autorouter_session_retention_period",
|
||||
"maximum_health_check_retention_period",
|
||||
"maximum_daily_tag_spend_retention_period",
|
||||
)
|
||||
)
|
||||
|
||||
|
|
@ -7623,7 +7630,10 @@ class ProxyConfig:
|
|||
db_values: Mapping[str, SettingsJsonValue],
|
||||
previous_retention_values: tuple[SettingsJsonValue | None, ...],
|
||||
) -> None:
|
||||
if previous_retention_values != self._resolved_retention_values():
|
||||
resolved: Final = self._resolved_retention_values()
|
||||
wants_job: Final = any(value is not None for value in resolved)
|
||||
has_job: Final = scheduler is not None and scheduler.get_job("spend_log_cleanup_job") is not None
|
||||
if previous_retention_values != resolved or wants_job != has_job:
|
||||
await self._reschedule_spend_log_cleanup_job()
|
||||
|
||||
async def _apply_ssrf_settings(self, db_values: Mapping[str, SettingsJsonValue]) -> None:
|
||||
|
|
@ -10484,6 +10494,7 @@ class ProxyStartupEvent:
|
|||
general_settings.get("maximum_spend_logs_retention_period") is not None
|
||||
or general_settings.get("maximum_autorouter_session_retention_period") is not None
|
||||
or general_settings.get("maximum_health_check_retention_period") is not None
|
||||
or general_settings.get("maximum_daily_tag_spend_retention_period") is not None
|
||||
):
|
||||
spend_log_cleanup: Final = SpendLogCleanup()
|
||||
cleanup_cron: Final = general_settings.get("maximum_spend_logs_cleanup_cron")
|
||||
|
|
@ -17697,6 +17708,7 @@ _GENERAL_SETTINGS_CONFIG_LIST_FIELD_TYPES: Final[Mapping[str, str]] = MappingPro
|
|||
"store_prompts_in_spend_logs": "Boolean",
|
||||
"maximum_spend_logs_retention_period": "String",
|
||||
"maximum_health_check_retention_period": "String",
|
||||
"maximum_daily_tag_spend_retention_period": "String",
|
||||
"maximum_spend_logs_cleanup_batch_size": "Integer",
|
||||
"maximum_spend_logs_cleanup_max_batches": "Integer",
|
||||
"maximum_spend_logs_cleanup_run_budget": "String",
|
||||
|
|
|
|||
247
tests/integration/spend/test_daily_tag_spend_retention.py
Normal file
247
tests/integration/spend/test_daily_tag_spend_retention.py
Normal file
|
|
@ -0,0 +1,247 @@
|
|||
import json
|
||||
import os
|
||||
import signal
|
||||
import uuid
|
||||
from datetime import datetime, timedelta, timezone
|
||||
from pathlib import Path
|
||||
from typing import Final
|
||||
|
||||
import psutil
|
||||
import psycopg
|
||||
import pytest
|
||||
import yaml
|
||||
from pydantic import JsonValue, TypeAdapter
|
||||
|
||||
from tests.integration._support.client import Gateway, eventually, string_value
|
||||
from tests.integration._support.database import read_rows
|
||||
from tests.integration._support.process import OwnedProxy, owned_proxy, owned_proxy_process
|
||||
|
||||
CLEANUP_EVERY_MINUTE: Final = "* * * * *"
|
||||
RETENTION_SETTING: Final = "maximum_daily_tag_spend_retention_period"
|
||||
_MAPPING: Final = TypeAdapter(dict[str, JsonValue])
|
||||
_SETTINGS: Final = TypeAdapter(list[dict[str, JsonValue]])
|
||||
|
||||
|
||||
def _day(days_ago: int) -> str:
|
||||
return (datetime.now(timezone.utc) - timedelta(days=days_ago)).strftime("%Y-%m-%d")
|
||||
|
||||
|
||||
def _seed_daily_tag_spend(tag: str, days: tuple[str, ...]) -> None:
|
||||
with psycopg.connect(os.environ["DATABASE_URL"], autocommit=True) as connection:
|
||||
for day in days:
|
||||
connection.execute(
|
||||
'INSERT INTO "LiteLLM_DailyTagSpend" (id, tag, date, api_key, model, spend, updated_at) '
|
||||
"VALUES (%s, %s, %s, %s, %s, 1.0, now())",
|
||||
(uuid.uuid4().hex, tag, day, f"integration-{tag}", "gpt-4o-mini"),
|
||||
)
|
||||
|
||||
|
||||
def _seed_old_spend_log(request_id: str, days_ago: int) -> None:
|
||||
with psycopg.connect(os.environ["DATABASE_URL"], autocommit=True) as connection:
|
||||
connection.execute(
|
||||
'INSERT INTO "LiteLLM_SpendLogs" (request_id, call_type, api_key, spend, "startTime", "endTime") '
|
||||
"VALUES (%s, 'acompletion', %s, 0, now() - make_interval(days => %s), now() - make_interval(days => %s))",
|
||||
(request_id, f"integration-{request_id}", str(days_ago), str(days_ago)),
|
||||
)
|
||||
|
||||
|
||||
def _delete_daily_tag_spend(tag: str) -> None:
|
||||
with psycopg.connect(os.environ["DATABASE_URL"], autocommit=True) as connection:
|
||||
connection.execute('DELETE FROM "LiteLLM_DailyTagSpend" WHERE tag = %s', (tag,))
|
||||
|
||||
|
||||
def _remaining_days(tag: str) -> tuple[str, ...]:
|
||||
rows: Final = read_rows('SELECT date FROM "LiteLLM_DailyTagSpend" WHERE tag = %s ORDER BY date', (tag,))
|
||||
return tuple(str(row["date"]) for row in rows)
|
||||
|
||||
|
||||
def _spend_log_present(request_id: str) -> bool:
|
||||
return bool(read_rows('SELECT request_id FROM "LiteLLM_SpendLogs" WHERE request_id = %s', (request_id,)))
|
||||
|
||||
|
||||
def _stored_retention_setting() -> JsonValue:
|
||||
rows: Final = read_rows(
|
||||
'SELECT param_value -> %s AS value FROM "LiteLLM_Config" WHERE param_name = %s',
|
||||
(RETENTION_SETTING, "general_settings"),
|
||||
)
|
||||
return rows[0]["value"] if rows else None
|
||||
|
||||
|
||||
def _store_retention_setting(value: JsonValue) -> None:
|
||||
with psycopg.connect(os.environ["DATABASE_URL"], autocommit=True) as connection:
|
||||
if value is None:
|
||||
connection.execute(
|
||||
'UPDATE "LiteLLM_Config" SET param_value = param_value - %s WHERE param_name = %s',
|
||||
(RETENTION_SETTING, "general_settings"),
|
||||
)
|
||||
return
|
||||
connection.execute(
|
||||
'UPDATE "LiteLLM_Config" SET param_value = jsonb_set(param_value, ARRAY[%s], %s::jsonb) '
|
||||
"WHERE param_name = %s",
|
||||
(RETENTION_SETTING, json.dumps(value), "general_settings"),
|
||||
)
|
||||
|
||||
|
||||
def _listening_workers(owned: OwnedProxy) -> tuple[psutil.Process, ...]:
|
||||
port: Final = owned.gateway.client.base_url.port
|
||||
return tuple(
|
||||
child
|
||||
for child in psutil.Process(owned.process.pid).children(recursive=True)
|
||||
if any(conn.status == psutil.CONN_LISTEN and conn.laddr.port == port for conn in child.net_connections("inet"))
|
||||
)
|
||||
|
||||
|
||||
def _listed_retention_value(gateway: Gateway) -> JsonValue:
|
||||
listed: Final = _SETTINGS.validate_json(
|
||||
gateway.request("GET", "/config/list", params={"config_type": "general_settings"}).content
|
||||
)
|
||||
matching: Final = tuple(entry for entry in listed if entry["field_name"] == RETENTION_SETTING)
|
||||
return matching[0]["field_value"] if matching else "not listed"
|
||||
|
||||
|
||||
def _completion_id(gateway: Gateway, model: str) -> str:
|
||||
return string_value(gateway.chat(model, text=f"retention audit {uuid.uuid4().hex}")["id"])
|
||||
|
||||
|
||||
def _cleanup_config(tmp_path: Path, retention: dict[str, JsonValue]) -> Path:
|
||||
base: Final = _MAPPING.validate_python(yaml.safe_load(Path("tests/integration/proxy_config.yaml").read_text()))
|
||||
config: Final = {
|
||||
**base,
|
||||
"general_settings": {
|
||||
**_MAPPING.validate_python(base["general_settings"]),
|
||||
**retention,
|
||||
"maximum_spend_logs_cleanup_cron": CLEANUP_EVERY_MINUTE,
|
||||
"scheduled_job_stagger": {"enabled": False},
|
||||
},
|
||||
}
|
||||
path: Final = tmp_path / "retention.yaml"
|
||||
path.write_text(yaml.safe_dump(config))
|
||||
return path
|
||||
|
||||
|
||||
@pytest.mark.covers("spend.daily_tag_spend.retention_prunes_rows_older_than_the_period_and_keeps_the_rest")
|
||||
def test_daily_tag_spend_retention_prunes_only_rows_older_than_the_period(gateway: Gateway, tmp_path: Path) -> None:
|
||||
tag: Final = f"integration-retention-{uuid.uuid4().hex}"
|
||||
expired, on_the_cutoff, today = _day(200), _day(30), _day(0)
|
||||
_seed_daily_tag_spend(tag, (expired, on_the_cutoff, today))
|
||||
try:
|
||||
config: Final = _cleanup_config(tmp_path, {"maximum_daily_tag_spend_retention_period": "30d"})
|
||||
with owned_proxy(gateway, tmp_path, {}, config=config):
|
||||
remaining: Final = eventually(
|
||||
lambda: _remaining_days(tag),
|
||||
lambda days: expired not in days,
|
||||
seconds=150,
|
||||
)
|
||||
assert remaining == (on_the_cutoff, today), remaining
|
||||
finally:
|
||||
_delete_daily_tag_spend(tag)
|
||||
|
||||
|
||||
@pytest.mark.covers("spend.daily_tag_spend.runtime_config_update_enables_cleanup_on_a_multi_worker_proxy")
|
||||
def test_config_update_turns_on_daily_tag_spend_cleanup_without_a_restart(gateway: Gateway, tmp_path: Path) -> None:
|
||||
tag: Final = f"integration-retention-{uuid.uuid4().hex}"
|
||||
expired, yesterday_of_cutoff, on_the_cutoff, today = _day(200), _day(31), _day(30), _day(0)
|
||||
_seed_daily_tag_spend(tag, (expired, yesterday_of_cutoff, on_the_cutoff, today))
|
||||
previously_stored: Final = _stored_retention_setting()
|
||||
_store_retention_setting(None)
|
||||
try:
|
||||
config: Final = _cleanup_config(tmp_path, {})
|
||||
with owned_proxy(gateway, tmp_path, {}, config=config, workers=2) as owned, owned.scenario() as scenario:
|
||||
model: Final = scenario.model()
|
||||
assert _listed_retention_value(owned) is None
|
||||
owned.post("/config/update", {"general_settings": {RETENTION_SETTING: "30d"}})
|
||||
assert _listed_retention_value(owned) == "30d"
|
||||
remaining: Final = eventually(
|
||||
lambda: _remaining_days(tag),
|
||||
lambda days: yesterday_of_cutoff not in days,
|
||||
seconds=150,
|
||||
)
|
||||
assert remaining == (on_the_cutoff, today), remaining
|
||||
assert _completion_id(owned, model).startswith("chatcmpl-")
|
||||
finally:
|
||||
_store_retention_setting(previously_stored)
|
||||
_delete_daily_tag_spend(tag)
|
||||
|
||||
|
||||
@pytest.mark.covers("spend.daily_tag_spend.invalid_retention_value_keeps_rows_and_leaves_the_proxy_serving")
|
||||
def test_unparseable_daily_tag_spend_retention_deletes_nothing_and_keeps_serving(
|
||||
gateway: Gateway, tmp_path: Path
|
||||
) -> None:
|
||||
tag: Final = f"integration-retention-{uuid.uuid4().hex}"
|
||||
request_id: Final = f"integration-retention-{uuid.uuid4().hex}"
|
||||
expired: Final = _day(200)
|
||||
_seed_daily_tag_spend(tag, (expired,))
|
||||
_seed_old_spend_log(request_id, days_ago=200)
|
||||
try:
|
||||
config: Final = _cleanup_config(
|
||||
tmp_path, {RETENTION_SETTING: "soon", "maximum_spend_logs_retention_period": "30d"}
|
||||
)
|
||||
with owned_proxy(gateway, tmp_path, {}, config=config) as owned, owned.scenario() as scenario:
|
||||
model: Final = scenario.model()
|
||||
eventually(lambda: _spend_log_present(request_id), lambda present: not present, seconds=150)
|
||||
assert _remaining_days(tag) == (expired,)
|
||||
assert _completion_id(owned, model).startswith("chatcmpl-")
|
||||
finally:
|
||||
_delete_daily_tag_spend(tag)
|
||||
|
||||
|
||||
@pytest.mark.covers("spend.daily_tag_spend.retention_horizon_is_independent_of_the_spend_log_horizon")
|
||||
def test_daily_tag_spend_keeps_days_the_shorter_spend_log_horizon_already_pruned(
|
||||
gateway: Gateway, tmp_path: Path
|
||||
) -> None:
|
||||
tag: Final = f"integration-retention-{uuid.uuid4().hex}"
|
||||
request_id: Final = f"integration-retention-{uuid.uuid4().hex}"
|
||||
expired, inside_tag_horizon = _day(200), _day(60)
|
||||
_seed_daily_tag_spend(tag, (expired, inside_tag_horizon))
|
||||
_seed_old_spend_log(request_id, days_ago=60)
|
||||
try:
|
||||
config: Final = _cleanup_config(
|
||||
tmp_path, {RETENTION_SETTING: "90d", "maximum_spend_logs_retention_period": "30d"}
|
||||
)
|
||||
with owned_proxy(gateway, tmp_path, {}, config=config):
|
||||
eventually(lambda: _spend_log_present(request_id), lambda present: not present, seconds=150)
|
||||
remaining: Final = eventually(lambda: _remaining_days(tag), lambda days: expired not in days, seconds=150)
|
||||
assert remaining == (inside_tag_horizon,), remaining
|
||||
finally:
|
||||
_delete_daily_tag_spend(tag)
|
||||
|
||||
|
||||
@pytest.mark.covers("spend.daily_tag_spend.cleanup_and_serving_survive_losing_one_of_two_workers")
|
||||
def test_daily_tag_spend_cleanup_completes_after_one_of_two_workers_is_killed(gateway: Gateway, tmp_path: Path) -> None:
|
||||
tag: Final = f"integration-retention-{uuid.uuid4().hex}"
|
||||
expired, today = _day(200), _day(0)
|
||||
_seed_daily_tag_spend(tag, (expired, today))
|
||||
try:
|
||||
config: Final = _cleanup_config(tmp_path, {RETENTION_SETTING: "30d"})
|
||||
with owned_proxy_process(gateway, tmp_path, {}, config=config, workers=2) as owned:
|
||||
with owned.gateway.scenario() as scenario:
|
||||
model: Final = scenario.model()
|
||||
workers: Final = eventually(
|
||||
lambda: _listening_workers(owned), lambda found: len(found) == 2, seconds=30
|
||||
)
|
||||
workers[0].send_signal(signal.SIGKILL)
|
||||
eventually(lambda: workers[0].is_running(), lambda alive: not alive, seconds=10)
|
||||
ids: Final = tuple(_completion_id(owned.gateway, model) for _ in range(6))
|
||||
assert len(set(ids)) == 6 and all(identity.startswith("chatcmpl-") for identity in ids), ids
|
||||
remaining: Final = eventually(
|
||||
lambda: _remaining_days(tag), lambda days: expired not in days, seconds=150
|
||||
)
|
||||
assert remaining == (today,), remaining
|
||||
finally:
|
||||
_delete_daily_tag_spend(tag)
|
||||
|
||||
|
||||
@pytest.mark.covers("spend.daily_tag_spend.unset_retention_never_deletes_even_while_spend_logs_are_pruned")
|
||||
def test_daily_tag_spend_is_kept_forever_when_its_retention_is_unset(gateway: Gateway, tmp_path: Path) -> None:
|
||||
tag: Final = f"integration-retention-{uuid.uuid4().hex}"
|
||||
request_id: Final = f"integration-retention-{uuid.uuid4().hex}"
|
||||
expired: Final = _day(200)
|
||||
_seed_daily_tag_spend(tag, (expired,))
|
||||
_seed_old_spend_log(request_id, days_ago=200)
|
||||
try:
|
||||
config: Final = _cleanup_config(tmp_path, {"maximum_spend_logs_retention_period": "30d"})
|
||||
with owned_proxy(gateway, tmp_path, {}, config=config):
|
||||
eventually(lambda: _spend_log_present(request_id), lambda present: not present, seconds=150)
|
||||
assert _remaining_days(tag) == (expired,)
|
||||
finally:
|
||||
_delete_daily_tag_spend(tag)
|
||||
|
|
@ -79,6 +79,7 @@ _PREVIOUSLY_DB_WINS: Final[tuple[str, ...]] = (
|
|||
"maximum_spend_logs_retention_period",
|
||||
"maximum_autorouter_session_retention_period",
|
||||
"maximum_health_check_retention_period",
|
||||
"maximum_daily_tag_spend_retention_period",
|
||||
"maximum_spend_logs_cleanup_batch_size",
|
||||
"maximum_spend_logs_cleanup_max_batches",
|
||||
"maximum_spend_logs_cleanup_run_budget",
|
||||
|
|
|
|||
|
|
@ -3789,6 +3789,52 @@ async def test_ProxyConfig__update_general_settings_updates_health_check_retenti
|
|||
reschedule.assert_awaited_once()
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_ProxyConfig__reschedule_spend_log_cleanup_job_daily_tag_spend_retention(monkeypatch):
|
||||
fake_scheduler = MagicMock()
|
||||
monkeypatch.setattr("litellm.proxy.proxy_server.scheduler", fake_scheduler)
|
||||
monkeypatch.setattr(
|
||||
"litellm.proxy.proxy_server.general_settings",
|
||||
{"maximum_daily_tag_spend_retention_period": "90d"},
|
||||
)
|
||||
monkeypatch.setattr("litellm.proxy.proxy_server.prisma_client", None)
|
||||
pc = ProxyConfig()
|
||||
await pc._reschedule_spend_log_cleanup_job()
|
||||
assert fake_scheduler.add_job.call_count == 1
|
||||
assert fake_scheduler.add_job.call_args.kwargs["id"] == "spend_log_cleanup_job"
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_ProxyConfig__update_general_settings_updates_daily_tag_spend_retention(monkeypatch):
|
||||
settings = {}
|
||||
monkeypatch.setattr("litellm.proxy.proxy_server.general_settings", settings)
|
||||
monkeypatch.setattr("litellm.proxy.proxy_server.scheduler", None)
|
||||
pc = ProxyConfig()
|
||||
reschedule = AsyncMock()
|
||||
monkeypatch.setattr(pc, "_reschedule_spend_log_cleanup_job", reschedule)
|
||||
await pc._update_general_settings({"maximum_daily_tag_spend_retention_period": "90d"})
|
||||
from litellm.proxy import proxy_server
|
||||
|
||||
assert proxy_server.general_settings["maximum_daily_tag_spend_retention_period"] == "90d"
|
||||
reschedule.assert_awaited_once()
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_ProxyConfig__update_general_settings_schedules_cleanup_when_db_row_was_already_applied(monkeypatch):
|
||||
"""A config reload applies the db row to the store before the side effects run, so the
|
||||
before/after snapshot is equal; the job must still be scheduled when none is running."""
|
||||
fake_scheduler = MagicMock()
|
||||
fake_scheduler.get_job.return_value = None
|
||||
monkeypatch.setattr("litellm.proxy.proxy_server.scheduler", fake_scheduler)
|
||||
monkeypatch.setattr("litellm.proxy.proxy_server.prisma_client", None)
|
||||
pc = ProxyConfig()
|
||||
pc.settings.apply_db_row("general_settings", {"maximum_daily_tag_spend_retention_period": "90d"})
|
||||
monkeypatch.setattr("litellm.proxy.proxy_server.general_settings", pc.settings)
|
||||
await pc._update_general_settings({"maximum_daily_tag_spend_retention_period": "90d"})
|
||||
assert fake_scheduler.add_job.call_count == 1
|
||||
assert fake_scheduler.add_job.call_args.kwargs["id"] == "spend_log_cleanup_job"
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# ProxyConfig._update_general_settings
|
||||
# ---------------------------------------------------------------------------
|
||||
|
|
@ -3920,6 +3966,7 @@ async def test_ProxyConfig__update_general_settings_skips_redundant_retention_re
|
|||
pc = ProxyConfig()
|
||||
reschedule: Final = AsyncMock()
|
||||
monkeypatch.setattr(proxy_server, "general_settings", {})
|
||||
monkeypatch.setattr(proxy_server, "scheduler", MagicMock())
|
||||
monkeypatch.setattr(pc, "_reschedule_spend_log_cleanup_job", reschedule)
|
||||
|
||||
await pc._update_general_settings({"maximum_health_check_retention_period": "30d"})
|
||||
|
|
|
|||
|
|
@ -826,6 +826,29 @@ async def test_health_check_retention_alone_cleans_only_the_health_check_table()
|
|||
assert abs((cutoff_date - expected_cutoff).total_seconds()) < 1
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_daily_tag_spend_retention_alone_prunes_only_that_table_by_calendar_day():
|
||||
client = _mock_prisma_for_retention([0])
|
||||
cleaner = SpendLogCleanup(general_settings={"maximum_daily_tag_spend_retention_period": "90d"})
|
||||
cleaner.pod_lock_manager = None
|
||||
await cleaner.cleanup_old_spend_logs(client)
|
||||
tables = [call[0][0] for call in client.db.execute_raw.call_args_list]
|
||||
assert len(tables) == 1
|
||||
assert '"LiteLLM_DailyTagSpend"' in tables[0]
|
||||
cutoff_day = client.db.execute_raw.call_args[0][1]
|
||||
assert cutoff_day == (datetime.now(timezone.utc) - timedelta(days=90)).date().isoformat()
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_spend_logs_retention_alone_keeps_daily_tag_spend_forever():
|
||||
client = _mock_prisma_for_retention([0, 0])
|
||||
cleaner = SpendLogCleanup(general_settings={"maximum_spend_logs_retention_period": "7d"})
|
||||
cleaner.pod_lock_manager = None
|
||||
await cleaner.cleanup_old_spend_logs(client)
|
||||
tables = [call[0][0] for call in client.db.execute_raw.call_args_list]
|
||||
assert not any('"LiteLLM_DailyTagSpend"' in sql for sql in tables)
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_each_retention_key_cuts_off_at_its_own_horizon():
|
||||
client = _mock_prisma_for_retention([0, 0, 0, 0, 0])
|
||||
|
|
|
|||
5
ui/litellm-dashboard/src/lib/http/schema.d.ts
generated
vendored
5
ui/litellm-dashboard/src/lib/http/schema.d.ts
generated
vendored
|
|
@ -27396,6 +27396,11 @@ export interface components {
|
|||
* @description Maximum retention period for auto-router benchmark session rollup rows (e.g., '365d'). Rows whose last turn is older than this are deleted by the spend log cleanup job, on that job's schedule. Unset means rollup rows are never deleted.
|
||||
*/
|
||||
maximum_autorouter_session_retention_period?: string | null;
|
||||
/**
|
||||
* Maximum Daily Tag Spend Retention Period
|
||||
* @description Maximum retention period for per-day tag spend aggregate rows (e.g., '90d'). Rows whose day is older than this are deleted by the spend log cleanup job, on that job's schedule. Unset means rows are never deleted. Only historical tag usage analytics are affected; tag budgets read the lifetime counter.
|
||||
*/
|
||||
maximum_daily_tag_spend_retention_period?: string | null;
|
||||
/**
|
||||
* Maximum Health Check Retention Period
|
||||
* @description Maximum retention period for health-check rows (e.g., '30d'). Rows whose checked_at is older than this are deleted by the spend log cleanup job, on that job's schedule. Unset means rows are never deleted. Set this well above health_check_interval because /health and the UI read the latest row per model.
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue