From 4f57b1a3f047706fb9f1dcab0961f6fc3206eb41 Mon Sep 17 00:00:00 2001 From: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> Date: Tue, 1 Sep 2026 23:28:00 +0000 Subject: [PATCH 1/7] feat(proxy): add maximum_daily_tag_spend_retention_period cleanup setting Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- litellm/proxy/_types.py | 9 ++ .../db_transaction_queue/spend_log_cleanup.py | 117 ++++++++++++++---- litellm/proxy/proxy_server.py | 9 +- .../proxy/test_spend_log_cleanup.py | 23 ++++ 4 files changed, 136 insertions(+), 22 deletions(-) diff --git a/litellm/proxy/_types.py b/litellm/proxy/_types.py index 54574ed64e3..ee25afa5216 100644 --- a/litellm/proxy/_types.py +++ b/litellm/proxy/_types.py @@ -2950,6 +2950,15 @@ 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 LiteLLM_DailyTagSpend rows (e.g., '90d'). Rows whose date is older than " + "this are deleted by the spend log cleanup job, on that job's schedule. The table only feeds usage " + "analytics (tag usage dashboards, /spend/tags), so deleting old rows truncates historical tag usage " + "charts but does not affect budget enforcement. Unset means rows are never deleted." + ), + ) 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.", diff --git a/litellm/proxy/db/db_transaction_queue/spend_log_cleanup.py b/litellm/proxy/db/db_transaction_queue/spend_log_cleanup.py index 85e19fa8a32..b64987412f4 100644 --- a/litellm/proxy/db/db_transaction_queue/spend_log_cleanup.py +++ b/litellm/proxy/db/db_transaction_queue/spend_log_cleanup.py @@ -277,7 +277,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: datetime | str, deadline: float ) -> int | None: """ Run one delete batch under a Postgres statement and lock timeout. @@ -296,11 +296,17 @@ class SpendLogCleanup: async with prisma_client.db.tx() as tx: await tx.execute_raw(f"SET LOCAL statement_timeout = {timeout_ms}") await tx.execute_raw(f"SET LOCAL lock_timeout = {timeout_ms}") - deleted_result: Final = await tx.execute_raw(delete_sql, cutoff_date, self.batch_size) + deleted_result: Final = await tx.execute_raw(delete_sql, cutoff, self.batch_size) 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: datetime | str, + table_name: str, + time_column: str, + time_cast: str, + deadline: float, ) -> int | None: """ Count expired rows still outstanding, stopping at a cap. @@ -313,7 +319,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::{time_cast} LIMIT $2 ) capped """ @@ -321,7 +327,7 @@ class SpendLogCleanup: async with prisma_client.db.tx() as tx: await tx.execute_raw(f"SET LOCAL statement_timeout = {self._timeout_ms(deadline)}") rows: Final = _REMAINING_ROWS.validate_python( - await tx.query_raw(count_sql, cutoff_date, SPEND_LOG_CLEANUP_REMAINING_COUNT_CAP) + await tx.query_raw(count_sql, cutoff, SPEND_LOG_CLEANUP_REMAINING_COUNT_CAP) ) except Exception as e: # noqa: BLE001 - an observability probe must never fail the cleanup run verbose_proxy_logger.warning("Could not count remaining %s rows: %s", table_name, e) @@ -331,11 +337,12 @@ class SpendLogCleanup: async def _delete_old_rows_batched( self, prisma_client: PrismaClient, - cutoff_date: datetime, + cutoff: datetime | str, table_name: str, key_columns: tuple[str, ...], time_column: str, deadline: float, + time_cast: str = "timestamptz", ) -> TableCleanupResult: """ Delete a table's rows older than the cutoff in batches. @@ -349,7 +356,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::{time_cast} LIMIT $2 ) """ @@ -364,19 +371,33 @@ class SpendLogCleanup: total_deleted, ) return await self._finish_table( - prisma_client, cutoff_date, table_name, time_column, total_deleted, "budget_exhausted", deadline + prisma_client, + cutoff, + table_name, + time_column, + time_cast, + total_deleted, + "budget_exhausted", + deadline, ) if run_count >= self.max_batches: verbose_proxy_logger.info( "Max batches reached for %s cleanup, remaining rows will be deleted in next run", table_name ) return await self._finish_table( - prisma_client, cutoff_date, table_name, time_column, total_deleted, "batch_cap_reached", deadline + prisma_client, + cutoff, + table_name, + time_column, + time_cast, + total_deleted, + "batch_cap_reached", + deadline, ) # Find rows and delete them in one go without fetching to application batch_started_at = time.monotonic() try: - batch_result = await self._execute_delete_batch(prisma_client, delete_sql, cutoff_date, deadline) + batch_result = await self._execute_delete_batch(prisma_client, delete_sql, cutoff, deadline) except Exception as batch_exc: if time.monotonic() >= deadline: # The statement timeout was clamped to the budget that was @@ -391,7 +412,14 @@ class SpendLogCleanup: total_deleted, ) return await self._finish_table( - prisma_client, cutoff_date, table_name, time_column, total_deleted, "budget_exhausted", deadline + prisma_client, + cutoff, + table_name, + time_column, + time_cast, + total_deleted, + "budget_exhausted", + deadline, ) # A single batch failure (e.g. Prisma/DB timeout) must not abort # the whole run — subsequent batches may still succeed. @@ -405,7 +433,7 @@ class SpendLogCleanup: run_count, consecutive_failures, self.batch_size, - cutoff_date.isoformat(), + cutoff.isoformat() if isinstance(cutoff, datetime) else cutoff, total_deleted, type(batch_exc).__name__, batch_exc, @@ -418,7 +446,7 @@ class SpendLogCleanup: total_deleted, ) return await self._finish_table( - prisma_client, cutoff_date, table_name, time_column, total_deleted, "aborted", deadline + prisma_client, cutoff, table_name, time_column, time_cast, total_deleted, "aborted", deadline ) await asyncio.sleep(SPEND_LOG_CLEANUP_BATCH_FAILURE_BACKOFF_SECONDS) continue @@ -429,7 +457,7 @@ class SpendLogCleanup: table_name, ) return await self._finish_table( - prisma_client, cutoff_date, table_name, time_column, total_deleted, "aborted", deadline + prisma_client, cutoff, table_name, time_column, time_cast, total_deleted, "aborted", deadline ) consecutive_failures = 0 @@ -440,7 +468,7 @@ class SpendLogCleanup: if deleted_count == 0: verbose_proxy_logger.info("No more %s rows to delete. Total deleted: %s", table_name, total_deleted) return await self._finish_table( - prisma_client, cutoff_date, table_name, time_column, total_deleted, "exhausted", deadline + prisma_client, cutoff, table_name, time_column, time_cast, total_deleted, "exhausted", deadline ) total_deleted += deleted_count @@ -453,9 +481,10 @@ class SpendLogCleanup: async def _finish_table( self, prisma_client: PrismaClient, - cutoff_date: datetime, + cutoff: datetime | str, table_name: str, time_column: str, + time_cast: str, rows_deleted: int, stop_reason: StopReason, deadline: float, @@ -473,7 +502,9 @@ class SpendLogCleanup: """ if time.monotonic() >= deadline: return TableCleanupResult(rows_deleted=rows_deleted, stop_reason=stop_reason) - remaining: Final = await self._count_remaining(prisma_client, cutoff_date, table_name, time_column, deadline) + remaining: Final = await self._count_remaining( + prisma_client, cutoff, table_name, time_column, time_cast, deadline + ) if remaining is not None: SpendLogCleanupMetrics.set_rows_remaining(table_name, remaining) return TableCleanupResult(rows_deleted=rows_deleted, stop_reason=stop_reason) @@ -540,6 +571,21 @@ class SpendLogCleanup: deadline=deadline, ) + async def _delete_old_daily_tag_spend_rows( + self, prisma_client: PrismaClient, cutoff_day: str, deadline: float + ) -> TableCleanupResult: + # "date" is a YYYY-MM-DD string, so lexicographic order matches chronological + # order and the comparison stays sargable on the existing date index. + return await self._delete_old_rows_batched( + prisma_client, + cutoff_day, + table_name="LiteLLM_DailyTagSpend", + key_columns=("id",), + time_column="date", + deadline=deadline, + time_cast="text", + ) + async def _clean_spend_log_tables( self, prisma_client: PrismaClient, deadline: float ) -> tuple[TableCleanupResult, ...]: @@ -621,6 +667,16 @@ class SpendLogCleanup: ) return (health_checks_result,) + async def _clean_daily_tag_spend( + self, prisma_client: PrismaClient, retention_seconds: int, deadline: float + ) -> tuple[TableCleanupResult, ...]: + cutoff_day: Final = (datetime.now(timezone.utc) - timedelta(seconds=float(retention_seconds))).strftime( + "%Y-%m-%d" + ) + tag_spend_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", tag_spend_result.rows_deleted) + return (tag_spend_result,) + @staticmethod def _run_outcome(results: tuple[TableCleanupResult, ...]) -> RunOutcome: """ @@ -657,10 +713,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 +752,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 +763,10 @@ 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) ) session_results: Final = ( await self._clean_session_rollup( @@ -714,18 +777,30 @@ class SpendLogCleanup: if autorouter_retention_seconds is not None else () ) + remaining_groups_after_sessions: Final = int(health_check_retention_seconds is not None) + int( + daily_tag_spend_retention_seconds is not None + ) health_check_results: Final = ( 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: diff --git a/litellm/proxy/proxy_server.py b/litellm/proxy/proxy_server.py index 92a75bf953a..f2337b504c3 100644 --- a/litellm/proxy/proxy_server.py +++ b/litellm/proxy/proxy_server.py @@ -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, ) @@ -10484,6 +10490,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") diff --git a/tests/test_litellm/proxy/test_spend_log_cleanup.py b/tests/test_litellm/proxy/test_spend_log_cleanup.py index e333da03950..49074c7126b 100644 --- a/tests/test_litellm/proxy/test_spend_log_cleanup.py +++ b/tests/test_litellm/proxy/test_spend_log_cleanup.py @@ -826,6 +826,23 @@ 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_cleans_only_the_daily_tag_spend_table(): + 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] + assert '"id"' in tables[0] + assert '"date"' in tables[0] + assert "$1::text" in tables[0] + cutoff_day = client.db.execute_raw.call_args[0][1] + expected_cutoff_day = (datetime.now(timezone.utc) - timedelta(days=90)).strftime("%Y-%m-%d") + assert cutoff_day == expected_cutoff_day + + @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]) @@ -834,6 +851,7 @@ async def test_each_retention_key_cuts_off_at_its_own_horizon(): "maximum_spend_logs_retention_period": "7d", "maximum_autorouter_session_retention_period": "365d", "maximum_health_check_retention_period": "30d", + "maximum_daily_tag_spend_retention_period": "90d", } ) cleaner.pod_lock_manager = None @@ -846,6 +864,8 @@ async def test_each_retention_key_cuts_off_at_its_own_horizon(): if '"LiteLLM_AutoRouterUserSession"' in call[0][0] else "LiteLLM_HealthCheckTable" if '"LiteLLM_HealthCheckTable"' in call[0][0] + else "LiteLLM_DailyTagSpend" + if '"LiteLLM_DailyTagSpend"' in call[0][0] else "logs" ): call[0][1] for call in client.db.execute_raw.call_args_list @@ -855,6 +875,7 @@ async def test_each_retention_key_cuts_off_at_its_own_horizon(): assert (now - cutoffs["LiteLLM_AutoRouterSession"]).days == 365 assert cutoffs["LiteLLM_AutoRouterUserSession"] == cutoffs["LiteLLM_AutoRouterSession"] assert (now - cutoffs["LiteLLM_HealthCheckTable"]).days == 30 + assert cutoffs["LiteLLM_DailyTagSpend"] == (now - timedelta(days=90)).strftime("%Y-%m-%d") @pytest.mark.asyncio @@ -1219,6 +1240,7 @@ async def test_the_outstanding_rows_probe_carries_a_statement_timeout(): datetime.now(timezone.utc) - timedelta(days=7), "LiteLLM_SpendLogs", "startTime", + "timestamptz", _far_deadline(), ) @@ -1285,6 +1307,7 @@ async def test_no_statement_is_issued_once_the_budget_is_spent(): datetime.now(timezone.utc) - timedelta(days=7), "LiteLLM_SpendLogs", "startTime", + "timestamptz", 123, "budget_exhausted", time.monotonic() - 1, From a0f1ee1c0aeacc33107660e4a6522ad4365fea0b Mon Sep 17 00:00:00 2001 From: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> Date: Tue, 1 Sep 2026 23:36:54 +0000 Subject: [PATCH 2/7] chore(ui): regenerate schema.d.ts for new retention setting Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- ui/litellm-dashboard/src/lib/http/schema.d.ts | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/ui/litellm-dashboard/src/lib/http/schema.d.ts b/ui/litellm-dashboard/src/lib/http/schema.d.ts index bdfd4aec316..1a01510cbfb 100644 --- a/ui/litellm-dashboard/src/lib/http/schema.d.ts +++ b/ui/litellm-dashboard/src/lib/http/schema.d.ts @@ -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 LiteLLM_DailyTagSpend rows (e.g., '90d'). Rows whose date is older than this are deleted by the spend log cleanup job, on that job's schedule. The table only feeds usage analytics (tag usage dashboards, /spend/tags), so deleting old rows truncates historical tag usage charts but does not affect budget enforcement. Unset means rows are never deleted. + */ + 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. From 57b27fcfc157c4d3a91eb8aeef17e2a62358c826 Mon Sep 17 00:00:00 2001 From: yucheng Date: Wed, 23 Sep 2026 07:00:59 +0000 Subject: [PATCH 3/7] feat(proxy): rebase daily tag spend retention onto the run-budgeted cleanup job Reworks the cleanup on top of the refactored SpendLogCleanup: the daily tag spend table is pruned through the shared batched delete with a text cutoff on the indexed ISO date column, the setting is picked up by /config/update and the scheduler registration, and an integration test proves rows older than the period are pruned while the cutoff day and unset retention are left alone Co-authored-by: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- litellm/proxy/_types.py | 7 +- .../db_transaction_queue/spend_log_cleanup.py | 105 +++++++----------- litellm/proxy/proxy_server.py | 2 + .../spend/test_daily_tag_spend_retention.py | 103 +++++++++++++++++ .../config_resolvers/test_settings_rules.py | 1 + .../proxy/proxy_server/test_proxy_config.py | 29 +++++ .../proxy/test_spend_log_cleanup.py | 25 +++-- ui/litellm-dashboard/src/lib/http/schema.d.ts | 2 +- 8 files changed, 192 insertions(+), 82 deletions(-) create mode 100644 tests/integration/spend/test_daily_tag_spend_retention.py diff --git a/litellm/proxy/_types.py b/litellm/proxy/_types.py index ee25afa5216..0eee0807ad6 100644 --- a/litellm/proxy/_types.py +++ b/litellm/proxy/_types.py @@ -2953,10 +2953,9 @@ class ConfigGeneralSettings(LiteLLMPydanticObjectBase): maximum_daily_tag_spend_retention_period: str | None = Field( None, description=( - "Maximum retention period for LiteLLM_DailyTagSpend rows (e.g., '90d'). Rows whose date is older than " - "this are deleted by the spend log cleanup job, on that job's schedule. The table only feeds usage " - "analytics (tag usage dashboards, /spend/tags), so deleting old rows truncates historical tag usage " - "charts but does not affect budget enforcement. Unset means rows are never deleted." + "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( diff --git a/litellm/proxy/db/db_transaction_queue/spend_log_cleanup.py b/litellm/proxy/db/db_transaction_queue/spend_log_cleanup.py index b64987412f4..d0239a1d8a4 100644 --- a/litellm/proxy/db/db_transaction_queue/spend_log_cleanup.py +++ b/litellm/proxy/db/db_transaction_queue/spend_log_cleanup.py @@ -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: datetime | str, 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. @@ -296,17 +307,11 @@ class SpendLogCleanup: async with prisma_client.db.tx() as tx: await tx.execute_raw(f"SET LOCAL statement_timeout = {timeout_ms}") await tx.execute_raw(f"SET LOCAL lock_timeout = {timeout_ms}") - deleted_result: Final = await tx.execute_raw(delete_sql, cutoff, self.batch_size) + deleted_result: Final = await tx.execute_raw(delete_sql, cutoff_date, self.batch_size) return deleted_result if isinstance(deleted_result, int) else None async def _count_remaining( - self, - prisma_client: PrismaClient, - cutoff: datetime | str, - table_name: str, - time_column: str, - time_cast: 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. @@ -319,7 +324,7 @@ class SpendLogCleanup: count_sql: Final = f""" SELECT count(*)::int AS remaining FROM ( SELECT 1 FROM "{table_name}" - WHERE "{time_column}" < $1::{time_cast} + WHERE "{time_column}" < $1::{_cutoff_cast(cutoff_date)} LIMIT $2 ) capped """ @@ -327,7 +332,7 @@ class SpendLogCleanup: async with prisma_client.db.tx() as tx: await tx.execute_raw(f"SET LOCAL statement_timeout = {self._timeout_ms(deadline)}") rows: Final = _REMAINING_ROWS.validate_python( - await tx.query_raw(count_sql, cutoff, SPEND_LOG_CLEANUP_REMAINING_COUNT_CAP) + await tx.query_raw(count_sql, cutoff_date, SPEND_LOG_CLEANUP_REMAINING_COUNT_CAP) ) except Exception as e: # noqa: BLE001 - an observability probe must never fail the cleanup run verbose_proxy_logger.warning("Could not count remaining %s rows: %s", table_name, e) @@ -337,12 +342,11 @@ class SpendLogCleanup: async def _delete_old_rows_batched( self, prisma_client: PrismaClient, - cutoff: datetime | str, + cutoff_date: Cutoff, table_name: str, key_columns: tuple[str, ...], time_column: str, deadline: float, - time_cast: str = "timestamptz", ) -> TableCleanupResult: """ Delete a table's rows older than the cutoff in batches. @@ -356,7 +360,7 @@ class SpendLogCleanup: DELETE FROM "{table_name}" WHERE ({key_list}) IN ( SELECT {key_list} FROM "{table_name}" - WHERE "{time_column}" < $1::{time_cast} + WHERE "{time_column}" < $1::{_cutoff_cast(cutoff_date)} LIMIT $2 ) """ @@ -371,33 +375,19 @@ class SpendLogCleanup: total_deleted, ) return await self._finish_table( - prisma_client, - cutoff, - table_name, - time_column, - time_cast, - total_deleted, - "budget_exhausted", - deadline, + prisma_client, cutoff_date, table_name, time_column, total_deleted, "budget_exhausted", deadline ) if run_count >= self.max_batches: verbose_proxy_logger.info( "Max batches reached for %s cleanup, remaining rows will be deleted in next run", table_name ) return await self._finish_table( - prisma_client, - cutoff, - table_name, - time_column, - time_cast, - total_deleted, - "batch_cap_reached", - deadline, + prisma_client, cutoff_date, table_name, time_column, total_deleted, "batch_cap_reached", deadline ) # Find rows and delete them in one go without fetching to application batch_started_at = time.monotonic() try: - batch_result = await self._execute_delete_batch(prisma_client, delete_sql, cutoff, deadline) + batch_result = await self._execute_delete_batch(prisma_client, delete_sql, cutoff_date, deadline) except Exception as batch_exc: if time.monotonic() >= deadline: # The statement timeout was clamped to the budget that was @@ -412,14 +402,7 @@ class SpendLogCleanup: total_deleted, ) return await self._finish_table( - prisma_client, - cutoff, - table_name, - time_column, - time_cast, - total_deleted, - "budget_exhausted", - deadline, + prisma_client, cutoff_date, table_name, time_column, total_deleted, "budget_exhausted", deadline ) # A single batch failure (e.g. Prisma/DB timeout) must not abort # the whole run — subsequent batches may still succeed. @@ -433,7 +416,7 @@ class SpendLogCleanup: run_count, consecutive_failures, self.batch_size, - cutoff.isoformat() if isinstance(cutoff, datetime) else cutoff, + _cutoff_text(cutoff_date), total_deleted, type(batch_exc).__name__, batch_exc, @@ -446,7 +429,7 @@ class SpendLogCleanup: total_deleted, ) return await self._finish_table( - prisma_client, cutoff, table_name, time_column, time_cast, total_deleted, "aborted", deadline + prisma_client, cutoff_date, table_name, time_column, total_deleted, "aborted", deadline ) await asyncio.sleep(SPEND_LOG_CLEANUP_BATCH_FAILURE_BACKOFF_SECONDS) continue @@ -457,7 +440,7 @@ class SpendLogCleanup: table_name, ) return await self._finish_table( - prisma_client, cutoff, table_name, time_column, time_cast, total_deleted, "aborted", deadline + prisma_client, cutoff_date, table_name, time_column, total_deleted, "aborted", deadline ) consecutive_failures = 0 @@ -468,7 +451,7 @@ class SpendLogCleanup: if deleted_count == 0: verbose_proxy_logger.info("No more %s rows to delete. Total deleted: %s", table_name, total_deleted) return await self._finish_table( - prisma_client, cutoff, table_name, time_column, time_cast, total_deleted, "exhausted", deadline + prisma_client, cutoff_date, table_name, time_column, total_deleted, "exhausted", deadline ) total_deleted += deleted_count @@ -481,10 +464,9 @@ class SpendLogCleanup: async def _finish_table( self, prisma_client: PrismaClient, - cutoff: datetime | str, + cutoff_date: Cutoff, table_name: str, time_column: str, - time_cast: str, rows_deleted: int, stop_reason: StopReason, deadline: float, @@ -502,9 +484,7 @@ class SpendLogCleanup: """ if time.monotonic() >= deadline: return TableCleanupResult(rows_deleted=rows_deleted, stop_reason=stop_reason) - remaining: Final = await self._count_remaining( - prisma_client, cutoff, table_name, time_column, time_cast, deadline - ) + remaining: Final = await self._count_remaining(prisma_client, cutoff_date, table_name, time_column, deadline) if remaining is not None: SpendLogCleanupMetrics.set_rows_remaining(table_name, remaining) return TableCleanupResult(rows_deleted=rows_deleted, stop_reason=stop_reason) @@ -574,8 +554,6 @@ class SpendLogCleanup: async def _delete_old_daily_tag_spend_rows( self, prisma_client: PrismaClient, cutoff_day: str, deadline: float ) -> TableCleanupResult: - # "date" is a YYYY-MM-DD string, so lexicographic order matches chronological - # order and the comparison stays sargable on the existing date index. return await self._delete_old_rows_batched( prisma_client, cutoff_day, @@ -583,7 +561,6 @@ class SpendLogCleanup: key_columns=("id",), time_column="date", deadline=deadline, - time_cast="text", ) async def _clean_spend_log_tables( @@ -670,12 +647,14 @@ class SpendLogCleanup: async def _clean_daily_tag_spend( self, prisma_client: PrismaClient, retention_seconds: int, deadline: float ) -> tuple[TableCleanupResult, ...]: - cutoff_day: Final = (datetime.now(timezone.utc) - timedelta(seconds=float(retention_seconds))).strftime( - "%Y-%m-%d" - ) - tag_spend_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", tag_spend_result.rows_deleted) - return (tag_spend_result,) + """ + 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: @@ -768,6 +747,9 @@ class SpendLogCleanup: + 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( prisma_client, @@ -777,9 +759,6 @@ class SpendLogCleanup: if autorouter_retention_seconds is not None else () ) - remaining_groups_after_sessions: Final = int(health_check_retention_seconds is not None) + int( - daily_tag_spend_retention_seconds is not None - ) health_check_results: Final = ( await self._clean_health_checks( prisma_client, @@ -790,11 +769,7 @@ class SpendLogCleanup: else () ) daily_tag_spend_results: Final = ( - await self._clean_daily_tag_spend( - prisma_client, - daily_tag_spend_retention_seconds, - deadline, - ) + await self._clean_daily_tag_spend(prisma_client, daily_tag_spend_retention_seconds, deadline) if daily_tag_spend_retention_seconds is not None else () ) diff --git a/litellm/proxy/proxy_server.py b/litellm/proxy/proxy_server.py index f2337b504c3..7cb2954c10e 100644 --- a/litellm/proxy/proxy_server.py +++ b/litellm/proxy/proxy_server.py @@ -7516,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", ) ) @@ -17704,6 +17705,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", diff --git a/tests/integration/spend/test_daily_tag_spend_retention.py b/tests/integration/spend/test_daily_tag_spend_retention.py new file mode 100644 index 00000000000..b231819962c --- /dev/null +++ b/tests/integration/spend/test_daily_tag_spend_retention.py @@ -0,0 +1,103 @@ +import os +import uuid +from datetime import datetime, timedelta, timezone +from pathlib import Path +from typing import Final + +import psycopg +import pytest +import yaml +from pydantic import JsonValue, TypeAdapter + +from tests.integration._support.client import Gateway, eventually +from tests.integration._support.database import read_rows +from tests.integration._support.process import owned_proxy + +CLEANUP_EVERY_MINUTE: Final = "* * * * *" +_CONFIG: Final = TypeAdapter(dict[str, 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 _cleanup_config(tmp_path: Path, retention: dict[str, JsonValue]) -> Path: + base: Final = _CONFIG.validate_python(yaml.safe_load(Path("tests/integration/proxy_config.yaml").read_text())) + config: Final = { + **base, + "general_settings": { + **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.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}" + _seed_daily_tag_spend(tag, (_day(200),)) + _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) == (_day(200),) + finally: + _delete_daily_tag_spend(tag) diff --git a/tests/test_litellm/proxy/config_resolvers/test_settings_rules.py b/tests/test_litellm/proxy/config_resolvers/test_settings_rules.py index ea5ebe6cf12..40e5870c804 100644 --- a/tests/test_litellm/proxy/config_resolvers/test_settings_rules.py +++ b/tests/test_litellm/proxy/config_resolvers/test_settings_rules.py @@ -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", diff --git a/tests/test_litellm/proxy/proxy_server/test_proxy_config.py b/tests/test_litellm/proxy/proxy_server/test_proxy_config.py index 89cb8356289..10e02634023 100644 --- a/tests/test_litellm/proxy/proxy_server/test_proxy_config.py +++ b/tests/test_litellm/proxy/proxy_server/test_proxy_config.py @@ -3789,6 +3789,35 @@ 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) + 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() + + # --------------------------------------------------------------------------- # ProxyConfig._update_general_settings # --------------------------------------------------------------------------- diff --git a/tests/test_litellm/proxy/test_spend_log_cleanup.py b/tests/test_litellm/proxy/test_spend_log_cleanup.py index 49074c7126b..612b24d7dd7 100644 --- a/tests/test_litellm/proxy/test_spend_log_cleanup.py +++ b/tests/test_litellm/proxy/test_spend_log_cleanup.py @@ -827,7 +827,7 @@ async def test_health_check_retention_alone_cleans_only_the_health_check_table() @pytest.mark.asyncio -async def test_daily_tag_spend_retention_alone_cleans_only_the_daily_tag_spend_table(): +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 @@ -835,12 +835,19 @@ async def test_daily_tag_spend_retention_alone_cleans_only_the_daily_tag_spend_t tables = [call[0][0] for call in client.db.execute_raw.call_args_list] assert len(tables) == 1 assert '"LiteLLM_DailyTagSpend"' in tables[0] - assert '"id"' in tables[0] - assert '"date"' in tables[0] - assert "$1::text" in tables[0] + assert '"date" < $1::text' in tables[0] cutoff_day = client.db.execute_raw.call_args[0][1] - expected_cutoff_day = (datetime.now(timezone.utc) - timedelta(days=90)).strftime("%Y-%m-%d") - assert cutoff_day == expected_cutoff_day + 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 @@ -851,7 +858,6 @@ async def test_each_retention_key_cuts_off_at_its_own_horizon(): "maximum_spend_logs_retention_period": "7d", "maximum_autorouter_session_retention_period": "365d", "maximum_health_check_retention_period": "30d", - "maximum_daily_tag_spend_retention_period": "90d", } ) cleaner.pod_lock_manager = None @@ -864,8 +870,6 @@ async def test_each_retention_key_cuts_off_at_its_own_horizon(): if '"LiteLLM_AutoRouterUserSession"' in call[0][0] else "LiteLLM_HealthCheckTable" if '"LiteLLM_HealthCheckTable"' in call[0][0] - else "LiteLLM_DailyTagSpend" - if '"LiteLLM_DailyTagSpend"' in call[0][0] else "logs" ): call[0][1] for call in client.db.execute_raw.call_args_list @@ -875,7 +879,6 @@ async def test_each_retention_key_cuts_off_at_its_own_horizon(): assert (now - cutoffs["LiteLLM_AutoRouterSession"]).days == 365 assert cutoffs["LiteLLM_AutoRouterUserSession"] == cutoffs["LiteLLM_AutoRouterSession"] assert (now - cutoffs["LiteLLM_HealthCheckTable"]).days == 30 - assert cutoffs["LiteLLM_DailyTagSpend"] == (now - timedelta(days=90)).strftime("%Y-%m-%d") @pytest.mark.asyncio @@ -1240,7 +1243,6 @@ async def test_the_outstanding_rows_probe_carries_a_statement_timeout(): datetime.now(timezone.utc) - timedelta(days=7), "LiteLLM_SpendLogs", "startTime", - "timestamptz", _far_deadline(), ) @@ -1307,7 +1309,6 @@ async def test_no_statement_is_issued_once_the_budget_is_spent(): datetime.now(timezone.utc) - timedelta(days=7), "LiteLLM_SpendLogs", "startTime", - "timestamptz", 123, "budget_exhausted", time.monotonic() - 1, diff --git a/ui/litellm-dashboard/src/lib/http/schema.d.ts b/ui/litellm-dashboard/src/lib/http/schema.d.ts index 1a01510cbfb..44e1b4b27f5 100644 --- a/ui/litellm-dashboard/src/lib/http/schema.d.ts +++ b/ui/litellm-dashboard/src/lib/http/schema.d.ts @@ -27398,7 +27398,7 @@ export interface components { maximum_autorouter_session_retention_period?: string | null; /** * Maximum Daily Tag Spend Retention Period - * @description Maximum retention period for LiteLLM_DailyTagSpend rows (e.g., '90d'). Rows whose date is older than this are deleted by the spend log cleanup job, on that job's schedule. The table only feeds usage analytics (tag usage dashboards, /spend/tags), so deleting old rows truncates historical tag usage charts but does not affect budget enforcement. Unset means rows are never deleted. + * @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; /** From 6cb3bc260904e7c1c6af7328789f64f9699e370d Mon Sep 17 00:00:00 2001 From: yucheng Date: Wed, 23 Sep 2026 07:46:04 +0000 Subject: [PATCH 4/7] fix(proxy): schedule the cleanup job when a retention db row lands before the side effects run A config reload applies the db row to the SettingsStore before _update_general_settings snapshots the previous retention values, so the before/after compare saw no change and a retention period first set through /config/update never scheduled the cleanup job. Also reschedule when the job is missing but a retention period is set Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- litellm/proxy/proxy_server.py | 5 ++++- .../proxy/proxy_server/test_proxy_config.py | 18 ++++++++++++++++++ 2 files changed, 22 insertions(+), 1 deletion(-) diff --git a/litellm/proxy/proxy_server.py b/litellm/proxy/proxy_server.py index 7cb2954c10e..92d98606b23 100644 --- a/litellm/proxy/proxy_server.py +++ b/litellm/proxy/proxy_server.py @@ -7630,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: diff --git a/tests/test_litellm/proxy/proxy_server/test_proxy_config.py b/tests/test_litellm/proxy/proxy_server/test_proxy_config.py index 10e02634023..6d7b89dfd09 100644 --- a/tests/test_litellm/proxy/proxy_server/test_proxy_config.py +++ b/tests/test_litellm/proxy/proxy_server/test_proxy_config.py @@ -3808,6 +3808,7 @@ async def test_ProxyConfig__reschedule_spend_log_cleanup_job_daily_tag_spend_ret 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) @@ -3818,6 +3819,22 @@ async def test_ProxyConfig__update_general_settings_updates_daily_tag_spend_rete 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 # --------------------------------------------------------------------------- @@ -3949,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"}) From a484a13917ddc7cadb65bf2bd12f295c056f34cf Mon Sep 17 00:00:00 2001 From: yucheng Date: Wed, 23 Sep 2026 07:52:18 +0000 Subject: [PATCH 5/7] test(integration): accept list-valued top-level keys in the base integration proxy config The shared tests/integration/proxy_config.yaml now carries list-valued top-level keys, so the retention config helper validates only the mapping it merges into. Also drops a SQL-shape assertion from the unit test in favor of the behavioral cutoff-day check Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- tests/integration/spend/test_daily_tag_spend_retention.py | 6 +++--- tests/test_litellm/proxy/test_spend_log_cleanup.py | 1 - 2 files changed, 3 insertions(+), 4 deletions(-) diff --git a/tests/integration/spend/test_daily_tag_spend_retention.py b/tests/integration/spend/test_daily_tag_spend_retention.py index b231819962c..5cbee96a2b1 100644 --- a/tests/integration/spend/test_daily_tag_spend_retention.py +++ b/tests/integration/spend/test_daily_tag_spend_retention.py @@ -14,7 +14,7 @@ from tests.integration._support.database import read_rows from tests.integration._support.process import owned_proxy CLEANUP_EVERY_MINUTE: Final = "* * * * *" -_CONFIG: Final = TypeAdapter(dict[str, dict[str, JsonValue]]) +_MAPPING: Final = TypeAdapter(dict[str, JsonValue]) def _day(days_ago: int) -> str: @@ -55,11 +55,11 @@ def _spend_log_present(request_id: str) -> bool: def _cleanup_config(tmp_path: Path, retention: dict[str, JsonValue]) -> Path: - base: Final = _CONFIG.validate_python(yaml.safe_load(Path("tests/integration/proxy_config.yaml").read_text())) + base: Final = _MAPPING.validate_python(yaml.safe_load(Path("tests/integration/proxy_config.yaml").read_text())) config: Final = { **base, "general_settings": { - **base["general_settings"], + **_MAPPING.validate_python(base["general_settings"]), **retention, "maximum_spend_logs_cleanup_cron": CLEANUP_EVERY_MINUTE, "scheduled_job_stagger": {"enabled": False}, diff --git a/tests/test_litellm/proxy/test_spend_log_cleanup.py b/tests/test_litellm/proxy/test_spend_log_cleanup.py index 612b24d7dd7..da78c96e8e3 100644 --- a/tests/test_litellm/proxy/test_spend_log_cleanup.py +++ b/tests/test_litellm/proxy/test_spend_log_cleanup.py @@ -835,7 +835,6 @@ async def test_daily_tag_spend_retention_alone_prunes_only_that_table_by_calenda tables = [call[0][0] for call in client.db.execute_raw.call_args_list] assert len(tables) == 1 assert '"LiteLLM_DailyTagSpend"' in tables[0] - assert '"date" < $1::text' 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() From 04f63dc235db66cb276291063382a077e577fba7 Mon Sep 17 00:00:00 2001 From: yucheng Date: Wed, 23 Sep 2026 08:54:02 +0000 Subject: [PATCH 6/7] test(integration): cover runtime update, invalid value, independent horizons and worker loss for daily tag spend retention Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- .../spend/test_daily_tag_spend_retention.py | 119 +++++++++++++++++- 1 file changed, 117 insertions(+), 2 deletions(-) diff --git a/tests/integration/spend/test_daily_tag_spend_retention.py b/tests/integration/spend/test_daily_tag_spend_retention.py index 5cbee96a2b1..475d7ecde11 100644 --- a/tests/integration/spend/test_daily_tag_spend_retention.py +++ b/tests/integration/spend/test_daily_tag_spend_retention.py @@ -1,20 +1,24 @@ 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 +from tests.integration._support.client import Gateway, eventually, string_value from tests.integration._support.database import read_rows -from tests.integration._support.process import owned_proxy +from tests.integration._support.process import 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: @@ -54,6 +58,26 @@ 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 _forget_stored_retention_setting() -> None: + with psycopg.connect(os.environ["DATABASE_URL"], autocommit=True) as connection: + connection.execute( + 'UPDATE "LiteLLM_Config" SET param_value = param_value - %s WHERE param_name = %s', + (RETENTION_SETTING, "general_settings"), + ) + + +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 = { @@ -88,6 +112,97 @@ def test_daily_tag_spend_retention_prunes_only_rows_older_than_the_period(gatewa _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)) + _forget_stored_retention_setting() + 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: + _forget_stored_retention_setting() + _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}" + _seed_daily_tag_spend(tag, (_day(200),)) + _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) == (_day(200),) + 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}" + _seed_daily_tag_spend(tag, (_day(200), _day(60))) + _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: _day(200) not in days, seconds=150) + assert remaining == (_day(60),), 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}" + _seed_daily_tag_spend(tag, (_day(200), _day(0))) + 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: psutil.Process(owned.process.pid).children(recursive=True), + lambda children: len(children) >= 2, + seconds=30, + ) + workers[0].send_signal(signal.SIGKILL) + 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: _day(200) not in days, seconds=150 + ) + assert remaining == (_day(0),), 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}" From 486465808169426ade8d867ec53dbd6b52db31f0 Mon Sep 17 00:00:00 2001 From: yucheng Date: Wed, 23 Sep 2026 09:27:25 +0000 Subject: [PATCH 7/7] test(integration): restore the shared retention setting, capture seeded days once and kill a listening worker Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- .../spend/test_daily_tag_spend_retention.py | 67 +++++++++++++------ 1 file changed, 48 insertions(+), 19 deletions(-) diff --git a/tests/integration/spend/test_daily_tag_spend_retention.py b/tests/integration/spend/test_daily_tag_spend_retention.py index 475d7ecde11..8642d073941 100644 --- a/tests/integration/spend/test_daily_tag_spend_retention.py +++ b/tests/integration/spend/test_daily_tag_spend_retention.py @@ -1,3 +1,4 @@ +import json import os import signal import uuid @@ -13,7 +14,7 @@ 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 owned_proxy, owned_proxy_process +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" @@ -58,14 +59,38 @@ 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 _forget_stored_retention_setting() -> None: +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 = param_value - %s WHERE param_name = %s', - (RETENTION_SETTING, "general_settings"), + '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 @@ -117,7 +142,8 @@ def test_config_update_turns_on_daily_tag_spend_cleanup_without_a_restart(gatewa 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)) - _forget_stored_retention_setting() + 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: @@ -133,7 +159,7 @@ def test_config_update_turns_on_daily_tag_spend_cleanup_without_a_restart(gatewa assert remaining == (on_the_cutoff, today), remaining assert _completion_id(owned, model).startswith("chatcmpl-") finally: - _forget_stored_retention_setting() + _store_retention_setting(previously_stored) _delete_daily_tag_spend(tag) @@ -143,7 +169,8 @@ def test_unparseable_daily_tag_spend_retention_deletes_nothing_and_keeps_serving ) -> None: tag: Final = f"integration-retention-{uuid.uuid4().hex}" request_id: Final = f"integration-retention-{uuid.uuid4().hex}" - _seed_daily_tag_spend(tag, (_day(200),)) + expired: Final = _day(200) + _seed_daily_tag_spend(tag, (expired,)) _seed_old_spend_log(request_id, days_ago=200) try: config: Final = _cleanup_config( @@ -152,7 +179,7 @@ def test_unparseable_daily_tag_spend_retention_deletes_nothing_and_keeps_serving 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) == (_day(200),) + assert _remaining_days(tag) == (expired,) assert _completion_id(owned, model).startswith("chatcmpl-") finally: _delete_daily_tag_spend(tag) @@ -164,7 +191,8 @@ def test_daily_tag_spend_keeps_days_the_shorter_spend_log_horizon_already_pruned ) -> None: tag: Final = f"integration-retention-{uuid.uuid4().hex}" request_id: Final = f"integration-retention-{uuid.uuid4().hex}" - _seed_daily_tag_spend(tag, (_day(200), _day(60))) + 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( @@ -172,8 +200,8 @@ def test_daily_tag_spend_keeps_days_the_shorter_spend_log_horizon_already_pruned ) 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: _day(200) not in days, seconds=150) - assert remaining == (_day(60),), remaining + 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) @@ -181,24 +209,24 @@ def test_daily_tag_spend_keeps_days_the_shorter_spend_log_horizon_already_pruned @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}" - _seed_daily_tag_spend(tag, (_day(200), _day(0))) + 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: psutil.Process(owned.process.pid).children(recursive=True), - lambda children: len(children) >= 2, - seconds=30, + 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: _day(200) not in days, seconds=150 + lambda: _remaining_days(tag), lambda days: expired not in days, seconds=150 ) - assert remaining == (_day(0),), remaining + assert remaining == (today,), remaining finally: _delete_daily_tag_spend(tag) @@ -207,12 +235,13 @@ def test_daily_tag_spend_cleanup_completes_after_one_of_two_workers_is_killed(ga 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}" - _seed_daily_tag_spend(tag, (_day(200),)) + 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) == (_day(200),) + assert _remaining_days(tag) == (expired,) finally: _delete_daily_tag_spend(tag)