diff --git a/agent/chat_completion_helpers.py b/agent/chat_completion_helpers.py index 946df82010..0a17862aae 100644 --- a/agent/chat_completion_helpers.py +++ b/agent/chat_completion_helpers.py @@ -2728,6 +2728,10 @@ class _StreamingCall(StreamingWaitMonitor): base_timeout, read_timeout, conn_cap = self._stream_timeouts() content_parts: list = [] reasoning_parts: list = [] + # OpenAI structured refusal (``delta.refusal``): the explanation streams here and + # ``delta.content`` stays empty, so an un-accumulated refusal looks like an empty + # stream and burns the empty-response retries (the non-streaming fix is #46013). + refusal_parts: list[str] = [] reasoning_details: list = [] # OpenRouter replay data (signatures, encrypted blocks) pending_text_parts: list[str] = [] tool_calls = _ToolCallAccumulator() @@ -2812,6 +2816,13 @@ class _StreamingCall(StreamingWaitMonitor): rd_delta = delta.model_extra.get("reasoning_details") for rd in rd_delta if isinstance(rd_delta, (list, tuple)) else (): append_streamed_reasoning_detail(reasoning_details, rd) + # 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) + if delta_refusal is None and isinstance(getattr(delta, "model_extra", None), dict): + delta_refusal = delta.model_extra.get("refusal") + if isinstance(delta_refusal, str) and delta_refusal: + refusal_parts.append(delta_refusal) # Text (list-of-blocks deltas flattened once); possible echoed SSE is # buffered until it can be judged. @@ -2847,7 +2858,8 @@ class _StreamingCall(StreamingWaitMonitor): return self._adopt_final_response(stream.final_response) return self._finish_chat_stream(stream, role, content_parts, reasoning_parts, tool_calls_acc, finish_reason, model_name, usage_obj, flush_pending=_flush_pending_stream_text, - response_id=response_id, upstream_provider=upstream_provider, reasoning_details=reasoning_details) + response_id=response_id, upstream_provider=upstream_provider, reasoning_details=reasoning_details, + refusal_parts=refusal_parts) def _adopt_final_response(self, final_response): """Adapter returned a completed response for ``stream=True``: switch the @@ -2896,7 +2908,8 @@ class _StreamingCall(StreamingWaitMonitor): return mock_tool_calls or None, has_truncated_tool_args def _finish_chat_stream(self, stream, role, content_parts, reasoning_parts, tool_calls_acc, finish_reason, - model_name, usage_obj, *, flush_pending, response_id=None, upstream_provider=None, reasoning_details=None): + model_name, usage_obj, *, flush_pending, response_id=None, upstream_provider=None, reasoning_details=None, + refusal_parts=None): """Assemble the non-streaming-shaped response after the chunk loop. A stream ending with no finish_reason is a drop, not a completion: return a partial-stream stub so the loop fails fast instead of executing empty @@ -2905,7 +2918,7 @@ class _StreamingCall(StreamingWaitMonitor): full_reasoning = "".join(reasoning_parts) or None mock_tool_calls, has_truncated_tool_args = self._assemble_tool_calls(tool_calls_acc, finish_reason) # Zero-chunk guard: nothing usable = upstream error / malformed SSE. - if finish_reason is None and not content_parts and not reasoning_parts and not tool_calls_acc: + if finish_reason is None and not content_parts and not reasoning_parts and not refusal_parts and not tool_calls_acc: raise EmptyStreamError( "Provider returned an empty stream with no finish_reason (possible upstream error or malformed SSE response).") if has_truncated_tool_args and finish_reason is None: @@ -2932,7 +2945,9 @@ class _StreamingCall(StreamingWaitMonitor): if provider_stream_error is not None: raise provider_stream_error flush_pending() - message = SimpleNamespace(role=role, content=full_content, tool_calls=mock_tool_calls, reasoning_content=full_reasoning) + message = SimpleNamespace(role=role, content=full_content, tool_calls=mock_tool_calls, reasoning_content=full_reasoning, + # ``normalize_response`` reads ``message.refusal`` — same contract as the non-streaming object. + refusal="".join(refusal_parts or ()) or None) if reasoning_details: # Only when present: _build_assistant_message's passthrough persists them # for replay, and non-reasoning providers keep the attribute absent. diff --git a/agent/chat_completion_helpers_relay.py b/agent/chat_completion_helpers_relay.py index 37dd8b850f..ef5923c819 100644 --- a/agent/chat_completion_helpers_relay.py +++ b/agent/chat_completion_helpers_relay.py @@ -34,6 +34,7 @@ class RelayChatAccumulator: def __init__(self) -> None: self._content: list[str] = [] self._reasoning: list[str] = [] + self._refusal: list[str] = [] # OpenAI ``delta.refusal`` — a refusal is content, not an empty stream self._tool_calls = _ToolCallAccumulator() self._model = self._usage = self._finish_reason = None self._role = "assistant" @@ -61,6 +62,9 @@ class RelayChatAccumulator: if reasoning: self._reasoning.append(separate_glued_reasoning_blocks( self._reasoning[-1] if self._reasoning else "", reasoning)) + refusal = delta.get("refusal") + if isinstance(refusal, str) and refusal: + self._refusal.append(refusal) for tc_delta in delta.get("tool_calls") or []: self._tool_calls.feed(_tool_call_delta_view(tc_delta)) @@ -68,6 +72,7 @@ class RelayChatAccumulator: acc = self._tool_calls.materialize() message = {"role": self._role, "content": "".join(self._content) or None, "reasoning_content": "".join(self._reasoning) or None, + "refusal": "".join(self._refusal) or None, "tool_calls": [acc[i] for i in sorted(acc)] or None} # "stop" also covers Nous Portal ``lastOne`` usage frames, which carry no finish_reason. return {"model": self._model, "usage": self._usage, diff --git a/agent/codex_responses_adapter.py b/agent/codex_responses_adapter.py index b98d83df73..f20079b6d6 100644 --- a/agent/codex_responses_adapter.py +++ b/agent/codex_responses_adapter.py @@ -912,8 +912,16 @@ def _text_chunks(parts: Any, types: Optional[set] = None) -> List[str]: def _extract_responses_message_text(item: Any) -> str: - """Extract assistant text from a Responses message output item.""" - return "".join(_text_chunks(getattr(item, "content", None), _OUTPUT_TEXT_TYPES)).strip() + """Assistant text from a Responses message output item. A ``refusal`` part carries the + model's explanation in ``refusal`` instead of ``text``; it is message text too, otherwise a + refusal-only turn reads as an empty response (sibling of chat_completions ``message.refusal``).""" + chunks = [] + for part in _as_list(_field(item, "content")): + ptype = _field(part, "type") + text = _field(part, "refusal") if ptype == "refusal" else (_field(part, "text") if ptype in _OUTPUT_TEXT_TYPES else None) + if _nonempty_str(text): + chunks.append(text) + return "".join(chunks).strip() def _extract_responses_reasoning_text(item: Any) -> str: diff --git a/agent/codex_runtime.py b/agent/codex_runtime.py index d3263b8f99..8cf7e2bb1e 100644 --- a/agent/codex_runtime.py +++ b/agent/codex_runtime.py @@ -656,6 +656,15 @@ class _CodexResponseAssembler: self._safe(self.on_first_delta, "on_first_delta") self._safe(self.on_text_delta, "on_text_delta", delta_text) + def _on_refusal_delta(self, event: Any, event_type: str) -> None: + # ``response.refusal.delta``: the model declined and streams its explanation on the refusal + # channel instead of output_text. It is answer text — a refusal-only stream must not end + # with zero content and "did not emit a terminal response". The done item's ``refusal`` + # part is read by the normalizer; the deltas cover backends that omit the done item. + refusal_text = _event_field(event, "delta", "") + if isinstance(refusal_text, str) and refusal_text: + self.text_deltas.append(refusal_text) + def _on_function_call(self, event: Any, event_type: str) -> None: self.has_tool_calls = True pending = self.pending_function_calls.get(str(_event_field(event, "item_id", ""))) @@ -723,6 +732,7 @@ class _CodexResponseAssembler: "error": lambda self, event, event_type: _raise_stream_error(event), "response.output_item.added": _on_item_added, "response.output_item.done": _on_item_done, "response.completed": _on_terminal, "response.incomplete": _on_terminal, "response.failed": _on_terminal, + "response.refusal.delta": _on_refusal_delta, } _FUZZY_HANDLERS = ( (lambda t: "output_text.delta" in t, _on_text_delta), (lambda t: "function_call" in t, _on_function_call), diff --git a/tests/agent/test_run_agent_codex_responses.py b/tests/agent/test_run_agent_codex_responses.py index ae7b03786e..b6ad9be33a 100644 --- a/tests/agent/test_run_agent_codex_responses.py +++ b/tests/agent/test_run_agent_codex_responses.py @@ -845,6 +845,55 @@ def test_consume_codex_stream_routes_commentary_phase_deltas_to_reasoning(monkey assert response.output_text == "" +def test_consume_codex_stream_collects_refusal_deltas_as_text(monkeypatch): + """A refusal-only Responses stream yields usable text, not RuntimeError. + + Port of anomalyco/opencode#43343: the model declines and streams the + explanation via ``response.refusal.delta`` with no output_text and (on + some compatible backends) no output_item.done — without collecting the + refusal the consumer sees zero usable content. + """ + from agent.codex_runtime import _consume_codex_event_stream + + response = _consume_codex_event_stream( + _FakeCreateStream([ + SimpleNamespace(type="response.created"), + SimpleNamespace(type="response.refusal.delta", delta="I can't"), + SimpleNamespace(type="response.refusal.delta", delta=" help with that."), + SimpleNamespace(type="response.completed", response=SimpleNamespace(status="completed")), + ]), + model="gpt-5-codex", + ) + + assert response.output_text == "I can't help with that." + # Synthesized message item so downstream normalization has content. + assert response.output and response.output[0].type == "message" + + +def test_extract_responses_message_text_reads_refusal_parts(): + """Refusal content parts inside a message item surface as text.""" + from agent.codex_responses_adapter import _extract_responses_message_text + + item = SimpleNamespace( + type="message", + role="assistant", + content=[ + SimpleNamespace(type="refusal", refusal="I must decline."), + ], + ) + assert _extract_responses_message_text(item) == "I must decline." + + mixed = SimpleNamespace( + type="message", + role="assistant", + content=[ + SimpleNamespace(type="output_text", text="Hello. "), + SimpleNamespace(type="refusal", refusal="But no more."), + ], + ) + assert _extract_responses_message_text(mixed) == "Hello. But no more." + + def test_consume_codex_stream_separates_commentary_from_analysis(monkeypatch): from agent.codex_runtime import _consume_codex_event_stream diff --git a/tests/agent/test_streaming.py b/tests/agent/test_streaming.py index 27d74323d0..a2be5abdbd 100644 --- a/tests/agent/test_streaming.py +++ b/tests/agent/test_streaming.py @@ -521,6 +521,111 @@ class TestStreamingAccumulator: # ── Test: Streaming Callbacks ──────────────────────────────────────────── + @patch("run_agent.AIAgent._create_request_openai_client") + @patch("run_agent.AIAgent._close_request_openai_client") + def test_streamed_refusal_accumulated(self, mock_close, mock_create): + """delta.refusal streams assemble onto message.refusal (opencode#43343). + + A refusal-only stream must not raise EmptyStreamError, and the + assembled message must expose ``refusal`` so the transport's + normalize_response promotes it to content + content_filter. + """ + from run_agent import AIAgent + + def _refusal_chunk(refusal, finish_reason=None): + delta = SimpleNamespace( + content=None, + tool_calls=None, + reasoning_content=None, + reasoning=None, + refusal=refusal, + ) + choice = SimpleNamespace(index=0, delta=delta, finish_reason=finish_reason) + return SimpleNamespace(choices=[choice], model="test-model", usage=None) + + chunks = [ + _refusal_chunk("I can't"), + _refusal_chunk(" help with that."), + _refusal_chunk(None, finish_reason="stop"), + ] + + mock_client = MagicMock() + mock_client.chat.completions.create.return_value = iter(chunks) + mock_create.return_value = mock_client + + agent = AIAgent( + api_key="test-key", + base_url="https://openrouter.ai/api/v1", + model="test/model", + quiet_mode=True, + skip_context_files=True, + skip_memory=True, + ) + agent.api_mode = "chat_completions" + agent._interrupt_requested = False + + response = agent._interruptible_streaming_api_call({}) + + message = response.choices[0].message + assert message.refusal == "I can't help with that." + assert message.content is None + + # The transport promotes a sole-payload refusal to a terminal + # content_filter (same contract as the non-streaming path, #46013). + from agent.transports.chat_completions import ChatCompletionsTransport + + normalized = ChatCompletionsTransport().normalize_response(response) + assert normalized.content == "I can't help with that." + assert normalized.finish_reason == "content_filter" + + @patch("run_agent.AIAgent._create_request_openai_client") + @patch("run_agent.AIAgent._close_request_openai_client") + def test_streamed_refusal_alongside_content_not_terminal(self, mock_close, mock_create): + """A refusal note next to real content stays a normal usable turn.""" + from run_agent import AIAgent + + def _chunk(content=None, refusal=None, finish_reason=None): + delta = SimpleNamespace( + content=content, + tool_calls=None, + reasoning_content=None, + reasoning=None, + refusal=refusal, + ) + choice = SimpleNamespace(index=0, delta=delta, finish_reason=finish_reason) + return SimpleNamespace(choices=[choice], model="test-model", usage=None) + + chunks = [ + _chunk(content="Partial answer."), + _chunk(refusal="But I won't do the rest."), + _chunk(finish_reason="stop"), + ] + + mock_client = MagicMock() + mock_client.chat.completions.create.return_value = iter(chunks) + mock_create.return_value = mock_client + + agent = AIAgent( + api_key="test-key", + base_url="https://openrouter.ai/api/v1", + model="test/model", + quiet_mode=True, + skip_context_files=True, + skip_memory=True, + ) + agent.api_mode = "chat_completions" + agent._interrupt_requested = False + + response = agent._interruptible_streaming_api_call({}) + + from agent.transports.chat_completions import ChatCompletionsTransport + + normalized = ChatCompletionsTransport().normalize_response(response) + assert normalized.content == "Partial answer." + assert normalized.finish_reason == "stop" + assert normalized.provider_data["refusal"] == "But I won't do the rest." + + class TestStreamingCallbacks: """Verify that delta callbacks fire correctly."""