parseSSEStream 解析 SSE 流
(body io.ReadCloser, chunks chan<- StreamChunk)
| 386 | |
| 387 | // parseSSEStream 解析 SSE 流 |
| 388 | func (p *GeminiProvider) parseSSEStream(body io.ReadCloser, chunks chan<- StreamChunk) { |
| 389 | defer func() { _ = body.Close() }() |
| 390 | defer close(chunks) |
| 391 | |
| 392 | scanner := bufio.NewScanner(body) |
| 393 | scanner.Split(bufio.ScanLines) |
| 394 | |
| 395 | for scanner.Scan() { |
| 396 | line := scanner.Text() |
| 397 | |
| 398 | // SSE 格式: "data: {json}" |
| 399 | if !strings.HasPrefix(line, "data: ") { |
| 400 | continue |
| 401 | } |
| 402 | |
| 403 | data := strings.TrimPrefix(line, "data: ") |
| 404 | data = strings.TrimSpace(data) |
| 405 | |
| 406 | // 解析 JSON |
| 407 | var chunk map[string]any |
| 408 | if err := json.Unmarshal([]byte(data), &chunk); err != nil { |
| 409 | continue |
| 410 | } |
| 411 | |
| 412 | // 解析 chunk 并转换为 StreamChunk |
| 413 | streamChunks := p.parseStreamChunk(chunk) |
| 414 | for _, sc := range streamChunks { |
| 415 | chunks <- sc |
| 416 | } |
| 417 | } |
| 418 | |
| 419 | if err := scanner.Err(); err != nil { |
| 420 | chunks <- StreamChunk{ |
| 421 | Type: string(ChunkTypeError), |
| 422 | Error: &StreamError{ |
| 423 | Code: "stream_error", |
| 424 | Message: err.Error(), |
| 425 | }, |
| 426 | } |
| 427 | } |
| 428 | } |
| 429 | |
| 430 | // parseStreamChunk 解析单个流式 chunk |
| 431 | func (p *GeminiProvider) parseStreamChunk(chunk map[string]any) []StreamChunk { |
no test coverage detected