Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
16 changes: 16 additions & 0 deletions gateway/platforms/api_server.py
Original file line number Diff line number Diff line change
Expand Up @@ -1707,6 +1707,7 @@ def _create_agent(
ephemeral_system_prompt: Optional[str] = None,
session_id: Optional[str] = None,
stream_delta_callback=None,
reasoning_delta_callback=None,
tool_progress_callback=None,
tool_start_callback=None,
tool_complete_callback=None,
Expand Down Expand Up @@ -1824,6 +1825,7 @@ def _create_agent(
session_id=session_id,
platform="api_server",
stream_delta_callback=stream_delta_callback,
reasoning_callback=reasoning_delta_callback,
tool_progress_callback=tool_progress_callback,
tool_start_callback=tool_start_callback,
tool_complete_callback=tool_complete_callback,
Expand Down Expand Up @@ -4786,6 +4788,19 @@ def _text_cb(delta: Optional[str]) -> None:
except Exception:
pass

def _reasoning_cb(delta: Optional[str]) -> None:
if not delta or run_id not in self._run_streams:
return
try:
loop.call_soon_threadsafe(_put_event_if_active, {
"event": "reasoning.delta",
"run_id": run_id,
"timestamp": time.time(),
"delta": delta,
})
except Exception:
pass

self._set_run_status(
run_id,
"queued",
Expand Down Expand Up @@ -4820,6 +4835,7 @@ async def _run_and_close():
ephemeral_system_prompt=ephemeral_system_prompt,
session_id=session_id,
stream_delta_callback=_text_cb,
reasoning_delta_callback=_reasoning_cb,
tool_progress_callback=event_cb,
gateway_session_key=gateway_session_key,
route=route,
Expand Down
40 changes: 40 additions & 0 deletions tests/gateway/test_api_server_runs.py
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@
"""

import asyncio
import json
import threading
import time
from unittest.mock import MagicMock, patch
Expand Down Expand Up @@ -297,6 +298,45 @@ async def test_status_requires_auth(self, auth_adapter):


class TestRunEvents:
@pytest.mark.asyncio
async def test_events_streams_reasoning_deltas_before_completion(self, adapter):
app = _create_runs_app(adapter)
async with TestClient(TestServer(app)) as cli:
def create_agent(**kwargs):
mock_agent = MagicMock()

def run_conversation(**_run_kwargs):
kwargs["reasoning_delta_callback"]("先分析需求")
kwargs["reasoning_delta_callback"](",再调用工具")
return {"final_response": "完成"}

mock_agent.run_conversation.side_effect = run_conversation
mock_agent.session_prompt_tokens = 10
mock_agent.session_completion_tokens = 5
mock_agent.session_total_tokens = 15
return mock_agent

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"]
events_resp = await cli.get(f"/v1/runs/{run_id}/events")
body = await events_resp.text()

events = [
json.loads(line.removeprefix("data: "))
for line in body.splitlines()
if line.startswith("data: ")
]
assert [event["event"] for event in events] == [
"reasoning.delta",
"reasoning.delta",
"run.completed",
]
assert [event["delta"] for event in events[:2]] == [
"先分析需求",
",再调用工具",
]

@pytest.mark.asyncio
async def test_events_stream_returns_completed(self, adapter):
"""Events stream should receive run.completed when agent finishes."""
Expand Down