diff --git a/agent/api_error_summary.py b/agent/api_error_summary.py index 3e907dbaf3..5a1e884c27 100644 --- a/agent/api_error_summary.py +++ b/agent/api_error_summary.py @@ -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) diff --git a/agent/chat_completion_helpers.py b/agent/chat_completion_helpers.py index e44353d0e0..e573e7068b 100644 --- a/agent/chat_completion_helpers.py +++ b/agent/chat_completion_helpers.py @@ -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: diff --git a/run_agent.py b/run_agent.py index e6601e440c..3174c413ed 100644 --- a/run_agent.py +++ b/run_agent.py @@ -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") diff --git a/tests/agent/test_streaming.py b/tests/agent/test_streaming.py index 52426a6a2b..c5e15fd548 100644 --- a/tests/agent/test_streaming.py +++ b/tests/agent/test_streaming.py @@ -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,