From 25c865876188959fd141435391ac0c0636a3e50d Mon Sep 17 00:00:00 2001 From: Harshit28j Date: Wed, 11 Mar 2026 15:34:24 +0530 Subject: [PATCH] Fix oversized-row handling, batch error resilience, and callback registration Vantage destination: - Skip individual CSV rows exceeding 2MB limit with a warning instead of uploading an oversized batch that Vantage would reject - Wrap each batch upload in try/except so remaining batches continue on failure; re-raise the first error after all batches are attempted Callback registration: - Use litellm.logging_callback_manager.add_litellm_callback() instead of raw litellm.callbacks.append() for DB-bootstrapped VantageLogger, ensuring proper dedup and manager visibility - Add "vantage" init handler in litellm_logging.py (both creation and lookup branches) so config.yaml string callbacks are properly resolved Co-Authored-By: Claude Opus 4.6 --- .../focus/destinations/vantage_destination.py | 49 +++++++++++++++---- litellm/litellm_core_utils/litellm_logging.py | 15 ++++++ litellm/proxy/proxy_server.py | 4 +- 3 files changed, 57 insertions(+), 11 deletions(-) diff --git a/litellm/integrations/focus/destinations/vantage_destination.py b/litellm/integrations/focus/destinations/vantage_destination.py index e2a21c773f4..59c4e0cdb9b 100644 --- a/litellm/integrations/focus/destinations/vantage_destination.py +++ b/litellm/integrations/focus/destinations/vantage_destination.py @@ -99,26 +99,41 @@ class FocusVantageDestination(FocusDestination): async def _upload_batched( self, client: httpx.AsyncClient, csv_bytes: bytes, filename: str ) -> None: - """Split the CSV into batches and upload each.""" + """Split the CSV into batches and upload each. + + Continues uploading remaining batches even if one fails, then raises + the first error encountered so callers know the export was partial. + """ lines = csv_bytes.split(b"\n") header = lines[0] data_lines = [line for line in lines[1:] if line.strip()] + first_error: Optional[Exception] = None batch_num = 0 for start in range(0, len(data_lines), VANTAGE_MAX_ROWS_PER_UPLOAD): batch_lines = data_lines[start : start + VANTAGE_MAX_ROWS_PER_UPLOAD] batch_csv = header + b"\n" + b"\n".join(batch_lines) + b"\n" - # If a single batch still exceeds 2 MB, split further by size - if len(batch_csv) > VANTAGE_MAX_BYTES_PER_UPLOAD: - await self._upload_size_limited( - client, header, batch_lines, filename, batch_num + try: + # If a single batch still exceeds 2 MB, split further by size + if len(batch_csv) > VANTAGE_MAX_BYTES_PER_UPLOAD: + await self._upload_size_limited( + client, header, batch_lines, filename, batch_num + ) + else: + batch_filename = f"{filename}.part{batch_num}" + await self._upload_csv(client, batch_csv, batch_filename) + except Exception as e: + verbose_logger.error( + "Vantage destination: batch %d failed: %s", batch_num, e ) - else: - batch_filename = f"{filename}.part{batch_num}" - await self._upload_csv(client, batch_csv, batch_filename) + if first_error is None: + first_error = e batch_num += 1 + if first_error is not None: + raise first_error + async def _upload_size_limited( self, client: httpx.AsyncClient, @@ -127,19 +142,33 @@ class FocusVantageDestination(FocusDestination): filename: str, batch_offset: int, ) -> None: - """Upload lines in chunks that stay under the 2 MB size limit.""" + """Upload lines in chunks that stay under the 2 MB size limit. + + Individual rows that exceed the limit on their own are skipped with + a warning — they cannot be split further. + """ current_chunk: list[bytes] = [] current_size = len(header) + 1 # header + newline sub_batch = 0 + header_size = len(header) + 1 for line in data_lines: line_size = len(line) + 1 # line + newline + + # Skip individual rows that exceed the limit on their own + if header_size + line_size > VANTAGE_MAX_BYTES_PER_UPLOAD: + verbose_logger.warning( + "Vantage destination: skipping oversized row (%d bytes)", + line_size, + ) + continue + if current_size + line_size > VANTAGE_MAX_BYTES_PER_UPLOAD and current_chunk: batch_csv = header + b"\n" + b"\n".join(current_chunk) + b"\n" batch_filename = f"{filename}.part{batch_offset}_{sub_batch}" await self._upload_csv(client, batch_csv, batch_filename) current_chunk = [] - current_size = len(header) + 1 + current_size = header_size sub_batch += 1 current_chunk.append(line) current_size += line_size diff --git a/litellm/litellm_core_utils/litellm_logging.py b/litellm/litellm_core_utils/litellm_logging.py index 6f587abcdf1..7157f4cb633 100644 --- a/litellm/litellm_core_utils/litellm_logging.py +++ b/litellm/litellm_core_utils/litellm_logging.py @@ -3873,6 +3873,15 @@ def _init_custom_logger_compatible_class( # noqa: PLR0915 focus_logger = FocusLogger() _in_memory_loggers.append(focus_logger) return focus_logger # type: ignore + elif logging_integration == "vantage": + from litellm.integrations.vantage.vantage_logger import VantageLogger + + for callback in _in_memory_loggers: + if isinstance(callback, VantageLogger): + return callback # type: ignore + vantage_logger = VantageLogger() + _in_memory_loggers.append(vantage_logger) + return vantage_logger # type: ignore elif logging_integration == "deepeval": for callback in _in_memory_loggers: if isinstance(callback, DeepEvalLogger): @@ -4243,6 +4252,12 @@ def get_custom_logger_compatible_class( # noqa: PLR0915 for callback in _in_memory_loggers: if isinstance(callback, FocusLogger): return callback + elif logging_integration == "vantage": + from litellm.integrations.vantage.vantage_logger import VantageLogger + + for callback in _in_memory_loggers: + if isinstance(callback, VantageLogger): + return callback elif logging_integration == "deepeval": for callback in _in_memory_loggers: if isinstance(callback, DeepEvalLogger): diff --git a/litellm/proxy/proxy_server.py b/litellm/proxy/proxy_server.py index 9bd4cd7f140..c71df3a2ce9 100644 --- a/litellm/proxy/proxy_server.py +++ b/litellm/proxy/proxy_server.py @@ -6152,7 +6152,9 @@ class ProxyStartupEvent: integration_token=db_settings.get("integration_token"), base_url=db_settings.get("base_url"), ) - litellm.callbacks.append(vantage_logger) # type: ignore[arg-type] + litellm.logging_callback_manager.add_litellm_callback( + vantage_logger + ) except Exception as e: verbose_proxy_logger.warning( "Failed to register VantageLogger from DB settings: %s", e