fix(api-server): /v1/runs emits message.interim; display.interim_assistant_messages gates all surfaces
POST /v1/runs + GET /v1/runs/{id}/events now carries mid-turn assistant
commentary as `message.interim` {text, already_streamed} — the same
contract the TUI gateway emits — so Runs clients can tell an active,
tool-heavy Codex turn from a stalled one instead of seeing only tool.*
events until run.completed (#67580).
APIServerAdapter._create_agent applies the same
`display.interim_assistant_messages` gate the messaging gateway and the
TUI apply (resolve_display_setting, per-platform override honoured):
when off, no callback is installed and nothing leaves the agent on any
of the three streaming surfaces. Dedup of repeated commentary stays in
agent/stream_delivery.py, where every surface already relies on it.
Tests: /v1/runs event contract; session SSE + /v1/responses item shape
plus the display gate (both red on the base commit). Docs: event
payloads on all three surfaces.
Co-authored-by: RoySRose <sungwook0115.kim@gmail.com>
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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.<status>`` 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(
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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).
|
||||
|
||||
|
||||
Reference in New Issue
Block a user