From 078d8322ad29e5b353f7c6328f448a539011d7b6 Mon Sep 17 00:00:00 2001 From: amasen02 Date: Fri, 4 Sep 2026 17:14:14 +0530 Subject: [PATCH 1/2] 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. --- .../db_transaction_queue/spend_log_cleanup.py | 18 ++++- .../proxy/test_spend_log_cleanup.py | 65 +++++++++++++++++++ 2 files changed, 82 insertions(+), 1 deletion(-) diff --git a/litellm/proxy/db/db_transaction_queue/spend_log_cleanup.py b/litellm/proxy/db/db_transaction_queue/spend_log_cleanup.py index e97e9f6e683..3068abbc1f4 100644 --- a/litellm/proxy/db/db_transaction_queue/spend_log_cleanup.py +++ b/litellm/proxy/db/db_transaction_queue/spend_log_cleanup.py @@ -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( diff --git a/tests/test_litellm/proxy/test_spend_log_cleanup.py b/tests/test_litellm/proxy/test_spend_log_cleanup.py index bf1538183ab..3819d2c5403 100644 --- a/tests/test_litellm/proxy/test_spend_log_cleanup.py +++ b/tests/test_litellm/proxy/test_spend_log_cleanup.py @@ -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() + From 913263d8429f3473b276338392c41a82dff0343f Mon Sep 17 00:00:00 2001 From: amasen02 Date: Fri, 4 Sep 2026 17:44:15 +0530 Subject: [PATCH 2/2] test(proxy): verify cancellation progress and lock release via prisma seam (#39730) --- .../proxy/test_spend_log_cleanup.py | 96 +++++++++++-------- 1 file changed, 58 insertions(+), 38 deletions(-) diff --git a/tests/test_litellm/proxy/test_spend_log_cleanup.py b/tests/test_litellm/proxy/test_spend_log_cleanup.py index 3819d2c5403..9175fc6a914 100644 --- a/tests/test_litellm/proxy/test_spend_log_cleanup.py +++ b/tests/test_litellm/proxy/test_spend_log_cleanup.py @@ -1424,45 +1424,20 @@ 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. + diagnostic containing progress metrics (including rows deleted and batches + completed prior to cancellation), 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"} + mock_db = MagicMock() + _wire_tx(mock_db) + # First batch deletes 150 rows. Second batch simulates cancellation mid-run. + mock_db.execute_raw = AsyncMock( + side_effect=[150, asyncio.CancelledError("simulated shutdown mid-run")] ) - 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] = [] + mock_prisma_client = MagicMock() + mock_prisma_client.db = mock_db - 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) @@ -1472,13 +1447,58 @@ async def test_cancelled_cleanup_releases_lock(): 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") + + 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) ) - prisma_client = MagicMock() with pytest.raises(asyncio.CancelledError): - await cleaner.cleanup_old_spend_logs(prisma_client) + await cleaner.cleanup_old_spend_logs(mock_prisma_client) + + assert recorded_runs == ["aborted"] + assert cleaner._run_rows_deleted == 150 + assert cleaner._run_batches == 1 + mock_lock_manager.release_lock.assert_awaited_once() + + 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=150" in rendered_log + assert "batches=1" 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_db = MagicMock() + _wire_tx(mock_db) + mock_db.execute_raw = AsyncMock( + side_effect=asyncio.CancelledError("simulated immediate shutdown") + ) + mock_prisma_client = MagicMock() + mock_prisma_client.db = mock_db + + 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 + + with pytest.raises(asyncio.CancelledError): + await cleaner.cleanup_old_spend_logs(mock_prisma_client) mock_lock_manager.release_lock.assert_awaited_once()