diff --git a/agent/chat_completion_helpers.py b/agent/chat_completion_helpers.py index 72e962c62b..48c7517827 100644 --- a/agent/chat_completion_helpers.py +++ b/agent/chat_completion_helpers.py @@ -46,7 +46,10 @@ from agent.message_sanitization import ( _sanitize_messages_surrogates, _sanitize_surrogates, _repair_tool_call_arguments, normalize_finish_reason as _normalize_finish_reason, sanitize_outbound_kwargs, strip_images_for_rejecting_model, ) -from agent.reasoning_summaries import append_streamed_reasoning_detail, separate_glued_reasoning_blocks +from agent.reasoning_summaries import ( + append_streamed_reasoning_detail, separate_glued_reasoning_blocks, + streamed_reasoning_detail_text, +) from agent.repetition_guard import is_repetition_dominated from agent.stream_single_writer import claim_stream_writer, stream_writer_is_current from tools.terminal_tool_lifecycle import is_persistent_env @@ -3230,7 +3233,6 @@ class _StreamingCall(StreamingWaitMonitor): reasoning_text = separate_glued_reasoning_blocks( reasoning_parts[-1] if reasoning_parts else "", reasoning_text) reasoning_parts.append(reasoning_text) - self._emit_reasoning(reasoning_text) # Structured reasoning_details deltas carry the provider's replay data; the # non-streaming path already keeps them, so dropping them here lost # reasoning continuity on nearly every turn. Pydantic parks unknown fields @@ -3238,8 +3240,16 @@ class _StreamingCall(StreamingWaitMonitor): rd_delta = getattr(delta, "reasoning_details", None) if rd_delta is None and isinstance(getattr(delta, "model_extra", None), dict): rd_delta = delta.model_extra.get("reasoning_details") + detail_text_parts = [] for rd in rd_delta if isinstance(rd_delta, (list, tuple)) else (): + detail_text_parts.append(streamed_reasoning_detail_text(rd)) append_streamed_reasoning_detail(reasoning_details, rd) + # Details may carry the full text while ordinary reasoning is only + # a sparse fragment or a mirror. Deliver one representation per + # chunk, without rewriting either persisted/replayed field. + display_reasoning = "".join(detail_text_parts) or reasoning_text + if display_reasoning: + self._emit_reasoning(display_reasoning) # Not routed to the live display: the transport promotes a sole-payload # refusal to content + ``content_filter`` and the loop surfaces it terminally. delta_refusal = getattr(delta, "refusal", None) diff --git a/agent/reasoning_summaries.py b/agent/reasoning_summaries.py index 72258ee48d..1ed51eb200 100644 --- a/agent/reasoning_summaries.py +++ b/agent/reasoning_summaries.py @@ -36,6 +36,16 @@ _MERGEABLE_DETAIL_TEXT_KEYS = {"reasoning.text": "text", "reasoning.summary": "s _BACKFILL_DETAIL_KEYS = ("signature", "id", "format", "index") +def streamed_reasoning_detail_text(detail: Any) -> str: + """Readable text from a detail delta; never expose opaque replay material.""" + dtype = detail.get("type") if isinstance(detail, dict) else getattr(detail, "type", None) + key = _MERGEABLE_DETAIL_TEXT_KEYS.get(dtype) if isinstance(dtype, str) else None + if key is None: + return "" + text = detail.get(key) if isinstance(detail, dict) else getattr(detail, key, None) + return text if isinstance(text, str) else "" + + def append_streamed_reasoning_detail(details_acc: list, detail: Any) -> None: """Accumulate one streamed ``reasoning_details`` delta entry into *details_acc*. diff --git a/tests/agent/test_streamed_reasoning_details.py b/tests/agent/test_streamed_reasoning_details.py index 4470630b87..b12fc8a47d 100644 --- a/tests/agent/test_streamed_reasoning_details.py +++ b/tests/agent/test_streamed_reasoning_details.py @@ -56,7 +56,20 @@ def test_streamed_details_land_on_final_message_and_persist(_mock_close, mock_cr mock_create.return_value = mock_client agent = _agent() + delivered = [] + agent.reasoning_callback = delivered.append + agent.stream_delta_callback = lambda text: None + + def streamed_chunks(): + yield chunks[0] + assert delivered == ["I should "] + yield chunks[1] + assert delivered == ["I should ", "answer."] + yield chunks[2] + + mock_client.chat.completions.create.return_value = streamed_chunks() response = agent._interruptible_streaming_api_call({}) + assert "".join(delivered) == "I should answer." msg = response.choices[0].message assert msg.content == "Hello!" assert msg.reasoning_details == [{"type": "reasoning.text", "text": "I should answer.", "signature": "sigZ"}] @@ -74,3 +87,46 @@ def test_no_details_leaves_attribute_absent(_mock_close, mock_create): mock_create.return_value = mock_client response = _agent()._interruptible_streaming_api_call({}) assert not hasattr(response.choices[0].message, "reasoning_details") + + +@patch("run_agent.AIAgent._create_request_openai_client") +@patch("run_agent.AIAgent._close_request_openai_client") +def test_live_details_keep_plain_fallback_and_opaque_replay(_mock_close, mock_create): + agent = _agent() + client = MagicMock() + mock_create.return_value = client + cases = [ + ([{"type": "reasoning.text", "text": "Complete thought"}], "C", "Complete thought"), + ([{"type": "reasoning.text", "text": "Same"}], "Same", "Same"), + ([SimpleNamespace(type="reasoning.summary", summary="Summary")], None, "Summary"), + ([{"type": "reasoning.encrypted", "data": "secret", "text": "not readable"}], "Plain", "Plain"), + ([{"type": "unknown", "text": "not readable"}], "Plain", "Plain"), + ([{"type": "reasoning.text", "text": ""}], "Plain", "Plain"), + ([], "Plain", "Plain"), + ([{"type": "reasoning.encrypted", "data": "secret"}], None, ""), + ] + for details, plain, expected in cases: + chunk = _make_chunk(content="Answer", finish_reason="stop") + chunk.choices[0].delta.reasoning = plain + chunk.choices[0].delta.model_extra = {"reasoning_details": details} + client.chat.completions.create.return_value = iter([chunk]) + delivered = [] + agent.reasoning_callback = delivered.append + response = agent._interruptible_streaming_api_call({}) + assert "".join(delivered) == expected + assert response.choices[0].message.content == "Answer" + assert response.choices[0].message.reasoning_content == plain + preserved = [] + for detail in details: + append_streamed_reasoning_detail(preserved, detail) + assert getattr(response.choices[0].message, "reasoning_details", []) == preserved + + def broken_callback(text): + raise RuntimeError("consumer failed") + + for callback in (None, broken_callback): + agent.reasoning_callback = callback + client.chat.completions.create.return_value = iter([ + _make_chunk(content="Answer", finish_reason="stop", reasoning_details=[ + {"type": "reasoning.text", "text": "Thought"}])]) + assert agent._interruptible_streaming_api_call({}).choices[0].message.content == "Answer"