fix(vertex_ai/gemini): asyncify transform only when extensionless gs:// may fetch GCS metadata

Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
S0ngRu1 2026-05-13 17:08:36 +08:00
parent f3e0fd11d6
commit ca448058ef
2 changed files with 290 additions and 6 deletions

View file

@ -241,6 +241,104 @@ def _is_valid_gcs_bucket_name(bucket: str) -> bool:
return True
def _gs_uri_requires_content_type_metadata(url: str) -> bool:
"""
True when _process_gemini_media would call _get_gcs_object_content_type
(extension-less gs:// and no explicit format passed into that helper).
"""
if "gs://" not in url:
return False
extension_with_dot = os.path.splitext(url)[-1]
extension = extension_with_dot[1:] if extension_with_dot else ""
return len(extension) == 0
def _image_url_payload_may_need_sync_gcs_metadata_fetch(
raw_image_url: Any,
) -> bool:
"""
True when this image_url value (content-part image_url or assistant ``images[]``
entry) can trigger a blocking GCS metadata read for MIME resolution.
"""
fmt: Optional[str] = None
url: Optional[str] = None
if isinstance(raw_image_url, dict):
url = raw_image_url.get("url") # type: ignore[assignment]
if not isinstance(url, str):
return False
fmt = (
raw_image_url.get("format")
or raw_image_url.get("mime_type")
or raw_image_url.get("content_type")
)
elif isinstance(raw_image_url, str):
url = raw_image_url
else:
return False
if "gs://" not in url or fmt:
return False
return _gs_uri_requires_content_type_metadata(url)
def _openai_messages_may_need_sync_gcs_metadata_fetch(
messages: List[AllMessageValues],
) -> bool:
"""
Heuristic: True if any message part can trigger a blocking GCS JSON
metadata read inside _transform_request_body (extension-less gs:// without
explicit MIME hints). Covers user/system ``content`` parts and assistant
``images`` (same paths as ``_gemini_convert_messages_with_history``). Used
to decide whether ``async_transform_request_body`` should offload the sync
transform via ``asyncify``.
"""
for raw in messages:
msg: Any = raw
if not isinstance(msg, dict) and hasattr(msg, "model_dump"):
msg = msg.model_dump(exclude_none=False)
if not isinstance(msg, dict):
continue
images_field = msg.get("images")
if isinstance(images_field, list):
for image_item in images_field:
if not isinstance(image_item, dict):
continue
if _image_url_payload_may_need_sync_gcs_metadata_fetch(
image_item.get("image_url")
):
return True
content = msg.get("content")
if not isinstance(content, list):
continue
for item in content:
if not isinstance(item, dict):
continue
itype = item.get("type")
if itype == "image_url":
if _image_url_payload_may_need_sync_gcs_metadata_fetch(
item.get("image_url")
):
return True
elif itype == "file":
file_obj = item.get("file")
if not isinstance(file_obj, dict):
continue
fmt = (
file_obj.get("format")
or file_obj.get("mime_type")
or file_obj.get("content_type")
)
passed = file_obj.get("file_id") or file_obj.get("file_data")
if (
isinstance(passed, str)
and "gs://" in passed
and not fmt
and _gs_uri_requires_content_type_metadata(passed)
):
return True
return False
def _get_gcs_object_content_type(
image_url: str,
vertex_project: Optional[str] = None,
@ -856,7 +954,11 @@ def _gemini_convert_messages_with_history( # noqa: PLR0915
image_url_obj = image_item.get("image_url")
if isinstance(image_url_obj, dict):
assistant_image_url = image_url_obj.get("url")
format = image_url_obj.get("format")
format = (
image_url_obj.get("format")
or image_url_obj.get("mime_type")
or image_url_obj.get("content_type")
)
detail = image_url_obj.get("detail")
media_resolution_enum = (
_convert_detail_to_media_resolution_enum(detail)
@ -1223,11 +1325,21 @@ async def async_transform_request_body(
vertex_auth_header=vertex_auth_header,
)
# _transform_request_body may issue a sync httpx.get (up to 5s timeout)
# via _get_gcs_object_content_type to fetch GCS object metadata. Run the
# whole sync transformation on a worker thread so it does not block the
# async event loop.
return await asyncify(_transform_request_body)(
if _openai_messages_may_need_sync_gcs_metadata_fetch(messages):
# _transform_request_body may issue a sync httpx.get (up to 5s timeout)
# via _get_gcs_object_content_type to fetch GCS object metadata. Run the
# whole sync transformation on a worker thread so it does not block the
# async event loop.
return await asyncify(_transform_request_body)(
messages=messages,
model=model,
custom_llm_provider=custom_llm_provider,
litellm_params=litellm_params,
cached_content=cached_content,
optional_params=optional_params,
)
return _transform_request_body(
messages=messages,
model=model,
custom_llm_provider=custom_llm_provider,

View file

@ -1734,6 +1734,178 @@ def test_async_transform_request_body_does_not_block_event_loop():
)
def test_openai_messages_may_need_sync_gcs_metadata_fetch_false_for_plain_text():
from litellm.llms.vertex_ai.gemini.transformation import (
_openai_messages_may_need_sync_gcs_metadata_fetch,
)
assert (
_openai_messages_may_need_sync_gcs_metadata_fetch(
[{"role": "user", "content": "hello"}]
)
is False
)
def test_openai_messages_may_need_sync_gcs_metadata_fetch_true_for_extensionless_gs_image_url():
from litellm.llms.vertex_ai.gemini.transformation import (
_openai_messages_may_need_sync_gcs_metadata_fetch,
)
messages = [
{
"role": "user",
"content": [
{
"type": "image_url",
"image_url": {"url": "gs://bucket/image-without-extension"},
}
],
}
]
assert _openai_messages_may_need_sync_gcs_metadata_fetch(messages) is True
def test_openai_messages_may_need_sync_gcs_metadata_fetch_false_when_gs_has_file_extension():
from litellm.llms.vertex_ai.gemini.transformation import (
_openai_messages_may_need_sync_gcs_metadata_fetch,
)
messages = [
{
"role": "user",
"content": [
{
"type": "image_url",
"image_url": {"url": "gs://bucket/image.png"},
}
],
}
]
assert _openai_messages_may_need_sync_gcs_metadata_fetch(messages) is False
def test_openai_messages_may_need_sync_gcs_metadata_fetch_false_when_extensionless_gs_has_mime_hint():
from litellm.llms.vertex_ai.gemini.transformation import (
_openai_messages_may_need_sync_gcs_metadata_fetch,
)
messages = [
{
"role": "user",
"content": [
{
"type": "image_url",
"image_url": {
"url": "gs://bucket/image-without-extension",
"mime_type": "image/png",
},
}
],
}
]
assert _openai_messages_may_need_sync_gcs_metadata_fetch(messages) is False
def test_openai_messages_may_need_sync_gcs_metadata_fetch_true_for_assistant_images_extensionless_gs():
from litellm.llms.vertex_ai.gemini.transformation import (
_openai_messages_may_need_sync_gcs_metadata_fetch,
)
messages = [
{
"role": "assistant",
"content": [],
"images": [
{"image_url": {"url": "gs://bucket/gen-without-extension"}},
],
}
]
assert _openai_messages_may_need_sync_gcs_metadata_fetch(messages) is True
def test_openai_messages_may_need_sync_gcs_metadata_fetch_false_for_assistant_images_gs_with_extension():
from litellm.llms.vertex_ai.gemini.transformation import (
_openai_messages_may_need_sync_gcs_metadata_fetch,
)
messages = [
{
"role": "assistant",
"content": [],
"images": [{"image_url": {"url": "gs://bucket/gen.png"}}],
}
]
assert _openai_messages_may_need_sync_gcs_metadata_fetch(messages) is False
def test_openai_messages_may_need_sync_gcs_metadata_fetch_false_for_assistant_images_with_mime_hint():
from litellm.llms.vertex_ai.gemini.transformation import (
_openai_messages_may_need_sync_gcs_metadata_fetch,
)
messages = [
{
"role": "assistant",
"content": [],
"images": [
{
"image_url": {
"url": "gs://bucket/gen-no-ext",
"mime_type": "image/png",
},
}
],
}
]
assert _openai_messages_may_need_sync_gcs_metadata_fetch(messages) is False
def test_async_transform_request_body_skips_asyncify_for_text_only_requests():
"""Hot path: no extension-less gs:// metadata fetch → run sync transform in-loop."""
import asyncio
from unittest.mock import MagicMock, patch
from litellm.llms.vertex_ai.gemini import transformation as gemini_transformation
async def fake_check_and_create_cache(self, **kwargs):
return kwargs["messages"], kwargs["optional_params"], None
async def run():
with patch(
"litellm.llms.vertex_ai.gemini.transformation.asyncify",
side_effect=AssertionError(
"asyncify must not run when no extensionless GCS metadata fetch is needed"
),
):
return await gemini_transformation.async_transform_request_body(
gemini_api_key=None,
messages=[{"role": "user", "content": "hello"}],
api_base=None,
model="gemini-2.5-flash",
client=None,
timeout=None,
extra_headers=None,
optional_params={},
logging_obj=MagicMock(),
custom_llm_provider="vertex_ai",
litellm_params={},
vertex_project=None,
vertex_location=None,
vertex_auth_header=None,
)
with patch(
"litellm.llms.vertex_ai.context_caching.vertex_ai_context_caching."
"ContextCachingEndpoints.async_check_and_create_cache",
new=fake_check_and_create_cache,
):
body = asyncio.run(run())
assert body is not None
assert "contents" in body
def test_get_image_mime_type_from_url():
"""Test the _get_image_mime_type_from_url function for different image URLs"""
from litellm.llms.vertex_ai.gemini.transformation import (