reiterability + test

This commit is contained in:
mubashir1osmani 2026-06-06 14:01:52 -07:00
parent 72347e7c90
commit 3fbae90839
No known key found for this signature in database
GPG key ID: AB055FF67D0B4D9A
3 changed files with 25 additions and 0 deletions

View file

@ -3217,6 +3217,7 @@ class BaseLLMHTTPHandler:
timeout=timeout,
)
except Exception as e:
verbose_logger.exception(f"Error creating file: {e}")
raise self._handle_error(e=e, provider_config=provider_config)
elif isinstance(transformed_request, str) or isinstance(
transformed_request, bytes
@ -3562,6 +3563,8 @@ class BaseLLMHTTPHandler:
resp = httpx_client.send(req, follow_redirects=False)
resp.read()
if resp.status_code not in ((200, 201) if is_final else (308,)):
# 4xx/5xx raise here; the ValueError catches an unexpected success
# status (e.g. a 200 where the protocol expects a 308 between chunks).
resp.raise_for_status()
raise ValueError(f"resumable upload: unexpected status {resp.status_code}")
return resp
@ -3642,6 +3645,8 @@ class BaseLLMHTTPHandler:
resp = await httpx_client.send(req, follow_redirects=False)
await resp.aread()
if resp.status_code not in ((200, 201) if is_final else (308,)):
# 4xx/5xx raise here; the ValueError catches an unexpected success
# status (e.g. a 200 where the protocol expects a 308 between chunks).
resp.raise_for_status()
raise ValueError(f"resumable upload: unexpected status {resp.status_code}")
return resp

View file

@ -228,6 +228,13 @@ def _iter_openai_jsonl_lines(openai_file_content: FileTypes) -> Iterator[str]:
return
if hasattr(content, "read"):
# Rewind seekable handles (BytesIO, temp files) so the body can be
# replayed on a retry; a non-seekable stream cannot be re-read.
if hasattr(content, "seek"):
try:
content.seek(0)
except (OSError, ValueError):
pass
for raw in content:
line = raw.decode("utf-8") if isinstance(raw, bytes) else raw
line = line.strip()

View file

@ -433,6 +433,19 @@ class TestResumableStreamBody:
second = b"".join(stream.iter_bytes())
assert first == second and len(first) > 0
def test_stream_is_reiterable_for_seekable_file_like_input(self):
# A seekable handle (BytesIO, temp file) must be rewound between calls;
# otherwise the first iter_bytes() exhausts it and a retry would upload
# an empty body silently.
cfg = VertexAIFilesConfig()
raw = _make_openai_jsonl_bytes(40)
stream = _OpenAIToVertexBatchUploadStream(
io.BytesIO(raw), cfg._map_openai_to_vertex_params
)
first = b"".join(stream.iter_bytes())
second = b"".join(stream.iter_bytes())
assert first == second and len(first) > 0
class TestResumableChunking:
def test_intermediate_chunks_are_exactly_chunk_size(self):