From 16578186ecb0fb335b5f4a8c4f5286530de71e19 Mon Sep 17 00:00:00 2001 From: EmBista Date: Thu, 20 Aug 2026 23:04:05 +1000 Subject: [PATCH 1/2] Fix non-stream Responses output aggregation --- chatmock/responses_api.py | 13 +++++ tests/test_routes.py | 100 ++++++++++++++++++++++++++++++++++++++ 2 files changed, 113 insertions(+) diff --git a/chatmock/responses_api.py b/chatmock/responses_api.py index ab66803..df6f439 100644 --- a/chatmock/responses_api.py +++ b/chatmock/responses_api.py @@ -1,5 +1,6 @@ from __future__ import annotations +import copy import json from dataclasses import dataclass from typing import Any, Dict, Iterable, Iterator, List @@ -171,6 +172,7 @@ def aggregate_response_from_sse( ) -> tuple[Dict[str, Any] | None, Dict[str, Any] | None]: response_obj: Dict[str, Any] | None = None error_obj: Dict[str, Any] | None = None + completed_output_items: List[Dict[str, Any]] = [] try: for evt in iter_sse_event_payloads(upstream): if callable(on_event): @@ -182,6 +184,10 @@ def aggregate_response_from_sse( if isinstance(response, dict): response_obj = response kind = evt.get("type") + if kind == "response.output_item.done": + item = evt.get("item") + if isinstance(item, dict): + completed_output_items.append(copy.deepcopy(item)) if kind == "response.failed": if isinstance(response, dict) and isinstance(response.get("error"), dict): error_obj = {"error": response.get("error")} @@ -189,6 +195,13 @@ def aggregate_response_from_sse( error_obj = {"error": {"message": "response.failed"}} break if kind == "response.completed": + if ( + isinstance(response_obj, dict) + and completed_output_items + and not response_obj.get("output") + ): + response_obj = dict(response_obj) + response_obj["output"] = completed_output_items break finally: upstream.close() diff --git a/tests/test_routes.py b/tests/test_routes.py index a490670..a380114 100644 --- a/tests/test_routes.py +++ b/tests/test_routes.py @@ -260,6 +260,106 @@ def test_responses_route_returns_completed_response_object(self, mock_start) -> self.assertEqual(outbound_payload["reasoning"]["effort"], "medium") self.assertIsInstance(outbound_payload["prompt_cache_key"], str) + @patch("chatmock.routes_openai.start_upstream_raw_request") + def test_responses_route_reconstructs_non_stream_output_from_item_events(self, mock_start) -> None: + output = [ + { + "type": "reasoning", + "id": "reasoning_1", + "summary": [{"type": "summary_text", "text": "Need the tool."}], + "encrypted_content": "encrypted", + }, + { + "type": "function_call", + "id": "fc_1", + "call_id": "call_1", + "name": "get_time", + "arguments": '{"city":"Paris"}', + "status": "completed", + }, + { + "type": "message", + "role": "assistant", + "id": "msg_1", + "status": "completed", + "content": [{"type": "output_text", "text": '{"city":"Paris"}'}], + }, + ] + events = [ + { + "type": "response.created", + "response": {"id": "resp_items", "object": "response", "status": "in_progress"}, + }, + *[ + {"type": "response.output_item.done", "output_index": index, "item": item} + for index, item in enumerate(output) + ], + { + "type": "response.completed", + "response": { + "id": "resp_items", + "object": "response", + "status": "completed", + "output": [], + }, + }, + ] + mock_start.return_value = ( + FakeUpstream(events, headers={"Content-Type": "text/event-stream"}), + None, + ) + + response = self.client.post( + "/v1/responses", + json={"model": "gpt-5.6-luna", "input": "Return structured output."}, + ) + + self.assertEqual(response.status_code, 200) + self.assertEqual(response.get_json()["output"], output) + + @patch("chatmock.routes_openai.start_upstream_raw_request") + def test_responses_route_keeps_output_from_completed_response(self, mock_start) -> None: + authoritative = { + "type": "message", + "role": "assistant", + "id": "msg_final", + "content": [{"type": "output_text", "text": "final"}], + } + mock_start.return_value = ( + FakeUpstream( + [ + { + "type": "response.output_item.done", + "item": { + "type": "message", + "role": "assistant", + "id": "msg_event", + "content": [{"type": "output_text", "text": "event"}], + }, + }, + { + "type": "response.completed", + "response": { + "id": "resp_final", + "object": "response", + "status": "completed", + "output": [authoritative], + }, + }, + ], + headers={"Content-Type": "text/event-stream"}, + ), + None, + ) + + response = self.client.post( + "/v1/responses", + json={"model": "gpt-5.6-luna", "input": "hello"}, + ) + + self.assertEqual(response.status_code, 200) + self.assertEqual(response.get_json()["output"], [authoritative]) + @patch("chatmock.routes_openai.start_upstream_raw_request") def test_responses_route_honors_debug_model_override(self, mock_start) -> None: app = create_app(debug_model="gpt-5.4", model_sync=False) From 83c35f10f6e9c9461cf9f2aca144169d8162465c Mon Sep 17 00:00:00 2001 From: EmBista Date: Fri, 21 Aug 2026 00:09:33 +1000 Subject: [PATCH 2/2] Preserve Responses output item ordering --- chatmock/responses_api.py | 17 ++++++++++++++--- tests/test_routes.py | 6 ++++-- 2 files changed, 18 insertions(+), 5 deletions(-) diff --git a/chatmock/responses_api.py b/chatmock/responses_api.py index df6f439..bbc28be 100644 --- a/chatmock/responses_api.py +++ b/chatmock/responses_api.py @@ -172,7 +172,8 @@ def aggregate_response_from_sse( ) -> tuple[Dict[str, Any] | None, Dict[str, Any] | None]: response_obj: Dict[str, Any] | None = None error_obj: Dict[str, Any] | None = None - completed_output_items: List[Dict[str, Any]] = [] + completed_output_items: Dict[int, Dict[str, Any]] = {} + unindexed_output_items = 0 try: for evt in iter_sse_event_payloads(upstream): if callable(on_event): @@ -187,7 +188,14 @@ def aggregate_response_from_sse( if kind == "response.output_item.done": item = evt.get("item") if isinstance(item, dict): - completed_output_items.append(copy.deepcopy(item)) + output_index = evt.get("output_index") + if not isinstance(output_index, int): + # Indexed items are the protocol norm. Keep malformed or + # older unindexed events deterministically after them, + # preserving their arrival order. + output_index = 1_000_000 + unindexed_output_items + unindexed_output_items += 1 + completed_output_items[output_index] = copy.deepcopy(item) if kind == "response.failed": if isinstance(response, dict) and isinstance(response.get("error"), dict): error_obj = {"error": response.get("error")} @@ -201,7 +209,10 @@ def aggregate_response_from_sse( and not response_obj.get("output") ): response_obj = dict(response_obj) - response_obj["output"] = completed_output_items + response_obj["output"] = [ + completed_output_items[index] + for index in sorted(completed_output_items) + ] break finally: upstream.close() diff --git a/tests/test_routes.py b/tests/test_routes.py index a380114..03c6c76 100644 --- a/tests/test_routes.py +++ b/tests/test_routes.py @@ -291,8 +291,10 @@ def test_responses_route_reconstructs_non_stream_output_from_item_events(self, m "response": {"id": "resp_items", "object": "response", "status": "in_progress"}, }, *[ - {"type": "response.output_item.done", "output_index": index, "item": item} - for index, item in enumerate(output) + {"type": "response.output_item.done", "output_index": index, "item": output[index]} + # Completion order is not output order; the protocol supplies + # output_index so non-stream aggregation can reconstruct it. + for index in (2, 0, 1) ], { "type": "response.completed",