diff --git a/litellm/llms/vertex_ai/files/transformation.py b/litellm/llms/vertex_ai/files/transformation.py index f30518bc7ca..1db77c39f44 100644 --- a/litellm/llms/vertex_ai/files/transformation.py +++ b/litellm/llms/vertex_ai/files/transformation.py @@ -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": } 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: diff --git a/tests/test_litellm/llms/vertex_ai/files/test_vertex_ai_binary_file_upload.py b/tests/test_litellm/llms/vertex_ai/files/test_vertex_ai_binary_file_upload.py index d4586134b13..071f89ac414 100644 --- a/tests/test_litellm/llms/vertex_ai/files/test_vertex_ai_binary_file_upload.py +++ b/tests/test_litellm/llms/vertex_ai/files/test_vertex_ai_binary_file_upload.py @@ -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", diff --git a/tests/test_litellm/llms/vertex_ai/files/test_vertex_ai_files_streaming.py b/tests/test_litellm/llms/vertex_ai/files/test_vertex_ai_files_streaming.py new file mode 100644 index 00000000000..6eac006a00a --- /dev/null +++ b/tests/test_litellm/llms/vertex_ai/files/test_vertex_ai_files_streaming.py @@ -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"