Coverage for agentos/core/streaming.py: 62%
56 statements
« prev ^ index » next coverage.py v7.14.3, created at 2026-07-06 23:17 +0800
« prev ^ index » next coverage.py v7.14.3, created at 2026-07-06 23:17 +0800
1"""
2AgentOS v0.20 流式输出系统。
3支持 SSE (Server-Sent Events) 格式流式传输。
4"""
6from __future__ import annotations
8from dataclasses import dataclass, field
9from enum import StrEnum
10from typing import Any
13class StreamEvent(StrEnum):
14 """流式事件。"""
16 START = "start"
17 STEP_START = "step_start"
18 THINKING = "thinking"
19 TOOL_CALL = "tool_call"
20 TOOL_RESULT = "tool_result"
21 TEXT = "text"
22 ERROR = "error"
23 COMPLETE = "complete"
24 CANCELLED = "cancelled"
27@dataclass
28class StreamChunk:
29 """流式输出的单块数据。"""
31 event: StreamEvent
32 data: dict[str, Any] = field(default_factory=dict)
33 timestamp: float = 0.0 # auto-filled by emitter
35 def to_sse(self) -> str:
36 import json
38 payload = {
39 "event": self.event.value,
40 **self.data,
41 "ts": self.timestamp,
42 }
43 return f"data: {json.dumps(payload, default=str)}\n\n"
45 @property
46 def is_terminal(self) -> bool:
47 return self.event in (StreamEvent.COMPLETE, StreamEvent.ERROR, StreamEvent.CANCELLED)
50class StreamEmitter:
51 """异步SSE发射器。"""
53 def __init__(self):
54 import time
56 self._start = time.time()
58 def emit(self, event: StreamEvent, **data) -> StreamChunk:
59 import time
61 chunk = StreamChunk(event=event, data=data)
62 chunk.timestamp = (time.time() - self._start) * 1000
63 return chunk
65 def thinking(self, text: str) -> StreamChunk:
66 return self.emit(StreamEvent.THINKING, text=text)
68 def text(self, text: str) -> StreamChunk:
69 return self.emit(StreamEvent.TEXT, text=text)
71 def tool_call(self, name: str, args: dict) -> StreamChunk:
72 return self.emit(StreamEvent.TOOL_CALL, name=name, arguments=args)
74 def tool_result(self, name: str, result: str) -> StreamChunk:
75 return self.emit(StreamEvent.TOOL_RESULT, name=name, result=result)
77 def error(self, message: str) -> StreamChunk:
78 return self.emit(StreamEvent.ERROR, error=message)
81class ResponseCollector:
82 """收集流式chunk并拼接为最终响应。"""
84 def __init__(self):
85 self.chunks: list[StreamChunk] = []
86 self._text_buf: list[str] = []
88 def feed(self, chunk: StreamChunk):
89 self.chunks.append(chunk)
90 if chunk.event == StreamEvent.TEXT:
91 self._text_buf.append(chunk.data.get("text", ""))
93 @property
94 def full_text(self) -> str:
95 return "".join(self._text_buf)