perf(spend-logs): bound retention cleanup so one run cannot saturate the database (#36594)

This commit is contained in:
Yassin Kortam 2026-08-13 18:56:37 -07:00 • committed by GitHub
parent 2b63919f67
commit 909a2e6232
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
14 changed files with 2050 additions and 205 deletions

View file

@ -1491,6 +1491,9 @@ SPEND_LOG_CLEANUP_MAX_CONSECUTIVE_BATCH_FAILURES = int(os.getenv("SPEND_LOG_CLEA
SPEND_LOG_CLEANUP_BATCH_FAILURE_BACKOFF_SECONDS: Final = float(
os.getenv("SPEND_LOG_CLEANUP_BATCH_FAILURE_BACKOFF_SECONDS", 0.5)
)
SPEND_LOG_CLEANUP_RUN_BUDGET_SECONDS: Final = float(os.getenv("SPEND_LOG_CLEANUP_RUN_BUDGET_SECONDS", "300"))
SPEND_LOG_CLEANUP_BATCH_TIMEOUT_SECONDS: Final = float(os.getenv("SPEND_LOG_CLEANUP_BATCH_TIMEOUT_SECONDS", "30"))
SPEND_LOG_CLEANUP_REMAINING_COUNT_CAP: Final = int(os.getenv("SPEND_LOG_CLEANUP_REMAINING_COUNT_CAP", "100000"))
TOOL_SPEND_TOP_TOOLS: Final = 100
SPEND_LOG_PARTITION_INTERVAL: Final = os.getenv("SPEND_LOG_PARTITION_INTERVAL", "day")
SPEND_LOG_PARTITION_PRECREATE_AHEAD: Final = int(os.getenv("SPEND_LOG_PARTITION_PRECREATE_AHEAD", 7))

View file

@ -2514,6 +2514,22 @@ class ConfigGeneralSettings(LiteLLMPydanticObjectBase):
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.",
)
maximum_spend_logs_cleanup_batch_size: int | None = Field(
None,
description="Rows deleted per DELETE statement by the spend log cleanup job. Defaults to 1000.",
)
maximum_spend_logs_cleanup_max_batches: int | None = Field(
None,
description="Maximum DELETE statements the spend log cleanup job issues per table per run. Defaults to 500.",
)
maximum_spend_logs_cleanup_run_budget: str | None = Field(
None,
description="Wall-clock budget for one spend log cleanup run (e.g. '5m'), shared across every table it prunes. A run that hits the budget stops and the next run resumes from where it left off. Defaults to '5m'.",
)
maximum_spend_logs_cleanup_batch_timeout: str | None = Field(
None,
description="Postgres statement_timeout and lock_timeout applied to each spend log cleanup delete batch (e.g. '30s'), so cleanup cannot hold row locks or a connection indefinitely. Defaults to '30s'.",
)
mcp_internal_ip_ranges: list[str] | None = Field(
None,
description="Custom CIDR ranges that define internal/private networks for MCP access control. When set, only these ranges are treated as internal. Defaults to RFC 1918 private ranges (10.0.0.0/8, 172.16.0.0/12, 192.168.0.0/16, 127.0.0.0/8).",

View file

@ -1,22 +1,60 @@
import asyncio
import time
from dataclasses import dataclass
from datetime import datetime, timedelta, timezone
from typing import Final
from typing import Final, Literal, TypeAlias
from pydantic import BaseModel, TypeAdapter
from litellm._logging import verbose_proxy_logger
from litellm.caching import RedisCache
from litellm.constants import (
SPEND_LOG_CLEANUP_BATCH_FAILURE_BACKOFF_SECONDS,
SPEND_LOG_CLEANUP_BATCH_SIZE,
SPEND_LOG_CLEANUP_BATCH_TIMEOUT_SECONDS,
SPEND_LOG_CLEANUP_JOB_NAME,
SPEND_LOG_CLEANUP_MAX_CONSECUTIVE_BATCH_FAILURES,
SPEND_LOG_CLEANUP_REMAINING_COUNT_CAP,
SPEND_LOG_CLEANUP_RUN_BUDGET_SECONDS,
SPEND_LOG_RUN_LOOPS,
)
from litellm.litellm_core_utils.duration_parser import duration_in_seconds
from litellm.proxy.db.db_transaction_queue.spend_log_cleanup_metrics import (
RunOutcome,
SpendLogCleanupMetrics,
)
from litellm.proxy.db.db_transaction_queue.spend_logs_partition_manager import (
RemainingTimeoutMs,
SpendLogsPartitionManager,
)
from litellm.proxy.utils import PrismaClient
StopReason: TypeAlias = Literal["exhausted", "budget_exhausted", "batch_cap_reached", "aborted"]
@dataclass(frozen=True, slots=True)
class TableCleanupResult:
"""Outcome of pruning one table, so the caller can report why a run ended."""
rows_deleted: int
stop_reason: StopReason
class _RemainingRow(BaseModel):
"""One row of the capped outstanding-rows probe, validated out of prisma's untyped result."""
remaining: int
_REMAINING_ROWS: Final = TypeAdapter(list[_RemainingRow])
SPEND_LOG_CLEANUP_BOUND_SETTINGS: Final = (
"maximum_spend_logs_cleanup_batch_size",
"maximum_spend_logs_cleanup_max_batches",
"maximum_spend_logs_cleanup_run_budget",
"maximum_spend_logs_cleanup_batch_timeout",
)
class SpendLogCleanup:
"""
@ -26,6 +64,24 @@ class SpendLogCleanup:
dropping whole partitions (instant, frees disk immediately). Otherwise it
falls back to deleting logs in batches.
Uses PodLockManager to ensure only one pod runs cleanup in multi-pod deployments.
Every run is bounded so it can never monopolise the database: a wall-clock
budget shared across all tables, a per-table batch cap, and a Postgres
statement/lock timeout on every statement the job issues, deletes and the
outstanding-rows probe alike. A run that hits a bound stops cleanly and the
next run resumes from where it left off, because the cutoff is recomputed
and deleted rows are gone.
The budget is a hard wall clock, not an advisory one. Every statement this
job issues, deletes, the outstanding-rows probe and partition DDL alike, is
issued with a timeout clamped to the budget that is still left, so one
started just under the deadline is cancelled by Postgres at the deadline
rather than running a further batch timeout past it. No statement is issued
at all once the budget is spent, which is why the probe is skipped on that
path. Partition DDL additionally carries a lock_timeout, because it takes an
ACCESS EXCLUSIVE lock and would otherwise queue behind a long-running reader
for as long as that reader lives; a partition this run cannot get is left
for the next one.
"""
def __init__(
@ -34,17 +90,88 @@ class SpendLogCleanup:
redis_cache: RedisCache | None = None,
partition_manager: SpendLogsPartitionManager | None = None,
):
self.batch_size = SPEND_LOG_CLEANUP_BATCH_SIZE
self.retention_seconds: int | None = None
self.partition_manager = partition_manager or SpendLogsPartitionManager()
from litellm.proxy.proxy_server import general_settings as default_settings
self.general_settings = general_settings or default_settings
self._refresh_bounds()
from litellm.proxy.proxy_server import proxy_logging_obj
pod_lock_manager: Final = proxy_logging_obj.db_spend_update_writer.pod_lock_manager
self.pod_lock_manager = pod_lock_manager
verbose_proxy_logger.info("SpendLogCleanup initialized with batch size: %s", self.batch_size)
verbose_proxy_logger.info(
"SpendLogCleanup initialized: batch_size=%s max_batches=%s run_budget=%ss batch_timeout=%ss",
self.batch_size,
self.max_batches,
self.run_budget_seconds,
self.batch_timeout_seconds,
)
def _refresh_bounds(self) -> None:
"""
Re-read every bound in SPEND_LOG_CLEANUP_BOUND_SETTINGS from settings.
The scheduler holds one long-lived instance, so a bound captured at
construction would never reflect a dashboard change. general_settings is
the same dict the periodic config reload mutates in place, so reading it
per run is what makes these knobs live. Every bound falls back to its
shipped default, so clearing a field restores that default.
"""
self.batch_size: int = self._positive_int_setting(
"maximum_spend_logs_cleanup_batch_size", SPEND_LOG_CLEANUP_BATCH_SIZE
)
self.max_batches: int = self._positive_int_setting(
"maximum_spend_logs_cleanup_max_batches", SPEND_LOG_RUN_LOOPS
)
self.run_budget_seconds: float = self._duration_setting(
"maximum_spend_logs_cleanup_run_budget", SPEND_LOG_CLEANUP_RUN_BUDGET_SECONDS
)
self.batch_timeout_seconds: float = self._duration_setting(
"maximum_spend_logs_cleanup_batch_timeout", SPEND_LOG_CLEANUP_BATCH_TIMEOUT_SECONDS
)
def _positive_int_setting(self, setting_name: str, default: int) -> int:
"""
Read a positive-integer knob, falling back to the default when unset or unusable.
"""
raw: Final = self.general_settings.get(setting_name)
if raw is None:
return default
try:
parsed: Final = int(raw)
except (TypeError, ValueError):
verbose_proxy_logger.warning("Invalid %s value: %s, using default %s", setting_name, raw, default)
return default
if parsed <= 0:
verbose_proxy_logger.warning("%s must be positive, got %s, using default %s", setting_name, parsed, default)
return default
return parsed
def _duration_setting(self, setting_name: str, default_seconds: float) -> float:
"""
Read a duration knob (e.g. '5m'), falling back to the default when unset or unusable.
The knob must never be able to remove the bound it exists to enforce, so
anything the parser rejects (including the non-finite spellings 'inf' and
'nan') and anything non-positive falls back rather than being honoured.
"""
raw: Final = self.general_settings.get(setting_name)
if raw is None:
return default_seconds
try:
parsed: Final = float(duration_in_seconds(str(raw)))
except (ValueError, TypeError) as e:
verbose_proxy_logger.warning(
"Invalid %s value: %s (%s), using default %ss", setting_name, raw, e, default_seconds
)
return default_seconds
if parsed <= 0:
verbose_proxy_logger.warning(
"%s must be a positive duration, got %s, using default %ss", setting_name, raw, default_seconds
)
return default_seconds
return parsed
def _retention_seconds_for(self, setting_name: str) -> int | None:
"""
@ -78,6 +205,91 @@ class SpendLogCleanup:
self.retention_seconds = self._retention_seconds_for("maximum_spend_logs_retention_period")
return self.retention_seconds is not None
def _timeout_ms(self, deadline: float) -> int:
"""
The per-statement bound in milliseconds: the batch timeout, or whatever
is left of the run budget, whichever is smaller.
Clamping to the remaining budget is what makes the budget a real
wall-clock bound rather than an advisory one. Postgres offers no "stop
at time T", only a per-statement duration, so a statement issued just
under the deadline would otherwise run a full batch timeout past it, and
with several tables those overruns stack.
Interpolating this into SQL is safe by construction: an int cannot carry
SQL, and SET does not accept a bind parameter.
"""
remaining_ms: Final = int((deadline - time.monotonic()) * 1000)
return max(1, min(int(self.batch_timeout_seconds * 1000), remaining_ms))
def _remaining_timeout_ms(self, deadline: float) -> RemainingTimeoutMs:
"""
The per-statement bound for work this job delegates, as a callable.
Partition maintenance issues one statement per partition, so handing it a
number would bound each statement by the budget that was left before the
FIRST one and never by what remains. Re-evaluating per statement is what
makes the loop itself bounded, and None tells the callee to stop rather
than issue a statement it has no budget for.
"""
def remaining() -> int | None:
return None if time.monotonic() >= deadline else self._timeout_ms(deadline)
return remaining
async def _execute_delete_batch(
self, prisma_client: PrismaClient, delete_sql: str, cutoff_date: datetime, deadline: float
) -> int | None:
"""
Run one delete batch under a Postgres statement and lock timeout.
The timeouts are what actually bound the work: a Prisma transaction
timeout cannot interrupt a statement that is already executing, so
without these a single batch blocked behind a lock would hold its
connection, and the row locks it already took, indefinitely. SET LOCAL
scopes both to this transaction so the pooled connection is unaffected.
Returns the row count, or None when the driver returned something that
is not a row count. That is a contract violation rather than a transient
fault, so the caller stops instead of retrying.
"""
timeout_ms: Final = self._timeout_ms(deadline)
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)
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
) -> int | None:
"""
Count expired rows still outstanding, stopping at a cap.
An uncapped COUNT(*) over an expired backlog would itself be the kind of
long scan this job exists to avoid, so the probe reads at most
SPEND_LOG_CLEANUP_REMAINING_COUNT_CAP index entries. A result equal to
the cap means "at least this many".
"""
count_sql: Final = f"""
SELECT count(*)::int AS remaining FROM (
SELECT 1 FROM "{table_name}"
WHERE "{time_column}" < $1::timestamptz
LIMIT $2
) capped
"""
try:
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)
)
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)
return None
return rows[0].remaining if rows else None
async def _delete_old_rows_batched(
self,
prisma_client: PrismaClient,
@ -85,10 +297,14 @@ class SpendLogCleanup:
table_name: str,
key_columns: tuple[str, ...],
time_column: str,
) -> int:
deadline: float,
) -> TableCleanupResult:
"""
Helper method to delete a table's rows older than the cutoff in batches.
Returns the total number of rows deleted.
Delete a table's rows older than the cutoff in batches.
Stops at whichever bound is reached first: the backlog running out, the
shared wall-clock deadline, the per-table batch cap, or too many
consecutive batch failures.
"""
key_list: Final = ", ".join(f'"{col}"' for col in key_columns)
delete_sql: Final = f"""
@ -103,23 +319,46 @@ class SpendLogCleanup:
run_count = 0
consecutive_failures = 0
while True:
if run_count > SPEND_LOG_RUN_LOOPS:
if time.monotonic() >= deadline:
verbose_proxy_logger.info(
"Run budget exhausted during %s cleanup after %d rows; the next run resumes from here",
table_name,
total_deleted,
)
return await self._finish_table(
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
)
break
# Step 1: Find rows and delete them in one go without fetching to application
# Delete in batches, limited by self.batch_size
try:
deleted_result = await prisma_client.db.execute_raw(
delete_sql,
cutoff_date,
self.batch_size,
return await self._finish_table(
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_date, deadline)
except Exception as batch_exc:
if time.monotonic() >= deadline:
# The statement timeout was clamped to the budget that was
# left, so this batch was cancelled by the deadline itself.
# That is the bound working, not a database fault, and
# counting it would both inflate the failure metric and push
# every budget-exhausted run toward the abort threshold.
verbose_proxy_logger.info(
"Run budget exhausted mid-batch during %s cleanup after %d rows; "
"the next run resumes from here",
table_name,
total_deleted,
)
return await self._finish_table(
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.
consecutive_failures += 1
SpendLogCleanupMetrics.record_batch_failure(table_name)
verbose_proxy_logger.exception(
"%s cleanup batch failed "
"(run_count=%d, consecutive_failures=%d, batch_size=%d, "
@ -140,28 +379,31 @@ class SpendLogCleanup:
consecutive_failures,
total_deleted,
)
break
return await self._finish_table(
prisma_client, cutoff_date, table_name, time_column, total_deleted, "aborted", deadline
)
await asyncio.sleep(SPEND_LOG_CLEANUP_BATCH_FAILURE_BACKOFF_SECONDS)
continue
consecutive_failures = 0
deleted_count = 0
if isinstance(deleted_result, int):
deleted_count = deleted_result
else:
if batch_result is None:
verbose_proxy_logger.error(
"Unexpected execute_raw return type for %s cleanup: %s; aborting cleanup to avoid infinite loop",
"Unexpected execute_raw return type for %s cleanup; aborting cleanup to avoid infinite loop",
table_name,
type(deleted_result),
)
break
return await self._finish_table(
prisma_client, cutoff_date, table_name, time_column, total_deleted, "aborted", deadline
)
consecutive_failures = 0
deleted_count = batch_result
SpendLogCleanupMetrics.record_batch(table_name, deleted_count, time.monotonic() - batch_started_at)
verbose_proxy_logger.info("Deleted %s %s rows in this batch", deleted_count, table_name)
if deleted_count == 0:
verbose_proxy_logger.info("No more %s rows to delete. Total deleted: %s", table_name, total_deleted)
break
return await self._finish_table(
prisma_client, cutoff_date, table_name, time_column, total_deleted, "exhausted", deadline
)
total_deleted += deleted_count
run_count += 1
@ -169,18 +411,49 @@ class SpendLogCleanup:
# Add a small sleep to prevent overwhelming the database
await asyncio.sleep(0.1)
return total_deleted
async def _finish_table(
self,
prisma_client: PrismaClient,
cutoff_date: datetime,
table_name: str,
time_column: str,
rows_deleted: int,
stop_reason: StopReason,
deadline: float,
) -> TableCleanupResult:
"""
Publish how much of this table is still outstanding, then report the run's result.
async def _delete_old_logs(self, prisma_client: PrismaClient, cutoff_date: datetime) -> int:
The probe is skipped once the budget is spent. It is the one piece of
work that would otherwise be ISSUED after the deadline, and every table
exits through here, including the ones a spent run never started, so
keeping it would put one more statement per table past the bound. A run
that ends this way already reports "budget_exhausted", which tells an
operator the backlog was not drained; the gauge simply keeps its value
from the last run that finished inside its budget.
"""
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)
if remaining is not None:
SpendLogCleanupMetrics.set_rows_remaining(table_name, remaining)
return TableCleanupResult(rows_deleted=rows_deleted, stop_reason=stop_reason)
async def _delete_old_logs(
self, prisma_client: PrismaClient, cutoff_date: datetime, deadline: float
) -> TableCleanupResult:
return await self._delete_old_rows_batched(
prisma_client,
cutoff_date,
table_name="LiteLLM_SpendLogs",
key_columns=("request_id", "startTime"),
time_column="startTime",
deadline=deadline,
)
async def _delete_old_tool_index_rows(self, prisma_client: PrismaClient, cutoff_date: datetime) -> int:
async def _delete_old_tool_index_rows(
self, prisma_client: PrismaClient, cutoff_date: datetime, deadline: float
) -> TableCleanupResult:
# SpendLogToolIndex rows are derived from spend logs, so they expire on the
# same cutoff; rows older than retention point at already-deleted logs.
return await self._delete_old_rows_batched(
@ -189,17 +462,87 @@ class SpendLogCleanup:
table_name="LiteLLM_SpendLogToolIndex",
key_columns=("request_id", "tool_name"),
time_column="start_time",
deadline=deadline,
)
async def _delete_old_autorouter_session_rows(self, prisma_client: PrismaClient, cutoff_date: datetime) -> int:
async def _delete_old_autorouter_session_rows(
self, prisma_client: PrismaClient, cutoff_date: datetime, deadline: float
) -> TableCleanupResult:
return await self._delete_old_rows_batched(
prisma_client,
cutoff_date,
table_name="LiteLLM_AutoRouterSession",
key_columns=("api_key", "session_id", "router_name"),
time_column="last_turn_at",
deadline=deadline,
)
async def _clean_spend_log_tables(
self, prisma_client: PrismaClient, deadline: float
) -> tuple[TableCleanupResult, ...]:
"""
Prune the spend logs and the tool index rows derived from them.
When the table is range-partitioned, whole expired partitions are dropped
first because that reclaims disk immediately. Expired rows can still sit in
the DEFAULT partition (backfill, coverage gaps) or in a partition that spans
the cutoff, so retention still deletes those stragglers row-wise.
"""
cutoff_date: Final = datetime.now(timezone.utc) - timedelta(seconds=float(self.retention_seconds or 0))
verbose_proxy_logger.info("Removing logs older than %s", cutoff_date.isoformat())
# Partition maintenance is DDL taking an ACCESS EXCLUSIVE lock, so it is
# only STARTED while the run still has budget, and each statement carries
# the same timeouts the batches do. Without those, a DROP would queue
# behind any long-running reader for as long as that reader lives, which
# is the one way this job could still outlast its budget without bound.
remaining_timeout_ms: Final = self._remaining_timeout_ms(deadline)
if time.monotonic() >= deadline:
verbose_proxy_logger.info("Run budget already spent, skipping partition maintenance this run")
elif self.general_settings.get(
"use_spend_logs_partitioning", False
) and await self.partition_manager.is_partitioned(prisma_client, remaining_timeout_ms):
await self.partition_manager.ensure_partitions(prisma_client, remaining_timeout_ms)
dropped: Final = await self.partition_manager.drop_partitions_older_than(
prisma_client, cutoff_date, remaining_timeout_ms
)
verbose_proxy_logger.info("Dropped %d expired spend-log partitions: %s", len(dropped), dropped)
logs_result: Final = await self._delete_old_logs(prisma_client, cutoff_date, deadline)
verbose_proxy_logger.info("Deleted %s logs", logs_result.rows_deleted)
index_result: Final = await self._delete_old_tool_index_rows(prisma_client, cutoff_date, deadline)
verbose_proxy_logger.info("Deleted %s expired tool index rows", index_result.rows_deleted)
return (logs_result, index_result)
async def _clean_session_rollup(
self, prisma_client: PrismaClient, retention_seconds: int, deadline: float
) -> tuple[TableCleanupResult, ...]:
"""
Prune auto-router session rollup rows, which carry their own retention horizon.
"""
session_cutoff: Final = datetime.now(timezone.utc) - timedelta(seconds=float(retention_seconds))
sessions_result: Final = await self._delete_old_autorouter_session_rows(prisma_client, session_cutoff, deadline)
verbose_proxy_logger.info("Deleted %s expired auto-router session rollup rows", sessions_result.rows_deleted)
return (sessions_result,)
@staticmethod
def _run_outcome(results: tuple[TableCleanupResult, ...]) -> RunOutcome:
"""
Report the most operationally significant reason the run stopped.
A bound that was hit matters more than a table that simply ran dry, so
those win over "completed", and an abort wins over everything.
"""
reasons: Final = frozenset(result.stop_reason for result in results)
if "aborted" in reasons:
return "aborted"
if "budget_exhausted" in reasons:
return "budget_exhausted"
if "batch_cap_reached" in reasons:
return "batch_cap_reached"
return "completed"
async def cleanup_old_spend_logs(self, prisma_client: PrismaClient) -> None:
"""
Main cleanup function. Deletes old spend logs in batches.
@ -209,16 +552,19 @@ class SpendLogCleanup:
lock_acquired = False
try:
verbose_proxy_logger.info("Cleanup job triggered at %s", datetime.now())
self._refresh_bounds()
delete_spend_logs: Final = self._should_delete_spend_logs()
autorouter_retention_seconds: Final = self._retention_seconds_for(
"maximum_autorouter_session_retention_period"
)
if not delete_spend_logs and autorouter_retention_seconds is None:
SpendLogCleanupMetrics.record_run("skipped_disabled")
return
if delete_spend_logs and self.retention_seconds is None:
verbose_proxy_logger.error("Retention seconds is None, cannot proceed with cleanup")
SpendLogCleanupMetrics.record_run("skipped_disabled")
return
# If we have a pod lock manager, try to acquire the lock
@ -235,43 +581,23 @@ class SpendLogCleanup:
if not lock_acquired:
verbose_proxy_logger.info("Another pod is already running cleanup")
SpendLogCleanupMetrics.record_run("skipped_locked")
return
if delete_spend_logs and self.retention_seconds is not None:
cutoff_date: Final = datetime.now(timezone.utc) - timedelta(seconds=float(self.retention_seconds))
verbose_proxy_logger.info("Removing logs older than %s", cutoff_date.isoformat())
deadline: Final = time.monotonic() + self.run_budget_seconds
if self.general_settings.get(
"use_spend_logs_partitioning", False
) and await self.partition_manager.is_partitioned(prisma_client):
await self.partition_manager.ensure_partitions(prisma_client)
dropped: Final = await self.partition_manager.drop_partitions_older_than(prisma_client, cutoff_date)
verbose_proxy_logger.info(
"Dropped %d expired spend-log partitions: %s",
len(dropped),
dropped,
)
# DROP only reclaims whole expired partitions. Expired rows can
# still sit in the DEFAULT partition (backfill, coverage gaps)
# or in a partition that spans the cutoff, so retention must
# also delete those stragglers row-wise.
total_deleted = await self._delete_old_logs(prisma_client, cutoff_date)
verbose_proxy_logger.info(
"Deleted %s expired logs not covered by dropped partitions", total_deleted
)
else:
total_deleted = await self._delete_old_logs(prisma_client, cutoff_date)
verbose_proxy_logger.info("Deleted %s logs", total_deleted)
spend_log_results: Final = (
await self._clean_spend_log_tables(prisma_client, deadline)
if delete_spend_logs and self.retention_seconds is not None
else ()
)
session_results: Final = (
await self._clean_session_rollup(prisma_client, autorouter_retention_seconds, deadline)
if autorouter_retention_seconds is not None
else ()
)
index_deleted: Final = await self._delete_old_tool_index_rows(prisma_client, cutoff_date)
verbose_proxy_logger.info("Deleted %s expired tool index rows", index_deleted)
if autorouter_retention_seconds is not None:
session_cutoff: Final = datetime.now(timezone.utc) - timedelta(
seconds=float(autorouter_retention_seconds)
)
sessions_deleted: Final = await self._delete_old_autorouter_session_rows(prisma_client, session_cutoff)
verbose_proxy_logger.info("Deleted %s expired auto-router session rollup rows", sessions_deleted)
SpendLogCleanupMetrics.record_run(self._run_outcome(spend_log_results + session_results))
except Exception as e:
# .exception() captures the traceback; str(e) alone on a Prisma/DB
@ -281,6 +607,7 @@ class SpendLogCleanup:
type(e).__name__,
e,
)
SpendLogCleanupMetrics.record_run("aborted")
return # Return after error handling
finally:
# Only release the lock if it was actually acquired

View file

@ -0,0 +1,122 @@
"""
Prometheus metrics for the spend-log retention cleanup job.
The job runs in the background on a single elected pod, so its cost is invisible
from request-path metrics. These instruments make a run's database footprint
observable: how much it deleted, how long each batch took, how much work is
still outstanding, and why a run stopped.
``prometheus_client`` is an optional dependency, so every recorder degrades to a
no-op when it is absent.
"""
from typing import TYPE_CHECKING, Final, Literal, TypeAlias
from litellm._logging import verbose_proxy_logger
if TYPE_CHECKING:
# aliased so the annotations below cannot be mistaken for collections.Counter
from prometheus_client import Counter as PrometheusCounter
from prometheus_client import Gauge as PrometheusGauge
from prometheus_client import Histogram as PrometheusHistogram
RunOutcome: TypeAlias = Literal[
"completed",
"budget_exhausted",
"batch_cap_reached",
"skipped_locked",
"skipped_disabled",
"aborted",
]
_BATCH_DURATION_BUCKETS: Final = (0.005, 0.025, 0.1, 0.25, 0.5, 1.0, 2.5, 5.0, 10.0, 30.0, 60.0)
_TABLE_LABEL: Final = ("table",)
_OUTCOME_LABEL: Final = ("outcome",)
class SpendLogCleanupMetrics:
"""
Lazily-registered Prometheus instruments for the retention cleanup job.
Registration is deferred to first use so that importing this module never
touches the Prometheus registry, which keeps it safe to import from the
proxy regardless of whether Prometheus is a configured callback.
"""
_initialized: bool = False
rows_deleted: "PrometheusCounter | None" = None
batch_duration: "PrometheusHistogram | None" = None
rows_remaining: "PrometheusGauge | None" = None
batch_failures: "PrometheusCounter | None" = None
runs: "PrometheusCounter | None" = None
@classmethod
def _ensure_initialized(cls) -> None:
if cls._initialized:
return
cls._initialized = True
try:
# prometheus_client is an optional extra, so it is resolved here rather
# than at module import: this module is reachable from proxy startup
# regardless of whether Prometheus is a configured callback.
from prometheus_client import Counter, Gauge, Histogram
cls.rows_deleted = Counter(
"litellm_spend_log_cleanup_rows_deleted_total",
"Rows deleted by the spend-log retention cleanup job",
labelnames=_TABLE_LABEL,
)
cls.batch_duration = Histogram(
"litellm_spend_log_cleanup_batch_duration_seconds",
"Wall-clock duration of one retention cleanup delete batch",
labelnames=_TABLE_LABEL,
buckets=_BATCH_DURATION_BUCKETS,
)
cls.rows_remaining = Gauge(
"litellm_spend_log_cleanup_rows_remaining",
"Expired rows still awaiting deletion, counted only up to "
"SPEND_LOG_CLEANUP_REMAINING_COUNT_CAP so the probe itself cannot scan a "
"large table; a value equal to that cap means at least that many remain",
labelnames=_TABLE_LABEL,
multiprocess_mode="livemax",
)
cls.batch_failures = Counter(
"litellm_spend_log_cleanup_batch_failures_total",
"Retention cleanup delete batches that raised",
labelnames=_TABLE_LABEL,
)
cls.runs = Counter(
"litellm_spend_log_cleanup_runs_total",
"Retention cleanup runs, labelled by why the run ended",
labelnames=_OUTCOME_LABEL,
)
except Exception as e: # noqa: BLE001 - a metrics problem must never fail the cleanup run
# Covers the extra being absent, a duplicate registration (repeated
# imports under a test runner), and registry misconfiguration alike.
verbose_proxy_logger.warning("Could not register spend-log cleanup metrics: %s", e)
@classmethod
def record_batch(cls, table_name: str, rows_deleted: int, duration_seconds: float) -> None:
cls._ensure_initialized()
if cls.rows_deleted is not None:
cls.rows_deleted.labels(table=table_name).inc(rows_deleted)
if cls.batch_duration is not None:
cls.batch_duration.labels(table=table_name).observe(duration_seconds)
@classmethod
def record_batch_failure(cls, table_name: str) -> None:
cls._ensure_initialized()
if cls.batch_failures is not None:
cls.batch_failures.labels(table=table_name).inc()
@classmethod
def set_rows_remaining(cls, table_name: str, remaining: int) -> None:
cls._ensure_initialized()
if cls.rows_remaining is not None:
cls.rows_remaining.labels(table=table_name).set(remaining)
@classmethod
def record_run(cls, outcome: RunOutcome) -> None:
cls._ensure_initialized()
if cls.runs is not None:
cls.runs.labels(outcome=outcome).inc()

View file

@ -14,8 +14,9 @@ keeps the batched-DELETE path, so existing deployments are untouched.
"""
import re
from collections.abc import Callable
from datetime import date, datetime, timedelta, timezone
from typing import Final
from typing import TYPE_CHECKING, Final, TypeAlias
from litellm._logging import verbose_proxy_logger
from litellm.constants import (
@ -23,8 +24,23 @@ from litellm.constants import (
SPEND_LOG_PARTITION_PRECREATE_AHEAD,
)
if TYPE_CHECKING:
from litellm.proxy.utils import PrismaClient
SPEND_LOGS_TABLE: Final = "LiteLLM_SpendLogs"
RemainingTimeoutMs: TypeAlias = Callable[[], "int | None"]
"""
The per-statement bound in milliseconds, or None once the caller's budget is
spent.
Injected rather than passed as a number so it is re-evaluated before EVERY
statement: a value read once at entry would let a loop issue N statements each
bounded by the budget that was left before the first of them, which is not a
bound on the loop at all. The caller owns the policy; this module only asks how
much time it may still use.
"""
PartitionInterval = str # "day" | "week" | "month"
VALID_PARTITION_INTERVALS: Final = {"day", "week", "month"}
@ -116,21 +132,26 @@ class SpendLogsPartitionManager:
self.interval = interval
self.precreate_ahead = precreate_ahead
async def is_partitioned(self, prisma_client) -> bool:
async def is_partitioned(self, prisma_client: "PrismaClient", remaining_timeout_ms: RemainingTimeoutMs) -> bool:
budget_ms: Final = remaining_timeout_ms()
if budget_ms is None:
return False
try:
rows: Final = await prisma_client.db.query_raw(
"""
SELECT EXISTS (
SELECT 1
FROM pg_partitioned_table pt
JOIN pg_class c ON c.oid = pt.partrelid
JOIN pg_namespace n ON n.oid = c.relnamespace
WHERE c.relname = $1
AND n.nspname = current_schema()
) AS partitioned
""",
SPEND_LOGS_TABLE,
)
async with prisma_client.db.tx() as tx:
await tx.execute_raw(f"SET LOCAL statement_timeout = {budget_ms}")
rows: Final = await tx.query_raw(
"""
SELECT EXISTS (
SELECT 1
FROM pg_partitioned_table pt
JOIN pg_class c ON c.oid = pt.partrelid
JOIN pg_namespace n ON n.oid = c.relnamespace
WHERE c.relname = $1
AND n.nspname = current_schema()
) AS partitioned
""",
SPEND_LOGS_TABLE,
)
except Exception as e:
verbose_proxy_logger.warning(
"Could not determine if %s is partitioned, assuming it is not: %s",
@ -140,7 +161,25 @@ class SpendLogsPartitionManager:
return False
return bool(rows and rows[0].get("partitioned"))
async def ensure_partitions(self, prisma_client) -> list[str]:
@staticmethod
async def _execute_bounded_ddl(prisma_client: "PrismaClient", statement: str, timeout_ms: int) -> None:
"""
Run one DDL statement under a Postgres statement and lock timeout.
Partition DDL takes an ACCESS EXCLUSIVE lock, so an unbounded statement
queues behind any long-running reader for as long as that reader lives,
and the caller's run budget cannot cut it short. lock_timeout bounds the
wait for the lock and statement_timeout bounds the work itself, so a
partition this run cannot get is simply left for the next one.
"""
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}")
await tx.execute_raw(statement)
async def ensure_partitions(
self, prisma_client: "PrismaClient", remaining_timeout_ms: RemainingTimeoutMs
) -> list[str]:
"""
Ensure the current and upcoming partitions exist, returning the names
now present. CREATE TABLE IF NOT EXISTS is a no-op for partitions that
@ -150,42 +189,61 @@ class SpendLogsPartitionManager:
for name, lower, upper in upcoming_partitions(
datetime.now(timezone.utc).date(), self.interval, self.precreate_ahead
):
budget_ms = remaining_timeout_ms()
if budget_ms is None:
verbose_proxy_logger.info("Run budget spent, leaving the remaining partitions for the next run")
break
try:
await prisma_client.db.execute_raw(
await self._execute_bounded_ddl(
prisma_client,
f'CREATE TABLE IF NOT EXISTS "{name}" '
f'PARTITION OF "{SPEND_LOGS_TABLE}" '
f"FOR VALUES FROM ('{lower.isoformat()}') TO ('{upper.isoformat()}')"
f"FOR VALUES FROM ('{lower.isoformat()}') TO ('{upper.isoformat()}')",
budget_ms,
)
ensured.append(name)
except Exception as e:
verbose_proxy_logger.warning("Failed to ensure spend-log partition %s: %s", name, e)
return ensured
async def _list_partitions(self, prisma_client) -> list[tuple[str, datetime | None]]:
rows: Final = await prisma_client.db.query_raw(
"""
SELECT c.relname AS name,
pg_get_expr(c.relpartbound, c.oid) AS bound
FROM pg_inherits i
JOIN pg_class c ON c.oid = i.inhrelid
JOIN pg_class p ON p.oid = i.inhparent
JOIN pg_namespace n ON n.oid = p.relnamespace
WHERE p.relname = $1
AND n.nspname = current_schema()
""",
SPEND_LOGS_TABLE,
)
async def _list_partitions(
self, prisma_client: "PrismaClient", timeout_ms: int
) -> list[tuple[str, datetime | None]]:
async with prisma_client.db.tx() as tx:
await tx.execute_raw(f"SET LOCAL statement_timeout = {timeout_ms}")
rows: Final = await tx.query_raw(
"""
SELECT c.relname AS name,
pg_get_expr(c.relpartbound, c.oid) AS bound
FROM pg_inherits i
JOIN pg_class c ON c.oid = i.inhrelid
JOIN pg_class p ON p.oid = i.inhparent
JOIN pg_namespace n ON n.oid = p.relnamespace
WHERE p.relname = $1
AND n.nspname = current_schema()
""",
SPEND_LOGS_TABLE,
)
return [(row["name"], parse_partition_upper_bound(row.get("bound") or "")) for row in rows]
async def drop_partitions_older_than(self, prisma_client, cutoff: datetime) -> list[str]:
async def drop_partitions_older_than(
self, prisma_client: "PrismaClient", cutoff: datetime, remaining_timeout_ms: RemainingTimeoutMs
) -> list[str]:
"""DROP every partition whose whole range is older than `cutoff`."""
list_budget_ms: Final = remaining_timeout_ms()
if list_budget_ms is None:
return []
cutoff_naive: Final = cutoff.astimezone(timezone.utc).replace(tzinfo=None)
partitions: Final = await self._list_partitions(prisma_client)
partitions: Final = await self._list_partitions(prisma_client, list_budget_ms)
to_drop: Final = select_partitions_to_drop(partitions, cutoff_naive)
dropped: Final[list[str]] = []
for name in to_drop:
budget_ms = remaining_timeout_ms()
if budget_ms is None:
verbose_proxy_logger.info("Run budget spent, leaving the remaining partitions for the next run")
break
try:
await prisma_client.db.execute_raw(f'DROP TABLE IF EXISTS "{name}"')
await self._execute_bounded_ddl(prisma_client, f'DROP TABLE IF EXISTS "{name}"', budget_ms)
dropped.append(name)
except Exception as e:
verbose_proxy_logger.warning("Failed to drop spend-log partition %s: %s", name, e)

View file

@ -367,7 +367,10 @@ from litellm.proxy.config_resolvers.alerting import (
)
from litellm.proxy.container_endpoints.endpoints import router as container_router
from litellm.proxy.credential_endpoints.endpoints import router as credential_router
from litellm.proxy.db.db_transaction_queue.spend_log_cleanup import SpendLogCleanup
from litellm.proxy.db.db_transaction_queue.spend_log_cleanup import (
SPEND_LOG_CLEANUP_BOUND_SETTINGS,
SpendLogCleanup,
)
from litellm.proxy.db.exception_handler import (
PrismaDBExceptionHandler,
call_with_db_reconnect_retry,
@ -4079,6 +4082,7 @@ class ProxyConfig:
# precedence over stale DB-cached values for these specific keys
# during periodic config reloads (_update_general_settings).
self._yaml_general_settings_keys: set[str] = set() # mutable-ok: populated once at startup, read-only thereafter # fmt: skip
self._yaml_spend_log_cleanup_bounds: dict[str, object] = {} # mutable-ok: snapshot of YAML bounds at load time # fmt: skip
def is_yaml(self, config_file_path: str) -> bool:
if not os.path.isfile(config_file_path):
@ -5015,6 +5019,12 @@ class ProxyConfig:
# These keys take precedence over DB-cached values during periodic
# reloads (see _update_general_settings).
self._yaml_general_settings_keys = set(general_settings.keys()) # mutable-ok: snapshot of YAML keys at load time # fmt: skip
# The VALUES matter for the cleanup bounds, not just which keys were
# set: clearing one from the dashboard has to fall back to what the
# YAML declared, and a set of names cannot answer that.
self._yaml_spend_log_cleanup_bounds = { # mutable-ok: snapshot of YAML bounds at load time # fmt: skip
key: general_settings[key] for key in SPEND_LOG_CLEANUP_BOUND_SETTINGS if key in general_settings
}
### LOAD KEY MANAGEMENT SETTINGS FIRST (needed for custom secret manager) ###
key_management_settings: Final = general_settings.get("key_management_settings", None)
@ -6299,6 +6309,18 @@ class ProxyConfig:
if old_session_value != new_session_value:
await self._reschedule_spend_log_cleanup_job()
## SPEND LOG CLEANUP BOUNDS ##
# The dashboard writes these straight to the DB, so without copying them
# here the running cleanup job never sees them. A key the DB no longer
# carries was cleared from the dashboard, and falls back to whatever
# config.yaml declared, or to None (the shipped default) when it declared
# nothing. Leaving the deleted DB value in memory would keep enforcing the
# bound the operator just removed.
for cleanup_key in SPEND_LOG_CLEANUP_BOUND_SETTINGS:
general_settings[cleanup_key] = _general_settings.get(
cleanup_key, self._yaml_spend_log_cleanup_bounds.get(cleanup_key)
)
for key in (
"user_url_allowed_hosts",
"user_url_validation",
@ -15529,6 +15551,10 @@ _GENERAL_SETTINGS_CONFIG_LIST_FIELD_TYPES: Final[Mapping[str, str]] = MappingPro
"store_model_in_db": "Boolean",
"store_prompts_in_spend_logs": "Boolean",
"maximum_spend_logs_retention_period": "String",
"maximum_spend_logs_cleanup_batch_size": "Integer",
"maximum_spend_logs_cleanup_max_batches": "Integer",
"maximum_spend_logs_cleanup_run_budget": "String",
"maximum_spend_logs_cleanup_batch_timeout": "String",
"mcp_internal_ip_ranges": "List",
"mcp_trusted_proxy_ranges": "List",
"mcp_xff_num_trusted_hops": "Integer",

View file

@ -3,6 +3,7 @@ Tests for SpendLogsPartitionManager: partition naming/bounds math, retention
selection, the non-partitioned no-op safety path, and the drop/ensure SQL flow.
"""
from contextlib import asynccontextmanager
from datetime import date, datetime, timezone
from unittest.mock import AsyncMock, MagicMock
@ -19,6 +20,46 @@ from litellm.proxy.db.db_transaction_queue.spend_logs_partition_manager import (
)
DDL_TIMEOUT_MS = 30000
def _budget(ms: "int | None" = DDL_TIMEOUT_MS):
"""The injected per-statement bound: a callable re-read before each statement."""
return lambda: ms
def _wire_tx(db) -> list[str]:
"""
Model the prisma seam the partition DDL uses.
Every statement this manager issues, DDL and catalog query alike, runs inside
db.tx() so it can carry SET LOCAL timeouts. Those SET LOCAL statements are
collected in the returned list rather than forwarded, so assertions on
db.execute_raw and db.query_raw still see only the real statements.
"""
session_settings: list[str] = []
@asynccontextmanager
async def _tx():
tx = MagicMock()
async def _execute_raw(sql, *args):
if sql.lstrip().upper().startswith("SET LOCAL"):
session_settings.append(sql.strip())
return 0
return await db.execute_raw(sql, *args)
async def _query_raw(sql, *args):
return await db.query_raw(sql, *args)
tx.execute_raw = _execute_raw
tx.query_raw = _query_raw
yield tx
db.tx = _tx
return session_settings
def test_period_start_per_interval():
d = date(2026, 6, 3) # a Wednesday
assert period_start(d, "day") == date(2026, 6, 3)
@ -78,11 +119,13 @@ async def test_is_partitioned_true_and_false():
client_true = MagicMock()
client_true.db.query_raw = AsyncMock(return_value=[{"partitioned": True}])
assert await mgr.is_partitioned(client_true) is True
_wire_tx(client_true.db)
assert await mgr.is_partitioned(client_true, _budget()) is True
client_false = MagicMock()
client_false.db.query_raw = AsyncMock(return_value=[{"partitioned": False}])
assert await mgr.is_partitioned(client_false) is False
_wire_tx(client_false.db)
assert await mgr.is_partitioned(client_false, _budget()) is False
@pytest.mark.asyncio
@ -94,13 +137,14 @@ async def test_catalog_queries_are_scoped_to_current_schema():
mgr = SpendLogsPartitionManager()
client = MagicMock()
client.db.query_raw = AsyncMock(return_value=[])
_wire_tx(client.db)
await mgr.is_partitioned(client)
await mgr.is_partitioned(client, _budget())
is_partitioned_sql = client.db.query_raw.call_args.args[0]
assert "pg_namespace" in is_partitioned_sql
assert "current_schema()" in is_partitioned_sql
await mgr._list_partitions(client)
await mgr._list_partitions(client, DDL_TIMEOUT_MS)
list_sql = client.db.query_raw.call_args.args[0]
assert "pg_namespace" in list_sql
assert "current_schema()" in list_sql
@ -112,7 +156,10 @@ async def test_is_partitioned_swallows_errors_and_returns_false():
mgr = SpendLogsPartitionManager()
client = MagicMock()
client.db.query_raw = AsyncMock(side_effect=Exception("db down"))
assert await mgr.is_partitioned(client) is False
# Wire the real seam: without it the async with itself raises, and the test
# would pass on the wrong exception.
_wire_tx(client.db)
assert await mgr.is_partitioned(client, _budget()) is False
@pytest.mark.asyncio
@ -133,9 +180,10 @@ async def test_drop_partitions_older_than_drops_expired_only():
]
)
client.db.execute_raw = AsyncMock(return_value=0)
_wire_tx(client.db)
cutoff = datetime(2026, 6, 5, 0, 0, 0, tzinfo=timezone.utc)
dropped = await mgr.drop_partitions_older_than(client, cutoff)
dropped = await mgr.drop_partitions_older_than(client, cutoff, _budget())
assert dropped == ["LiteLLM_SpendLogs_p20260601"]
executed = " ".join(call.args[0] for call in client.db.execute_raw.call_args_list)
@ -149,8 +197,9 @@ async def test_ensure_partitions_issues_create_for_each_period():
mgr = SpendLogsPartitionManager(interval="day", precreate_ahead=2)
client = MagicMock()
client.db.execute_raw = AsyncMock(return_value=0)
_wire_tx(client.db)
created = await mgr.ensure_partitions(client)
created = await mgr.ensure_partitions(client, _budget())
assert len(created) == 3 # current + 2 ahead
assert client.db.execute_raw.await_count == 3
@ -159,6 +208,105 @@ async def test_ensure_partitions_issues_create_for_each_period():
assert "CREATE TABLE IF NOT EXISTS" in first_sql
@pytest.mark.asyncio
async def test_partition_ddl_carries_a_statement_and_lock_timeout():
"""
Partition DDL takes an ACCESS EXCLUSIVE lock, so an unbounded DROP queues
behind any long-running reader for as long as that reader lives. That is the
one path by which cleanup could outlast its run budget without bound, and
lock_timeout is what bounds the wait rather than only the work.
"""
mgr = SpendLogsPartitionManager(interval="day", precreate_ahead=0)
client = MagicMock()
client.db.execute_raw = AsyncMock(return_value=0)
client.db.query_raw = AsyncMock(
return_value=[
{
"name": "LiteLLM_SpendLogs_p20260601",
"bound": "FOR VALUES FROM ('2026-06-01 00:00:00') TO ('2026-06-02 00:00:00')",
}
]
)
session_settings = _wire_tx(client.db)
await mgr.ensure_partitions(client, _budget(7000))
await mgr.drop_partitions_older_than(client, datetime(2026, 6, 5, tzinfo=timezone.utc), _budget(7000))
# Three statements were issued: the CREATE, the catalog list the drop needs,
# and the DROP. All three carry a statement timeout; only the two that take
# a lock also carry a lock timeout, since the catalog read takes none.
assert session_settings.count("SET LOCAL statement_timeout = 7000") == 3
assert session_settings.count("SET LOCAL lock_timeout = 7000") == 2
@pytest.mark.asyncio
async def test_catalog_queries_carry_a_statement_timeout():
"""
Bounding only the DDL leaves the two catalog lookups as statements this job
issues with no bound at all, so a run could still outlast its budget waiting
on one. Every statement the manager issues carries the caller's timeout.
"""
mgr = SpendLogsPartitionManager()
client = MagicMock()
client.db.query_raw = AsyncMock(return_value=[])
session_settings = _wire_tx(client.db)
await mgr.is_partitioned(client, _budget(4000))
assert session_settings == ["SET LOCAL statement_timeout = 4000"], (
f"is_partitioned issued no statement timeout: {session_settings}"
)
session_settings.clear()
await mgr._list_partitions(client, 4000)
assert session_settings == ["SET LOCAL statement_timeout = 4000"], (
f"_list_partitions issued no statement timeout: {session_settings}"
)
@pytest.mark.asyncio
async def test_partition_loops_stop_when_the_budget_runs_out_mid_way():
"""
Each loop issues one statement per partition, so a bound read once at entry
would let N statements each run for the budget that was left before the
first of them. The bound is re-read per statement and the loop stops.
"""
mgr = SpendLogsPartitionManager(interval="day", precreate_ahead=4)
client = MagicMock()
client.db.execute_raw = AsyncMock(return_value=0)
_wire_tx(client.db)
# Budget for two statements, then spent.
calls = {"n": 0}
def budget() -> "int | None":
calls["n"] += 1
return 5000 if calls["n"] <= 2 else None
created = await mgr.ensure_partitions(client, budget)
assert len(created) == 2, f"the loop ran past its budget and created {len(created)}"
assert client.db.execute_raw.await_count == 2
@pytest.mark.asyncio
async def test_partition_maintenance_issues_nothing_when_the_budget_is_already_spent():
"""A run with no budget left must not issue even the catalog lookups."""
mgr = SpendLogsPartitionManager(interval="day", precreate_ahead=2)
client = MagicMock()
client.db.execute_raw = AsyncMock(return_value=0)
client.db.query_raw = AsyncMock(return_value=[])
_wire_tx(client.db)
spent = _budget(None)
assert await mgr.is_partitioned(client, spent) is False
assert await mgr.ensure_partitions(client, spent) == []
assert await mgr.drop_partitions_older_than(client, datetime(2026, 6, 5, tzinfo=timezone.utc), spent) == []
client.db.execute_raw.assert_not_awaited()
client.db.query_raw.assert_not_awaited()
def test_unsupported_interval_raises():
with pytest.raises(ValueError):
period_start(date(2026, 6, 1), "year")
@ -178,8 +326,9 @@ async def test_ensure_partitions_continues_when_one_create_fails():
mgr = SpendLogsPartitionManager(interval="day", precreate_ahead=2)
client = MagicMock()
client.db.execute_raw = AsyncMock(side_effect=[0, Exception("overlap"), 0])
_wire_tx(client.db)
created = await mgr.ensure_partitions(client)
created = await mgr.ensure_partitions(client, _budget())
# the failed partition is skipped, the others still created
assert len(created) == 2
@ -202,8 +351,9 @@ async def test_invalid_interval_does_not_abort_ensure_partitions():
mgr = SpendLogsPartitionManager(interval="fortnight", precreate_ahead=1)
client = MagicMock()
client.db.execute_raw = AsyncMock(return_value=0)
_wire_tx(client.db)
created = await mgr.ensure_partitions(client)
created = await mgr.ensure_partitions(client, _budget())
assert len(created) == 2 # current + 1 ahead, day-based fallback
@ -225,9 +375,10 @@ async def test_drop_partitions_continues_when_one_drop_fails():
]
)
client.db.execute_raw = AsyncMock(side_effect=[Exception("locked"), 0])
_wire_tx(client.db)
cutoff = datetime(2026, 6, 10, 0, 0, 0, tzinfo=timezone.utc)
dropped = await mgr.drop_partitions_older_than(client, cutoff)
dropped = await mgr.drop_partitions_older_than(client, cutoff, _budget())
# both were eligible; the first drop failed so only the second is reported
assert dropped == ["LiteLLM_SpendLogs_p20260602"]

View file

@ -6872,6 +6872,91 @@ async def test_update_general_settings_propagates_apply_user_budget_to_team_keys
assert ps.general_settings["apply_user_budget_to_team_keys"] is True
@pytest.mark.asyncio
async def test_update_general_settings_propagates_spend_log_cleanup_bounds():
"""The dashboard writes the cleanup bounds straight to the DB config, so
without runtime propagation the scheduled job never sees them and the knobs
do nothing until the process restarts."""
from litellm.proxy.db.db_transaction_queue.spend_log_cleanup import (
SPEND_LOG_CLEANUP_BOUND_SETTINGS,
)
from litellm.proxy.proxy_server import ProxyConfig
proxy_config = ProxyConfig()
db_settings = {
"maximum_spend_logs_cleanup_batch_size": 2000,
"maximum_spend_logs_cleanup_max_batches": 250,
"maximum_spend_logs_cleanup_run_budget": "90s",
"maximum_spend_logs_cleanup_batch_timeout": "10s",
}
assert set(db_settings) == set(SPEND_LOG_CLEANUP_BOUND_SETTINGS)
with patch("litellm.proxy.proxy_server.general_settings", {}):
await proxy_config._update_general_settings(db_general_settings=db_settings)
import litellm.proxy.proxy_server as ps
assert {key: ps.general_settings.get(key) for key in db_settings} == db_settings
@pytest.mark.asyncio
async def test_update_general_settings_clears_a_spend_log_cleanup_bound_dropped_from_the_db():
"""Blanking the field in the dashboard deletes the key outright, so leaving
the last value in memory would keep a bound the operator just removed."""
from litellm.proxy.proxy_server import ProxyConfig
proxy_config = ProxyConfig()
with patch(
"litellm.proxy.proxy_server.general_settings",
{"maximum_spend_logs_cleanup_run_budget": "90s", "maximum_spend_logs_cleanup_batch_timeout": "10s"},
):
await proxy_config._update_general_settings(
db_general_settings={"maximum_spend_logs_cleanup_batch_timeout": "10s"}
)
import litellm.proxy.proxy_server as ps
assert ps.general_settings["maximum_spend_logs_cleanup_run_budget"] is None
assert ps.general_settings["maximum_spend_logs_cleanup_batch_timeout"] == "10s"
@pytest.mark.asyncio
async def test_update_general_settings_keeps_a_yaml_set_spend_log_cleanup_bound():
"""A YAML-set bound never appears in the DB object, so treating its absence
as a dashboard clear would discard the deployed config on every reload."""
from litellm.proxy.proxy_server import ProxyConfig
proxy_config = ProxyConfig()
proxy_config._yaml_spend_log_cleanup_bounds = {"maximum_spend_logs_cleanup_run_budget": "90s"}
with patch("litellm.proxy.proxy_server.general_settings", {"maximum_spend_logs_cleanup_run_budget": "90s"}):
await proxy_config._update_general_settings(db_general_settings={"store_model_in_db": True})
import litellm.proxy.proxy_server as ps
assert ps.general_settings["maximum_spend_logs_cleanup_run_budget"] == "90s"
@pytest.mark.asyncio
async def test_update_general_settings_clearing_a_db_override_falls_back_to_the_yaml_bound():
"""Clearing a dashboard override of a YAML-declared bound must restore the
YAML value. Leaving the deleted override in memory would keep enforcing the
bound the operator just removed, until the process restarted."""
from litellm.proxy.proxy_server import ProxyConfig
proxy_config = ProxyConfig()
proxy_config._yaml_spend_log_cleanup_bounds = {"maximum_spend_logs_cleanup_run_budget": "90s"}
# Memory currently holds the dashboard override, and the DB no longer carries it.
with patch("litellm.proxy.proxy_server.general_settings", {"maximum_spend_logs_cleanup_run_budget": "30s"}):
await proxy_config._update_general_settings(db_general_settings={"store_model_in_db": True})
import litellm.proxy.proxy_server as ps
assert ps.general_settings["maximum_spend_logs_cleanup_run_budget"] == "90s"
@pytest.mark.asyncio
async def test_update_general_settings_apply_user_budget_to_team_keys_yaml_wins():
"""A DB value must not silently override an explicit YAML setting on reload."""

View file

@ -2,12 +2,66 @@
Test cases for spend log cleanup functionality
"""
import asyncio
import math
import time
from contextlib import asynccontextmanager
from datetime import datetime, timedelta, timezone
from unittest.mock import AsyncMock, MagicMock
import pytest
from litellm.proxy.db.db_transaction_queue.spend_log_cleanup import SpendLogCleanup
from litellm.constants import (
SPEND_LOG_CLEANUP_BATCH_SIZE,
SPEND_LOG_CLEANUP_REMAINING_COUNT_CAP,
SPEND_LOG_CLEANUP_RUN_BUDGET_SECONDS,
)
from litellm.proxy.db.db_transaction_queue.spend_log_cleanup import (
SPEND_LOG_CLEANUP_BOUND_SETTINGS,
SpendLogCleanup,
TableCleanupResult,
)
from litellm.proxy.db.db_transaction_queue.spend_log_cleanup_metrics import (
SpendLogCleanupMetrics,
)
def _far_deadline() -> float:
"""A run deadline far enough out that only the other bounds can stop a batch loop."""
return time.monotonic() + 3600
def _wire_tx(db):
"""
Model the prisma seam the cleanup job actually uses.
Every statement the job issues runs inside db.tx() so it can carry a SET
LOCAL statement_timeout. Batch and probe statements are forwarded to
db.execute_raw and db.query_raw, which is what tests configure and assert
on, while the SET LOCAL statements are answered here so they neither consume
a side_effect entry nor show up in the recorded call list. Lookup is
deferred to call time so this can be wired before a test assigns its own
execute_raw.
"""
@asynccontextmanager
async def _tx():
tx = MagicMock()
async def _execute_raw(sql, *args):
if sql.lstrip().upper().startswith("SET LOCAL"):
return 0
return await db.execute_raw(sql, *args)
async def _query_raw(sql, *args):
return await db.query_raw(sql, *args)
tx.execute_raw = _execute_raw
tx.query_raw = _query_raw
yield tx
db.tx = _tx
db.query_raw = AsyncMock(return_value=[{"remaining": 0}])
def test_spend_log_cleanup_cron_scheduling():
@ -49,6 +103,7 @@ def test_spend_log_cleanup_cron_scheduler_integration():
# Mock scheduler
mock_scheduler = MagicMock()
mock_prisma_client = MagicMock()
_wire_tx(mock_prisma_client.db)
mock_cleanup_instance = MagicMock()
# Test Case 1: Cron-based scheduling
@ -155,7 +210,9 @@ async def test_cleanup_old_spend_logs_batch_deletion():
# Setup Prisma client
mock_prisma_client = MagicMock()
_wire_tx(mock_prisma_client.db)
mock_db = MagicMock()
_wire_tx(mock_db)
# Mock execute_raw to return deleted counts (3 spend-log batches, then the
# tool-index cleanup's first batch returning 0)
@ -207,7 +264,9 @@ async def test_cleanup_old_spend_logs_retention_period_cutoff():
"""
# Setup Prisma client
mock_prisma_client = MagicMock()
_wire_tx(mock_prisma_client.db)
mock_db = MagicMock()
_wire_tx(mock_db)
mock_db.execute_raw = AsyncMock(return_value=0)
mock_prisma_client.db = mock_db
@ -244,6 +303,7 @@ async def test_cleanup_drops_partitions_when_enabled_and_partitioned():
from unittest.mock import AsyncMock, MagicMock
mock_prisma_client = MagicMock()
_wire_tx(mock_prisma_client.db)
mock_prisma_client.db.execute_raw = AsyncMock(return_value=0)
partition_manager = MagicMock()
@ -285,6 +345,7 @@ async def test_cleanup_uses_delete_when_partitioning_not_enabled():
from unittest.mock import AsyncMock, MagicMock
mock_prisma_client = MagicMock()
_wire_tx(mock_prisma_client.db)
mock_prisma_client.db.execute_raw = AsyncMock(side_effect=[10, 0, 0])
partition_manager = MagicMock()
@ -316,6 +377,7 @@ async def test_cleanup_uses_delete_when_not_partitioned():
from unittest.mock import AsyncMock, MagicMock
mock_prisma_client = MagicMock()
_wire_tx(mock_prisma_client.db)
mock_prisma_client.db.execute_raw = AsyncMock(side_effect=[10, 0, 0])
partition_manager = MagicMock()
@ -346,6 +408,7 @@ async def test_cleanup_old_spend_logs_no_retention_period():
Test that no logs are deleted when no retention period is set
"""
mock_prisma_client = MagicMock()
_wire_tx(mock_prisma_client.db)
mock_prisma_client.db.execute_raw = AsyncMock()
cleaner = SpendLogCleanup(general_settings={}) # no retention
@ -361,6 +424,7 @@ async def test_lock_not_released_when_not_acquired():
before the lock is ever acquired.
"""
mock_prisma_client = MagicMock()
_wire_tx(mock_prisma_client.db)
mock_prisma_client.db.execute_raw = AsyncMock()
mock_redis_cache = MagicMock()
@ -418,7 +482,9 @@ async def test_delete_old_logs_aborts_on_non_int_execute_raw_return():
"""should abort deletion loop immediately when execute_raw returns a non-int
(e.g. None or dict), preventing an infinite loop."""
mock_prisma_client = MagicMock()
_wire_tx(mock_prisma_client.db)
mock_db = MagicMock()
_wire_tx(mock_db)
mock_db.execute_raw = AsyncMock(return_value=None)
mock_prisma_client.db = mock_db
@ -427,17 +493,19 @@ async def test_delete_old_logs_aborts_on_non_int_execute_raw_return():
)
cutoff_date = datetime.now(timezone.utc) - timedelta(days=7)
total_deleted = await cleaner._delete_old_logs(mock_prisma_client, cutoff_date)
result = await cleaner._delete_old_logs(mock_prisma_client, cutoff_date, _far_deadline())
assert mock_db.execute_raw.call_count == 1
assert total_deleted == 0
assert result.rows_deleted == 0
@pytest.mark.asyncio
async def test_delete_old_logs_continues_on_valid_int_return():
"""should continue deletion loop across batches when execute_raw returns valid int counts."""
mock_prisma_client = MagicMock()
_wire_tx(mock_prisma_client.db)
mock_db = MagicMock()
_wire_tx(mock_db)
mock_db.execute_raw = AsyncMock(side_effect=[500, 300, 0])
mock_prisma_client.db = mock_db
@ -446,35 +514,37 @@ async def test_delete_old_logs_continues_on_valid_int_return():
)
cutoff_date = datetime.now(timezone.utc) - timedelta(days=7)
total_deleted = await cleaner._delete_old_logs(mock_prisma_client, cutoff_date)
result = await cleaner._delete_old_logs(mock_prisma_client, cutoff_date, _far_deadline())
assert mock_db.execute_raw.call_count == 3
assert total_deleted == 800
assert result.rows_deleted == 800
@pytest.mark.asyncio
async def test_delete_old_rows_stops_at_max_batches(monkeypatch):
"""The run-loop backstop must halt a cleanup that keeps finding rows, so a
huge backlog is spread across scheduled runs instead of one unbounded loop."""
import litellm.proxy.db.db_transaction_queue.spend_log_cleanup as cleanup_module
monkeypatch.setattr(cleanup_module, "SPEND_LOG_RUN_LOOPS", 2)
async def test_delete_old_rows_stops_at_max_batches():
"""The batch cap must halt a cleanup that keeps finding rows, so a huge
backlog is spread across scheduled runs instead of one unbounded loop, and
the operator-facing knob must mean exactly the number of statements it names."""
mock_prisma_client = MagicMock()
_wire_tx(mock_prisma_client.db)
mock_db = MagicMock()
_wire_tx(mock_db)
mock_db.execute_raw = AsyncMock(return_value=1000)
mock_prisma_client.db = mock_db
cleaner = SpendLogCleanup(
general_settings={"maximum_spend_logs_retention_period": "7d"}
general_settings={
"maximum_spend_logs_retention_period": "7d",
"maximum_spend_logs_cleanup_max_batches": 2,
}
)
cutoff_date = datetime.now(timezone.utc) - timedelta(days=7)
total_deleted = await cleaner._delete_old_logs(mock_prisma_client, cutoff_date)
result = await cleaner._delete_old_logs(mock_prisma_client, cutoff_date, _far_deadline())
# run_count exceeds the cap only after 3 full batches (0, 1, 2)
assert mock_db.execute_raw.call_count == 3
assert total_deleted == 3000
assert mock_db.execute_raw.call_count == 2
assert result.rows_deleted == 2000
assert result.stop_reason == "batch_cap_reached"
@pytest.mark.asyncio
@ -482,7 +552,9 @@ async def test_delete_old_tool_index_rows_deletes_on_composite_key():
"""Tool index rows are derived from spend logs and expire on the same cutoff;
the delete must match on the table's composite primary key."""
mock_prisma_client = MagicMock()
_wire_tx(mock_prisma_client.db)
mock_db = MagicMock()
_wire_tx(mock_db)
mock_db.execute_raw = AsyncMock(side_effect=[5, 0])
mock_prisma_client.db = mock_db
@ -491,9 +563,9 @@ async def test_delete_old_tool_index_rows_deletes_on_composite_key():
)
cutoff_date = datetime.now(timezone.utc) - timedelta(days=7)
total_deleted = await cleaner._delete_old_tool_index_rows(mock_prisma_client, cutoff_date)
result = await cleaner._delete_old_tool_index_rows(mock_prisma_client, cutoff_date, _far_deadline())
assert total_deleted == 5
assert result.rows_deleted == 5
delete_sql = mock_db.execute_raw.call_args_list[0][0][0]
assert 'DELETE FROM "LiteLLM_SpendLogToolIndex"' in delete_sql
assert 'WHERE ("request_id", "tool_name") IN' in delete_sql
@ -513,7 +585,9 @@ async def test_delete_old_logs_continues_after_single_batch_failure(monkeypatch)
)
mock_prisma_client = MagicMock()
_wire_tx(mock_prisma_client.db)
mock_db = MagicMock()
_wire_tx(mock_db)
# batch 1 succeeds, batch 2 raises (one-off DB timeout), batches 3-4 succeed,
# batch 5 returns 0 → loop exits naturally.
mock_db.execute_raw = AsyncMock(
@ -526,11 +600,11 @@ async def test_delete_old_logs_continues_after_single_batch_failure(monkeypatch)
)
cutoff_date = datetime.now(timezone.utc) - timedelta(days=7)
total_deleted = await cleaner._delete_old_logs(mock_prisma_client, cutoff_date)
result = await cleaner._delete_old_logs(mock_prisma_client, cutoff_date, _far_deadline())
# All 5 batches should have been attempted; 100 + 200 + 50 = 350 deleted.
assert mock_db.execute_raw.call_count == 5
assert total_deleted == 350
assert result.rows_deleted == 350
@pytest.mark.asyncio
@ -548,7 +622,9 @@ async def test_delete_old_logs_aborts_after_consecutive_failures(monkeypatch):
)
mock_prisma_client = MagicMock()
_wire_tx(mock_prisma_client.db)
mock_db = MagicMock()
_wire_tx(mock_db)
# Every batch raises — must abort after exactly 3 attempts, not loop forever.
mock_db.execute_raw = AsyncMock(
side_effect=ConnectionError("simulated persistent DB outage")
@ -560,10 +636,10 @@ async def test_delete_old_logs_aborts_after_consecutive_failures(monkeypatch):
)
cutoff_date = datetime.now(timezone.utc) - timedelta(days=7)
total_deleted = await cleaner._delete_old_logs(mock_prisma_client, cutoff_date)
result = await cleaner._delete_old_logs(mock_prisma_client, cutoff_date, _far_deadline())
assert mock_db.execute_raw.call_count == 3
assert total_deleted == 0
assert result.rows_deleted == 0
@pytest.mark.asyncio
@ -580,7 +656,9 @@ async def test_delete_old_logs_resets_consecutive_failures_on_success(monkeypatc
)
mock_prisma_client = MagicMock()
_wire_tx(mock_prisma_client.db)
mock_db = MagicMock()
_wire_tx(mock_db)
# Pattern: fail, fail, success (resets counter), fail, fail, success, done.
# Without reset, three of these would trip abort; with reset, they don't.
mock_db.execute_raw = AsyncMock(
@ -601,10 +679,10 @@ async def test_delete_old_logs_resets_consecutive_failures_on_success(monkeypatc
)
cutoff_date = datetime.now(timezone.utc) - timedelta(days=7)
total_deleted = await cleaner._delete_old_logs(mock_prisma_client, cutoff_date)
result = await cleaner._delete_old_logs(mock_prisma_client, cutoff_date, _far_deadline())
assert mock_db.execute_raw.call_count == 7
assert total_deleted == 150
assert result.rows_deleted == 150
@pytest.mark.asyncio
@ -617,6 +695,7 @@ async def test_cleanup_uses_logger_exception_for_full_traceback(monkeypatch):
monkeypatch.setattr(cleanup_module, "verbose_proxy_logger", mock_logger)
mock_prisma_client = MagicMock()
_wire_tx(mock_prisma_client.db)
# Force the outer try/except to fire by making _should_delete_spend_logs raise.
cleaner = cleanup_module.SpendLogCleanup(
general_settings={"maximum_spend_logs_retention_period": "7d"}
@ -653,7 +732,9 @@ async def test_cleanup_releases_lock_after_persistent_batch_failures(monkeypatch
)
mock_prisma_client = MagicMock()
_wire_tx(mock_prisma_client.db)
mock_db = MagicMock()
_wire_tx(mock_db)
mock_db.execute_raw = AsyncMock(side_effect=TimeoutError("DB down"))
mock_prisma_client.db = mock_db
@ -698,6 +779,7 @@ def _mock_prisma_for_retention(side_effect: list) -> "MagicMock":
from unittest.mock import AsyncMock, MagicMock
client = MagicMock()
_wire_tx(client.db)
client.db.execute_raw = AsyncMock(side_effect=side_effect)
return client
@ -753,3 +835,536 @@ async def test_no_retention_keys_means_no_cleanup_at_all():
cleaner.pod_lock_manager = None
await cleaner.cleanup_old_spend_logs(client)
assert client.db.execute_raw.await_count == 0
@pytest.mark.asyncio
async def test_run_budget_stops_the_loop_and_leaves_the_backlog_for_the_next_run():
"""
The wall-clock budget is the bound that keeps a large backlog from turning
into one multi-hour run. With rows always available, the loop must stop on
the deadline rather than on the batch cap, and must report that reason so
operators can tell a budgeted stop from a drained table.
"""
mock_prisma_client = MagicMock()
_wire_tx(mock_prisma_client.db)
mock_db = MagicMock()
_wire_tx(mock_db)
mock_db.execute_raw = AsyncMock(return_value=1000)
mock_prisma_client.db = mock_db
cleaner = SpendLogCleanup(
general_settings={
"maximum_spend_logs_retention_period": "7d",
# Comfortably more batches than a sub-second budget can reach (each
# batch sleeps 0.1s), but small enough that a broken deadline fails
# this test in seconds instead of hanging it
"maximum_spend_logs_cleanup_max_batches": 50,
}
)
cutoff_date = datetime.now(timezone.utc) - timedelta(days=7)
started_at = time.monotonic()
result = await cleaner._delete_old_logs(mock_prisma_client, cutoff_date, time.monotonic() + 0.25)
elapsed = time.monotonic() - started_at
assert result.stop_reason == "budget_exhausted"
assert elapsed < 3, f"budgeted run overran its deadline: {elapsed}s"
assert mock_db.execute_raw.call_count < 50
assert result.rows_deleted > 0
@pytest.mark.asyncio
async def test_run_budget_is_shared_across_tables_not_granted_per_table():
"""
A per-table budget would let a run take N times the configured bound. The
deadline is computed once per run, so once it is spent on the first table
the later tables must stop immediately rather than each getting a fresh one.
"""
mock_prisma_client = MagicMock()
_wire_tx(mock_prisma_client.db)
mock_db = MagicMock()
_wire_tx(mock_db)
mock_db.execute_raw = AsyncMock(return_value=1000)
mock_prisma_client.db = mock_db
cleaner = SpendLogCleanup(
general_settings={
"maximum_spend_logs_retention_period": "7d",
"maximum_autorouter_session_retention_period": "365d",
# Comfortably more batches than a sub-second budget can reach (each
# batch sleeps 0.1s), but small enough that a broken deadline fails
# this test in seconds instead of hanging it
"maximum_spend_logs_cleanup_max_batches": 50,
"maximum_spend_logs_cleanup_run_budget": "1s",
}
)
cleaner.pod_lock_manager = None
started_at = time.monotonic()
await cleaner.cleanup_old_spend_logs(mock_prisma_client)
elapsed = time.monotonic() - started_at
# three tables are eligible; a per-table budget would push this past 3s
assert elapsed < 2.5, f"budget was granted per table, not per run: {elapsed}s"
tables_touched = {call[0][0].split('"')[1] for call in mock_db.execute_raw.call_args_list}
assert "LiteLLM_SpendLogs" in tables_touched
@pytest.mark.asyncio
async def test_each_batch_carries_a_statement_and_lock_timeout():
"""
A Prisma transaction timeout cannot interrupt a statement already running,
so the Postgres statement_timeout and lock_timeout are the only things
stopping one batch from holding row locks and a pooled connection
indefinitely. Both must be set, inside the batch's own transaction, and
scoped with SET LOCAL so the pooled connection is left unchanged.
"""
recorded: list[str] = []
mock_prisma_client = MagicMock()
mock_db = MagicMock()
@asynccontextmanager
async def _tx():
tx = MagicMock()
async def _execute_raw(sql, *args):
recorded.append(sql.strip())
return 0
tx.execute_raw = _execute_raw
yield tx
mock_db.tx = _tx
mock_db.query_raw = AsyncMock(return_value=[{"remaining": 0}])
mock_prisma_client.db = mock_db
cleaner = SpendLogCleanup(
general_settings={
"maximum_spend_logs_retention_period": "7d",
"maximum_spend_logs_cleanup_batch_timeout": "12s",
}
)
await cleaner._delete_old_logs(
mock_prisma_client, datetime.now(timezone.utc) - timedelta(days=7), _far_deadline()
)
assert "SET LOCAL statement_timeout = 12000" in recorded
assert "SET LOCAL lock_timeout = 12000" in recorded
# the timeouts must precede the delete they are meant to bound
assert recorded.index("SET LOCAL statement_timeout = 12000") < next(
i for i, sql in enumerate(recorded) if sql.startswith("DELETE")
)
@pytest.mark.parametrize(
"setting_value",
["inf", "-inf", "nan", "1e400", "0s", "-5m", "not-a-duration"],
)
def test_a_non_finite_or_non_positive_budget_falls_back_to_the_default(setting_value):
"""
The knob must not be able to remove the bound it exists to enforce.
'inf', 'nan' and '1e400' are the spellings that would turn the deadline
into no deadline at all, and '0s' and '-5m' would make every run stop before
deleting anything. All of them must land on the default rather than being
honoured, and the resulting budget must be usable arithmetic.
"""
cleaner = SpendLogCleanup(
general_settings={
"maximum_spend_logs_retention_period": "7d",
"maximum_spend_logs_cleanup_run_budget": setting_value,
}
)
assert cleaner.run_budget_seconds == SPEND_LOG_CLEANUP_RUN_BUDGET_SECONDS
assert math.isfinite(cleaner.run_budget_seconds)
assert cleaner.run_budget_seconds > 0
@pytest.mark.parametrize("setting_value", [0, -1, "abc", "", 2.9])
def test_a_bad_batch_size_falls_back_to_the_default(setting_value):
"""A zero or negative batch size would make every DELETE a no-op and the
loop spin, so unusable values must fall back rather than be honoured."""
cleaner = SpendLogCleanup(
general_settings={
"maximum_spend_logs_retention_period": "7d",
"maximum_spend_logs_cleanup_batch_size": setting_value,
}
)
assert cleaner.batch_size >= 1
def test_operator_knobs_override_the_env_defaults():
"""The knobs are meant to be reachable from general_settings (and therefore
from the admin UI), not only from environment variables."""
cleaner = SpendLogCleanup(
general_settings={
"maximum_spend_logs_retention_period": "7d",
"maximum_spend_logs_cleanup_batch_size": 250,
"maximum_spend_logs_cleanup_max_batches": 7,
"maximum_spend_logs_cleanup_run_budget": "90s",
"maximum_spend_logs_cleanup_batch_timeout": "2m",
}
)
assert cleaner.batch_size == 250
assert cleaner.max_batches == 7
assert cleaner.run_budget_seconds == 90
assert cleaner.batch_timeout_seconds == 120
_BOUND_SETTING_CASES = (
("maximum_spend_logs_cleanup_batch_size", 137, "batch_size", 137),
("maximum_spend_logs_cleanup_max_batches", 9, "max_batches", 9),
("maximum_spend_logs_cleanup_run_budget", "45s", "run_budget_seconds", 45.0),
("maximum_spend_logs_cleanup_batch_timeout", "8s", "batch_timeout_seconds", 8.0),
)
@pytest.mark.parametrize("setting_name, setting_value, attribute, expected", _BOUND_SETTING_CASES)
@pytest.mark.asyncio
async def test_a_bound_changed_after_construction_reaches_the_next_run(
setting_name, setting_value, attribute, expected
):
"""The scheduler holds one long-lived instance and the config reload mutates
general_settings in place, so a bound captured at construction would leave
every dashboard change inert until the process restarts."""
settings = {"maximum_spend_logs_retention_period": "7d"}
cleaner = SpendLogCleanup(general_settings=settings)
cleaner.pod_lock_manager = None
assert getattr(cleaner, attribute) != expected
settings[setting_name] = setting_value
await cleaner.cleanup_old_spend_logs(_mock_prisma_for_retention([0, 0]))
assert getattr(cleaner, attribute) == expected
@pytest.mark.parametrize("cleared_to_none", [True, False])
@pytest.mark.asyncio
async def test_a_bound_cleared_after_construction_falls_back_to_its_default(cleared_to_none):
"""Blanking the field in the dashboard has to restore the shipped default
rather than leave the operator's old bound in force, whether the reload
spells the clear as an explicit None or as an absent key."""
settings = {"maximum_spend_logs_retention_period": "7d", "maximum_spend_logs_cleanup_batch_size": 137}
cleaner = SpendLogCleanup(general_settings=settings)
cleaner.pod_lock_manager = None
assert cleaner.batch_size == 137
if cleared_to_none:
settings["maximum_spend_logs_cleanup_batch_size"] = None
else:
del settings["maximum_spend_logs_cleanup_batch_size"]
await cleaner.cleanup_old_spend_logs(_mock_prisma_for_retention([0, 0]))
assert cleaner.batch_size == SPEND_LOG_CLEANUP_BATCH_SIZE
def test_every_declared_bound_setting_is_covered_by_a_live_reread_case():
"""A bound added to the declared set without a live-reread case would be
propagated by the proxy and then ignored by the running job."""
assert {case[0] for case in _BOUND_SETTING_CASES} == set(SPEND_LOG_CLEANUP_BOUND_SETTINGS)
@pytest.mark.asyncio
async def test_remaining_rows_probe_is_capped_so_it_cannot_scan_the_table():
"""The remaining-eligible-rows metric must never itself become the long
scan this job exists to avoid, so its probe carries a LIMIT."""
mock_prisma_client = MagicMock()
_wire_tx(mock_prisma_client.db)
mock_db = MagicMock()
_wire_tx(mock_db)
mock_db.execute_raw = AsyncMock(return_value=0)
mock_prisma_client.db = mock_db
cleaner = SpendLogCleanup(general_settings={"maximum_spend_logs_retention_period": "7d"})
await cleaner._delete_old_logs(
mock_prisma_client, datetime.now(timezone.utc) - timedelta(days=7), _far_deadline()
)
count_sql = mock_db.query_raw.call_args[0][0]
assert "count(*)" in count_sql
assert "LIMIT $2" in count_sql
assert mock_db.query_raw.call_args[0][2] == SPEND_LOG_CLEANUP_REMAINING_COUNT_CAP
@pytest.mark.asyncio
async def test_a_run_skipped_because_another_pod_holds_the_lock_is_reported():
"""Operators need to tell "nothing to do" apart from "someone else is doing
it", so a lock-skipped run is recorded under its own outcome."""
recorded: list[str] = []
original_record_run = SpendLogCleanupMetrics.record_run
mock_prisma_client = MagicMock()
_wire_tx(mock_prisma_client.db)
cleaner = SpendLogCleanup(general_settings={"maximum_spend_logs_retention_period": "7d"})
cleaner.pod_lock_manager = MagicMock()
cleaner.pod_lock_manager.redis_cache = MagicMock()
cleaner.pod_lock_manager.acquire_lock = AsyncMock(return_value=False)
cleaner.pod_lock_manager.release_lock = AsyncMock()
SpendLogCleanupMetrics.record_run = classmethod(lambda cls, outcome: recorded.append(outcome))
try:
await cleaner.cleanup_old_spend_logs(mock_prisma_client)
finally:
SpendLogCleanupMetrics.record_run = original_record_run
assert recorded == ["skipped_locked"]
cleaner.pod_lock_manager.release_lock.assert_not_awaited()
@pytest.mark.asyncio
async def test_the_outstanding_rows_probe_carries_a_statement_timeout():
"""
The probe is a statement like any other, so if it were issued bare a slow one
would hold a connection past the budget the job advertises, which is exactly
what the bounds exist to prevent. With budget to spare it carries the same
per-statement timeout the delete batches do.
"""
recorded: list[str] = []
mock_prisma_client = MagicMock()
mock_db = MagicMock()
@asynccontextmanager
async def _tx():
tx = MagicMock()
async def _execute_raw(sql, *args):
recorded.append(sql.strip())
return 0
async def _query_raw(sql, *args):
recorded.append(sql.strip())
return [{"remaining": 7}]
tx.execute_raw = _execute_raw
tx.query_raw = _query_raw
yield tx
mock_db.tx = _tx
mock_prisma_client.db = mock_db
cleaner = SpendLogCleanup(
general_settings={
"maximum_spend_logs_retention_period": "7d",
"maximum_spend_logs_cleanup_batch_timeout": "8s",
}
)
remaining = await cleaner._count_remaining(
mock_prisma_client,
datetime.now(timezone.utc) - timedelta(days=7),
"LiteLLM_SpendLogs",
"startTime",
_far_deadline(),
)
assert remaining == 7
count_index = next(i for i, sql in enumerate(recorded) if sql.startswith("SELECT count(*)"))
assert "SET LOCAL statement_timeout = 8000" in recorded[:count_index], (
f"the probe ran without a statement timeout: {recorded}"
)
@pytest.mark.asyncio
async def test_a_statement_timeout_is_clamped_to_the_budget_that_is_left():
"""
Postgres has no 'stop at time T', only a per-statement duration, so a batch
issued just under the deadline would run a whole batch timeout past it and
the run budget would be advisory. Clamping the timeout to the remaining
budget is what makes the budget a real wall clock.
"""
recorded: list[str] = []
client = MagicMock()
@asynccontextmanager
async def _tx():
tx = MagicMock()
async def _execute_raw(sql, *args):
recorded.append(sql.strip())
return 0
tx.execute_raw = _execute_raw
tx.query_raw = AsyncMock(return_value=[{"remaining": 0}])
yield tx
client.db.tx = _tx
cleaner = SpendLogCleanup(
general_settings={
"maximum_spend_logs_retention_period": "7d",
"maximum_spend_logs_cleanup_batch_timeout": "30s",
}
)
# Only 2s of budget left against a 30s batch timeout.
await cleaner._execute_delete_batch(client, "DELETE FROM x", datetime.now(timezone.utc), time.monotonic() + 2)
timeouts = [sql for sql in recorded if "statement_timeout" in sql]
assert timeouts, f"no statement timeout was issued: {recorded}"
issued_ms = int(timeouts[0].split("=")[1].strip())
assert issued_ms <= 2000, f"the batch was given {issued_ms}ms with only 2000ms of budget left"
@pytest.mark.asyncio
async def test_no_statement_is_issued_once_the_budget_is_spent():
"""
Every table exits through _finish_table, including the ones a spent run never
started, so an unconditional probe there would put one more statement per
table past the bound.
"""
client = _mock_prisma_for_retention([0, 0])
cleaner = SpendLogCleanup(general_settings={"maximum_spend_logs_retention_period": "7d"})
result = await cleaner._finish_table(
client,
datetime.now(timezone.utc) - timedelta(days=7),
"LiteLLM_SpendLogs",
"startTime",
123,
"budget_exhausted",
time.monotonic() - 1,
)
assert result.rows_deleted == 123
assert result.stop_reason == "budget_exhausted"
client.db.query_raw.assert_not_called()
@pytest.mark.asyncio
async def test_a_batch_cancelled_by_the_deadline_is_budget_exhaustion_not_a_failure(monkeypatch):
"""
Clamping the timeout means the last batch of a budget-exhausted run is
cancelled by the deadline itself. Counting that as a batch failure would
inflate the failure metric on every such run and walk it toward the abort
threshold, so it has to be classified as the bound working.
"""
failures: list[str] = []
client = MagicMock()
_wire_tx(client.db)
# The deadline has to pass DURING the batch, not before it: a deadline
# already spent is caught by the loop's own check and no batch is ever
# issued, which would exercise none of the classification under test.
async def _cancelled_after_the_deadline(sql, *args):
await asyncio.sleep(0.05)
raise Exception("canceling statement due to statement timeout")
client.db.execute_raw = _cancelled_after_the_deadline
cleaner = SpendLogCleanup(general_settings={"maximum_spend_logs_retention_period": "7d"})
monkeypatch.setattr(SpendLogCleanupMetrics, "record_batch_failure", lambda table: failures.append(table))
result = await cleaner._delete_old_logs(
client, datetime.now(timezone.utc) - timedelta(days=7), time.monotonic() + 0.02
)
assert result.stop_reason == "budget_exhausted"
assert failures == [], f"a deadline cancellation was recorded as a batch failure: {failures}"
@pytest.mark.asyncio
async def test_partition_maintenance_is_skipped_once_the_run_budget_is_spent():
"""
Dropping a partition is DDL holding an ACCESS EXCLUSIVE lock, and unlike a
delete batch it cannot be cut short once it has started. A run whose budget is
already gone must therefore not start it at all; the next tick picks it up.
"""
mock_prisma_client = MagicMock()
_wire_tx(mock_prisma_client.db)
mock_db = MagicMock()
_wire_tx(mock_db)
mock_db.execute_raw = AsyncMock(return_value=0)
mock_prisma_client.db = mock_db
partition_manager = MagicMock()
partition_manager.is_partitioned = AsyncMock(return_value=True)
partition_manager.ensure_partitions = AsyncMock()
partition_manager.drop_partitions_older_than = AsyncMock(return_value=[])
cleaner = SpendLogCleanup(
general_settings={
"maximum_spend_logs_retention_period": "7d",
"use_spend_logs_partitioning": True,
},
partition_manager=partition_manager,
)
cleaner._should_delete_spend_logs()
# a deadline already in the past is what a run that spent its budget on an
# earlier table looks like
await cleaner._clean_spend_log_tables(mock_prisma_client, time.monotonic() - 1)
partition_manager.ensure_partitions.assert_not_awaited()
partition_manager.drop_partitions_older_than.assert_not_awaited()
@pytest.mark.asyncio
async def test_partition_maintenance_still_runs_while_the_run_has_budget():
"""The skip above must be caused by the spent budget, not by breaking the
partition path outright."""
mock_prisma_client = MagicMock()
_wire_tx(mock_prisma_client.db)
mock_db = MagicMock()
_wire_tx(mock_db)
mock_db.execute_raw = AsyncMock(return_value=0)
mock_prisma_client.db = mock_db
partition_manager = MagicMock()
partition_manager.is_partitioned = AsyncMock(return_value=True)
partition_manager.ensure_partitions = AsyncMock()
partition_manager.drop_partitions_older_than = AsyncMock(return_value=["LiteLLM_SpendLogs_p20260601"])
cleaner = SpendLogCleanup(
general_settings={
"maximum_spend_logs_retention_period": "7d",
"use_spend_logs_partitioning": True,
},
partition_manager=partition_manager,
)
cleaner._should_delete_spend_logs()
await cleaner._clean_spend_log_tables(mock_prisma_client, _far_deadline())
partition_manager.ensure_partitions.assert_awaited_once()
partition_manager.drop_partitions_older_than.assert_awaited_once()
@pytest.mark.parametrize(
"stop_reasons, expected",
[
(("exhausted",), "completed"),
(("exhausted", "exhausted"), "completed"),
(("exhausted", "batch_cap_reached"), "batch_cap_reached"),
(("batch_cap_reached", "exhausted"), "batch_cap_reached"),
(("exhausted", "budget_exhausted"), "budget_exhausted"),
(("budget_exhausted", "exhausted"), "budget_exhausted"),
(("batch_cap_reached", "budget_exhausted"), "budget_exhausted"),
(("budget_exhausted", "batch_cap_reached"), "budget_exhausted"),
(("exhausted", "aborted"), "aborted"),
(("aborted", "exhausted"), "aborted"),
(("budget_exhausted", "aborted"), "aborted"),
(("aborted", "budget_exhausted"), "aborted"),
(("aborted", "budget_exhausted", "batch_cap_reached"), "aborted"),
],
)
def test_the_reported_run_outcome_is_the_most_significant_reason_in_any_order(stop_reasons, expected):
"""
The run outcome answers "why did this run stop", so a table that merely ran
dry must never mask one that hit a bound, and an abort must outrank both.
Both orders of every pair are covered because this folds several per-table
results into one answer: a first-match-wins implementation would pass on
whichever order happened to be written and fail on its mirror.
"""
results = tuple(TableCleanupResult(rows_deleted=0, stop_reason=reason) for reason in stop_reasons)
assert SpendLogCleanup._run_outcome(results) == expected

View file

@ -17,6 +17,10 @@ export enum ConfigType {
*/
export enum GeneralSettingsFieldName {
MAXIMUM_SPEND_LOGS_RETENTION_PERIOD = "maximum_spend_logs_retention_period",
MAXIMUM_SPEND_LOGS_CLEANUP_BATCH_SIZE = "maximum_spend_logs_cleanup_batch_size",
MAXIMUM_SPEND_LOGS_CLEANUP_MAX_BATCHES = "maximum_spend_logs_cleanup_max_batches",
MAXIMUM_SPEND_LOGS_CLEANUP_RUN_BUDGET = "maximum_spend_logs_cleanup_run_budget",
MAXIMUM_SPEND_LOGS_CLEANUP_BATCH_TIMEOUT = "maximum_spend_logs_cleanup_batch_timeout",
// Add more field names here as needed
}

View file

@ -6,6 +6,10 @@ import { proxyConfigKeys } from "../proxyConfig/useProxyConfig";
export interface StoreRequestInSpendLogsParams {
store_prompts_in_spend_logs: boolean;
maximum_spend_logs_retention_period?: string;
maximum_spend_logs_cleanup_batch_size?: number;
maximum_spend_logs_cleanup_max_batches?: number;
maximum_spend_logs_cleanup_run_budget?: string;
maximum_spend_logs_cleanup_batch_timeout?: string;
}
export interface StoreRequestInSpendLogsResponse {
@ -19,6 +23,8 @@ const performStoreRequestInSpendLogs = async (
const proxyBaseUrl = getProxyBaseUrl();
const url = proxyBaseUrl ? `${proxyBaseUrl}/config/update` : `/config/update`;
const { store_prompts_in_spend_logs, ...optionalSettings } = params;
const response = await fetch(url, {
method: "POST",
headers: {
@ -27,10 +33,8 @@ const performStoreRequestInSpendLogs = async (
},
body: JSON.stringify({
general_settings: {
store_prompts_in_spend_logs: params.store_prompts_in_spend_logs,
...(params.maximum_spend_logs_retention_period && {
maximum_spend_logs_retention_period: params.maximum_spend_logs_retention_period,
}),
store_prompts_in_spend_logs,
...optionalSettings,
},
}),
});

View file

@ -1,4 +1,8 @@
import { useDeleteProxyConfigField, useProxyConfig } from "@/app/(dashboard)/hooks/proxyConfig/useProxyConfig";
import {
DeleteProxyConfigFieldRequest,
useDeleteProxyConfigField,
useProxyConfig,
} from "@/app/(dashboard)/hooks/proxyConfig/useProxyConfig";
import { useStoreRequestInSpendLogs } from "@/app/(dashboard)/hooks/storeRequestInSpendLogs/useStoreRequestInSpendLogs";
import NotificationsManager from "@/components/molecules/notifications_manager";
import { parseErrorMessage } from "@/components/shared/errorUtils";
@ -40,6 +44,64 @@ describe("LoggingSettings", () => {
const mockDeleteField = vi.fn();
const mockRefetch = vi.fn();
const clearedFieldNames = (): string[] =>
mockDeleteField.mock.calls.map((call) => (call[0] as DeleteProxyConfigFieldRequest).field_name);
// Every optional knob already persisted. Clearing is only ever issued for a
// field that has a stored value, so any test about the clear path has to say
// so; the default mock below is an empty config, which is a proxy that has
// never saved these settings and therefore has nothing to clear.
const withEveryOptionalFieldStored = () =>
mockUseProxyConfig.mockReturnValue({
data: [
{
field_name: "maximum_spend_logs_retention_period",
field_type: "string",
field_description: "Maximum retention period",
field_value: "30d",
stored_in_db: true,
},
{
field_name: "maximum_spend_logs_cleanup_batch_size",
field_type: "Integer",
field_description: "Rows per delete",
field_value: 2000,
stored_in_db: true,
},
{
field_name: "maximum_spend_logs_cleanup_max_batches",
field_type: "Integer",
field_description: "Deletes per table per run",
field_value: 50,
stored_in_db: true,
},
{
field_name: "maximum_spend_logs_cleanup_run_budget",
field_type: "string",
field_description: "Wall clock budget per run",
field_value: "90s",
stored_in_db: true,
},
{
field_name: "maximum_spend_logs_cleanup_batch_timeout",
field_type: "string",
field_description: "Statement and lock timeout per batch",
field_value: "10s",
stored_in_db: true,
},
],
isLoading: false,
refetch: mockRefetch,
} as unknown as ReturnType<typeof useProxyConfig>);
// Blank every optional input the form rendered from stored values, which is
// what an admin does to reset a knob to its default.
const blankEveryOptionalField = async (user: ReturnType<typeof userEvent.setup>) => {
for (const placeholder of ["e.g., 7d, 30d", "e.g., 1000", "e.g., 500", "e.g., 5m", "e.g., 30s"]) {
await user.clear(screen.getByPlaceholderText(placeholder));
}
};
beforeEach(() => {
vi.resetAllMocks();
mockUseStoreRequestInSpendLogs.mockReturnValue({
@ -68,6 +130,19 @@ describe("LoggingSettings", () => {
expect(screen.getByRole("button", { name: "Save Settings" })).toBeInTheDocument();
});
it("should render a control for every spend logs cleanup knob", () => {
renderWithProviders(<LoggingSettings />);
expect(screen.getByLabelText("Spend Logs Cleanup Batch Size (Optional)")).toBeInTheDocument();
expect(screen.getByLabelText("Spend Logs Cleanup Max Batches (Optional)")).toBeInTheDocument();
expect(screen.getByLabelText("Spend Logs Cleanup Run Budget (Optional)")).toBeInTheDocument();
expect(screen.getByLabelText("Spend Logs Cleanup Batch Timeout (Optional)")).toBeInTheDocument();
expect(screen.getByPlaceholderText("e.g., 1000")).toBeInTheDocument();
expect(screen.getByPlaceholderText("e.g., 500")).toBeInTheDocument();
expect(screen.getByPlaceholderText("e.g., 5m")).toBeInTheDocument();
expect(screen.getByPlaceholderText("e.g., 30s")).toBeInTheDocument();
});
it("should toggle store prompts switch", async () => {
const user = userEvent.setup();
renderWithProviders(<LoggingSettings />);
@ -94,6 +169,9 @@ describe("LoggingSettings", () => {
it("should submit form with store prompts enabled and retention period", async () => {
const user = userEvent.setup();
mockDeleteField.mockImplementation((_params, options) => {
options?.onSettled?.();
});
mockMutate.mockImplementation((_params, options) => {
options?.onSuccess?.();
});
@ -110,7 +188,6 @@ describe("LoggingSettings", () => {
await user.click(saveButton);
await waitFor(() => {
expect(mockDeleteField).not.toHaveBeenCalled();
expect(mockMutate).toHaveBeenCalledWith(
{
store_prompts_in_spend_logs: true,
@ -119,9 +196,42 @@ describe("LoggingSettings", () => {
expect.any(Object),
);
});
expect(clearedFieldNames()).not.toContain("maximum_spend_logs_retention_period");
});
it("should delete retention period field when left empty on submit", async () => {
it("should submit every spend logs cleanup setting that has a value", async () => {
const user = userEvent.setup();
mockMutate.mockImplementation((_params, options) => {
options?.onSuccess?.();
});
renderWithProviders(<LoggingSettings />);
await user.click(screen.getByRole("switch"));
await user.type(screen.getByPlaceholderText("e.g., 7d, 30d"), "30d");
await user.type(screen.getByPlaceholderText("e.g., 1000"), "2000");
await user.type(screen.getByPlaceholderText("e.g., 500"), "50");
await user.type(screen.getByPlaceholderText("e.g., 5m"), "90s");
await user.type(screen.getByPlaceholderText("e.g., 30s"), "10s");
await user.click(screen.getByRole("button", { name: "Save Settings" }));
const expectedParams = {
store_prompts_in_spend_logs: true,
maximum_spend_logs_retention_period: "30d",
maximum_spend_logs_cleanup_batch_size: 2000,
maximum_spend_logs_cleanup_max_batches: 50,
maximum_spend_logs_cleanup_run_budget: "90s",
maximum_spend_logs_cleanup_batch_timeout: "10s",
};
await waitFor(() => {
expect(mockMutate).toHaveBeenCalledWith(expectedParams, expect.any(Object));
});
expect(mockDeleteField).not.toHaveBeenCalled();
});
it("should omit blank cleanup settings from the save payload instead of sending empty values", async () => {
const user = userEvent.setup();
mockDeleteField.mockImplementation((_params, options) => {
options?.onSettled?.();
@ -132,11 +242,70 @@ describe("LoggingSettings", () => {
renderWithProviders(<LoggingSettings />);
await user.type(screen.getByPlaceholderText("e.g., 1000"), "2000");
await user.click(screen.getByRole("button", { name: "Save Settings" }));
await waitFor(() => {
expect(mockMutate).toHaveBeenCalled();
});
const submittedParams = mockMutate.mock.calls[0][0];
expect(submittedParams).not.toHaveProperty("maximum_spend_logs_retention_period");
expect(submittedParams).not.toHaveProperty("maximum_spend_logs_cleanup_max_batches");
expect(submittedParams).not.toHaveProperty("maximum_spend_logs_cleanup_run_budget");
expect(submittedParams).not.toHaveProperty("maximum_spend_logs_cleanup_batch_timeout");
expect(submittedParams).toEqual({
store_prompts_in_spend_logs: false,
maximum_spend_logs_cleanup_batch_size: 2000,
});
});
it("should clear the stored value of every cleanup setting left blank", async () => {
const user = userEvent.setup();
withEveryOptionalFieldStored();
mockDeleteField.mockImplementation((_params, options) => {
options?.onSettled?.();
});
mockMutate.mockImplementation((_params, options) => {
options?.onSuccess?.();
});
renderWithProviders(<LoggingSettings />);
await blankEveryOptionalField(user);
await user.type(screen.getByPlaceholderText("e.g., 5m"), "10m");
await user.click(screen.getByRole("button", { name: "Save Settings" }));
await waitFor(() => {
expect(mockMutate).toHaveBeenCalled();
});
expect(clearedFieldNames().sort()).toEqual([
"maximum_spend_logs_cleanup_batch_size",
"maximum_spend_logs_cleanup_batch_timeout",
"maximum_spend_logs_cleanup_max_batches",
"maximum_spend_logs_retention_period",
]);
});
it("should delete retention period field when left empty on submit", async () => {
const user = userEvent.setup();
withEveryOptionalFieldStored();
mockDeleteField.mockImplementation((_params, options) => {
options?.onSettled?.();
});
mockMutate.mockImplementation((_params, options) => {
options?.onSuccess?.();
});
renderWithProviders(<LoggingSettings />);
await blankEveryOptionalField(user);
const saveButton = screen.getByRole("button", { name: "Save Settings" });
await user.click(saveButton);
await waitFor(() => {
expect(mockDeleteField).toHaveBeenCalled();
expect(clearedFieldNames()).toContain("maximum_spend_logs_retention_period");
expect(mockMutate).toHaveBeenCalledWith(
{
store_prompts_in_spend_logs: false,
@ -234,6 +403,38 @@ describe("LoggingSettings", () => {
stored_in_db: true,
field_default_value: undefined,
},
{
field_name: "maximum_spend_logs_cleanup_batch_size",
field_type: "Integer",
field_description: "Rows per delete",
field_value: 2000,
stored_in_db: true,
field_default_value: 1000,
},
{
field_name: "maximum_spend_logs_cleanup_max_batches",
field_type: "Integer",
field_description: "Deletes per table per run",
field_value: 50,
stored_in_db: true,
field_default_value: 500,
},
{
field_name: "maximum_spend_logs_cleanup_run_budget",
field_type: "string",
field_description: "Wall clock budget per run",
field_value: "90s",
stored_in_db: true,
field_default_value: "5m",
},
{
field_name: "maximum_spend_logs_cleanup_batch_timeout",
field_type: "string",
field_description: "Statement and lock timeout per batch",
field_value: "10s",
stored_in_db: true,
field_default_value: "30s",
},
],
isLoading: false,
refetch: mockRefetch,
@ -246,6 +447,10 @@ describe("LoggingSettings", () => {
expect(switchElement).toBeChecked();
expect(retentionInput).toHaveValue("30d");
expect(screen.getByPlaceholderText("e.g., 1000")).toHaveDisplayValue("2000");
expect(screen.getByPlaceholderText("e.g., 500")).toHaveDisplayValue("50");
expect(screen.getByPlaceholderText("e.g., 5m")).toHaveValue("90s");
expect(screen.getByPlaceholderText("e.g., 30s")).toHaveValue("10s");
});
it("should reflect persisted values that arrive after the initial loading render", async () => {
@ -307,11 +512,11 @@ describe("LoggingSettings", () => {
expect(skeletons.length).toBeGreaterThan(0);
});
it("should continue with update even if deleteField fails", async () => {
it("should report an error and not claim success when clearing a field fails", async () => {
const user = userEvent.setup();
const deleteError = new Error("Field does not exist");
withEveryOptionalFieldStored();
mockDeleteField.mockImplementation((_params, options) => {
options?.onError?.(deleteError);
options?.onError?.(new Error("Field does not exist"));
options?.onSettled?.();
});
mockMutate.mockImplementation((_params, options) => {
@ -320,19 +525,52 @@ describe("LoggingSettings", () => {
renderWithProviders(<LoggingSettings />);
await blankEveryOptionalField(user);
const saveButton = screen.getByRole("button", { name: "Save Settings" });
await user.click(saveButton);
await waitFor(() => {
expect(mockNotificationsManager.fromBackend).toHaveBeenCalled();
});
// the old value is still in force server side, so an unqualified success
// notification would tell the admin the opposite of what happened
expect(mockNotificationsManager.success).not.toHaveBeenCalled();
});
it("should clear fields one at a time, never concurrently", async () => {
const user = userEvent.setup();
withEveryOptionalFieldStored();
let inFlight = 0;
let maxInFlight = 0;
mockDeleteField.mockImplementation((_params, options) => {
inFlight += 1;
maxInFlight = Math.max(maxInFlight, inFlight);
// Settle on a microtask rather than synchronously, so a parallel
// implementation genuinely overlaps: Promise.all would issue every call
// before any of them settles, driving inFlight to the number of fields.
void Promise.resolve().then(() => {
inFlight -= 1;
options?.onSettled?.();
});
});
mockMutate.mockImplementation((_params, options) => {
options?.onSuccess?.();
});
renderWithProviders(<LoggingSettings />);
await blankEveryOptionalField(user);
const saveButton = screen.getByRole("button", { name: "Save Settings" });
await user.click(saveButton);
await waitFor(() => {
expect(mockDeleteField).toHaveBeenCalled();
expect(mockMutate).toHaveBeenCalledWith(
{
store_prompts_in_spend_logs: false,
},
expect.any(Object),
);
expect(mockNotificationsManager.success).toHaveBeenCalled();
});
// /config/field/delete rewrites the whole general_settings object, so two of
// them in flight at once means the later write restores what the earlier cleared
expect(maxInFlight).toBe(1);
expect(mockDeleteField.mock.calls.length).toBeGreaterThan(1);
});
it("should submit with only store prompts enabled when retention is empty", async () => {
@ -353,7 +591,6 @@ describe("LoggingSettings", () => {
await user.click(saveButton);
await waitFor(() => {
expect(mockDeleteField).toHaveBeenCalled();
expect(mockMutate).toHaveBeenCalledWith(
{
store_prompts_in_spend_logs: true,
@ -361,5 +598,57 @@ describe("LoggingSettings", () => {
expect.any(Object),
);
});
// nothing is stored for the blank fields, so there is nothing to clear
expect(mockDeleteField).not.toHaveBeenCalled();
});
it("should save on a proxy that has never stored these settings, without clearing anything", async () => {
// The first save on a new deployment: no general_settings row exists, so
// /config/field/delete answers 400 for every blank field. Issuing those
// clears anyway failed the whole save and persisted nothing.
const user = userEvent.setup();
mockDeleteField.mockImplementation((_params, options) => {
options?.onError?.(new Error("Field name=... not in config"));
options?.onSettled?.();
});
mockMutate.mockImplementation((_params, options) => {
options?.onSuccess?.();
});
renderWithProviders(<LoggingSettings />);
const switchElement = screen.getByRole("switch");
await user.click(switchElement);
await user.click(screen.getByRole("button", { name: "Save Settings" }));
await waitFor(() => {
expect(mockNotificationsManager.success).toHaveBeenCalled();
});
expect(mockDeleteField).not.toHaveBeenCalled();
expect(mockMutate).toHaveBeenCalledWith({ store_prompts_in_spend_logs: true }, expect.any(Object));
expect(mockNotificationsManager.fromBackend).not.toHaveBeenCalled();
});
it("should still clear a field that does have a stored value", async () => {
// The guard above must not turn into "never clear anything": a field the
// admin blanks out that IS stored still has to be deleted server side.
const user = userEvent.setup();
withEveryOptionalFieldStored();
mockDeleteField.mockImplementation((_params, options) => {
options?.onSettled?.();
});
mockMutate.mockImplementation((_params, options) => {
options?.onSuccess?.();
});
renderWithProviders(<LoggingSettings />);
await user.clear(screen.getByPlaceholderText("e.g., 5m"));
await user.click(screen.getByRole("button", { name: "Save Settings" }));
await waitFor(() => {
expect(mockNotificationsManager.success).toHaveBeenCalled();
});
expect(clearedFieldNames()).toContain("maximum_spend_logs_cleanup_run_budget");
});
});

View file

@ -13,43 +13,165 @@ import {
import NotificationsManager from "@/components/molecules/notifications_manager";
import { parseErrorMessage } from "@/components/shared/errorUtils";
import { ClockCircleOutlined } from "@ant-design/icons";
import { Button, Card, Form, Input, Skeleton, Space, Switch, Typography } from "antd";
import React, { useMemo } from "react";
import { Button, Card, Form, Input, InputNumber, Skeleton, Space, Switch, Typography } from "antd";
import React, { useCallback, useMemo } from "react";
const STORE_PROMPTS_FIELD_NAME = "store_prompts_in_spend_logs";
interface OptionalField {
readonly name: GeneralSettingsFieldName;
readonly kind: "duration" | "count";
readonly label: string;
readonly placeholder: string;
readonly fallbackTooltip: string;
}
const OPTIONAL_FIELDS: readonly OptionalField[] = [
{
name: GeneralSettingsFieldName.MAXIMUM_SPEND_LOGS_RETENTION_PERIOD,
kind: "duration",
label: "Maximum Spend Logs Retention Period (Optional)",
placeholder: "e.g., 7d, 30d",
fallbackTooltip:
"Set the maximum retention period for spend logs (e.g., '7d' for 7 days, '30d' for 30 days). Leave empty for no limit.",
},
{
name: GeneralSettingsFieldName.MAXIMUM_SPEND_LOGS_CLEANUP_BATCH_SIZE,
kind: "count",
label: "Spend Logs Cleanup Batch Size (Optional)",
placeholder: "e.g., 1000",
fallbackTooltip: "Rows deleted per DELETE statement during cleanup. Leave empty to use the default of 1000.",
},
{
name: GeneralSettingsFieldName.MAXIMUM_SPEND_LOGS_CLEANUP_MAX_BATCHES,
kind: "count",
label: "Spend Logs Cleanup Max Batches (Optional)",
placeholder: "e.g., 500",
fallbackTooltip:
"Maximum number of DELETE statements run per table per cleanup run. Leave empty to use the default of 500.",
},
{
name: GeneralSettingsFieldName.MAXIMUM_SPEND_LOGS_CLEANUP_RUN_BUDGET,
kind: "duration",
label: "Spend Logs Cleanup Run Budget (Optional)",
placeholder: "e.g., 5m",
fallbackTooltip:
"Wall-clock budget for a whole cleanup run, shared across every table it cleans (e.g., '5m'). Leave empty to use the default of 5m.",
},
{
name: GeneralSettingsFieldName.MAXIMUM_SPEND_LOGS_CLEANUP_BATCH_TIMEOUT,
kind: "duration",
label: "Spend Logs Cleanup Batch Timeout (Optional)",
placeholder: "e.g., 30s",
fallbackTooltip:
"Postgres statement and lock timeout applied to each cleanup batch, so cleanup never monopolizes a connection (e.g., '30s'). Leave empty to use the default of 30s.",
},
];
interface LoggingSettingsFormValues {
store_prompts_in_spend_logs: boolean;
maximum_spend_logs_retention_period?: string | null;
maximum_spend_logs_cleanup_batch_size?: number | null;
maximum_spend_logs_cleanup_max_batches?: number | null;
maximum_spend_logs_cleanup_run_budget?: string | null;
maximum_spend_logs_cleanup_batch_timeout?: string | null;
}
const hasDuration = (value: string | null | undefined): value is string =>
typeof value === "string" && value.trim() !== "";
const hasCount = (value: number | null | undefined): value is number =>
typeof value === "number" && Number.isFinite(value);
const buildUpdateParams = (formValues: LoggingSettingsFormValues): StoreRequestInSpendLogsParams => ({
store_prompts_in_spend_logs: formValues.store_prompts_in_spend_logs,
...(hasDuration(formValues.maximum_spend_logs_retention_period) && {
maximum_spend_logs_retention_period: formValues.maximum_spend_logs_retention_period,
}),
...(hasCount(formValues.maximum_spend_logs_cleanup_batch_size) && {
maximum_spend_logs_cleanup_batch_size: formValues.maximum_spend_logs_cleanup_batch_size,
}),
...(hasCount(formValues.maximum_spend_logs_cleanup_max_batches) && {
maximum_spend_logs_cleanup_max_batches: formValues.maximum_spend_logs_cleanup_max_batches,
}),
...(hasDuration(formValues.maximum_spend_logs_cleanup_run_budget) && {
maximum_spend_logs_cleanup_run_budget: formValues.maximum_spend_logs_cleanup_run_budget,
}),
...(hasDuration(formValues.maximum_spend_logs_cleanup_batch_timeout) && {
maximum_spend_logs_cleanup_batch_timeout: formValues.maximum_spend_logs_cleanup_batch_timeout,
}),
});
// A blank field only needs clearing when something is actually stored for it.
// Asking the proxy to clear a field it has no value for is a 400 whenever no
// general_settings row exists at all, which is the state of every deployment
// that has never saved one, so clearing unconditionally would fail the first
// save on a new proxy and take the rest of the form down with it.
const omittedFieldNames = (
updateParams: StoreRequestInSpendLogsParams,
isStored: (name: GeneralSettingsFieldName) => boolean,
): readonly GeneralSettingsFieldName[] =>
OPTIONAL_FIELDS.map((field) => field.name).filter((name) => !(name in updateParams) && isStored(name));
const LoggingSettings: React.FC = () => {
const [form] = Form.useForm();
const [form] = Form.useForm<LoggingSettingsFormValues>();
const { mutate, isPending } = useStoreRequestInSpendLogs();
const { mutate: deleteField, isPending: isDeletingField } = useDeleteProxyConfigField();
const { data: proxyConfigData, isLoading: isLoadingConfig } = useProxyConfig(ConfigType.GENERAL_SETTINGS);
const describeField = (name: string, fallback: string) =>
proxyConfigData?.find((field) => field.field_name === name)?.field_description || fallback;
const storedValue = useCallback(
(name: string) => proxyConfigData?.find((field) => field.field_name === name)?.field_value,
[proxyConfigData],
);
const isStored = (name: GeneralSettingsFieldName) => {
const value = storedValue(name);
return value !== null && value !== undefined;
};
const initialValues = useMemo(() => {
if (!proxyConfigData) {
return {
store_prompts_in_spend_logs: false,
maximum_spend_logs_retention_period: undefined,
};
}
const storePromptsField = proxyConfigData.find((field) => field.field_name === "store_prompts_in_spend_logs");
const retentionPeriodField = proxyConfigData.find(
(field) => field.field_name === "maximum_spend_logs_retention_period",
);
return {
store_prompts_in_spend_logs: storePromptsField?.field_value ?? false,
maximum_spend_logs_retention_period: retentionPeriodField?.field_value ?? undefined,
store_prompts_in_spend_logs: storedValue(STORE_PROMPTS_FIELD_NAME) ?? false,
...Object.fromEntries(OPTIONAL_FIELDS.map((field) => [field.name, storedValue(field.name)])),
};
}, [proxyConfigData]);
}, [storedValue]);
const handleFormSubmit = (formValues: StoreRequestInSpendLogsParams) => {
const retentionPeriodValue = formValues.maximum_spend_logs_retention_period;
const hasRetentionPeriod = typeof retentionPeriodValue === "string" && retentionPeriodValue.trim() !== "";
// Resolves to the field name when clearing it failed, or null when it worked.
const clearStoredField = (fieldName: GeneralSettingsFieldName) =>
new Promise<GeneralSettingsFieldName | null>((resolve) => {
let failed = false;
deleteField(
{ config_type: ConfigType.GENERAL_SETTINGS, field_name: fieldName },
{
onError: () => {
failed = true;
},
onSettled: () => resolve(failed ? fieldName : null),
},
);
});
const updateParams: StoreRequestInSpendLogsParams = {
store_prompts_in_spend_logs: formValues.store_prompts_in_spend_logs,
...(hasRetentionPeriod && { maximum_spend_logs_retention_period: retentionPeriodValue }),
};
// Clearing a field rewrites the whole stored general_settings object server
// side, so these must run one at a time: in parallel the last write back wins
// and silently restores the fields the earlier ones just cleared.
const clearStoredFieldsInSequence = async (
fieldNames: readonly GeneralSettingsFieldName[],
): Promise<readonly GeneralSettingsFieldName[]> => {
const failed: GeneralSettingsFieldName[] = [];
for (const fieldName of fieldNames) {
const failure = await clearStoredField(fieldName);
if (failure !== null) {
failed.push(failure);
}
}
return failed;
};
const handleFormSubmit = (formValues: LoggingSettingsFormValues) => {
const updateParams = buildUpdateParams(formValues);
const submitUpdate = () =>
mutate(updateParams, {
onSuccess: () => NotificationsManager.success("Spend logs settings updated successfully"),
@ -57,21 +179,21 @@ const LoggingSettings: React.FC = () => {
NotificationsManager.fromBackend("Failed to save spend logs settings: " + parseErrorMessage(error)),
});
if (hasRetentionPeriod) {
const fieldsToClear = omittedFieldNames(updateParams, isStored);
if (fieldsToClear.length === 0) {
submitUpdate();
return;
}
deleteField(
{
config_type: ConfigType.GENERAL_SETTINGS,
field_name: GeneralSettingsFieldName.MAXIMUM_SPEND_LOGS_RETENTION_PERIOD,
},
{
onError: (deleteError) => console.warn("Failed to delete retention period field (may not exist):", deleteError),
onSettled: submitUpdate,
},
);
void clearStoredFieldsInSequence(fieldsToClear).then((failed) => {
if (failed.length > 0) {
// Reporting an unqualified success here would tell the admin a setting
// was reset to its default while the old value is still in force.
NotificationsManager.fromBackend(`Failed to clear saved value for: ${failed.join(", ")}`);
return;
}
submitUpdate();
});
};
return (
@ -87,27 +209,30 @@ const LoggingSettings: React.FC = () => {
<Form form={form} layout="vertical" onFinish={handleFormSubmit} initialValues={initialValues}>
<Form.Item
label="Store Prompts in Spend Logs"
name="store_prompts_in_spend_logs"
tooltip={
proxyConfigData?.find((f) => f.field_name === "store_prompts_in_spend_logs")?.field_description ||
"When enabled, prompts will be stored in spend logs for tracking and analysis purposes."
}
name={STORE_PROMPTS_FIELD_NAME}
tooltip={describeField(
STORE_PROMPTS_FIELD_NAME,
"When enabled, prompts will be stored in spend logs for tracking and analysis purposes.",
)}
valuePropName="checked"
>
<Switch />
</Form.Item>
<Form.Item
label="Maximum Spend Logs Retention Period (Optional)"
name="maximum_spend_logs_retention_period"
tooltip={
proxyConfigData?.find((f) => f.field_name === "maximum_spend_logs_retention_period")
?.field_description ||
"Set the maximum retention period for spend logs (e.g., '7d' for 7 days, '30d' for 30 days). Leave empty for no limit."
}
>
<Input placeholder="e.g., 7d, 30d" prefix={<ClockCircleOutlined />} />
</Form.Item>
{OPTIONAL_FIELDS.map((field) => (
<Form.Item
key={field.name}
label={field.label}
name={field.name}
tooltip={describeField(field.name, field.fallbackTooltip)}
>
{field.kind === "duration" ? (
<Input placeholder={field.placeholder} prefix={<ClockCircleOutlined />} />
) : (
<InputNumber min={1} precision={0} placeholder={field.placeholder} style={{ width: "100%" }} />
)}
</Form.Item>
))}
<Form.Item>
<Button type="primary" htmlType="submit" loading={isPending || isDeletingField}>

View file

@ -23675,6 +23675,26 @@ 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 Spend Logs Cleanup Batch Size
* @description Rows deleted per DELETE statement by the spend log cleanup job. Defaults to 1000.
*/
maximum_spend_logs_cleanup_batch_size?: number | null;
/**
* Maximum Spend Logs Cleanup Batch Timeout
* @description Postgres statement_timeout and lock_timeout applied to each spend log cleanup delete batch (e.g. '30s'), so cleanup cannot hold row locks or a connection indefinitely. Defaults to '30s'.
*/
maximum_spend_logs_cleanup_batch_timeout?: string | null;
/**
* Maximum Spend Logs Cleanup Max Batches
* @description Maximum DELETE statements the spend log cleanup job issues per table per run. Defaults to 500.
*/
maximum_spend_logs_cleanup_max_batches?: number | null;
/**
* Maximum Spend Logs Cleanup Run Budget
* @description Wall-clock budget for one spend log cleanup run (e.g. '5m'), shared across every table it prunes. A run that hits the budget stops and the next run resumes from where it left off. Defaults to '5m'.
*/
maximum_spend_logs_cleanup_run_budget?: string | null;
/**
* Maximum Spend Logs Retention Period
* @description Maximum retention period for spend logs (e.g., '7d' for 7 days). Logs older than this will be deleted.