diff --git a/headroom/proxy/handlers/openai.py b/headroom/proxy/handlers/openai.py index 39926a430..f7bbd8eec 100644 --- a/headroom/proxy/handlers/openai.py +++ b/headroom/proxy/handlers/openai.py @@ -957,31 +957,71 @@ def _dedup_responses_output_items( def _openai_responses_to_sse(response: dict[str, Any]) -> list[bytes]: - """Convert a complete Responses API JSON body into a minimal SSE stream. + """Convert a complete Responses API JSON body into an SSE stream. - Used only for the buffered-CCR path: the client asked for - ``stream: true`` but we forced a non-streaming upstream call so CCR - retrieval could be resolved server-side. This reconstructs just enough - of the real event sequence (``response.created`` + ``response.completed``) - for Responses API clients that key off the terminal event's full - response object — it does not replay incremental output-item/text - deltas. Mirrors the equivalent simplification in - ``StreamingMixin._response_to_sse`` for the Anthropic buffered path. + Used only for the buffered-CCR path: the client asked for ``stream: true`` + but we forced a non-streaming upstream call so CCR retrieval could be + resolved server-side. We then have to replay the response as SSE. + + Some Responses clients read the whole answer off the terminal + ``response.completed`` event, but others (OpenCode / the Vercel AI SDK) + render output only from the *incremental* item/text events and show nothing + when they are absent (#2410). So reconstruct the real event sequence: + ``response.created`` -> ``response.in_progress`` -> per output item + (``response.output_item.added``, and for message items the + ``response.content_part.added`` / ``response.output_text.delta`` / + ``response.output_text.done`` / ``response.content_part.done`` sequence) -> + ``response.output_item.done`` -> ``response.completed`` -> ``[DONE]``. """ - created_response = {**response, "status": "in_progress", "output": []} events: list[bytes] = [] - for seq, (event_type, event_response) in enumerate( - ( - ("response.created", created_response), - ("response.completed", response), - ) - ): - payload = { - "type": event_type, - "sequence_number": seq, - "response": event_response, - } + seq = 0 + + def _emit(event_type: str, extra: dict[str, Any]) -> None: + nonlocal seq + payload = {"type": event_type, "sequence_number": seq, **extra} events.append(f"event: {event_type}\ndata: {json.dumps(payload)}\n\n".encode()) + seq += 1 + + output_items = response.get("output") or [] + created_response = {**response, "status": "in_progress", "output": []} + _emit("response.created", {"response": created_response}) + _emit("response.in_progress", {"response": created_response}) + + for out_idx, item in enumerate(output_items): + if not isinstance(item, dict): + continue + item_id = item.get("id", f"item_{out_idx}") + + # ``output_item.added`` carries the item shell; message content streams + # via the content-part events below, so start it empty there. + if item.get("type") == "message": + added_item = {k: v for k, v in item.items() if k != "content"} + added_item["content"] = [] + else: + added_item = item + _emit("response.output_item.added", {"output_index": out_idx, "item": added_item}) + + content = item.get("content") + if item.get("type") == "message" and isinstance(content, list): + for c_idx, part in enumerate(content): + if not isinstance(part, dict): + continue + loc = {"item_id": item_id, "output_index": out_idx, "content_index": c_idx} + if part.get("type") in ("output_text", "text"): + text = part.get("text", "") or "" + _emit("response.content_part.added", {**loc, "part": {**part, "text": ""}}) + if text: + _emit("response.output_text.delta", {**loc, "delta": text}) + _emit("response.output_text.done", {**loc, "text": text}) + _emit("response.content_part.done", {**loc, "part": part}) + else: + # Non-text part (e.g. refusal): add + done with the full part. + _emit("response.content_part.added", {**loc, "part": part}) + _emit("response.content_part.done", {**loc, "part": part}) + + _emit("response.output_item.done", {"output_index": out_idx, "item": item}) + + _emit("response.completed", {"response": response}) events.append(b"data: [DONE]\n\n") return events diff --git a/tests/test_openai_responses_buffered_sse.py b/tests/test_openai_responses_buffered_sse.py new file mode 100644 index 000000000..a3271f5dc --- /dev/null +++ b/tests/test_openai_responses_buffered_sse.py @@ -0,0 +1,109 @@ +"""Regression for #2410: the buffered-CCR Responses -> SSE reconstruction must +replay the incremental output-item/text events, not just response.created + +response.completed, so AI-SDK / OpenCode clients render the output.""" + +from __future__ import annotations + +import json + +from headroom.proxy.handlers.openai import _openai_responses_to_sse + + +def _parse(events: list[bytes]) -> list[dict]: + out: list[dict] = [] + for e in events: + s = e.decode() + if s.startswith("data: [DONE]"): + out.append({"type": "[DONE]"}) + continue + out.append(json.loads(s.split("data: ", 1)[1])) + return out + + +def test_responses_sse_replays_incremental_output_text() -> None: + resp = { + "id": "resp_1", + "object": "response", + "status": "completed", + "model": "gpt-5.3-codex", + "output": [ + {"type": "reasoning", "id": "rs_1", "summary": []}, + { + "type": "message", + "id": "msg_1", + "status": "completed", + "role": "assistant", + "content": [{"type": "output_text", "text": "Hello world", "annotations": []}], + }, + ], + "usage": {"input_tokens": 10, "output_tokens": 3}, + } + + parsed = _parse(_openai_responses_to_sse(resp)) + types = [p["type"] for p in parsed] + + assert types[0] == "response.created" + assert types[1] == "response.in_progress" + assert types[-2] == "response.completed" + assert types[-1] == "[DONE]" + + # The visible assistant text is streamed as an output_text.delta. + deltas = [p for p in parsed if p["type"] == "response.output_text.delta"] + assert len(deltas) == 1 + assert deltas[0]["delta"] == "Hello world" + assert deltas[0]["output_index"] == 1 + assert deltas[0]["content_index"] == 0 + + # The message item gets the full content-part sequence; the reasoning item + # gets add/done with no content parts. + assert types.count("response.output_item.added") == 2 + assert types.count("response.output_item.done") == 2 + assert "response.content_part.added" in types + assert "response.output_text.done" in types + assert "response.content_part.done" in types + + # created / in_progress carry an empty output; completed carries the full one. + created = next(p for p in parsed if p["type"] == "response.created") + assert created["response"]["output"] == [] + completed = next(p for p in parsed if p["type"] == "response.completed") + assert completed["response"]["output"] == resp["output"] + + # Sequence numbers are contiguous from 0. + seqs = [p["sequence_number"] for p in parsed if p["type"] != "[DONE]"] + assert seqs == list(range(len(seqs))) + + +def test_responses_sse_empty_output_still_valid() -> None: + resp = {"id": "resp_2", "status": "completed", "output": [], "usage": {}} + types = [p["type"] for p in _parse(_openai_responses_to_sse(resp))] + assert types == ["response.created", "response.in_progress", "response.completed", "[DONE]"] + + +def test_responses_sse_non_message_item_added_and_done() -> None: + resp = { + "id": "resp_3", + "status": "completed", + "output": [ + { + "type": "function_call", + "id": "fc_1", + "call_id": "c1", + "name": "grep", + "arguments": "{}", + } + ], + "usage": {}, + } + parsed = _parse(_openai_responses_to_sse(resp)) + types = [p["type"] for p in parsed] + assert types == [ + "response.created", + "response.in_progress", + "response.output_item.added", + "response.output_item.done", + "response.completed", + "[DONE]", + ] + # The function_call item is preserved whole on added and done. + done = next(p for p in parsed if p["type"] == "response.output_item.done") + assert done["item"]["name"] == "grep"