mirror of
https://github.com/BerriAI/litellm.git
synced 2026-10-03 02:22:24 +00:00
add soem changes
This commit is contained in:
parent
118176f21a
commit
e384cb771a
3 changed files with 455 additions and 148 deletions
|
|
@ -3,7 +3,7 @@ import json
|
|||
import os
|
||||
import re
|
||||
import time
|
||||
from typing import Any, Callable, Dict, List, Optional, Tuple, Union
|
||||
from typing import Any, Callable, Dict, Iterator, List, Optional, Tuple, Union, cast
|
||||
|
||||
import httpx
|
||||
from httpx import Headers, Response
|
||||
|
|
@ -137,42 +137,150 @@ def _get_litellm_batch_custom_id_from_labels(labels: Dict[str, Any]) -> str:
|
|||
return str(labels.get("litellm_custom_id", "unknown"))
|
||||
|
||||
|
||||
def _openai_batch_jsonl_entries_to_vertex_wrapped_requests(
|
||||
openai_jsonl_content: List[Dict[str, Any]],
|
||||
def _openai_batch_jsonl_entry_to_vertex_wrapped_request(
|
||||
openai_entry: Dict[str, Any],
|
||||
map_openai_to_vertex_params: Callable[[Dict[str, Any]], Dict[str, Any]],
|
||||
) -> List[Dict[str, Any]]:
|
||||
) -> Dict[str, Any]:
|
||||
"""
|
||||
Transforms OpenAI JSONL batch entries to Vertex AI JSONL lines.
|
||||
Transforms a single OpenAI JSONL batch entry into its Vertex wrapped request.
|
||||
|
||||
jsonl body for vertex is {"request": <request_body>}
|
||||
Example Vertex jsonl
|
||||
{"request":{"contents": [{"role": "user", "parts": [{"text": "What is the relation between the following video and image samples?"}, {"fileData": {"fileUri": "gs://cloud-samples-data/generative-ai/video/animals.mp4", "mimeType": "video/mp4"}}, {"fileData": {"fileUri": "gs://cloud-samples-data/generative-ai/image/cricket.jpeg", "mimeType": "image/jpeg"}}]}]}}
|
||||
{"request":{"contents": [{"role": "user", "parts": [{"text": "Describe what is happening in this video."}, {"fileData": {"fileUri": "gs://cloud-samples-data/generative-ai/video/another_video.mov", "mimeType": "video/mov"}}]}]}}
|
||||
"""
|
||||
openai_request_body = openai_entry.get("body") or {}
|
||||
vertex_request_body = _transform_request_body(
|
||||
messages=openai_request_body.get("messages", []),
|
||||
model=openai_request_body.get("model", ""),
|
||||
optional_params=map_openai_to_vertex_params(openai_request_body),
|
||||
custom_llm_provider="vertex_ai",
|
||||
litellm_params={},
|
||||
cached_content=None,
|
||||
)
|
||||
|
||||
vertex_jsonl_content = []
|
||||
for _openai_jsonl_content in openai_jsonl_content:
|
||||
openai_request_body = _openai_jsonl_content.get("body") or {}
|
||||
vertex_request_body = _transform_request_body(
|
||||
messages=openai_request_body.get("messages", []),
|
||||
model=openai_request_body.get("model", ""),
|
||||
optional_params=map_openai_to_vertex_params(openai_request_body),
|
||||
custom_llm_provider="vertex_ai",
|
||||
litellm_params={},
|
||||
cached_content=None,
|
||||
custom_id = openai_entry.get("custom_id")
|
||||
if custom_id is not None:
|
||||
if "labels" not in vertex_request_body:
|
||||
vertex_request_body["labels"] = {}
|
||||
_set_litellm_batch_custom_id_labels(vertex_request_body["labels"], custom_id)
|
||||
|
||||
return {"request": vertex_request_body}
|
||||
|
||||
|
||||
def _openai_batch_jsonl_entries_to_vertex_wrapped_requests(
|
||||
openai_jsonl_content: List[Dict[str, Any]],
|
||||
map_openai_to_vertex_params: Callable[[Dict[str, Any]], Dict[str, Any]],
|
||||
) -> List[Dict[str, Any]]:
|
||||
return [
|
||||
_openai_batch_jsonl_entry_to_vertex_wrapped_request(
|
||||
entry, map_openai_to_vertex_params
|
||||
)
|
||||
for entry in openai_jsonl_content
|
||||
]
|
||||
|
||||
# Add custom_id as a label for correlation in batch outputs
|
||||
custom_id = _openai_jsonl_content.get("custom_id")
|
||||
if custom_id is not None:
|
||||
if "labels" not in vertex_request_body:
|
||||
vertex_request_body["labels"] = {}
|
||||
_set_litellm_batch_custom_id_labels(
|
||||
vertex_request_body["labels"], custom_id
|
||||
|
||||
def _iter_openai_jsonl_lines(openai_file_content: FileTypes) -> Iterator[str]:
|
||||
"""
|
||||
Yield non-empty JSONL lines one at a time without materializing the whole
|
||||
payload, so peak memory stays bounded regardless of payload size. Mirrors
|
||||
``str.splitlines()`` + ``line.strip()`` for ``\\n`` / ``\\r\\n`` delimited
|
||||
JSONL.
|
||||
"""
|
||||
content: Any = openai_file_content
|
||||
if isinstance(content, tuple):
|
||||
content = content[1]
|
||||
|
||||
if isinstance(content, PathLike):
|
||||
with open(str(content), "rb") as handle:
|
||||
for raw in handle:
|
||||
line = raw.decode("utf-8").strip()
|
||||
if line:
|
||||
yield line
|
||||
return
|
||||
|
||||
if isinstance(content, (bytes, bytearray)):
|
||||
newline = ord("\n")
|
||||
start, length = 0, len(content)
|
||||
while start < length:
|
||||
idx = content.find(newline, start)
|
||||
if idx == -1:
|
||||
chunk, start = content[start:], length
|
||||
else:
|
||||
chunk, start = content[start:idx], idx + 1
|
||||
line = chunk.decode("utf-8").strip()
|
||||
if line:
|
||||
yield line
|
||||
return
|
||||
|
||||
if isinstance(content, str):
|
||||
str_start, str_length = 0, len(content)
|
||||
while str_start < str_length:
|
||||
str_idx = content.find("\n", str_start)
|
||||
if str_idx == -1:
|
||||
str_chunk, str_start = content[str_start:], str_length
|
||||
else:
|
||||
str_chunk, str_start = content[str_start:str_idx], str_idx + 1
|
||||
str_line = str_chunk.strip()
|
||||
if str_line:
|
||||
yield str_line
|
||||
return
|
||||
|
||||
if hasattr(content, "read"):
|
||||
for raw in content:
|
||||
line = raw.decode("utf-8") if isinstance(raw, bytes) else raw
|
||||
line = line.strip()
|
||||
if line:
|
||||
yield line
|
||||
return
|
||||
|
||||
raise ValueError("Unsupported file content type")
|
||||
|
||||
|
||||
def _iter_openai_jsonl_entries(
|
||||
openai_file_content: FileTypes,
|
||||
) -> Iterator[Dict[str, Any]]:
|
||||
for line in _iter_openai_jsonl_lines(openai_file_content):
|
||||
yield json.loads(line)
|
||||
|
||||
|
||||
def _stream_openai_jsonl_to_vertex(
|
||||
openai_file_content: FileTypes,
|
||||
map_openai_to_vertex_params: Callable[[Dict[str, Any]], Dict[str, Any]],
|
||||
as_bytes: bool,
|
||||
) -> Tuple[Union[str, bytes], Optional[Dict[str, Any]]]:
|
||||
"""
|
||||
Stream the OpenAI -> Vertex JSONL transform entry-by-entry.
|
||||
|
||||
Returns the joined Vertex JSONL payload and the first parsed OpenAI entry,
|
||||
so callers can derive the GCS object name without re-parsing. ``as_bytes``
|
||||
returns bytes (the upload path ships them to httpx without a str->bytes
|
||||
re-encode); otherwise a str is returned.
|
||||
|
||||
The bytes path extends one ``bytearray`` in place instead of collecting a
|
||||
list of encoded lines and joining, which would hold both the list and the
|
||||
joined result at once and double peak output memory.
|
||||
"""
|
||||
first_entry: Optional[Dict[str, Any]] = None
|
||||
byte_buf = bytearray()
|
||||
str_parts: List[str] = []
|
||||
for entry in _iter_openai_jsonl_entries(openai_file_content):
|
||||
if first_entry is None:
|
||||
first_entry = entry
|
||||
line = json.dumps(
|
||||
_openai_batch_jsonl_entry_to_vertex_wrapped_request(
|
||||
entry, map_openai_to_vertex_params
|
||||
)
|
||||
)
|
||||
if as_bytes:
|
||||
if byte_buf:
|
||||
byte_buf.extend(b"\n")
|
||||
byte_buf.extend(line.encode("utf-8"))
|
||||
else:
|
||||
str_parts.append(line)
|
||||
|
||||
vertex_jsonl_content.append({"request": vertex_request_body})
|
||||
return vertex_jsonl_content
|
||||
if as_bytes:
|
||||
return bytes(byte_buf), first_entry
|
||||
return "\n".join(str_parts), first_entry
|
||||
|
||||
|
||||
class VertexAIFilesConfig(VertexBase, BaseFilesConfig):
|
||||
|
|
@ -208,43 +316,6 @@ class VertexAIFilesConfig(VertexBase, BaseFilesConfig):
|
|||
headers["Authorization"] = f"Bearer {api_key}"
|
||||
return headers
|
||||
|
||||
def _get_content_from_openai_file(self, openai_file_content: FileTypes) -> str:
|
||||
"""
|
||||
Helper to extract content from various OpenAI file types and return as string.
|
||||
|
||||
Handles:
|
||||
- Direct content (str, bytes, IO[bytes])
|
||||
- Tuple formats: (filename, content, [content_type], [headers])
|
||||
- PathLike objects
|
||||
"""
|
||||
content: Union[str, bytes] = b""
|
||||
# Extract file content from tuple if necessary
|
||||
if isinstance(openai_file_content, tuple):
|
||||
# Take the second element which is always the file content
|
||||
file_content = openai_file_content[1]
|
||||
else:
|
||||
file_content = openai_file_content
|
||||
|
||||
# Handle different file content types
|
||||
if isinstance(file_content, str):
|
||||
# String content can be used directly
|
||||
content = file_content
|
||||
elif isinstance(file_content, bytes):
|
||||
# Bytes content can be decoded
|
||||
content = file_content
|
||||
elif isinstance(file_content, PathLike): # PathLike
|
||||
with open(str(file_content), "rb") as f:
|
||||
content = f.read()
|
||||
elif hasattr(file_content, "read"): # IO[bytes]
|
||||
# File-like objects need to be read
|
||||
content = file_content.read()
|
||||
|
||||
# Ensure content is string
|
||||
if isinstance(content, bytes):
|
||||
content = content.decode("utf-8")
|
||||
|
||||
return content
|
||||
|
||||
def _get_gcs_object_name_from_batch_jsonl(
|
||||
self,
|
||||
openai_jsonl_content: List[Dict[str, Any]],
|
||||
|
|
@ -273,17 +344,12 @@ class VertexAIFilesConfig(VertexBase, BaseFilesConfig):
|
|||
raise ValueError("file content is required")
|
||||
|
||||
if purpose == "batch":
|
||||
## 1. If jsonl, check if there's a model name
|
||||
file_content = self._get_content_from_openai_file(
|
||||
extracted_file_data_content
|
||||
## 1. If jsonl, derive the object name from the first entry's model
|
||||
first_entry = next(
|
||||
_iter_openai_jsonl_entries(extracted_file_data_content), None
|
||||
)
|
||||
|
||||
# Split into lines and parse each line as JSON
|
||||
openai_jsonl_content = [
|
||||
json.loads(line) for line in file_content.splitlines() if line.strip()
|
||||
]
|
||||
if len(openai_jsonl_content) > 0:
|
||||
return self._get_gcs_object_name_from_batch_jsonl(openai_jsonl_content)
|
||||
if first_entry is not None:
|
||||
return self._get_gcs_object_name_from_batch_jsonl([first_entry])
|
||||
|
||||
## 2. If not jsonl, store under a server-generated managed object name
|
||||
filename = extracted_file_data.get("filename")
|
||||
|
|
@ -399,21 +465,12 @@ class VertexAIFilesConfig(VertexBase, BaseFilesConfig):
|
|||
create_file_data=create_file_data,
|
||||
extracted_file_data=extracted_file_data,
|
||||
):
|
||||
## 1. If jsonl, check if there's a model name
|
||||
file_content = self._get_content_from_openai_file(
|
||||
extracted_file_data_content
|
||||
vertex_jsonl_bytes, _ = _stream_openai_jsonl_to_vertex(
|
||||
extracted_file_data_content,
|
||||
self._map_openai_to_vertex_params,
|
||||
as_bytes=True,
|
||||
)
|
||||
|
||||
# Split into lines and parse each line as JSON
|
||||
openai_jsonl_content = [
|
||||
json.loads(line) for line in file_content.splitlines() if line.strip()
|
||||
]
|
||||
vertex_jsonl_content = (
|
||||
self._transform_openai_jsonl_content_to_vertex_ai_jsonl_content(
|
||||
openai_jsonl_content
|
||||
)
|
||||
)
|
||||
return "\n".join(json.dumps(item) for item in vertex_jsonl_content)
|
||||
return vertex_jsonl_bytes
|
||||
elif isinstance(extracted_file_data_content, bytes):
|
||||
return extracted_file_data_content
|
||||
else:
|
||||
|
|
@ -811,25 +868,16 @@ class VertexAIJsonlFilesTransformation(VertexGeminiConfig):
|
|||
|
||||
if openai_file_content is None:
|
||||
raise ValueError("contents of file are None")
|
||||
# Read the content of the file
|
||||
file_content = self._get_content_from_openai_file(openai_file_content)
|
||||
|
||||
# Split into lines and parse each line as JSON
|
||||
openai_jsonl_content = [
|
||||
json.loads(line) for line in file_content.splitlines() if line.strip()
|
||||
]
|
||||
vertex_jsonl_content = (
|
||||
self._transform_openai_jsonl_content_to_vertex_ai_jsonl_content(
|
||||
openai_jsonl_content
|
||||
)
|
||||
vertex_jsonl_string, first_entry = _stream_openai_jsonl_to_vertex(
|
||||
openai_file_content,
|
||||
self._map_openai_to_vertex_params,
|
||||
as_bytes=False,
|
||||
)
|
||||
vertex_jsonl_string = "\n".join(
|
||||
json.dumps(item) for item in vertex_jsonl_content
|
||||
)
|
||||
object_name = self._get_gcs_object_name(
|
||||
openai_jsonl_content=openai_jsonl_content
|
||||
)
|
||||
return vertex_jsonl_string, object_name
|
||||
if first_entry is None:
|
||||
raise ValueError("contents of file are empty")
|
||||
object_name = self._get_gcs_object_name(openai_jsonl_content=[first_entry])
|
||||
return cast(str, vertex_jsonl_string), object_name
|
||||
|
||||
def _transform_openai_jsonl_content_to_vertex_ai_jsonl_content(
|
||||
self, openai_jsonl_content: List[Dict[str, Any]]
|
||||
|
|
@ -871,43 +919,6 @@ class VertexAIJsonlFilesTransformation(VertexGeminiConfig):
|
|||
)
|
||||
return vertex_params
|
||||
|
||||
def _get_content_from_openai_file(self, openai_file_content: FileTypes) -> str:
|
||||
"""
|
||||
Helper to extract content from various OpenAI file types and return as string.
|
||||
|
||||
Handles:
|
||||
- Direct content (str, bytes, IO[bytes])
|
||||
- Tuple formats: (filename, content, [content_type], [headers])
|
||||
- PathLike objects
|
||||
"""
|
||||
content: Union[str, bytes] = b""
|
||||
# Extract file content from tuple if necessary
|
||||
if isinstance(openai_file_content, tuple):
|
||||
# Take the second element which is always the file content
|
||||
file_content = openai_file_content[1]
|
||||
else:
|
||||
file_content = openai_file_content
|
||||
|
||||
# Handle different file content types
|
||||
if isinstance(file_content, str):
|
||||
# String content can be used directly
|
||||
content = file_content
|
||||
elif isinstance(file_content, bytes):
|
||||
# Bytes content can be decoded
|
||||
content = file_content
|
||||
elif isinstance(file_content, PathLike): # PathLike
|
||||
with open(str(file_content), "rb") as f:
|
||||
content = f.read()
|
||||
elif hasattr(file_content, "read"): # IO[bytes]
|
||||
# File-like objects need to be read
|
||||
content = file_content.read()
|
||||
|
||||
# Ensure content is string
|
||||
if isinstance(content, bytes):
|
||||
content = content.decode("utf-8")
|
||||
|
||||
return content
|
||||
|
||||
def transform_gcs_bucket_response_to_openai_file_object(
|
||||
self, create_file_data: CreateFileRequest, gcs_upload_response: Dict[str, Any]
|
||||
) -> OpenAIFileObject:
|
||||
|
|
|
|||
|
|
@ -8,6 +8,7 @@ Regression test for: UTF-8 codec error when uploading binary files
|
|||
"""
|
||||
|
||||
import io
|
||||
import json
|
||||
import pytest
|
||||
from unittest.mock import AsyncMock, MagicMock, patch
|
||||
|
||||
|
|
@ -137,11 +138,12 @@ class TestVertexAIBinaryFileUpload:
|
|||
), "Binary file data should remain as bytes"
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_jsonl_file_upload_returns_string(self):
|
||||
async def test_jsonl_file_upload_returns_utf8_bytes(self):
|
||||
"""
|
||||
Test that JSONL files (text) are correctly transformed to strings.
|
||||
Test that JSONL batch files are transformed to UTF-8 bytes.
|
||||
|
||||
This ensures we handle both binary and text files correctly.
|
||||
The transform emits bytes so the upload ships the payload to GCS without
|
||||
a str->bytes re-encode, avoiding a full extra copy of the payload.
|
||||
"""
|
||||
# Create mock JSONL content
|
||||
mock_jsonl_content = (
|
||||
|
|
@ -164,10 +166,14 @@ class TestVertexAIBinaryFileUpload:
|
|||
litellm_params={},
|
||||
)
|
||||
|
||||
# JSONL files should be transformed to string
|
||||
assert isinstance(
|
||||
transformed_request, str
|
||||
), f"Expected string for JSONL file, got {type(transformed_request)}"
|
||||
transformed_request, bytes
|
||||
), f"Expected bytes for JSONL file, got {type(transformed_request)}"
|
||||
|
||||
decoded = json.loads(transformed_request.decode("utf-8"))
|
||||
assert (
|
||||
"request" in decoded
|
||||
), "JSONL transform must wrap each row in {'request': ...}"
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_mixed_file_types_in_sequence(self):
|
||||
|
|
@ -208,7 +214,7 @@ class TestVertexAIBinaryFileUpload:
|
|||
optional_params={},
|
||||
litellm_params={},
|
||||
)
|
||||
assert isinstance(result2, str)
|
||||
assert isinstance(result2, bytes)
|
||||
|
||||
# Test 3: Upload another binary file
|
||||
binary_content2 = b"\xc4\xe5\xf2\xe5\xeb"
|
||||
|
|
@ -251,7 +257,7 @@ class TestVertexAIBinaryFileUpload:
|
|||
},
|
||||
"text_files": {
|
||||
"input_type": "str or bytes",
|
||||
"output_type": "str",
|
||||
"output_type": "bytes",
|
||||
"examples": ["JSONL", "CSV", "TXT"],
|
||||
"http_method": "POST",
|
||||
"encoding": "UTF-8",
|
||||
|
|
|
|||
|
|
@ -0,0 +1,290 @@
|
|||
"""
|
||||
Tests for the streaming OpenAI -> Vertex JSONL batch transform.
|
||||
|
||||
The transform converts batch uploads entry-by-entry rather than materializing
|
||||
the payload in full intermediate lists (decoded str, parsed dicts, transformed
|
||||
dicts, joined output), which keeps peak memory bounded on large uploads.
|
||||
|
||||
These tests lock in the behaviour that would regress if the streaming path were
|
||||
replaced by a list-based pipeline:
|
||||
1. Byte-for-byte output parity with a list pipeline (wire format).
|
||||
2. The streaming transform peaks at a clear fraction of a list pipeline on the
|
||||
same input (relative differential, robust to GC noise).
|
||||
3. ``get_object_name`` only parses the first JSONL row, so a payload whose
|
||||
later rows are not valid JSON does not raise.
|
||||
4. A tuple-wrapped file handle uploaded through the real create_file ordering
|
||||
keeps every row, including entry 0 (no partial upload from a consumed
|
||||
cursor).
|
||||
"""
|
||||
|
||||
import gc
|
||||
import io
|
||||
import json
|
||||
import tracemalloc
|
||||
|
||||
import pytest
|
||||
|
||||
from litellm.llms.vertex_ai.files.transformation import (
|
||||
VertexAIFilesConfig,
|
||||
VertexAIJsonlFilesTransformation,
|
||||
_get_litellm_batch_custom_id_from_labels,
|
||||
_iter_openai_jsonl_entries,
|
||||
_iter_openai_jsonl_lines,
|
||||
_stream_openai_jsonl_to_vertex,
|
||||
)
|
||||
from litellm.types.llms.openai import CreateFileRequest
|
||||
|
||||
|
||||
def _make_openai_jsonl_bytes(n_rows: int, padding: int = 400) -> bytes:
|
||||
pad = "x" * padding
|
||||
rows = []
|
||||
for i in range(n_rows):
|
||||
rows.append(
|
||||
json.dumps(
|
||||
{
|
||||
"custom_id": f"request-{i}",
|
||||
"method": "POST",
|
||||
"url": "/v1/chat/completions",
|
||||
"body": {
|
||||
"model": "gemini-2.5-flash",
|
||||
"messages": [{"role": "user", "content": f"{pad} {i}"}],
|
||||
"max_tokens": 4,
|
||||
},
|
||||
}
|
||||
)
|
||||
)
|
||||
return ("\n".join(rows)).encode("utf-8")
|
||||
|
||||
|
||||
def _legacy_vertex_jsonl_string(cfg: VertexAIFilesConfig, content: str) -> str:
|
||||
"""A list-based pipeline, reconstructed for parity comparison."""
|
||||
entries = [json.loads(line) for line in content.splitlines() if line.strip()]
|
||||
vertex = cfg._transform_openai_jsonl_content_to_vertex_ai_jsonl_content(entries)
|
||||
return "\n".join(json.dumps(item) for item in vertex)
|
||||
|
||||
|
||||
class TestStreamingOutputParity:
|
||||
def test_streaming_bytes_match_legacy_pipeline(self):
|
||||
cfg = VertexAIFilesConfig()
|
||||
raw = _make_openai_jsonl_bytes(500)
|
||||
|
||||
streamed, first_entry = _stream_openai_jsonl_to_vertex(
|
||||
raw, cfg._map_openai_to_vertex_params, as_bytes=True
|
||||
)
|
||||
legacy = _legacy_vertex_jsonl_string(cfg, raw.decode("utf-8"))
|
||||
|
||||
assert isinstance(streamed, bytes)
|
||||
assert streamed.decode("utf-8") == legacy
|
||||
assert first_entry is not None
|
||||
assert first_entry["custom_id"] == "request-0"
|
||||
|
||||
def test_transform_create_file_request_returns_bytes_parity(self):
|
||||
cfg = VertexAIFilesConfig()
|
||||
raw = _make_openai_jsonl_bytes(300)
|
||||
request: CreateFileRequest = {
|
||||
"file": ("batch.jsonl", raw, "application/jsonl"),
|
||||
"purpose": "batch",
|
||||
}
|
||||
|
||||
out = cfg.transform_create_file_request(
|
||||
model="", create_file_data=request, optional_params={}, litellm_params={}
|
||||
)
|
||||
|
||||
assert isinstance(out, bytes)
|
||||
assert out.decode("utf-8") == _legacy_vertex_jsonl_string(
|
||||
cfg, raw.decode("utf-8")
|
||||
)
|
||||
|
||||
|
||||
class TestFileLikeInputNotPartiallyConsumed:
|
||||
"""
|
||||
In ``llm_http_handler.create_file`` the object-name step
|
||||
(get_complete_file_url -> get_object_name) runs before
|
||||
transform_create_file_request, and each independently calls
|
||||
``extract_file_data`` on the same create_file_data. When the file is a
|
||||
tuple-wrapped open handle, the streaming reader must still emit every row
|
||||
including entry 0: ``extract_file_data`` materializes the handle to bytes and
|
||||
rewinds it (seek(0)), so neither step consumes the other's cursor. A partial
|
||||
upload missing the first request would be silent and hard to catch, so this
|
||||
locks the full-payload invariant in.
|
||||
"""
|
||||
|
||||
def test_filehandle_create_file_keeps_first_entry(self):
|
||||
cfg = VertexAIFilesConfig()
|
||||
n_rows = 25
|
||||
raw = _make_openai_jsonl_bytes(n_rows)
|
||||
create_file_data: CreateFileRequest = {
|
||||
"file": ("batch.jsonl", io.BytesIO(raw), "application/jsonl"),
|
||||
"purpose": "batch",
|
||||
}
|
||||
|
||||
# Object-name step first (as the handler does), then the transform, both
|
||||
# reading the same live BytesIO handle.
|
||||
cfg.get_complete_file_url(
|
||||
api_base=None,
|
||||
api_key=None,
|
||||
model="",
|
||||
optional_params={},
|
||||
litellm_params={"bucket_name": "test-bucket"},
|
||||
data=create_file_data,
|
||||
)
|
||||
out = cfg.transform_create_file_request(
|
||||
model="",
|
||||
create_file_data=create_file_data,
|
||||
optional_params={},
|
||||
litellm_params={},
|
||||
)
|
||||
|
||||
assert isinstance(out, bytes)
|
||||
lines = out.decode("utf-8").splitlines()
|
||||
assert len(lines) == n_rows, "no batch row may be dropped from the upload"
|
||||
first_labels = json.loads(lines[0])["request"]["labels"]
|
||||
assert _get_litellm_batch_custom_id_from_labels(first_labels) == "request-0"
|
||||
|
||||
|
||||
class TestStreamingLineIterator:
|
||||
def test_skips_blank_and_whitespace_lines(self):
|
||||
content = b'{"a": 1}\n\n \n{"b": 2}\n'
|
||||
assert list(_iter_openai_jsonl_lines(content)) == ['{"a": 1}', '{"b": 2}']
|
||||
|
||||
def test_handles_crlf_and_missing_trailing_newline(self):
|
||||
content = b'{"a": 1}\r\n{"b": 2}'
|
||||
assert [json.loads(line) for line in _iter_openai_jsonl_lines(content)] == [
|
||||
{"a": 1},
|
||||
{"b": 2},
|
||||
]
|
||||
|
||||
def test_accepts_str_bytes_tuple_and_filelike(self):
|
||||
expected = [{"a": 1}, {"b": 2}]
|
||||
text = '{"a": 1}\n{"b": 2}\n'
|
||||
for source in (
|
||||
text,
|
||||
text.encode("utf-8"),
|
||||
("name.jsonl", text.encode("utf-8"), "application/jsonl"),
|
||||
io.BytesIO(text.encode("utf-8")),
|
||||
):
|
||||
assert list(_iter_openai_jsonl_entries(source)) == expected
|
||||
|
||||
def test_str_input_without_trailing_newline(self):
|
||||
assert list(_iter_openai_jsonl_lines('{"a": 1}\n{"b": 2}')) == [
|
||||
'{"a": 1}',
|
||||
'{"b": 2}',
|
||||
]
|
||||
|
||||
def test_pathlike_input_is_read_line_by_line(self, tmp_path):
|
||||
path = tmp_path / "batch.jsonl"
|
||||
path.write_bytes(b'{"a": 1}\n{"b": 2}\n')
|
||||
assert list(_iter_openai_jsonl_entries(path)) == [{"a": 1}, {"b": 2}]
|
||||
|
||||
def test_unsupported_content_type_raises(self):
|
||||
with pytest.raises(ValueError, match="Unsupported file content type"):
|
||||
list(_iter_openai_jsonl_lines(12345)) # type: ignore[arg-type]
|
||||
|
||||
def test_is_lazy_does_not_parse_past_first_entry(self):
|
||||
# Second row is invalid JSON; pulling only the first entry must not raise.
|
||||
content = b'{"custom_id": "first"}\nnot-json-at-all\n'
|
||||
gen = _iter_openai_jsonl_entries(content)
|
||||
assert next(gen)["custom_id"] == "first"
|
||||
with pytest.raises(json.JSONDecodeError):
|
||||
next(gen)
|
||||
|
||||
|
||||
class TestLegacyHandlerPathStreaming:
|
||||
def test_returns_str_with_object_name_from_first_row(self):
|
||||
transformer = VertexAIJsonlFilesTransformation()
|
||||
raw = _make_openai_jsonl_bytes(50)
|
||||
|
||||
vertex_str, object_name = (
|
||||
transformer.transform_openai_file_content_to_vertex_ai_file_content(raw)
|
||||
)
|
||||
|
||||
assert isinstance(vertex_str, str)
|
||||
assert "gemini-2.5-flash" in object_name
|
||||
# First line is a valid Vertex-wrapped request.
|
||||
assert "request" in json.loads(vertex_str.splitlines()[0])
|
||||
|
||||
def test_empty_payload_raises(self):
|
||||
transformer = VertexAIJsonlFilesTransformation()
|
||||
with pytest.raises(ValueError, match="empty"):
|
||||
transformer.transform_openai_file_content_to_vertex_ai_file_content(b"\n\n")
|
||||
|
||||
|
||||
class TestGetObjectNameLazyParse:
|
||||
def test_only_parses_first_row_for_model(self):
|
||||
cfg = VertexAIFilesConfig()
|
||||
# Tail rows are deliberately not valid JSON. Parsing the whole payload
|
||||
# would raise here; a first-row-only parse must not.
|
||||
raw = (
|
||||
b'{"custom_id": "r-0", "body": {"model": "gemini-2.5-flash"}}\n'
|
||||
b"garbage line that is not json\n"
|
||||
)
|
||||
from litellm.litellm_core_utils.prompt_templates.common_utils import (
|
||||
extract_file_data,
|
||||
)
|
||||
|
||||
extracted = extract_file_data(("batch.jsonl", raw, "application/jsonl"))
|
||||
object_name = cfg.get_object_name(extracted, purpose="batch")
|
||||
assert "gemini-2.5-flash" in object_name
|
||||
|
||||
|
||||
class TestStreamingPeakMemory:
|
||||
"""
|
||||
Differential guard: the streaming transform must stay well under the peak
|
||||
that a list pipeline incurs on the same input. If the hot path builds full
|
||||
intermediate lists, the streaming assertion fails.
|
||||
|
||||
The assertion that matters is the *relative* one: ``streaming_peak`` must be
|
||||
a clear fraction of ``list_peak`` on the identical input. Absolute
|
||||
``tracemalloc`` ratios drift with GC timing and the live set carried in from
|
||||
earlier tests, so they make poor CI gates; the relative comparison cancels
|
||||
that shared noise and is exactly what regresses (toward 1.0) when the hot
|
||||
path builds full intermediate lists. ``gc.collect()`` before each
|
||||
measurement removes any garbage the previous run left behind.
|
||||
"""
|
||||
|
||||
def _measure(self, fn):
|
||||
gc.collect()
|
||||
tracemalloc.start()
|
||||
try:
|
||||
fn()
|
||||
_, peak = tracemalloc.get_traced_memory()
|
||||
finally:
|
||||
tracemalloc.stop()
|
||||
return peak
|
||||
|
||||
def test_streaming_peak_well_below_list_pipeline(self):
|
||||
cfg = VertexAIFilesConfig()
|
||||
raw = _make_openai_jsonl_bytes(8000)
|
||||
content_str = raw.decode("utf-8")
|
||||
|
||||
streaming_peak = self._measure(
|
||||
lambda: _stream_openai_jsonl_to_vertex(
|
||||
raw, cfg._map_openai_to_vertex_params, as_bytes=True
|
||||
)
|
||||
)
|
||||
list_peak = self._measure(lambda: _legacy_vertex_jsonl_string(cfg, content_str))
|
||||
|
||||
# Core guard: streaming peaks at well under two-thirds of the list
|
||||
# pipeline. Building full intermediate lists in the hot path pushes this
|
||||
# ratio back toward 1.0 and fails the test.
|
||||
assert streaming_peak < list_peak * 0.6, (
|
||||
f"streaming peak {streaming_peak} not a clear win over list pipeline "
|
||||
f"{list_peak} (ratio {streaming_peak / list_peak:.2f})"
|
||||
)
|
||||
|
||||
def test_get_object_name_does_not_scale_with_payload(self):
|
||||
cfg = VertexAIFilesConfig()
|
||||
from litellm.litellm_core_utils.prompt_templates.common_utils import (
|
||||
extract_file_data,
|
||||
)
|
||||
|
||||
raw = _make_openai_jsonl_bytes(8000)
|
||||
extracted = extract_file_data(("batch.jsonl", raw, "application/jsonl"))
|
||||
|
||||
# The payload bytes already exist before measurement starts, so a lazy
|
||||
# first-row parse should allocate only a small fraction of the payload;
|
||||
# parsing every row would blow past this bound.
|
||||
peak = self._measure(lambda: cfg.get_object_name(extracted, purpose="batch"))
|
||||
assert (
|
||||
peak / len(raw) < 2.0
|
||||
), "get_object_name should not copy the whole payload"
|
||||
Loading…
Add table
Reference in a new issue