fix(vertex_ai): raise error on mid-stream 429/error chunks instead of silently swallowing

- Detect inline error chunks (e.g. 429 RESOURCE_EXHAUSTED) in Vertex AI
  streaming responses and raise VertexAIError so the Router can retry/fallback
- Add isinstance guard for non-dict error payloads to avoid AttributeError
- Cast error code to int for type safety against malformed responses
This commit is contained in:
Kris Xia 2026-03-17 12:27:36 +08:00 • committed by S0ngRu1
parent a9c9bee6c9
commit 25bc066dba
2 changed files with 314 additions and 0 deletions

View file

@ -3139,6 +3139,25 @@ class ModelResponseIterator:
self.cumulative_tool_call_index: int = 0
self.has_seen_tool_calls: bool = False
@staticmethod
def _check_streaming_error(chunk: dict) -> None:
"""检测流式响应中内嵌的错误(如 429 RESOURCE_EXHAUSTED),有错误时抛出 VertexAIError。"""
if "error" not in chunk:
return
error_data = chunk["error"]
if not isinstance(error_data, dict):
raise VertexAIError(
status_code=500,
message=f"VertexAIError: unexpected error format: {error_data}",
)
error_code = int(error_data.get("code", 500))
error_message = error_data.get("message", "Unknown error")
error_status = error_data.get("status", "UNKNOWN")
raise VertexAIError(
status_code=error_code,
message=f"VertexAIError: {error_status} - {error_message}",
)
def _apply_stream_candidates(
self,
_candidates: List[Candidates],
@ -3256,6 +3275,11 @@ class ModelResponseIterator:
def chunk_parser(self, chunk: dict) -> Optional["ModelResponseStream"]:
try:
verbose_logger.debug(f"RAW GEMINI CHUNK: {chunk}")
# 检测流式响应中内嵌的错误(如 429 RESOURCE_EXHAUSTED)。
# 这类错误以 HTTP 200 返回,但 SSE body 中包含 error JSON。
self._check_streaming_error(chunk)
from litellm.types.utils import ModelResponseStream
processed_chunk = GenerateContentResponseBody(**chunk) # type: ignore

View file

@ -4290,3 +4290,293 @@ def test_transform_response_does_not_leak_body_on_parse_failure():
msg = str(exc_info.value)
assert "secret content" not in msg
assert "Error converting to valid response block" in msg
def test_chunk_parser_raises_on_429_error_chunk():
"""Test chunk_parser raises VertexAIError on 429 RESOURCE_EXHAUSTED error chunk"""
from unittest.mock import Mock
from litellm.llms.vertex_ai.common_utils import VertexAIError
from litellm.llms.vertex_ai.gemini.vertex_and_google_ai_studio_gemini import (
ModelResponseIterator,
)
error_chunk = {
"error": {
"code": 429,
"message": "Resource exhausted. Please try again later. Please refer to https://cloud.google.com/vertex-ai/generative-ai/docs/error-code-429 for more details.",
"status": "RESOURCE_EXHAUSTED",
}
}
logging_obj = Mock()
logging_obj.optional_params = {}
streaming_obj = ModelResponseIterator(
streaming_response=iter([]),
sync_stream=True,
logging_obj=logging_obj,
)
with pytest.raises(VertexAIError) as exc_info:
streaming_obj.chunk_parser(error_chunk)
assert exc_info.value.status_code == 429
assert "RESOURCE_EXHAUSTED" in str(exc_info.value.message)
assert "Resource exhausted" in str(exc_info.value.message)
def test_chunk_parser_raises_on_500_error_chunk():
"""Test chunk_parser raises VertexAIError on 500 INTERNAL error chunk"""
from unittest.mock import Mock
from litellm.llms.vertex_ai.common_utils import VertexAIError
from litellm.llms.vertex_ai.gemini.vertex_and_google_ai_studio_gemini import (
ModelResponseIterator,
)
error_chunk = {
"error": {
"code": 500,
"message": "Internal error encountered.",
"status": "INTERNAL",
}
}
logging_obj = Mock()
logging_obj.optional_params = {}
streaming_obj = ModelResponseIterator(
streaming_response=iter([]),
sync_stream=True,
logging_obj=logging_obj,
)
with pytest.raises(VertexAIError) as exc_info:
streaming_obj.chunk_parser(error_chunk)
assert exc_info.value.status_code == 500
assert "INTERNAL" in str(exc_info.value.message)
def test_chunk_parser_raises_on_error_chunk_with_minimal_fields():
"""Test chunk_parser handles error chunks with missing optional fields"""
from unittest.mock import Mock
from litellm.llms.vertex_ai.common_utils import VertexAIError
from litellm.llms.vertex_ai.gemini.vertex_and_google_ai_studio_gemini import (
ModelResponseIterator,
)
error_chunk = {
"error": {
"code": 429,
"message": "Resource exhausted.",
}
}
logging_obj = Mock()
logging_obj.optional_params = {}
streaming_obj = ModelResponseIterator(
streaming_response=iter([]),
sync_stream=True,
logging_obj=logging_obj,
)
with pytest.raises(VertexAIError) as exc_info:
streaming_obj.chunk_parser(error_chunk)
assert exc_info.value.status_code == 429
def test_chunk_parser_normal_chunk_unaffected_by_error_check():
"""Test that normal streaming chunks still work correctly after error check addition"""
from unittest.mock import Mock
from litellm.llms.vertex_ai.gemini.vertex_and_google_ai_studio_gemini import (
ModelResponseIterator,
)
normal_chunk = {
"candidates": [
{
"content": {
"role": "model",
"parts": [{"text": "Hello"}],
},
"index": 0,
}
],
"usageMetadata": {
"promptTokenCount": 5,
"candidatesTokenCount": 1,
"totalTokenCount": 6,
},
}
logging_obj = Mock()
logging_obj.optional_params = {}
streaming_obj = ModelResponseIterator(
streaming_response=iter([]),
sync_stream=True,
logging_obj=logging_obj,
)
result = streaming_obj.chunk_parser(normal_chunk)
assert result is not None
assert len(result.choices) > 0
assert result.choices[0].delta.content == "Hello"
def test_chunk_parser_raises_on_non_dict_error():
"""Test chunk_parser raises VertexAIError when chunk['error'] is not a dict"""
from unittest.mock import Mock
from litellm.llms.vertex_ai.common_utils import VertexAIError
from litellm.llms.vertex_ai.gemini.vertex_and_google_ai_studio_gemini import (
ModelResponseIterator,
)
error_chunk = {"error": "something went wrong"}
logging_obj = Mock()
logging_obj.optional_params = {}
streaming_obj = ModelResponseIterator(
streaming_response=iter([]),
sync_stream=True,
logging_obj=logging_obj,
)
with pytest.raises(VertexAIError) as exc_info:
streaming_obj.chunk_parser(error_chunk)
assert exc_info.value.status_code == 500
assert "unexpected error format" in str(exc_info.value.message)
def test_chunk_parser_raises_on_string_error_code():
"""Test chunk_parser correctly converts string error code to int"""
from unittest.mock import Mock
from litellm.llms.vertex_ai.common_utils import VertexAIError
from litellm.llms.vertex_ai.gemini.vertex_and_google_ai_studio_gemini import (
ModelResponseIterator,
)
# code 字段为字符串类型 "429" 而非 int
error_chunk = {
"error": {
"code": "429",
"message": "Resource exhausted.",
"status": "RESOURCE_EXHAUSTED",
}
}
logging_obj = Mock()
logging_obj.optional_params = {}
streaming_obj = ModelResponseIterator(
streaming_response=iter([]),
sync_stream=True,
logging_obj=logging_obj,
)
with pytest.raises(VertexAIError) as exc_info:
streaming_obj.chunk_parser(error_chunk)
assert exc_info.value.status_code == 429
assert isinstance(exc_info.value.status_code, int)
def test_mid_stream_429_error_raises_during_iteration():
"""
Simulate a full streaming scenario: normal thinking chunks arrive first,
then a 429 RESOURCE_EXHAUSTED error chunk arrives mid-stream.
Verify that ModelResponseIterator raises VertexAIError during iteration.
"""
import json
from unittest.mock import Mock
from litellm.llms.vertex_ai.common_utils import VertexAIError
from litellm.llms.vertex_ai.gemini.vertex_and_google_ai_studio_gemini import (
ModelResponseIterator,
)
# Simulate Vertex AI SSE stream: normal chunks followed by a 429 error chunk
normal_chunk_1 = json.dumps(
{
"candidates": [
{
"content": {
"role": "model",
"parts": [{"text": "Let me think about this...", "thought": True}],
},
"index": 0,
}
],
"usageMetadata": {
"promptTokenCount": 10,
"candidatesTokenCount": 5,
"totalTokenCount": 15,
},
"modelVersion": "gemini-3.1-flash-image-preview",
}
)
normal_chunk_2 = json.dumps(
{
"candidates": [
{
"content": {
"role": "model",
"parts": [{"text": "I'll generate the image now.", "thought": True}],
},
"index": 0,
}
],
"usageMetadata": {
"promptTokenCount": 10,
"candidatesTokenCount": 12,
"totalTokenCount": 22,
},
}
)
error_chunk = json.dumps(
{
"error": {
"code": 429,
"message": "Resource exhausted. Please try again later. Please refer to https://cloud.google.com/vertex-ai/generative-ai/docs/error-code-429 for more details.",
"status": "RESOURCE_EXHAUSTED",
}
}
)
# Build a mock SSE stream (lines returned by iter_lines)
sse_lines = iter([normal_chunk_1, normal_chunk_2, error_chunk])
logging_obj = Mock()
logging_obj.optional_params = {}
streaming_obj = ModelResponseIterator(
streaming_response=sse_lines,
sync_stream=True,
logging_obj=logging_obj,
)
# Iterate the stream: first chunks should succeed, then 429 error should be raised
results = []
with pytest.raises(VertexAIError) as exc_info:
for chunk in streaming_obj:
if chunk is not None:
results.append(chunk)
# Verify: received normal chunks before the error
assert len(results) >= 1, "Should have received at least 1 normal chunk before the error"
# Verify: 429 error is properly raised
assert exc_info.value.status_code == 429
assert "RESOURCE_EXHAUSTED" in str(exc_info.value.message)