Files

86 lines
1.9 KiB
Go

package httpclient
import (
"bufio"
"context"
"encoding/json"
"io"
"strings"
"github.com/gogf/gf/v2/frame/g"
)
// SSE 常量
const (
ssePrefixData = "data:"
ssePrefixEvent = "event:"
ssePrefixComment = ":"
sseStreamDone = "[DONE]"
scanBufInitSize = 64 * 1024 // 64KB
scanMaxLineSize = 1024 * 1024 // 单行最大 1MB
)
// ParseSSEStream 标准 SSE 流式解析,逐分片回调,支持多行data、上下文取消
func ParseSSEStream(ctx context.Context, respBody io.Reader, onChunk func(ctx context.Context, chunk map[string]any) error) {
scanner := bufio.NewScanner(respBody)
scanner.Buffer(make([]byte, 0, scanBufInitSize), scanMaxLineSize)
var dataBuilder strings.Builder
for scanner.Scan() {
// 监听上下文取消,及时终止
select {
case <-ctx.Done():
g.Log().Infof(ctx, "[SSE] 上下文取消,终止流读取: %v", ctx.Err())
return
default:
}
line := scanner.Text()
// 跳过注释、事件行
if strings.HasPrefix(line, ssePrefixComment) || strings.HasPrefix(line, ssePrefixEvent) {
continue
}
lineTrim := strings.TrimSpace(line)
// 空行 = 一个SSE事件结束
if lineTrim == "" {
if dataBuilder.Len() == 0 {
continue
}
dataStr := dataBuilder.String()
dataBuilder.Reset()
if dataStr == sseStreamDone {
continue
}
var chunk map[string]any
if err := json.Unmarshal([]byte(dataStr), &chunk); err != nil {
g.Log().Debugf(ctx, "[SSE] JSON解析失败: %s, err: %v", dataStr, err)
continue
}
if onChunk != nil {
onChunk(ctx, chunk)
}
continue
}
// 拼接多行 data 数据
if strings.HasPrefix(line, ssePrefixData) {
raw := strings.TrimPrefix(line, ssePrefixData)
dataBuilder.WriteString(strings.TrimSpace(raw))
}
}
// 捕获读取异常
if err := scanner.Err(); err != nil {
g.Log().Errorf(ctx, "[SSE] 流读取异常: %v", err)
return
}
g.Log().Infof(ctx, "[SSE] 流式读取正常结束")
}