From ecfeebc79c9c0b56c49b76acbc323fbbbdc2dc34 Mon Sep 17 00:00:00 2001 From: Praveen Ghuge Date: Sun, 26 Apr 2026 22:02:16 +0530 Subject: [PATCH] =?UTF-8?q?fix(mavvrik):=20two=20P1s=20=E2=80=94=20delete?= =?UTF-8?q?=20for=20env-var=20deployments=20+=20GCS=20session=20leak?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit P1: Service.delete() fails for env-var-only deployments settings.delete() raised LookupError (no DB row) before remove_job() was reached — the scheduler kept running with no way to stop it via API. Fix: deregister scheduler job first (independent of DB), then only delete the DB row when credentials are NOT from env vars. P1: GCS resumable session left open on _put_chunk failure A mid-stream exception abandoned the session URI; GCS held it open for up to 1 week, accumulating stale sessions on repeated transient failures. Fix: wrap the chunk loop in try/except — on any exception, send DELETE to the session URI (best-effort via contextlib.suppress) then re-raise. Co-Authored-By: Claude Sonnet 4.6 (1M context) --- litellm/integrations/mavvrik/__init__.py | 16 +++++-- litellm/integrations/mavvrik/uploader.py | 57 ++++++++++++++---------- 2 files changed, 45 insertions(+), 28 deletions(-) diff --git a/litellm/integrations/mavvrik/__init__.py b/litellm/integrations/mavvrik/__init__.py index 7d9a8d58571..2a0442963c6 100644 --- a/litellm/integrations/mavvrik/__init__.py +++ b/litellm/integrations/mavvrik/__init__.py @@ -259,17 +259,21 @@ class Service: # ------------------------------------------------------------------ async def delete(self) -> dict: - """Remove all Mavvrik settings and deregister the scheduler job. + """Deregister the scheduler job and remove DB settings if present. + + Works for both env-var-only and DB-backed deployments: + - Scheduler job is always deregistered (independent of DB). + - DB row is only deleted when credentials come from the DB; env-var + deployments have no row so this step is skipped. Raises: - LookupError: when no settings exist in the database. + LookupError: only when DB is connected and no settings row exists. """ from litellm.constants import MAVVRIK_EXPORT_USAGE_DATA_JOB_NAME import litellm.proxy.proxy_server as _pserver - await self._settings.delete() - + # Deregister scheduler first — independent of whether creds are in DB or env vars. _scheduler = getattr(_pserver, "scheduler", None) if _scheduler is not None: try: @@ -277,6 +281,10 @@ class Service: except Exception: pass # job may not exist if scheduler was restarted + # Only delete the DB row for DB-backed deployments; env-var-only has no row. + if not self._settings.has_env_vars: + await self._settings.delete() # raises LookupError if no row found + verbose_proxy_logger.info("mavvrik settings deleted") return {"message": "Mavvrik settings deleted successfully", "status": "success"} diff --git a/litellm/integrations/mavvrik/uploader.py b/litellm/integrations/mavvrik/uploader.py index 342c875fa84..af553c37c5b 100644 --- a/litellm/integrations/mavvrik/uploader.py +++ b/litellm/integrations/mavvrik/uploader.py @@ -19,6 +19,7 @@ GCS resumable upload protocol reference: https://cloud.google.com/storage/docs/resumable-uploads """ +import contextlib import gzip import io from typing import TYPE_CHECKING, Any, AsyncIterator @@ -184,35 +185,43 @@ class Uploader: session_uri: str = "" has_data = False - async for csv_chunk in pages: - if not csv_chunk: - continue + try: + async for csv_chunk in pages: + if not csv_chunk: + continue + + if not has_data: + signed_url = await self._client.get_signed_url(date_str) + session_uri = await self._initiate_resumable_upload(signed_url) + has_data = True + + gz.write(csv_chunk.encode("utf-8")) + gz.flush() + gz_buffer.extend(raw_buf.getvalue()) + raw_buf.seek(0) + raw_buf.truncate(0) + + while len(gz_buffer) >= _GCS_CHUNK_SIZE: + chunk = bytes(gz_buffer[:_GCS_CHUNK_SIZE]) + gz_buffer = gz_buffer[_GCS_CHUNK_SIZE:] + await self._put_chunk(session_uri, chunk, offset=offset, final=False) + offset += len(chunk) if not has_data: - signed_url = await self._client.get_signed_url(date_str) - session_uri = await self._initiate_resumable_upload(signed_url) - has_data = True + verbose_proxy_logger.debug("uploader: no data to stream, skipping upload") + return 0 - gz.write(csv_chunk.encode("utf-8")) - gz.flush() + gz.close() gz_buffer.extend(raw_buf.getvalue()) - raw_buf.seek(0) - raw_buf.truncate(0) + total = offset + len(gz_buffer) + await self._put_chunk(session_uri, bytes(gz_buffer), offset=offset, final=True) - while len(gz_buffer) >= _GCS_CHUNK_SIZE: - chunk = bytes(gz_buffer[:_GCS_CHUNK_SIZE]) - gz_buffer = gz_buffer[_GCS_CHUNK_SIZE:] - await self._put_chunk(session_uri, chunk, offset=offset, final=False) - offset += len(chunk) - - if not has_data: - verbose_proxy_logger.debug("uploader: no data to stream, skipping upload") - return 0 - - gz.close() - gz_buffer.extend(raw_buf.getvalue()) - total = offset + len(gz_buffer) - await self._put_chunk(session_uri, bytes(gz_buffer), offset=offset, final=True) + except Exception: + # Cancel the open GCS session so it doesn't linger for up to 1 week. + if session_uri: + with contextlib.suppress(Exception): + await http_request("DELETE", session_uri, timeout=10.0, label="cancel") + raise verbose_proxy_logger.info( "uploader: stream upload complete — %d bytes for date %s", total, date_str