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 <noreply@anthropic.com>
This commit is contained in:
Harshit28j 2026-03-11 15:34:24 +05:30
parent 9e623ceaff
commit 25c8658761
3 changed files with 57 additions and 11 deletions

View file

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

View file

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

View file

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