fix: cancel the agent run when a streaming query is abandoned

The Agents SDK starts the run eagerly; breaking out of the stream (or a
client disconnect in a server handler) left the full agent loop running
to completion in the background — further LLM calls, billed, with no way
to stop them. The cloud query_stream already handled this; the local
path did not.
This commit is contained in:
Ray 2026-07-21 17:57:24 +08:00
parent 7805169bff
commit 9fc7b680e2

View file

@ -114,23 +114,29 @@ class QueryStream:
from openai.types.responses import ResponseTextDeltaEvent
streamed_run = Runner.run_streamed(self._agent, self._question)
async for event in streamed_run.stream_events():
if isinstance(event, RawResponsesStreamEvent):
if isinstance(event.data, ResponseTextDeltaEvent):
yield QueryEvent(type="text_delta", data=event.data.delta)
elif isinstance(event, RunItemStreamEvent):
item = event.item
if item.type == "tool_call_item":
raw = item.raw_item
yield QueryEvent(type="tool_call", data={
"name": raw.name, "args": getattr(raw, "arguments", "{}"),
})
elif item.type == "tool_call_output_item":
yield QueryEvent(type="tool_result", data=str(item.output))
elif item.type == "message_output_item":
text = ItemHelpers.text_message_output(item)
if text:
yield QueryEvent(type="text_done", data=text)
# cancel() in finally: the SDK starts the run eagerly, and abandoning
# this generator (consumer breaks / client disconnects) would otherwise
# leave the agent loop running — and billing LLM calls — to completion.
try:
async for event in streamed_run.stream_events():
if isinstance(event, RawResponsesStreamEvent):
if isinstance(event.data, ResponseTextDeltaEvent):
yield QueryEvent(type="text_delta", data=event.data.delta)
elif isinstance(event, RunItemStreamEvent):
item = event.item
if item.type == "tool_call_item":
raw = item.raw_item
yield QueryEvent(type="tool_call", data={
"name": raw.name, "args": getattr(raw, "arguments", "{}"),
})
elif item.type == "tool_call_output_item":
yield QueryEvent(type="tool_result", data=str(item.output))
elif item.type == "message_output_item":
text = ItemHelpers.text_message_output(item)
if text:
yield QueryEvent(type="text_done", data=text)
finally:
streamed_run.cancel()
def __aiter__(self):
return self.stream_events()