fix(api-server): /v1/responses stream carries commentary as a phase-tagged message item

Codex phase="commentary" progress (and any mid-turn assistant text) is
emitted on POST /v1/responses stream=true as its own completed
`message` output item with `"phase": "commentary"` — added/done pair,
listed in response.completed output — and kept out of the final answer
item so clients can render live progress without polluting the final
text (#67580). Text the agent already streamed as output_text.delta is
not re-emitted.

Hand-reapplied from #67593 onto the _ResponsesStream emitter.
This commit is contained in:
PRATHAMESH75
2026-09-19 00:56:29 -07:00
committed by teknium1
parent b44ccdcc0e
commit 968f99e2eb

View File

@@ -221,6 +221,18 @@ class _ResponsesStream:
"output_index": self.message_output_index, "content_index": 0, "delta": delta_text,
"logprobs": []})
async def emit_commentary(self, text: str) -> None:
"""Mid-turn assistant commentary as its own completed ``message`` item carrying
``"phase": "commentary"`` — never appended to ``final_text_parts``, so the final answer
item stays clean (#67580)."""
item = {"id": f"msg_{uuid.uuid4().hex[:24]}", "status": "completed", "phase": "commentary",
**_message_item(text)}
idx = self.output_index
self.output_index += 1
self.emitted_items.append({"phase": "commentary", **_message_item(text)})
for event in ("response.output_item.added", "response.output_item.done"):
await self.write_event(event, {"type": event, "output_index": idx, "item": item})
async def emit_tool_started(self, payload: Dict[str, Any]) -> None:
"""function_call ``output_item.added``; the agent's tool_call_id beats a generated call id."""
self.call_counter += 1
@@ -274,6 +286,8 @@ class _ResponsesStream:
await self.emit_tool_started(payload)
elif tag == "__tool_completed__":
await self.emit_tool_completed(payload)
elif tag == "__commentary__":
await self.emit_commentary(payload["text"])
elif isinstance(item, str):
self._batch_buf.append(item)
if self._batch_timer is None:
@@ -888,10 +902,16 @@ class OpenAICompatRoutesMixin:
_stream_q.put_threadsafe(("__tool_completed__", {
"tool_call_id": tool_call_id, "name": function_name,
"arguments": function_args or {}, "result": function_result}))
def _on_commentary(text, *, already_streamed: bool = False):
# Already-streamed text went out as output_text.delta of the final item; a second
# copy as a commentary item would duplicate it.
if not already_streamed and isinstance(text, str) and text.strip():
_stream_q.put_threadsafe(("__commentary__", {"text": text}))
agent_task, agent_ref = self._spawn_stream_agent(
_stream_q, tool_progress_callback=_on_tool_progress,
tool_start_callback=_on_tool_start, tool_complete_callback=_on_tool_complete,
**run_kwargs)
interim_assistant_callback=_on_commentary, **run_kwargs)
return await self._write_sse_responses(
request=request, response_id=f"resp_{uuid.uuid4().hex[:28]}",
model=body.get("model", self._model_name), created_at=int(time.time()),