From 8e69c5b053f2f6625441318a7b2bfeecca4b995d Mon Sep 17 00:00:00 2001 From: adaofeliz Date: Thu, 30 Jul 2026 08:54:14 +0100 Subject: [PATCH] feat(api-server): surface reasoning_content on /v1/chat/completions MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit /v1/runs and /v1/responses already relay reasoning.available previews as structured events, but /v1/chat/completions never wired a callback to receive them, so reasoning/thinking text was silently dropped for that endpoint — even though it's the one most OpenAI-SDK clients use. Add a reasoning.available-filtered tool_progress_callback in both the streaming and non-streaming paths: - Streaming: emits delta.reasoning_content chunks (same field name DeepSeek/Moonshot/OpenRouter use), separate from delta.content and the existing hermes.tool.progress tool-lifecycle events. - Non-streaming: accumulates previews into message.reasoning_content on the final response. Tool-call events are unaffected — the filter only forwards reasoning.available, so tool_start_callback/tool_complete_callback still own the tool-lifecycle channel exactly as before. --- gateway/platforms/api_server.py | 36 +++++++++++--- tests/gateway/test_api_server.py | 85 ++++++++++++++++++++++++++++++++ 2 files changed, 115 insertions(+), 6 deletions(-) diff --git a/gateway/platforms/api_server.py b/gateway/platforms/api_server.py index 62dff2878fa59..3dc96005cf5e0 100644 --- a/gateway/platforms/api_server.py +++ b/gateway/platforms/api_server.py @@ -3916,14 +3916,15 @@ def _on_tool_complete(tool_call_id, function_name, function_args, function_resul "status": "completed", })) + # Relay ``reasoning.available`` previews as ``delta.reasoning_content`` + # chunks (see ``_emit``). Filtered so tool-lifecycle events (already + # covered by ``_on_tool_start``/``_on_tool_complete``) aren't duplicated. + def _on_reasoning(event_type, tool_name=None, preview=None, args=None, **kwargs): + if event_type == "reasoning.available" and preview: + _stream_q.put(("__reasoning__", preview)) + # Start agent in background. agent_ref is a mutable container # so the SSE writer can interrupt the agent on client disconnect. - # - # ``tool_progress_callback`` is intentionally not wired here: - # it would duplicate every emit because ``run_agent`` fires it - # side-by-side with ``tool_start_callback``/``tool_complete_callback``. - # The structured callbacks are strictly richer (they carry - # the tool_call id), so they own the chat-completions SSE channel. agent_ref = [None] agent_task = asyncio.ensure_future(self._run_agent( user_message=user_message, @@ -3931,6 +3932,7 @@ def _on_tool_complete(tool_call_id, function_name, function_args, function_resul ephemeral_system_prompt=system_prompt, session_id=session_id, stream_delta_callback=_on_delta, + tool_progress_callback=_on_reasoning, tool_start_callback=_on_tool_start, tool_complete_callback=_on_tool_complete, agent_ref=agent_ref, @@ -3949,6 +3951,12 @@ def _on_tool_complete(tool_call_id, function_name, function_args, function_resul ) # Non-streaming: run the agent (with optional Idempotency-Key) + _reasoning_parts: List[str] = [] + + def _on_reasoning(event_type, tool_name=None, preview=None, args=None, **kwargs): + if event_type == "reasoning.available" and preview: + _reasoning_parts.append(preview) + async def _compute_completion(): return await self._run_agent( user_message=user_message, @@ -3956,6 +3964,7 @@ async def _compute_completion(): ephemeral_system_prompt=system_prompt, session_id=session_id, gateway_session_key=gateway_session_key, + tool_progress_callback=_on_reasoning, **agent_overrides, route=route, ) @@ -4049,6 +4058,9 @@ async def _compute_completion(): "total_tokens": usage.get("total_tokens", 0), }, } + reasoning_text = "\n\n".join(_reasoning_parts).strip() + if reasoning_text: + response_data["choices"][0]["message"]["reasoning_content"] = reasoning_text if is_partial or is_failed or not completed: response_data["hermes"] = { "completed": completed, @@ -4117,12 +4129,24 @@ async def _emit(item): frontends can display them without storing the markers in conversation history. See #6972 for the original event, #16588 for the ``toolCallId``/``status`` lifecycle fields. + Tagged tuples ``("__reasoning__", text)`` are sent as a + standard ``delta.reasoning_content`` chunk — the same + field name DeepSeek/Moonshot/OpenRouter use for thinking + traces — so OpenAI-SDK clients can opt into rendering the + model's reasoning without it leaking into ``content``. """ if isinstance(item, tuple) and len(item) == 2 and item[0] == "__tool_progress__": event_data = json.dumps(item[1]) await response.write( f"event: hermes.tool.progress\ndata: {event_data}\n\n".encode() ) + elif isinstance(item, tuple) and len(item) == 2 and item[0] == "__reasoning__": + reasoning_chunk = { + "id": completion_id, "object": "chat.completion.chunk", + "created": created, "model": model, + "choices": [{"index": 0, "delta": {"reasoning_content": item[1]}, "finish_reason": None}], + } + await response.write(f"data: {json.dumps(reasoning_chunk)}\n\n".encode()) else: content_chunk = { "id": completion_id, "object": "chat.completion.chunk", diff --git a/tests/gateway/test_api_server.py b/tests/gateway/test_api_server.py index 86c3b03a77d06..9092af4490022 100644 --- a/tests/gateway/test_api_server.py +++ b/tests/gateway/test_api_server.py @@ -1051,6 +1051,91 @@ async def _mock_run_agent(**kwargs): assert '"status": "running"' not in body assert '"status": "completed"' not in body + @pytest.mark.asyncio + async def test_stream_includes_reasoning_content(self, adapter): + """``reasoning.available`` previews surface as ``delta.reasoning_content`` + chunks, not mixed into ``delta.content`` — mirrors the DeepSeek/ + Moonshot/OpenRouter ``reasoning_content`` field so OpenAI-SDK + clients can opt into rendering the model's reasoning.""" + import asyncio + import json as _json + + app = _create_app(adapter) + async with TestClient(TestServer(app)) as cli: + async def _mock_run_agent(**kwargs): + cb = kwargs.get("stream_delta_callback") + reasoning_cb = kwargs.get("tool_progress_callback") + if reasoning_cb: + reasoning_cb("reasoning.available", "_thinking", "Let me check the calendar.", None) + if cb: + await asyncio.sleep(0.05) + cb("Here's your schedule.") + return ( + {"final_response": "Here's your schedule.", "messages": [], "api_calls": 1}, + {"input_tokens": 10, "output_tokens": 5, "total_tokens": 15}, + ) + + with patch.object(adapter, "_run_agent", side_effect=_mock_run_agent): + resp = await cli.post( + "/v1/chat/completions", + json={ + "model": "test", + "messages": [{"role": "user", "content": "what's on my calendar"}], + "stream": True, + }, + ) + assert resp.status == 200 + body = await resp.text() + + reasoning_chunks = [] + for line in body.splitlines(): + if line.startswith("data: ") and line.strip() != "data: [DONE]": + try: + chunk = _json.loads(line[len("data: "):]) + except _json.JSONDecodeError: + continue + for choice in chunk.get("choices", []): + delta = choice.get("delta", {}) + if "reasoning_content" in delta: + reasoning_chunks.append(delta["reasoning_content"]) + # Reasoning text must never leak into delta.content. + assert "Let me check the calendar." not in delta.get("content", "") + + assert reasoning_chunks == ["Let me check the calendar."] + assert "Here's your schedule." in body + + @pytest.mark.asyncio + async def test_non_streaming_includes_reasoning_content(self, adapter): + """Non-streaming ``/v1/chat/completions`` surfaces the model's + reasoning as ``message.reasoning_content``, matching the + streaming behaviour above.""" + app = _create_app(adapter) + + async def _mock_run_agent(**kwargs): + reasoning_cb = kwargs.get("tool_progress_callback") + if reasoning_cb: + reasoning_cb("reasoning.available", "_thinking", "Checking inbox for follow-ups.", None) + return ( + {"final_response": "You have two emails needing a reply.", "messages": [], "api_calls": 1}, + {"input_tokens": 10, "output_tokens": 5, "total_tokens": 15}, + ) + + async with TestClient(TestServer(app)) as cli: + with patch.object(adapter, "_run_agent", side_effect=_mock_run_agent): + resp = await cli.post( + "/v1/chat/completions", + json={ + "model": "test", + "messages": [{"role": "user", "content": "check my email"}], + }, + ) + assert resp.status == 200 + data = await resp.json() + + message = data["choices"][0]["message"] + assert message["reasoning_content"] == "Checking inbox for follow-ups." + assert message["content"] == "You have two emails needing a reply." + # --------------------------------------------------------------------------- # _derive_chat_session_id unit tests