Coverage for agentos/core/streaming.py: 62%

56 statements  

« prev     ^ index     » next       coverage.py v7.14.3, created at 2026-07-10 07:44 +0800

1""" 

2AgentOS v0.20 流式输出系统。 

3支持 SSE (Server-Sent Events) 格式流式传输。 

4""" 

5 

6from __future__ import annotations 

7 

8from dataclasses import dataclass, field 

9from enum import StrEnum 

10from typing import Any 

11 

12 

13class StreamEvent(StrEnum): 

14 """流式事件。""" 

15 

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" 

25 

26 

27@dataclass 

28class StreamChunk: 

29 """流式输出的单块数据。""" 

30 

31 event: StreamEvent 

32 data: dict[str, Any] = field(default_factory=dict) 

33 timestamp: float = 0.0 # auto-filled by emitter 

34 

35 def to_sse(self) -> str: 

36 import json 

37 

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" 

44 

45 @property 

46 def is_terminal(self) -> bool: 

47 return self.event in (StreamEvent.COMPLETE, StreamEvent.ERROR, StreamEvent.CANCELLED) 

48 

49 

50class StreamEmitter: 

51 """异步SSE发射器。""" 

52 

53 def __init__(self): 

54 import time 

55 

56 self._start = time.time() 

57 

58 def emit(self, event: StreamEvent, **data) -> StreamChunk: 

59 import time 

60 

61 chunk = StreamChunk(event=event, data=data) 

62 chunk.timestamp = (time.time() - self._start) * 1000 

63 return chunk 

64 

65 def thinking(self, text: str) -> StreamChunk: 

66 return self.emit(StreamEvent.THINKING, text=text) 

67 

68 def text(self, text: str) -> StreamChunk: 

69 return self.emit(StreamEvent.TEXT, text=text) 

70 

71 def tool_call(self, name: str, args: dict) -> StreamChunk: 

72 return self.emit(StreamEvent.TOOL_CALL, name=name, arguments=args) 

73 

74 def tool_result(self, name: str, result: str) -> StreamChunk: 

75 return self.emit(StreamEvent.TOOL_RESULT, name=name, result=result) 

76 

77 def error(self, message: str) -> StreamChunk: 

78 return self.emit(StreamEvent.ERROR, error=message) 

79 

80 

81class ResponseCollector: 

82 """收集流式chunk并拼接为最终响应。""" 

83 

84 def __init__(self): 

85 self.chunks: list[StreamChunk] = [] 

86 self._text_buf: list[str] = [] 

87 

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", "")) 

92 

93 @property 

94 def full_text(self) -> str: 

95 return "".join(self._text_buf)