fix(agent): preserve streamed refusals as text (port of anomalyco/opencode#43343)

A model that declines mid-stream delivers the explanation on the
structured refusal channel (chat_completions delta.refusal; Responses
response.refusal.delta / refusal content parts) and leaves content
empty. The streaming accumulators dropped that channel entirely, so a
streamed refusal assembled into an empty message and fell into the
empty/invalid-response retry loops - burning paid retries reproducing a
deterministic refusal - while the non-streaming path had already fixed
this class in #46013.

- chat_completions streaming: accumulate delta.refusal (incl.
  model_extra), expose message.refusal on the assembled mock response so
  ChatCompletionsTransport.normalize_response applies the existing
  sole-payload -> content_filter promotion; count refusal deltas in the
  zero-chunk guard; carry refusal in the Relay final-response dict.
- Codex Responses stream consumer: collect response.refusal.delta as
  answer text so a refusal-only stream no longer raises 'did not emit a
  terminal response' with zero usable content.
- Responses normalizer: read type=refusal content parts in
  _extract_responses_message_text (attr and dict shapes).

Sabotage-verified: each new test fails with its wiring line disabled.
E2E: refusal-only stream -> terminal content_filter with explanation;
refusal-alongside-content stays a normal usable turn; plain-text
streams unchanged.
This commit is contained in:
Hermes Agent
2026-08-20 17:21:10 -07:00
committed by Teknium
parent 021ab58a23
commit e27f16365d
6 changed files with 198 additions and 6 deletions
+19 -4
View File
@@ -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.
+5
View File
@@ -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,
+10 -2
View File
@@ -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:
+10
View File
@@ -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),
@@ -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
+105
View File
@@ -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."""