Coverage for .venv/lib/python3.13/site-packages/litellm/proxy/db/db_transaction_queue/spend_log_cleanup.py: 16%
273 statements
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-10 12:01 +0000
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-10 12:01 +0000
1import asyncio
2import time
3from contextvars import ContextVar
4from dataclasses import dataclass
5from datetime import datetime, timedelta, timezone
6from typing import Final, Literal, TypeAlias
8from pydantic import BaseModel, TypeAdapter
10from litellm._logging import verbose_proxy_logger
11from litellm.caching import RedisCache
12from litellm.constants import (
13 SPEND_LOG_CLEANUP_BATCH_FAILURE_BACKOFF_SECONDS,
14 SPEND_LOG_CLEANUP_BATCH_SIZE,
15 SPEND_LOG_CLEANUP_BATCH_TIMEOUT_SECONDS,
16 SPEND_LOG_CLEANUP_JOB_NAME,
17 SPEND_LOG_CLEANUP_MAX_CONSECUTIVE_BATCH_FAILURES,
18 SPEND_LOG_CLEANUP_REMAINING_COUNT_CAP,
19 SPEND_LOG_CLEANUP_RUN_BUDGET_SECONDS,
20 SPEND_LOG_RUN_LOOPS,
21)
22from litellm.litellm_core_utils.duration_parser import duration_in_seconds
23from litellm.proxy.db.db_transaction_queue.spend_log_cleanup_metrics import (
24 RunOutcome,
25 SpendLogCleanupMetrics,
26)
27from litellm.proxy.db.db_transaction_queue.spend_logs_partition_manager import (
28 RemainingTimeoutMs,
29 SpendLogsPartitionManager,
30)
31from litellm.proxy.utils import PrismaClient
33StopReason: TypeAlias = Literal["exhausted", "budget_exhausted", "batch_cap_reached", "aborted"]
36@dataclass(frozen=True, slots=True)
37class TableCleanupResult:
38 """Outcome of pruning one table, so the caller can report why a run ended."""
40 rows_deleted: int
41 stop_reason: StopReason
44class _RunProgress:
45 """How far one cleanup run has got, reported if that run is cancelled"""
47 def __init__(self) -> None:
48 self.rows_deleted: int = 0
49 self.batches: int = 0
51 def record_batch(self, rows_deleted: int) -> None:
52 self.rows_deleted += rows_deleted
53 self.batches += 1
56_run_progress: ContextVar[_RunProgress] = ContextVar("spend_log_cleanup_run_progress")
59def _record_run_batch(rows_deleted: int) -> None:
60 """Count a batch towards the run in progress, if a run is what issued it"""
61 progress: Final = _run_progress.get(None)
62 if progress is not None:
63 progress.record_batch(rows_deleted)
66class _RemainingRow(BaseModel):
67 """One row of the capped outstanding-rows probe, validated out of prisma's untyped result."""
69 remaining: int
72_REMAINING_ROWS: Final = TypeAdapter(list[_RemainingRow])
74SPEND_LOG_CLEANUP_BOUND_SETTINGS: Final = (
75 "maximum_spend_logs_cleanup_batch_size",
76 "maximum_spend_logs_cleanup_max_batches",
77 "maximum_spend_logs_cleanup_run_budget",
78 "maximum_spend_logs_cleanup_batch_timeout",
79)
82class SpendLogCleanup:
83 """
84 Handles cleaning up old spend logs based on maximum retention period.
86 When LiteLLM_SpendLogs is range-partitioned, expired data is reclaimed by
87 dropping whole partitions (instant, frees disk immediately). Otherwise it
88 falls back to deleting logs in batches.
89 Uses PodLockManager to ensure only one pod runs cleanup in multi-pod deployments.
91 Every run is bounded so it can never monopolise the database: a wall-clock
92 budget shared across all tables, a per-table batch cap, and a Postgres
93 statement/lock timeout on every statement the job issues, deletes and the
94 outstanding-rows probe alike. A run that hits a bound stops cleanly and the
95 next run resumes from where it left off, because the cutoff is recomputed
96 and deleted rows are gone.
98 The budget is a hard wall clock, not an advisory one. Every statement this
99 job issues, deletes, the outstanding-rows probe and partition DDL alike, is
100 issued with a timeout clamped to the budget that is still left, so one
101 started just under the deadline is cancelled by Postgres at the deadline
102 rather than running a further batch timeout past it. No statement is issued
103 at all once the budget is spent, which is why the probe is skipped on that
104 path. Partition DDL additionally carries a lock_timeout, because it takes an
105 ACCESS EXCLUSIVE lock and would otherwise queue behind a long-running reader
106 for as long as that reader lives; a partition this run cannot get is left
107 for the next one.
108 """
110 def __init__(
111 self,
112 general_settings=None,
113 redis_cache: RedisCache | None = None,
114 partition_manager: SpendLogsPartitionManager | None = None,
115 ):
116 self.retention_seconds: int | None = None
117 self.partition_manager = partition_manager or SpendLogsPartitionManager()
118 from litellm.proxy.proxy_server import general_settings as default_settings
120 self.general_settings = general_settings or default_settings
121 self._refresh_bounds()
122 from litellm.proxy.proxy_server import proxy_logging_obj
124 pod_lock_manager: Final = proxy_logging_obj.db_spend_update_writer.pod_lock_manager
125 self.pod_lock_manager = pod_lock_manager
126 verbose_proxy_logger.info(
127 "SpendLogCleanup initialized: batch_size=%s max_batches=%s run_budget=%ss batch_timeout=%ss",
128 self.batch_size,
129 self.max_batches,
130 self.run_budget_seconds,
131 self.batch_timeout_seconds,
132 )
134 def _refresh_bounds(self) -> None:
135 """
136 Re-read every bound in SPEND_LOG_CLEANUP_BOUND_SETTINGS from settings.
138 The scheduler holds one long-lived instance, so a bound captured at
139 construction would never reflect a dashboard change. general_settings is
140 the same dict the periodic config reload mutates in place, so reading it
141 per run is what makes these knobs live. Every bound falls back to its
142 shipped default, so clearing a field restores that default.
143 """
144 self.batch_size: int = self._positive_int_setting(
145 "maximum_spend_logs_cleanup_batch_size", SPEND_LOG_CLEANUP_BATCH_SIZE
146 )
147 self.max_batches: int = self._positive_int_setting(
148 "maximum_spend_logs_cleanup_max_batches", SPEND_LOG_RUN_LOOPS
149 )
150 self.run_budget_seconds: float = self._duration_setting(
151 "maximum_spend_logs_cleanup_run_budget", SPEND_LOG_CLEANUP_RUN_BUDGET_SECONDS
152 )
153 self.batch_timeout_seconds: float = self._duration_setting(
154 "maximum_spend_logs_cleanup_batch_timeout", SPEND_LOG_CLEANUP_BATCH_TIMEOUT_SECONDS
155 )
157 def _positive_int_setting(self, setting_name: str, default: int) -> int:
158 """
159 Read a positive-integer knob, falling back to the default when unset or unusable.
160 """
161 raw: Final = self.general_settings.get(setting_name)
162 if raw is None:
163 return default
164 try:
165 parsed: Final = int(raw)
166 except (TypeError, ValueError):
167 verbose_proxy_logger.warning("Invalid %s value: %s, using default %s", setting_name, raw, default)
168 return default
169 if parsed <= 0:
170 verbose_proxy_logger.warning("%s must be positive, got %s, using default %s", setting_name, parsed, default)
171 return default
172 return parsed
174 def _duration_setting(self, setting_name: str, default_seconds: float) -> float:
175 """
176 Read a duration knob (e.g. '5m'), falling back to the default when unset or unusable.
178 The knob must never be able to remove the bound it exists to enforce, so
179 anything the parser rejects (including the non-finite spellings 'inf' and
180 'nan') and anything non-positive falls back rather than being honoured.
181 """
182 raw: Final = self.general_settings.get(setting_name)
183 if raw is None:
184 return default_seconds
185 try:
186 parsed: Final = float(duration_in_seconds(str(raw)))
187 except (ValueError, TypeError) as e:
188 verbose_proxy_logger.warning(
189 "Invalid %s value: %s (%s), using default %ss", setting_name, raw, e, default_seconds
190 )
191 return default_seconds
192 if parsed <= 0:
193 verbose_proxy_logger.warning(
194 "%s must be a positive duration, got %s, using default %ss", setting_name, raw, default_seconds
195 )
196 return default_seconds
197 return parsed
199 def _retention_seconds_for(self, setting_name: str) -> int | None:
200 """
201 Parse one retention setting into seconds, or None when unset or invalid.
202 """
203 retention_setting = self.general_settings.get(setting_name)
204 verbose_proxy_logger.info("Checking %s: %s", setting_name, retention_setting)
206 if retention_setting is None:
207 return None
209 try:
210 if isinstance(retention_setting, int):
211 verbose_proxy_logger.warning(
212 "%s is an integer (%s); treating as days. Use a string like '3d' to be explicit.",
213 setting_name,
214 retention_setting,
215 )
216 retention_setting = f"{retention_setting}d"
217 retention_seconds: Final = duration_in_seconds(retention_setting)
218 except ValueError as e:
219 verbose_proxy_logger.warning("Invalid %s value: %s, error: %s", setting_name, retention_setting, e)
220 return None
221 verbose_proxy_logger.info("%s set to %s seconds", setting_name, retention_seconds)
222 return retention_seconds
224 def _should_delete_spend_logs(self) -> bool:
225 """
226 Determines if logs should be deleted based on the max retention period in settings.
227 """
228 self.retention_seconds = self._retention_seconds_for("maximum_spend_logs_retention_period")
229 return self.retention_seconds is not None
231 def _timeout_ms(self, deadline: float) -> int:
232 """
233 The per-statement bound in milliseconds: the batch timeout, or whatever
234 is left of the run budget, whichever is smaller.
236 Clamping to the remaining budget is what makes the budget a real
237 wall-clock bound rather than an advisory one. Postgres offers no "stop
238 at time T", only a per-statement duration, so a statement issued just
239 under the deadline would otherwise run a full batch timeout past it, and
240 with several tables those overruns stack.
242 Interpolating this into SQL is safe by construction: an int cannot carry
243 SQL, and SET does not accept a bind parameter.
244 """
245 remaining_ms: Final = int((deadline - time.monotonic()) * 1000)
246 return max(1, min(int(self.batch_timeout_seconds * 1000), remaining_ms))
248 @staticmethod
249 def _group_deadline(overall_deadline: float, groups_remaining: int) -> float:
250 """
251 Give each pending cleanup group an equal share of the time left.
253 A single group keeps the whole run budget, while a persistent backlog
254 on an earlier group cannot starve a later group.
255 """
256 if groups_remaining == 1:
257 return overall_deadline
258 current_time: Final = time.monotonic()
259 if current_time >= overall_deadline:
260 return overall_deadline
261 return current_time + (overall_deadline - current_time) / groups_remaining
263 def _remaining_timeout_ms(self, deadline: float) -> RemainingTimeoutMs:
264 """
265 The per-statement bound for work this job delegates, as a callable.
267 Partition maintenance issues one statement per partition, so handing it a
268 number would bound each statement by the budget that was left before the
269 FIRST one and never by what remains. Re-evaluating per statement is what
270 makes the loop itself bounded, and None tells the callee to stop rather
271 than issue a statement it has no budget for.
272 """
274 def remaining() -> int | None:
275 return None if time.monotonic() >= deadline else self._timeout_ms(deadline)
277 return remaining
279 async def _execute_delete_batch(
280 self, prisma_client: PrismaClient, delete_sql: str, cutoff_date: datetime, deadline: float
281 ) -> int | None:
282 """
283 Run one delete batch under a Postgres statement and lock timeout.
285 The timeouts are what actually bound the work: a Prisma transaction
286 timeout cannot interrupt a statement that is already executing, so
287 without these a single batch blocked behind a lock would hold its
288 connection, and the row locks it already took, indefinitely. SET LOCAL
289 scopes both to this transaction so the pooled connection is unaffected.
291 Returns the row count, or None when the driver returned something that
292 is not a row count. That is a contract violation rather than a transient
293 fault, so the caller stops instead of retrying.
294 """
295 timeout_ms: Final = self._timeout_ms(deadline)
296 async with prisma_client.db.tx() as tx:
297 await tx.execute_raw(f"SET LOCAL statement_timeout = {timeout_ms}")
298 await tx.execute_raw(f"SET LOCAL lock_timeout = {timeout_ms}")
299 deleted_result: Final = await tx.execute_raw(delete_sql, cutoff_date, self.batch_size)
300 return deleted_result if isinstance(deleted_result, int) else None
302 async def _count_remaining(
303 self, prisma_client: PrismaClient, cutoff_date: datetime, table_name: str, time_column: str, deadline: float
304 ) -> int | None:
305 """
306 Count expired rows still outstanding, stopping at a cap.
308 An uncapped COUNT(*) over an expired backlog would itself be the kind of
309 long scan this job exists to avoid, so the probe reads at most
310 SPEND_LOG_CLEANUP_REMAINING_COUNT_CAP index entries. A result equal to
311 the cap means "at least this many".
312 """
313 count_sql: Final = f"""
314 SELECT count(*)::int AS remaining FROM (
315 SELECT 1 FROM "{table_name}"
316 WHERE "{time_column}" < $1::timestamptz
317 LIMIT $2
318 ) capped
319 """
320 try:
321 async with prisma_client.db.tx() as tx:
322 await tx.execute_raw(f"SET LOCAL statement_timeout = {self._timeout_ms(deadline)}")
323 rows: Final = _REMAINING_ROWS.validate_python(
324 await tx.query_raw(count_sql, cutoff_date, SPEND_LOG_CLEANUP_REMAINING_COUNT_CAP)
325 )
326 except Exception as e: # noqa: BLE001 - an observability probe must never fail the cleanup run
327 verbose_proxy_logger.warning("Could not count remaining %s rows: %s", table_name, e)
328 return None
329 return rows[0].remaining if rows else None
331 async def _delete_old_rows_batched(
332 self,
333 prisma_client: PrismaClient,
334 cutoff_date: datetime,
335 table_name: str,
336 key_columns: tuple[str, ...],
337 time_column: str,
338 deadline: float,
339 ) -> TableCleanupResult:
340 """
341 Delete a table's rows older than the cutoff in batches.
343 Stops at whichever bound is reached first: the backlog running out, the
344 shared wall-clock deadline, the per-table batch cap, or too many
345 consecutive batch failures.
346 """
347 key_list: Final = ", ".join(f'"{col}"' for col in key_columns)
348 delete_sql: Final = f"""
349 DELETE FROM "{table_name}"
350 WHERE ({key_list}) IN (
351 SELECT {key_list} FROM "{table_name}"
352 WHERE "{time_column}" < $1::timestamptz
353 LIMIT $2
354 )
355 """
356 total_deleted = 0
357 run_count = 0
358 consecutive_failures = 0
359 while True:
360 if time.monotonic() >= deadline:
361 verbose_proxy_logger.info(
362 "Run budget exhausted during %s cleanup after %d rows; the next run resumes from here",
363 table_name,
364 total_deleted,
365 )
366 return await self._finish_table(
367 prisma_client, cutoff_date, table_name, time_column, total_deleted, "budget_exhausted", deadline
368 )
369 if run_count >= self.max_batches:
370 verbose_proxy_logger.info(
371 "Max batches reached for %s cleanup, remaining rows will be deleted in next run", table_name
372 )
373 return await self._finish_table(
374 prisma_client, cutoff_date, table_name, time_column, total_deleted, "batch_cap_reached", deadline
375 )
376 # Find rows and delete them in one go without fetching to application
377 batch_started_at = time.monotonic()
378 try:
379 batch_result = await self._execute_delete_batch(prisma_client, delete_sql, cutoff_date, deadline)
380 except Exception as batch_exc:
381 if time.monotonic() >= deadline:
382 # The statement timeout was clamped to the budget that was
383 # left, so this batch was cancelled by the deadline itself.
384 # That is the bound working, not a database fault, and
385 # counting it would both inflate the failure metric and push
386 # every budget-exhausted run toward the abort threshold.
387 verbose_proxy_logger.info(
388 "Run budget exhausted mid-batch during %s cleanup after %d rows; "
389 "the next run resumes from here",
390 table_name,
391 total_deleted,
392 )
393 return await self._finish_table(
394 prisma_client, cutoff_date, table_name, time_column, total_deleted, "budget_exhausted", deadline
395 )
396 # A single batch failure (e.g. Prisma/DB timeout) must not abort
397 # the whole run — subsequent batches may still succeed.
398 consecutive_failures += 1
399 SpendLogCleanupMetrics.record_batch_failure(table_name)
400 verbose_proxy_logger.exception(
401 "%s cleanup batch failed "
402 "(run_count=%d, consecutive_failures=%d, batch_size=%d, "
403 "cutoff=%s, total_deleted_so_far=%d): %s: %s",
404 table_name,
405 run_count,
406 consecutive_failures,
407 self.batch_size,
408 cutoff_date.isoformat(),
409 total_deleted,
410 type(batch_exc).__name__,
411 batch_exc,
412 )
413 if consecutive_failures >= SPEND_LOG_CLEANUP_MAX_CONSECUTIVE_BATCH_FAILURES:
414 verbose_proxy_logger.error(
415 "Aborting %s cleanup after %d consecutive batch failures; total deleted before abort: %d",
416 table_name,
417 consecutive_failures,
418 total_deleted,
419 )
420 return await self._finish_table(
421 prisma_client, cutoff_date, table_name, time_column, total_deleted, "aborted", deadline
422 )
423 await asyncio.sleep(SPEND_LOG_CLEANUP_BATCH_FAILURE_BACKOFF_SECONDS)
424 continue
426 if batch_result is None:
427 verbose_proxy_logger.error(
428 "Unexpected execute_raw return type for %s cleanup; aborting cleanup to avoid infinite loop",
429 table_name,
430 )
431 return await self._finish_table(
432 prisma_client, cutoff_date, table_name, time_column, total_deleted, "aborted", deadline
433 )
435 consecutive_failures = 0
436 deleted_count = batch_result
437 SpendLogCleanupMetrics.record_batch(table_name, deleted_count, time.monotonic() - batch_started_at)
438 verbose_proxy_logger.info("Deleted %s %s rows in this batch", deleted_count, table_name)
440 if deleted_count == 0:
441 verbose_proxy_logger.info("No more %s rows to delete. Total deleted: %s", table_name, total_deleted)
442 return await self._finish_table(
443 prisma_client, cutoff_date, table_name, time_column, total_deleted, "exhausted", deadline
444 )
446 total_deleted += deleted_count
447 run_count += 1
448 _record_run_batch(deleted_count)
450 # Add a small sleep to prevent overwhelming the database
451 await asyncio.sleep(0.1)
453 async def _finish_table(
454 self,
455 prisma_client: PrismaClient,
456 cutoff_date: datetime,
457 table_name: str,
458 time_column: str,
459 rows_deleted: int,
460 stop_reason: StopReason,
461 deadline: float,
462 ) -> TableCleanupResult:
463 """
464 Publish how much of this table is still outstanding, then report the run's result.
466 The probe is skipped once the budget is spent. It is the one piece of
467 work that would otherwise be ISSUED after the deadline, and every table
468 exits through here, including the ones a spent run never started, so
469 keeping it would put one more statement per table past the bound. A run
470 that ends this way already reports "budget_exhausted", which tells an
471 operator the backlog was not drained; the gauge simply keeps its value
472 from the last run that finished inside its budget.
473 """
474 if time.monotonic() >= deadline:
475 return TableCleanupResult(rows_deleted=rows_deleted, stop_reason=stop_reason)
476 remaining: Final = await self._count_remaining(prisma_client, cutoff_date, table_name, time_column, deadline)
477 if remaining is not None:
478 SpendLogCleanupMetrics.set_rows_remaining(table_name, remaining)
479 return TableCleanupResult(rows_deleted=rows_deleted, stop_reason=stop_reason)
481 async def _delete_old_logs(
482 self, prisma_client: PrismaClient, cutoff_date: datetime, deadline: float
483 ) -> TableCleanupResult:
484 return await self._delete_old_rows_batched(
485 prisma_client,
486 cutoff_date,
487 table_name="LiteLLM_SpendLogs",
488 key_columns=("request_id", "startTime"),
489 time_column="startTime",
490 deadline=deadline,
491 )
493 async def _delete_old_tool_index_rows(
494 self, prisma_client: PrismaClient, cutoff_date: datetime, deadline: float
495 ) -> TableCleanupResult:
496 # SpendLogToolIndex rows are derived from spend logs, so they expire on the
497 # same cutoff; rows older than retention point at already-deleted logs.
498 return await self._delete_old_rows_batched(
499 prisma_client,
500 cutoff_date,
501 table_name="LiteLLM_SpendLogToolIndex",
502 key_columns=("request_id", "tool_name"),
503 time_column="start_time",
504 deadline=deadline,
505 )
507 async def _delete_old_autorouter_session_rows(
508 self, prisma_client: PrismaClient, cutoff_date: datetime, deadline: float
509 ) -> TableCleanupResult:
510 return await self._delete_old_rows_batched(
511 prisma_client,
512 cutoff_date,
513 table_name="LiteLLM_AutoRouterSession",
514 key_columns=("api_key", "session_id", "router_name"),
515 time_column="last_turn_at",
516 deadline=deadline,
517 )
519 async def _delete_old_autorouter_user_session_rows(
520 self, prisma_client: PrismaClient, cutoff_date: datetime, deadline: float
521 ) -> TableCleanupResult:
522 return await self._delete_old_rows_batched(
523 prisma_client,
524 cutoff_date,
525 table_name="LiteLLM_AutoRouterUserSession",
526 key_columns=("user_id", "api_key", "session_id", "router_name"),
527 time_column="last_turn_at",
528 deadline=deadline,
529 )
531 async def _delete_old_health_check_rows(
532 self, prisma_client: PrismaClient, cutoff_date: datetime, deadline: float
533 ) -> TableCleanupResult:
534 return await self._delete_old_rows_batched(
535 prisma_client,
536 cutoff_date,
537 table_name="LiteLLM_HealthCheckTable",
538 key_columns=("health_check_id",),
539 time_column="checked_at",
540 deadline=deadline,
541 )
543 async def _clean_spend_log_tables(
544 self, prisma_client: PrismaClient, deadline: float
545 ) -> tuple[TableCleanupResult, ...]:
546 """
547 Prune the spend logs and the tool index rows derived from them.
549 When the table is range-partitioned, whole expired partitions are dropped
550 first because that reclaims disk immediately. Expired rows can still sit in
551 the DEFAULT partition (backfill, coverage gaps) or in a partition that spans
552 the cutoff, so retention still deletes those stragglers row-wise.
553 """
554 cutoff_date: Final = datetime.now(timezone.utc) - timedelta(seconds=float(self.retention_seconds or 0))
555 verbose_proxy_logger.info("Removing logs older than %s", cutoff_date.isoformat())
557 # Partition maintenance is DDL taking an ACCESS EXCLUSIVE lock, so it is
558 # only STARTED while the run still has budget, and each statement carries
559 # the same timeouts the batches do. Without those, a DROP would queue
560 # behind any long-running reader for as long as that reader lives, which
561 # is the one way this job could still outlast its budget without bound.
562 remaining_timeout_ms: Final = self._remaining_timeout_ms(deadline)
563 if time.monotonic() >= deadline:
564 verbose_proxy_logger.info("Run budget already spent, skipping partition maintenance this run")
565 elif self.general_settings.get(
566 "use_spend_logs_partitioning", False
567 ) and await self.partition_manager.is_partitioned(prisma_client, remaining_timeout_ms):
568 await self.partition_manager.ensure_partitions(prisma_client, remaining_timeout_ms)
569 dropped: Final = await self.partition_manager.drop_partitions_older_than(
570 prisma_client, cutoff_date, remaining_timeout_ms
571 )
572 verbose_proxy_logger.info("Dropped %d expired spend-log partitions: %s", len(dropped), dropped)
574 logs_result: Final = await self._delete_old_logs(prisma_client, cutoff_date, deadline)
575 verbose_proxy_logger.info("Deleted %s logs", logs_result.rows_deleted)
577 index_result: Final = await self._delete_old_tool_index_rows(prisma_client, cutoff_date, deadline)
578 verbose_proxy_logger.info("Deleted %s expired tool index rows", index_result.rows_deleted)
579 return (logs_result, index_result)
581 async def _clean_session_rollup(
582 self, prisma_client: PrismaClient, retention_seconds: int, deadline: float
583 ) -> tuple[TableCleanupResult, ...]:
584 """
585 Prune auto-router session rollup rows, which carry their own retention horizon.
586 """
587 session_cutoff: Final = datetime.now(timezone.utc) - timedelta(seconds=float(retention_seconds))
588 from litellm.proxy.db.baseline_accounting import BaselineAccountingStore
590 if remaining_ms := self._remaining_timeout_ms(deadline)():
591 try:
592 await BaselineAccountingStore.for_client(prisma_client).retire_before(
593 session_cutoff,
594 self.batch_size,
595 remaining_ms,
596 )
597 except Exception: # noqa: BLE001 # retained observations are retried by the next cleanup job
598 verbose_proxy_logger.warning("Auto-router baseline retention remains pending")
599 sessions_result: Final = await self._delete_old_autorouter_session_rows(
600 prisma_client, session_cutoff, self._group_deadline(deadline, 2)
601 )
602 verbose_proxy_logger.info("Deleted %s expired auto-router session rollup rows", sessions_result.rows_deleted)
603 user_sessions_result: Final = await self._delete_old_autorouter_user_session_rows(
604 prisma_client, session_cutoff, deadline
605 )
606 verbose_proxy_logger.info(
607 "Deleted %s expired auto-router user session rollup rows", user_sessions_result.rows_deleted
608 )
609 return (sessions_result, user_sessions_result)
611 async def _clean_health_checks(
612 self, prisma_client: PrismaClient, retention_seconds: int, deadline: float
613 ) -> tuple[TableCleanupResult, ...]:
614 health_check_cutoff: Final = datetime.now(timezone.utc) - timedelta(seconds=float(retention_seconds))
615 health_checks_result: Final = await self._delete_old_health_check_rows(
616 prisma_client, health_check_cutoff, deadline
617 )
618 verbose_proxy_logger.info(
619 "Deleted %s expired health-check rows",
620 health_checks_result.rows_deleted,
621 )
622 return (health_checks_result,)
624 @staticmethod
625 def _run_outcome(results: tuple[TableCleanupResult, ...]) -> RunOutcome:
626 """
627 Report the most operationally significant reason the run stopped.
629 A bound that was hit matters more than a table that simply ran dry, so
630 those win over "completed", and an abort wins over everything.
631 """
632 reasons: Final = frozenset(result.stop_reason for result in results)
633 if "aborted" in reasons:
634 return "aborted"
635 if "budget_exhausted" in reasons:
636 return "budget_exhausted"
637 if "batch_cap_reached" in reasons:
638 return "batch_cap_reached"
639 return "completed"
641 async def cleanup_old_spend_logs(self, prisma_client: PrismaClient) -> None:
642 """
643 Main cleanup function. Deletes old spend logs in batches.
644 If pod_lock_manager is available, ensures only one pod runs cleanup.
645 If no pod_lock_manager, runs cleanup without distributed locking.
646 """
647 lock_acquired = False
648 run_started_at: Final = time.monotonic()
649 progress: Final = _RunProgress()
650 progress_token: Final = _run_progress.set(progress)
651 try:
652 verbose_proxy_logger.info("Cleanup job triggered at %s", datetime.now())
653 self._refresh_bounds()
655 delete_spend_logs: Final = self._should_delete_spend_logs()
656 autorouter_retention_seconds: Final = self._retention_seconds_for(
657 "maximum_autorouter_session_retention_period"
658 )
659 health_check_retention_seconds: Final = self._retention_seconds_for("maximum_health_check_retention_period")
660 if (
661 not delete_spend_logs
662 and autorouter_retention_seconds is None
663 and health_check_retention_seconds is None
664 ):
665 SpendLogCleanupMetrics.record_run("skipped_disabled")
666 return
668 if delete_spend_logs and self.retention_seconds is None:
669 verbose_proxy_logger.error("Retention seconds is None, cannot proceed with cleanup")
670 SpendLogCleanupMetrics.record_run("skipped_disabled")
671 return
673 # If we have a pod lock manager, try to acquire the lock
674 if self.pod_lock_manager and self.pod_lock_manager.redis_cache:
675 lock_acquired = (
676 await self.pod_lock_manager.acquire_lock(
677 cronjob_id=SPEND_LOG_CLEANUP_JOB_NAME,
678 )
679 or False
680 )
681 verbose_proxy_logger.info(
682 "Lock acquisition attempt: %s at %s", "successful" if lock_acquired else "failed", datetime.now()
683 )
685 if not lock_acquired:
686 verbose_proxy_logger.info("Another pod is already running cleanup")
687 SpendLogCleanupMetrics.record_run("skipped_locked")
688 return
690 deadline: Final = time.monotonic() + self.run_budget_seconds
691 configured_group_count: Final = (
692 int(delete_spend_logs and self.retention_seconds is not None)
693 + int(autorouter_retention_seconds is not None)
694 + int(health_check_retention_seconds is not None)
695 )
697 spend_log_results: Final = (
698 await self._clean_spend_log_tables(
699 prisma_client,
700 self._group_deadline(deadline, configured_group_count),
701 )
702 if delete_spend_logs and self.retention_seconds is not None
703 else ()
704 )
705 remaining_groups_after_spend_logs: Final = int(autorouter_retention_seconds is not None) + int(
706 health_check_retention_seconds is not None
707 )
708 session_results: Final = (
709 await self._clean_session_rollup(
710 prisma_client,
711 autorouter_retention_seconds,
712 self._group_deadline(deadline, remaining_groups_after_spend_logs),
713 )
714 if autorouter_retention_seconds is not None
715 else ()
716 )
717 health_check_results: Final = (
718 await self._clean_health_checks(
719 prisma_client,
720 health_check_retention_seconds,
721 deadline,
722 )
723 if health_check_retention_seconds is not None
724 else ()
725 )
727 SpendLogCleanupMetrics.record_run(
728 self._run_outcome(spend_log_results + session_results + health_check_results)
729 )
731 except asyncio.CancelledError:
732 verbose_proxy_logger.error(
733 "Spend log cleanup cancelled after %.2fs (rows_deleted=%d, batches=%d); the next run resumes from here",
734 time.monotonic() - run_started_at,
735 progress.rows_deleted,
736 progress.batches,
737 )
738 SpendLogCleanupMetrics.record_run("aborted")
739 raise
740 except Exception as e:
741 # .exception() captures the traceback; str(e) alone on a Prisma/DB
742 # timeout is often empty and gives operators no signal to diagnose.
743 verbose_proxy_logger.exception(
744 "Error during spend log cleanup: %s: %s",
745 type(e).__name__,
746 e,
747 )
748 SpendLogCleanupMetrics.record_run("aborted")
749 return # Return after error handling
750 finally:
751 _run_progress.reset(progress_token)
752 # Only release the lock if it was actually acquired
753 if lock_acquired and self.pod_lock_manager and self.pod_lock_manager.redis_cache:
754 await self.pod_lock_manager.release_lock(cronjob_id=SPEND_LOG_CLEANUP_JOB_NAME)
755 verbose_proxy_logger.info("Released cleanup lock")