diff --git a/agent/chat_completion_helpers.py b/agent/chat_completion_helpers.py index b830a4b99d..18f5a2c259 100644 --- a/agent/chat_completion_helpers.py +++ b/agent/chat_completion_helpers.py @@ -3897,6 +3897,34 @@ def interruptible_streaming_api_call(agent, api_kwargs: dict, *, on_first_delta= if agent._interrupt_requested: return None + + def _tool_use_dropped_mid_stream(message) -> bool: + """True when the stream died mid tool call (#80498 sibling). + + Mirror of the chat_completions zero-byte/truncated-args gate: a + legitimate completion always carries a ``stop_reason`` + (``tool_use``/``end_turn``/...), so a message that contains a + ``tool_use`` block but NO stop_reason means the SSE closed after + ``content_block_start`` and before ``message_delta`` — the + block's ``input`` is whatever partial state the SDK snapshot + accumulated (typically ``{}`` when no ``input_json_delta`` ever + arrived). Without this gate the empty-input call passed the + empty-stream guards (content is non-empty) and executed the tool + with no arguments and no retry. Raising EmptyStreamError blocks + that execution on every path; when no assistant text streamed + before the drop it additionally rides the bounded stream-retry + the eventless case uses (probe-verified recovery), while a + drop after streamed preamble text degrades to the + partial-stream-stub/continuation path instead — still never an + empty-args execution. + """ + if getattr(message, "stop_reason", None) is not None: + return False + for block in getattr(message, "content", None) or []: + if getattr(block, "type", None) == "tool_use": + return True + return False + if ( base_final_message is not None and not getattr(base_final_message, "content", None) @@ -3907,6 +3935,12 @@ def interruptible_streaming_api_call(agent, api_kwargs: dict, *, on_first_delta= "(possible upstream error or malformed event stream)." ) if base_final_message is not None and not stream.output_modified: + if _tool_use_dropped_mid_stream(base_final_message): + raise EmptyStreamError( + "Stream ended with no stop_reason while a tool_use " + "block was still incomplete; treating as a " + "mid-tool-call stream drop (#80498)." + ) return base_final_message final_message = accumulator.response(base_final_message) if ( @@ -3917,6 +3951,12 @@ def interruptible_streaming_api_call(agent, api_kwargs: dict, *, on_first_delta= "Provider returned an empty stream with no stop_reason " "(possible upstream error or malformed event stream)." ) + if _tool_use_dropped_mid_stream(final_message): + raise EmptyStreamError( + "Stream ended with no stop_reason while a tool_use " + "block was still incomplete; treating as a " + "mid-tool-call stream drop (#80498)." + ) return final_message def _call(): diff --git a/tests/run_agent/test_anthropic_mid_tool_call_drop.py b/tests/run_agent/test_anthropic_mid_tool_call_drop.py new file mode 100644 index 0000000000..f1bcc305c9 --- /dev/null +++ b/tests/run_agent/test_anthropic_mid_tool_call_drop.py @@ -0,0 +1,119 @@ +"""The Anthropic streaming path must not accept a tool_use block whose stream +died before its input arrived (#80498 sibling). + +The chat_completions accumulator flags zero-byte tool-call arguments on a +clean no-finish_reason stream end and routes them through the +partial-stream-stub retry path. The Anthropic path had the same gap in a +different shape: a clean SSE close after ``content_block_start(tool_use)`` +but before any ``input_json_delta``/``message_delta`` yields an SDK +final-message snapshot whose content is NON-empty (the tool_use block is +there, ``input={}``) and whose ``stop_reason`` is None — which sailed past +both empty-stream guards and executed the tool with empty input, no retry. + +The fix raises EmptyStreamError for a tool_use-bearing message with no +stop_reason, riding the same bounded stream-retry the eventless case uses. +""" +from types import SimpleNamespace +from unittest.mock import MagicMock + +import pytest + + +def _make_anthropic_agent(**kwargs): + from run_agent import AIAgent + + defaults = dict( + api_key="test-key", + base_url="https://example.com/v1", + model="claude-opus-4-7", + quiet_mode=True, + skip_context_files=True, + skip_memory=True, + ) + defaults.update(kwargs) + agent = AIAgent(**defaults) + agent.api_mode = "anthropic_messages" + agent._anthropic_client = MagicMock() + agent._anthropic_api_key = "test-anthropic-key" + agent._create_request_anthropic_client = lambda *a, **k: agent._anthropic_client + return agent + + +def _stream_cm(final_message, events=()): + cm = MagicMock() + stream = MagicMock() + stream.__iter__ = MagicMock(return_value=iter(list(events))) + stream.get_final_message = MagicMock(return_value=final_message) + cm.__enter__ = MagicMock(return_value=stream) + cm.__exit__ = MagicMock(return_value=False) + return cm + + +def _tool_use_block(name="write_file", input_obj=None): + return SimpleNamespace( + type="tool_use", + id="toolu_x", + name=name, + input=input_obj if input_obj is not None else {}, + ) + + +def _tool_use_start_event(name="write_file"): + return SimpleNamespace( + type="content_block_start", + content_block=SimpleNamespace(type="tool_use", name=name), + ) + + +class TestAnthropicMidToolCallStreamDrop: + def test_tool_use_without_stop_reason_raises_empty_stream(self): + """The #80498-sibling shape: tool_use block present, stop_reason None.""" + from agent.chat_completion_helpers import EmptyStreamError + + dropped = MagicMock() + dropped.content = [_tool_use_block()] + dropped.stop_reason = None + dropped.usage = SimpleNamespace(input_tokens=10, output_tokens=2) + + agent = _make_anthropic_agent() + agent._anthropic_client.messages.stream = MagicMock( + return_value=_stream_cm(dropped, events=[_tool_use_start_event()]) + ) + + with pytest.raises(EmptyStreamError, match="tool_use"): + agent._interruptible_streaming_api_call({"model": "claude-opus-4-7"}) + + def test_completed_tool_use_with_stop_reason_passes(self): + """A legitimate tool_use completion (stop_reason set) is untouched.""" + done = MagicMock() + done.content = [_tool_use_block(input_obj={"path": "a.txt"})] + done.stop_reason = "tool_use" + done.usage = SimpleNamespace(input_tokens=10, output_tokens=5) + + agent = _make_anthropic_agent() + agent._anthropic_client.messages.stream = MagicMock( + return_value=_stream_cm(done, events=[_tool_use_start_event()]) + ) + + response = agent._interruptible_streaming_api_call( + {"model": "claude-opus-4-7"} + ) + assert response is done + + def test_text_only_message_without_stop_reason_passes(self): + """No tool_use block -> the new gate stays out of the way (text-only + no-stop_reason handling keeps its pre-existing behavior).""" + text_only = MagicMock() + text_only.content = [SimpleNamespace(type="text", text="partial answer")] + text_only.stop_reason = None + text_only.usage = SimpleNamespace(input_tokens=10, output_tokens=5) + + agent = _make_anthropic_agent() + agent._anthropic_client.messages.stream = MagicMock( + return_value=_stream_cm(text_only) + ) + + response = agent._interruptible_streaming_api_call( + {"model": "claude-opus-4-7"} + ) + assert response is text_only