86 lines
1.9 KiB
Go
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] 流式读取正常结束")
|
|
}
|