fix(agent): deliver readable reasoning details during streaming
The streaming reasoning display read only reasoning_content/reasoning deltas; reasoning_details deltas were accumulated for replay continuity but their text was never routed to the live display, so models that stream thought solely via reasoning_details (xiaomi/mimo-v2.6-pro via OpenRouter) showed no thinking display at all. The non-streaming extract_reasoning already reads detail text — this closes the asymmetry. One representation per chunk is emitted: detail text when the details carry it, otherwise the plain reasoning field. Persisted/replayed fields are not rewritten; opaque replay material (encrypted/unknown types) stays unexposed. A callback that raises (or is None) no longer breaks the stream. Fixes #118851 Salvaged from #118888 by Wenfengcheng
This commit is contained in:
@@ -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)
|
||||
|
||||
@@ -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*.
|
||||
|
||||
|
||||
@@ -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"
|
||||
|
||||
Reference in New Issue
Block a user