MCPcopy Create free account
hub / github.com/astercloud/aster / parseSSEStream

Method parseSSEStream

pkg/provider/gemini.go:388–428  ·  view source on GitHub ↗

parseSSEStream 解析 SSE 流

(body io.ReadCloser, chunks chan<- StreamChunk)

Source from the content-addressed store, hash-verified

386
387// parseSSEStream 解析 SSE 流
388func (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
431func (p *GeminiProvider) parseStreamChunk(chunk map[string]any) []StreamChunk {

Callers 1

StreamMethod · 0.95

Calls 3

parseStreamChunkMethod · 0.95
CloseMethod · 0.65
ErrorMethod · 0.65

Tested by

no test coverage detected