fix(mavvrik): two P1s — delete for env-var deployments + GCS session leak

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) <noreply@anthropic.com>
This commit is contained in:
Praveen Ghuge 2026-04-26 22:02:16 +05:30
parent 4962c1d62a
commit ecfeebc79c
2 changed files with 45 additions and 28 deletions

View file

@ -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"}

View file

@ -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