feat(api_server): stream agent status lines as hermes.status SSE events
The OpenAI-compatible SSE writers carried tool progress and reasoning but no
lifecycle status, so an API client waiting through a provider outage (now the
auto-recovery ladder) saw a silent socket with no way to tell "waiting on the
provider" from "hung". _spawn_stream_agent wires status_callback into the
agent and both writers (/v1/chat/completions, /v1/responses) emit
`event: hermes.status` with {kind, text}, redacted like every other API-bound
error text. The Responses writer's tag dispatch moves to a table so the new tag
does not grow an if/elif ladder.
This commit is contained in:
@@ -2167,7 +2167,7 @@ class APIServerAdapter(OpenAICompatRoutesMixin, BasePlatformAdapter):
|
||||
self, ephemeral_system_prompt: Optional[str] = None, session_id: Optional[str] = None,
|
||||
stream_delta_callback=None, tool_progress_callback=None, tool_start_callback=None,
|
||||
tool_complete_callback=None, interim_assistant_callback=None, reasoning_callback=None,
|
||||
gateway_session_key: Optional[str] = None,
|
||||
status_callback=None, gateway_session_key: Optional[str] = None,
|
||||
requested_model: Optional[str] = None, requested_provider: Optional[str] = None,
|
||||
model_options: Optional[Dict[str, Any]] = None, route: Optional[Dict[str, Any]] = None,
|
||||
session_model: Optional[str] = None, confirmed_runtime_lock: bool = False,
|
||||
@@ -2227,6 +2227,7 @@ class APIServerAdapter(OpenAICompatRoutesMixin, BasePlatformAdapter):
|
||||
"tool_complete_callback": tool_complete_callback,
|
||||
"interim_assistant_callback": interim_assistant_callback,
|
||||
"reasoning_callback": reasoning_callback,
|
||||
"status_callback": status_callback,
|
||||
"session_db": self._ensure_session_db(),
|
||||
# Same fallback provider chain as Telegram/Discord/Slack.
|
||||
"fallback_model": None if confirmed_runtime_lock else GatewayRunner._load_fallback_model(),
|
||||
@@ -3767,7 +3768,7 @@ class APIServerAdapter(OpenAICompatRoutesMixin, BasePlatformAdapter):
|
||||
ephemeral_system_prompt: Optional[str] = None, session_id: Optional[str] = None,
|
||||
stream_delta_callback=None, tool_progress_callback=None, tool_start_callback=None,
|
||||
tool_complete_callback=None, interim_assistant_callback=None, reasoning_callback=None,
|
||||
agent_ref: Optional[list] = None, active_run_id: Optional[str] = None,
|
||||
status_callback=None, agent_ref: Optional[list] = None, active_run_id: Optional[str] = None,
|
||||
gateway_session_key: Optional[str] = None, requested_model: Optional[str] = None,
|
||||
requested_provider: Optional[str] = None, model_options: Optional[Dict[str, Any]] = None,
|
||||
route: Optional[Dict[str, Any]] = None, session_model: Optional[str] = None,
|
||||
@@ -3808,7 +3809,7 @@ class APIServerAdapter(OpenAICompatRoutesMixin, BasePlatformAdapter):
|
||||
stream_delta_callback=stream_delta_callback, tool_progress_callback=tool_progress_callback,
|
||||
tool_start_callback=tool_start_callback, tool_complete_callback=tool_complete_callback,
|
||||
interim_assistant_callback=interim_assistant_callback,
|
||||
reasoning_callback=reasoning_callback,
|
||||
reasoning_callback=reasoning_callback, status_callback=status_callback,
|
||||
gateway_session_key=gateway_session_key, requested_model=requested_model,
|
||||
requested_provider=requested_provider, model_options=model_options, route=route,
|
||||
session_model=session_model, confirmed_runtime_lock=confirmed_runtime_lock)
|
||||
|
||||
@@ -416,19 +416,29 @@ class _ResponsesStream:
|
||||
for event in ("response.output_item.added", "response.output_item.done"):
|
||||
await self.write_event(event, {"type": event, "output_index": idx, "item": output_item})
|
||||
|
||||
async def emit_status(self, payload: Dict[str, Any]) -> None:
|
||||
"""Lifecycle/warning status (provider wait, auto-recovery countdown, fallback switch) as a
|
||||
``hermes.status`` custom event; not a Responses output item."""
|
||||
await self.response.write(self._api._sse_frame(payload, event="hermes.status"))
|
||||
|
||||
# queue tag -> (method name, payload adapter)
|
||||
_TAG_HANDLERS = {
|
||||
"__tool_started__": ("emit_tool_started", lambda p: p),
|
||||
"__tool_completed__": ("emit_tool_completed", lambda p: p),
|
||||
"__commentary__": ("emit_commentary", lambda p: p["text"]),
|
||||
"__reasoning__": ("emit_reasoning_delta", lambda p: p),
|
||||
"__status__": ("emit_status", lambda p: p),
|
||||
}
|
||||
|
||||
async def dispatch(self, item: Any) -> None:
|
||||
"""Route one queue item: tool tuples emit immediately, strings are batched, others dropped."""
|
||||
"""Route one queue item: tagged tuples emit immediately, strings are batched, others dropped."""
|
||||
if isinstance(item, tuple) and len(item) == 2 and isinstance(item[0], str):
|
||||
tag, payload = item
|
||||
await self.flush_batch()
|
||||
if tag == "__tool_started__":
|
||||
await self.emit_tool_started(payload)
|
||||
elif tag == "__tool_completed__":
|
||||
await self.emit_tool_completed(payload)
|
||||
elif tag == "__commentary__":
|
||||
await self.emit_commentary(payload["text"])
|
||||
elif tag == "__reasoning__":
|
||||
await self.emit_reasoning_delta(payload)
|
||||
handler = self._TAG_HANDLERS.get(tag)
|
||||
if handler is not None:
|
||||
method, adapt = handler
|
||||
await getattr(self, method)(adapt(payload))
|
||||
elif isinstance(item, str):
|
||||
self._batch_buf.append(item)
|
||||
if self._batch_timer is None:
|
||||
@@ -563,10 +573,17 @@ class OpenAICompatRoutesMixin:
|
||||
# keep them distinct from answer text.
|
||||
if text:
|
||||
stream_q.put_threadsafe(("__reasoning__", text))
|
||||
def _on_status(kind, message=None):
|
||||
# Lifecycle/warning status (provider wait, auto-recovery countdown, fallback switch) as a
|
||||
# ``hermes.status`` event, so a client sees why the stream is silent instead of a dead socket.
|
||||
from gateway.platforms.api_server import _redact_api_error_text
|
||||
text = _redact_api_error_text(message if message is not None else kind or "").strip()
|
||||
if text:
|
||||
stream_q.put_threadsafe(("__status__", {"kind": str(kind), "text": text}))
|
||||
agent_ref = [None]
|
||||
agent_task = asyncio.ensure_future(self._run_agent(
|
||||
stream_delta_callback=_on_delta, reasoning_callback=_on_reasoning, agent_ref=agent_ref,
|
||||
**run_kwargs))
|
||||
stream_delta_callback=_on_delta, reasoning_callback=_on_reasoning, status_callback=_on_status,
|
||||
agent_ref=agent_ref, **run_kwargs))
|
||||
agent_task.add_done_callback(lambda _fut: stream_q.put_nowait(None))
|
||||
return agent_task, agent_ref
|
||||
|
||||
@@ -845,6 +862,8 @@ class OpenAICompatRoutesMixin:
|
||||
# DeepSeek-style ``delta.reasoning_content`` (#99552), the field Open WebUI,
|
||||
# opencode and the Vercel AI SDK render as a thinking block.
|
||||
await response.write(_sse_frame(_chunk({"reasoning_content": delta[1]})))
|
||||
elif isinstance(delta, tuple) and len(delta) == 2 and delta[0] == "__status__":
|
||||
await response.write(_sse_frame(delta[1], event="hermes.status"))
|
||||
else:
|
||||
await response.write(_sse_frame(_chunk({"content": delta})))
|
||||
# The agent can fail after the queue drains (task raises / result flagged failed or
|
||||
|
||||
Reference in New Issue
Block a user