fix(anthropic): recover malformed streamed tool JSON
This commit is contained in:
@@ -142,7 +142,9 @@ class ApiErrorSummaryMixin:
|
||||
)
|
||||
current = current.__cause__ or current.__context__
|
||||
|
||||
if isinstance(error, ValueError) and "expected ident at line" in raw.lower():
|
||||
if isinstance(error, ValueError) and any(
|
||||
marker in raw.lower() for marker in ("expected ident at line", "expected value at line")
|
||||
):
|
||||
return f"Malformed provider streaming response: {raw[:300]}"
|
||||
|
||||
prefix = _http_prefix(error)
|
||||
|
||||
@@ -3019,6 +3019,7 @@ class _StreamingCall(StreamingWaitMonitor):
|
||||
self._writer_token = None
|
||||
_stream_context = {"manager": None, "stream": None}
|
||||
base_final_message = None
|
||||
stream_parse_error = None
|
||||
|
||||
from agent import relay_llm
|
||||
from agent.anthropic_adapter import sanitize_anthropic_kwargs
|
||||
@@ -3074,6 +3075,17 @@ class _StreamingCall(StreamingWaitMonitor):
|
||||
raise EmptyStreamError(
|
||||
"Provider returned an empty stream with no events (possible upstream error or malformed event stream).") from None
|
||||
raise
|
||||
except ValueError as exc:
|
||||
# Fine-grained tool streaming exposes raw partial JSON. The Anthropic SDK
|
||||
# may fail while incrementally decoding it before get_final_message() can
|
||||
# return its buffered/repaired representation. Only recover after a tool
|
||||
# block started; unrelated ValueErrors retain their normal handling.
|
||||
_error_text = str(exc).strip().lower()
|
||||
if not has_tool_use or not any(
|
||||
marker in _error_text for marker in ("expected value at line", "expected ident at line")
|
||||
):
|
||||
raise
|
||||
stream_parse_error = exc
|
||||
finally:
|
||||
try:
|
||||
self._close_managed_stream()
|
||||
@@ -3084,6 +3096,19 @@ class _StreamingCall(StreamingWaitMonitor):
|
||||
|
||||
if self.agent._interrupt_requested:
|
||||
return None
|
||||
if stream_parse_error is not None:
|
||||
fallback_kwargs = dict(self.api_kwargs)
|
||||
fallback_kwargs.pop("stream", None)
|
||||
sanitize_anthropic_kwargs(
|
||||
fallback_kwargs, log_prefix=getattr(self.agent, "log_prefix", "")
|
||||
)
|
||||
logger.warning(
|
||||
"%sAnthropic tool stream JSON parsing failed (%s); retrying once with buffered messages.create()",
|
||||
getattr(self.agent, "log_prefix", ""),
|
||||
stream_parse_error,
|
||||
)
|
||||
base_final_message = request_client.messages.create(**fallback_kwargs)
|
||||
return self._check_anthropic_message(base_final_message)
|
||||
if base_final_message is not None:
|
||||
self._check_anthropic_message(base_final_message, tool_drop=False)
|
||||
if not stream.output_modified:
|
||||
|
||||
@@ -485,7 +485,8 @@ class AIAgent(
|
||||
that is wire trouble, not local validation, so it follows the truncated-JSON retry path."""
|
||||
return (getattr(self, "api_mode", None) == "anthropic_messages" and isinstance(error, ValueError)
|
||||
and not isinstance(error, (UnicodeEncodeError, json.JSONDecodeError))
|
||||
and "expected ident at line" in str(error).strip().lower())
|
||||
and any(marker in str(error).strip().lower() for marker in (
|
||||
"expected ident at line", "expected value at line")))
|
||||
|
||||
_log_stream_retry = _forward("agent.stream_diag", "log_stream_retry")
|
||||
_emit_stream_drop = _forward("agent.stream_diag", "emit_stream_drop")
|
||||
|
||||
@@ -1136,6 +1136,52 @@ class TestAnthropicStreamCallbacks:
|
||||
assert mock_rebuild.call_count == 0
|
||||
assert agent._anthropic_client.close.call_count >= 1
|
||||
|
||||
def test_anthropic_malformed_tool_json_falls_back_to_buffered_message(self):
|
||||
"""Malformed fine-grained tool JSON uses Anthropic's buffered response."""
|
||||
from run_agent import AIAgent
|
||||
|
||||
agent = AIAgent(
|
||||
api_key="test-key",
|
||||
base_url="https://api.anthropic.com",
|
||||
model="claude-sonnet-4-5",
|
||||
quiet_mode=True,
|
||||
skip_context_files=True,
|
||||
skip_memory=True,
|
||||
)
|
||||
agent.api_mode = "anthropic_messages"
|
||||
agent._interrupt_requested = False
|
||||
|
||||
class _MalformedToolStream:
|
||||
response = None
|
||||
|
||||
def __enter__(self):
|
||||
return self
|
||||
|
||||
def __exit__(self, *_args):
|
||||
return False
|
||||
|
||||
def __iter__(self):
|
||||
yield SimpleNamespace(
|
||||
type="content_block_start",
|
||||
content_block=SimpleNamespace(type="tool_use", name="cronjob_manage"),
|
||||
)
|
||||
raise ValueError("expected value at line 1 column 11")
|
||||
|
||||
repaired_message = SimpleNamespace(
|
||||
content=[SimpleNamespace(type="tool_use", name="cronjob_manage", input={"names": "cronjob_manage"})],
|
||||
stop_reason="tool_use",
|
||||
)
|
||||
agent._anthropic_client = MagicMock()
|
||||
agent._anthropic_client.messages.stream.return_value = _MalformedToolStream()
|
||||
agent._anthropic_client.messages.create.return_value = repaired_message
|
||||
agent._create_request_anthropic_client = lambda *a, **k: agent._anthropic_client
|
||||
|
||||
response = agent._interruptible_streaming_api_call({"model": agent.model})
|
||||
|
||||
assert response is repaired_message
|
||||
assert agent._anthropic_client.messages.stream.call_count == 1
|
||||
assert agent._anthropic_client.messages.create.call_count == 1
|
||||
|
||||
@patch("run_agent.AIAgent._replace_primary_openai_client")
|
||||
def test_generic_anthropic_valueerror_still_propagates_without_stream_retry(
|
||||
self, mock_replace, monkeypatch,
|
||||
|
||||
Reference in New Issue
Block a user