SSE 增量解析器。 核心思路:维护一个字符串缓冲区。 每次收到新数据就追加到缓冲区,然后尝试从中提取完整的事件。 对应源码: api/sse.rs:5-61 (SseParser) runtime/sse.rs:12-97 (IncrementalSseParser)
| 127 | # 4. 没找到 → 继续等下一块数据 |
| 128 | |
| 129 | class SseParser: |
| 130 | """ |
| 131 | SSE 增量解析器。 |
| 132 | |
| 133 | 核心思路:维护一个字符串缓冲区。 |
| 134 | 每次收到新数据就追加到缓冲区,然后尝试从中提取完整的事件。 |
| 135 | |
| 136 | 对应源码: api/sse.rs:5-61 (SseParser) |
| 137 | runtime/sse.rs:12-97 (IncrementalSseParser) |
| 138 | """ |
| 139 | |
| 140 | def __init__(self): |
| 141 | self.buffer = "" # 缓冲区:存放还没解析完的数据 |
| 142 | self._event_name = None # 当前正在解析的事件名 |
| 143 | self._data_lines = [] # 当前事件的 data 行 |
| 144 | self._id = None # 当前事件的 id |
| 145 | self._retry = None # 当前事件的 retry |
| 146 | |
| 147 | def push_chunk(self, chunk: str) -> list[SseEvent]: |
| 148 | """ |
| 149 | 推入一块新收到的数据,返回解析出的完整事件列表。 |
| 150 | |
| 151 | 可能返回 0 个事件(数据不够完整) |
| 152 | 可能返回 1 个事件(刚好凑齐一个) |
| 153 | 可能返回多个事件(一块数据里包含了好几个事件) |
| 154 | |
| 155 | 对应源码: runtime/sse.rs:26-39 (push_chunk) |
| 156 | """ |
| 157 | self.buffer += chunk |
| 158 | events = [] |
| 159 | |
| 160 | # 逐行处理缓冲区 |
| 161 | while "\n" in self.buffer: |
| 162 | # 找到第一个换行符的位置 |
| 163 | index = self.buffer.index("\n") |
| 164 | # 提取这一行(不包含换行符) |
| 165 | line = self.buffer[:index] |
| 166 | # 从缓冲区移除已处理的部分(包含换行符) |
| 167 | self.buffer = self.buffer[index + 1:] |
| 168 | # 处理 \r\n 的情况(Windows 风格换行) |
| 169 | line = line.rstrip("\r") |
| 170 | # 处理这一行 |
| 171 | self._process_line(line, events) |
| 172 | |
| 173 | return events |
| 174 | |
| 175 | def finish(self) -> list[SseEvent]: |
| 176 | """ |
| 177 | 流结束时调用。处理缓冲区中剩余的数据。 |
| 178 | |
| 179 | 对应源码: runtime/sse.rs:44-53 (finish) |
| 180 | """ |
| 181 | events = [] |
| 182 | if self.buffer: |
| 183 | line = self.buffer.rstrip("\r") |
| 184 | self.buffer = "" |
| 185 | self._process_line(line, events) |
| 186 | # 如果还有未完成的事件,提取出来 |
no outgoing calls
no test coverage detected