处理流式响应
(self, response: AsyncIterable[bytes])
| 125 | self.show_think = think |
| 126 | |
| 127 | async def process(self, response: AsyncIterable[bytes]) -> AsyncGenerator[str, None]: |
| 128 | """处理流式响应""" |
| 129 | try: |
| 130 | async for line in response: |
| 131 | if not line: |
| 132 | continue |
| 133 | try: |
| 134 | data = orjson.loads(line) |
| 135 | except orjson.JSONDecodeError: |
| 136 | continue |
| 137 | |
| 138 | resp = data.get("result", {}).get("response", {}) |
| 139 | |
| 140 | # 元数据 |
| 141 | if (llm := resp.get("llmInfo")) and not self.fingerprint: |
| 142 | self.fingerprint = llm.get("modelHash", "") |
| 143 | if rid := resp.get("responseId"): |
| 144 | self.response_id = rid |
| 145 | |
| 146 | # 首次发送 role |
| 147 | if not self.role_sent: |
| 148 | yield self._sse(role="assistant") |
| 149 | self.role_sent = True |
| 150 | |
| 151 | # 图像生成进度 |
| 152 | if img := resp.get("streamingImageGenerationResponse"): |
| 153 | if self.show_think: |
| 154 | if not self.think_opened: |
| 155 | yield self._sse("<think>\n") |
| 156 | self.think_opened = True |
| 157 | idx = img.get('imageIndex', 0) + 1 |
| 158 | progress = img.get('progress', 0) |
| 159 | yield self._sse(f"正在生成第{idx}张图片中,当前进度{progress}%\n") |
| 160 | continue |
| 161 | |
| 162 | # modelResponse |
| 163 | if mr := resp.get("modelResponse"): |
| 164 | if self.think_opened and self.show_think: |
| 165 | if msg := mr.get("message"): |
| 166 | yield self._sse(msg + "\n") |
| 167 | yield self._sse("</think>\n") |
| 168 | self.think_opened = False |
| 169 | |
| 170 | # 处理生成的图片 |
| 171 | for url in mr.get("generatedImageUrls", []): |
| 172 | parts = url.split("/") |
| 173 | img_id = parts[-2] if len(parts) >= 2 else "image" |
| 174 | |
| 175 | if self.image_format == "base64": |
| 176 | dl_service = self._get_dl() |
| 177 | base64_data = await dl_service.to_base64(url, self.token, "image") |
| 178 | if base64_data: |
| 179 | yield self._sse(f"\n") |
| 180 | else: |
| 181 | final_url = await self.process_url(url, "image") |
| 182 | yield self._sse(f"\n") |
| 183 | else: |
| 184 | final_url = await self.process_url(url, "image") |
no test coverage detected