diff --git a/gateway/platforms/api_server.py b/gateway/platforms/api_server.py index f02af4bacc..8183521d86 100644 --- a/gateway/platforms/api_server.py +++ b/gateway/platforms/api_server.py @@ -116,6 +116,7 @@ except ImportError: web = None # type: ignore[assignment] from gateway.config import Platform, PlatformConfig +from gateway.display_config import resolve_display_setting from gateway.platforms import api_server_room_dispatch as _room_dispatch from gateway.platforms import api_server_room_grants as _room_grants from gateway.platforms import api_server_runs as _api_runs @@ -2181,6 +2182,10 @@ class APIServerAdapter(OpenAICompatRoutesMixin, BasePlatformAdapter): gateway_session_key=gateway_session_key, session_id=session_id) user_config = _load_gateway_config() enabled_toolsets = sorted(_get_platform_tools(user_config, "api_server")) + # Same gate the messaging gateway and TUI apply: ``display.interim_assistant_messages`` + # off means no callback is installed, so mid-turn commentary never leaves the agent. + if not resolve_display_setting(user_config, "api_server", "interim_assistant_messages", True): + interim_assistant_callback = None max_iterations = _current_max_iterations() if room_dispatch is not None: from gateway.hosted_room_execution_policy import RoomExecutionPolicy diff --git a/gateway/platforms/api_server_runs.py b/gateway/platforms/api_server_runs.py index ea7efb7f1a..ae290cf44f 100644 --- a/gateway/platforms/api_server_runs.py +++ b/gateway/platforms/api_server_runs.py @@ -728,6 +728,16 @@ async def _execute_run(self, run: _RunLaunch, *, _api_server) -> None: with suppress(Exception): loop.call_soon_threadsafe(run.put_event, _run_event(run_id, "message.delta", delta=delta)) + def _interim_cb(text: str, *, already_streamed: bool = False) -> None: + # Mid-turn assistant commentary (Codex ``phase="commentary"``, text beside tool calls), + # same ``message.interim`` contract as the TUI gateway; reasoning never reaches this + # callback and the final answer still arrives via ``run.completed`` (#67580). + if not isinstance(text, str) or not text.strip() or run_id not in self._run_streams: + return + with suppress(Exception): + loop.call_soon_threadsafe(run.put_event, _run_event( + run_id, "message.interim", text=text, already_streamed=bool(already_streamed))) + def _finish(status: str, extra: Optional[dict] = None, **fields: Any) -> None: """Terminal status, then best-effort ``run.`` event; key order is wire shape.""" extra = extra or {} @@ -752,7 +762,7 @@ async def _execute_run(self, run: _RunLaunch, *, _api_server) -> None: with self._profile_scope(run.request_profile): agent = self._create_agent( stream_delta_callback=_text_cb, tool_progress_callback=self._make_run_event_callback(run_id, loop), - **run.agent_kwargs) + interim_assistant_callback=_interim_cb, **run.agent_kwargs) self._active_run_agents[run_id] = agent approval_notify = _make_approval_notify(self, run, _api_server=_api_server) result, usage = await loop.run_in_executor( diff --git a/tests/gateway/test_api_server_runs.py b/tests/gateway/test_api_server_runs.py index 279cab2d5b..a85b235e0f 100644 --- a/tests/gateway/test_api_server_runs.py +++ b/tests/gateway/test_api_server_runs.py @@ -337,6 +337,39 @@ class TestStartRun: assert adapter._run_statuses == {} mock_create.assert_not_called() + @pytest.mark.asyncio + async def test_events_stream_forwards_interim_commentary(self, adapter): + """Mid-turn assistant commentary (Codex ``phase="commentary"``) reaches /v1/runs clients + as ``message.interim`` {text, already_streamed}; the final answer is unchanged (#67580).""" + import json + + app = _create_runs_app(adapter) + + def create_agent(**kwargs): + interim = kwargs["interim_assistant_callback"] + agent = MagicMock() + + def run_conversation(**_kw): + interim("Checking the docs first.", already_streamed=False) + interim("Applying the fix.", already_streamed=True) + return {"final_response": "Done."} + + agent.run_conversation.side_effect = run_conversation + agent.session_prompt_tokens = agent.session_completion_tokens = agent.session_total_tokens = 0 + return agent + + async with TestClient(TestServer(app)) as cli: + with patch.object(adapter, "_create_agent", side_effect=create_agent): + resp = await cli.post("/v1/runs", json={"input": "hello"}) + run_id = (await resp.json())["run_id"] + body = await (await cli.get(f"/v1/runs/{run_id}/events")).text() + + events = [json.loads(line[6:]) for line in body.splitlines() if line.startswith("data: ")] + interim = [(e["text"], e["already_streamed"]) for e in events if e["event"] == "message.interim"] + assert interim == [("Checking the docs first.", False), ("Applying the fix.", True)] + completed = next(e for e in events if e["event"] == "run.completed") + assert completed["output"] == "Done." + @pytest.mark.asyncio async def test_start_passes_request_model_provider_options_to_create_agent(self, adapter): app = _create_runs_app(adapter) diff --git a/tests/gateway/test_session_api.py b/tests/gateway/test_session_api.py index 1973d947e0..e97b128de8 100644 --- a/tests/gateway/test_session_api.py +++ b/tests/gateway/test_session_api.py @@ -1155,3 +1155,69 @@ async def test_session_stream_records_reply_text_for_post_disconnect_recovery( response = await adapter._handle_get_run(get_request) assert response.status == 200 assert "the answer worth keeping" in response.text + + +@pytest.mark.asyncio +async def test_interim_commentary_reaches_session_sse_and_responses_stream(adapter, session_db, monkeypatch): + """Codex commentary / mid-turn assistant text is a typed ``assistant.commentary`` event on the + session SSE endpoint and a ``phase: commentary`` message item on /v1/responses, never part of + the final answer; ``display.interim_assistant_messages: false`` installs no callback (#67580).""" + import json as _json + + session_id = session_db.create_session("commentary-session", "api_server") + + def fake_create_agent(**kwargs): + interim = kwargs["interim_assistant_callback"] + + class FakeAgent: + provider, model = "openai-codex", "gpt-5" + session_prompt_tokens = session_completion_tokens = session_total_tokens = 0 + + def run_conversation(self, **_kw): + interim("Checking the docs first.", already_streamed=False) + return {"final_response": "Done.", "messages": [], "api_calls": 1} + + return FakeAgent() + + app = _create_session_app(adapter) + app.router.add_post("/v1/responses", adapter._handle_responses) + with patch.object(adapter, "_create_agent", side_effect=fake_create_agent): + async with TestClient(TestServer(app)) as cli: + sse = await (await cli.post(f"/api/sessions/{session_id}/chat/stream", json={"message": "go"})).text() + responses = await (await cli.post( + "/v1/responses", json={"model": "hermes-agent", "input": "go", "stream": True})).text() + + def _events(body): + out = [] + for block in body.split("\n\n"): + lines = block.splitlines() + name = next((ln[7:] for ln in lines if ln.startswith("event: ")), None) + data = next((ln[6:] for ln in lines if ln.startswith("data: ")), None) + if data: + out.append((name, _json.loads(data))) + return out + + sse_events = _events(sse) + commentary = [d for n, d in sse_events if n == "assistant.commentary"] + assert [(d["text"], d["already_streamed"]) for d in commentary] == [("Checking the docs first.", False)] + assert next(d for n, d in sse_events if n == "assistant.completed")["content"] == "Done." + + done_items = [d["item"] for n, d in _events(responses) if n == "response.output_item.done"] + assert [(i.get("phase"), i["content"][0]["text"]) for i in done_items if i["type"] == "message"] == [ + ("commentary", "Checking the docs first."), (None, "Done.")] + + # Display gate: the callback is dropped before it reaches AIAgent, like the gateway/TUI. + _patch_api_server_runtime(monkeypatch) + monkeypatch.setattr( + "gateway.run._load_gateway_config", lambda: {"display": {"interim_assistant_messages": False}}) + captured = {} + + class CapturingAgent: + provider, model = "openrouter", "global/model" + + def __init__(self, **kwargs): + captured.update(kwargs) + + monkeypatch.setattr("run_agent.AIAgent", CapturingAgent) + adapter._create_agent(session_id="gated", interim_assistant_callback=lambda *_a, **_k: None) + assert captured["interim_assistant_callback"] is None diff --git a/website/docs/user-guide/features/api-server.md b/website/docs/user-guide/features/api-server.md index 9e8fc51379..752e6e732b 100644 --- a/website/docs/user-guide/features/api-server.md +++ b/website/docs/user-guide/features/api-server.md @@ -146,6 +146,8 @@ OpenAI Responses API format. Supports server-side conversation state via `previo Tool calls in the `output` array were already executed server-side by the Hermes agent — they are replayed with `"status": "completed"` for structured tool UI, never as pending calls for the client to execute. +With `"stream": true`, mid-turn assistant commentary (the `openai-codex` backend's `phase="commentary"` progress preambles, or text a model writes alongside its tool calls) arrives as its own completed `message` output item carrying `"phase": "commentary"` (`response.output_item.added` + `response.output_item.done`, also listed in `response.completed`). It is never merged into the final answer item, so clients can render it as live progress and skip it when assembling the reply. Private reasoning never reaches this item. `display.interim_assistant_messages: false` (or the `display.platforms.api_server` override) suppresses it on every API-server surface. + **Inline image input:** `input[].content` can contain `input_text` and `input_image` parts. Both remote URLs and `data:image/...` URLs are supported: ```json @@ -475,6 +477,13 @@ Statuses are retained briefly after terminal states (`completed`, `failed`, `can Server-Sent Events stream of the run's tool-call progress, token deltas, and lifecycle events. Designed for dashboards and thick clients that want to attach/detach without losing state. +Mid-turn assistant commentary — the `openai-codex` backend's `phase="commentary"` progress +preambles, or text a model writes alongside its tool calls — arrives as `message.interim` +(`text`, `already_streamed`), the same contract the TUI gateway uses. `already_streamed: true` +means the text also went out as `message.delta`, so clients that render deltas can skip it. +The final answer still arrives only in `run.completed`; private reasoning never becomes +`message.interim`. Gate: `display.interim_assistant_messages` (default `true`). + Tool lifecycle events carry `tool.started` (`tool`, `preview` of the arguments) and `tool.completed` (`tool`, `duration` in seconds, `error`, and a `preview` of the result). The `error` flag reflects the tool's own outcome — a non-zero terminal `exit_code`, a structured @@ -590,7 +599,7 @@ External UIs can manage Hermes sessions over REST without standing up the dashbo | `GET` | `/api/sessions/{id}/messages` | Message history for a session | | `POST` | `/api/sessions/{id}/fork` | Branch the session via `SessionDB` lineage (matches CLI `/branch` semantics) | | `POST` | `/api/sessions/{id}/chat` | Run one synchronous agent turn | -| `POST` | `/api/sessions/{id}/chat/stream` | SSE wrapper over a single turn — emits `assistant.delta`, `tool.started`, `tool.completed`, then a terminal `run.completed` / `run.failed` / `run.cancelled` event that matches how the turn ended (see [Terminal run status](../../developer-guide/programmatic-integration.md#terminal-run-status)) | +| `POST` | `/api/sessions/{id}/chat/stream` | SSE wrapper over a single turn — emits `assistant.delta`, `assistant.commentary` (mid-turn commentary: `message_id`, `text`, `already_streamed`; never folded into `assistant.completed`), `tool.started`, `tool.completed`, then a terminal `run.completed` / `run.failed` / `run.cancelled` event that matches how the turn ended (see [Terminal run status](../../developer-guide/programmatic-integration.md#terminal-run-status)) | `/v1/capabilities` advertises the full surface via `session_*` feature flags and `endpoints.session_*` entries so external UIs can detect support and fall back safely. Inline images are supported in `chat` and `chat/stream` payloads (multimodal-aware path).