fix(proxy): record aborted outcome and log progress on cancelled spend-log cleanup (#39730)

When `SpendLogCleanup.cleanup_old_spend_logs()` is cancelled during shutdown
or pod termination, `asyncio.CancelledError` (deriving from `BaseException`)
bypassed the top-level `except Exception:` block. As a result, the run stopped
without recording `outcome="aborted"` in `SpendLogCleanupMetrics` and without
emitting any failure diagnostic before propagating.

- Catch `asyncio.CancelledError` explicitly in `cleanup_old_spend_logs()`.
- Log an exception diagnostic with elapsed time, rows deleted, and batch counts.
- Record `SpendLogCleanupMetrics.record_run("aborted")`.
- Re-raise `CancelledError` to preserve asyncio cancellation semantics while letting `finally:` cleanly release any held pod lock.
- Track in-flight `_run_rows_deleted` and `_run_batches` on `SpendLogCleanup`.
- Add unit tests covering cancellation metrics/logging and lock release.
This commit is contained in:
amasen02 2026-09-04 17:14:14 +05:30
parent c8635ecc67
commit 078d8322ad
2 changed files with 82 additions and 1 deletions

View file

@ -96,6 +96,8 @@ class SpendLogCleanup:
self.general_settings = general_settings or default_settings
self._refresh_bounds()
self._run_rows_deleted: int = 0
self._run_batches: int = 0
from litellm.proxy.proxy_server import proxy_logging_obj
pod_lock_manager: Final = proxy_logging_obj.db_spend_update_writer.pod_lock_manager
@ -422,6 +424,8 @@ class SpendLogCleanup:
total_deleted += deleted_count
run_count += 1
self._run_rows_deleted += deleted_count
self._run_batches += 1
# Add a small sleep to prevent overwhelming the database
await asyncio.sleep(0.1)
@ -590,6 +594,9 @@ class SpendLogCleanup:
If no pod_lock_manager, runs cleanup without distributed locking.
"""
lock_acquired = False
run_start: Final = time.monotonic()
self._run_rows_deleted = 0
self._run_batches = 0
try:
verbose_proxy_logger.info("Cleanup job triggered at %s", datetime.now())
self._refresh_bounds()
@ -670,7 +677,16 @@ class SpendLogCleanup:
self._run_outcome(spend_log_results + session_results + health_check_results)
)
except Exception as e:
except asyncio.CancelledError:
verbose_proxy_logger.exception(
"Spend log cleanup cancelled (elapsed=%.2fs, rows_deleted=%d, batches=%d)",
time.monotonic() - run_start,
self._run_rows_deleted,
self._run_batches,
)
SpendLogCleanupMetrics.record_run("aborted")
raise
except Exception as e: # noqa: BLE001 - top-level cleanup job exception handler
# .exception() captures the traceback; str(e) alone on a Prisma/DB
# timeout is often empty and gives operators no signal to diagnose.
verbose_proxy_logger.exception(

View file

@ -1417,3 +1417,68 @@ def test_the_reported_run_outcome_is_the_most_significant_reason_in_any_order(st
"""
results = tuple(TableCleanupResult(rows_deleted=0, stop_reason=reason) for reason in stop_reasons)
assert SpendLogCleanup._run_outcome(results) == expected
@pytest.mark.asyncio
async def test_cancelled_cleanup_records_and_logs_aborted_run(monkeypatch):
"""
When cleanup_old_spend_logs is cancelled (e.g. during pod shutdown),
asyncio.CancelledError must record outcome='aborted', emit an error-level
diagnostic containing progress metrics, and re-raise CancelledError.
"""
import litellm.proxy.db.db_transaction_queue.spend_log_cleanup as cleanup_module
cleaner = SpendLogCleanup(
general_settings={"maximum_spend_logs_retention_period": "7d"}
)
cleaner.pod_lock_manager = None
cleaner._clean_spend_log_tables = AsyncMock(
side_effect=asyncio.CancelledError("simulated shutdown")
)
prisma_client = MagicMock()
logger = MagicMock()
recorded_runs: list[str] = []
monkeypatch.setattr(cleanup_module, "verbose_proxy_logger", logger)
monkeypatch.setattr(
cleanup_module.SpendLogCleanupMetrics, "record_run", lambda outcome: recorded_runs.append(outcome)
)
with pytest.raises(asyncio.CancelledError):
await cleaner.cleanup_old_spend_logs(prisma_client)
assert recorded_runs == ["aborted"]
logger.exception.assert_called_once()
log_call = logger.exception.call_args[0]
rendered_log = log_call[0] % log_call[1:]
assert "cancelled" in rendered_log.lower()
assert "elapsed=" in rendered_log
assert "rows_deleted=" in rendered_log
assert "batches=" in rendered_log
@pytest.mark.asyncio
async def test_cancelled_cleanup_releases_lock():
"""
When a cleanup run holds a distributed pod lock and is cancelled,
the finally block must release the lock before propagating CancelledError.
"""
mock_lock_manager = MagicMock()
mock_lock_manager.redis_cache = MagicMock()
mock_lock_manager.acquire_lock = AsyncMock(return_value=True)
mock_lock_manager.release_lock = AsyncMock()
cleaner = SpendLogCleanup(
general_settings={"maximum_spend_logs_retention_period": "7d"}
)
cleaner.pod_lock_manager = mock_lock_manager
cleaner._clean_spend_log_tables = AsyncMock(
side_effect=asyncio.CancelledError("shutdown")
)
prisma_client = MagicMock()
with pytest.raises(asyncio.CancelledError):
await cleaner.cleanup_old_spend_logs(prisma_client)
mock_lock_manager.release_lock.assert_awaited_once()