refactor: 重构模型 HTTP 客户端并移除冗余代码

This commit is contained in:
2026-08-29 11:53:02 +08:00
parent e079f4ea28
commit 4a9ae2d412
43 changed files with 719 additions and 4295 deletions
+167
View File
@@ -0,0 +1,167 @@
// Package httpclient 模型网关的传输层:模型 HTTP 请求(含瞬时网络错误重试)与 SSE 流式解析。
// 纯基础设施,不依赖 session/task/call 等业务逻辑;业务代码只通过三个导出函数使用。
package httpclient
import (
"context"
"errors"
"fmt"
"io"
"net/http"
"strings"
"time"
commonHttp "gitea.redpowerfuture.com/red-future/common/http"
"github.com/gogf/gf/v2/frame/g"
"github.com/gogf/gf/v2/net/gclient"
"github.com/gogf/gf/v2/util/gconv"
)
// modelCallHeaderTimeout 模型响应头等待超时。
// commonHttp 底层 gclient 默认 ResponseHeaderTimeout 只有 30s,模型生成首字节
// (尤其非流式、大 max_tokens)经常超过 30s,导致 http2: timeout awaiting response
// headers。模型调用必须用独立 client 并把该超时调大,与模型配置的超时保持一致。
const modelCallHeaderTimeout = 30 * time.Minute
// modelHTTPClient 构建模型调用专用 HTTP client
// 克隆 commonHttp 客户端(保留 ContentJson、header 注入等行为),但把
// ResponseHeaderTimeout 从默认 30s 调大到 modelCallHeaderTimeout。
func modelHTTPClient() *gclient.Client {
client := commonHttp.Httpclient.Clone()
if tr, ok := client.Transport.(*http.Transport); ok {
tr = tr.Clone() // 独立拷贝,避免改动全局共享 transport
tr.ResponseHeaderTimeout = modelCallHeaderTimeout
client.Transport = tr
}
return client
}
// modelNetRetryTimes 模型请求瞬时网络错误最大重试次数(不含首次);modelNetRetryBackoff 为退避基数。
// 模型域名 DNS 解析失败(Docker 内 127.0.0.11 偶发 no such host)是瞬时错误,短退避重试即可恢复。
// 重试在 HTTP 层完成,覆盖同步/异步/流式全部调用路径;流式场景发生在写 SSE 响应头之前,重试安全。
const (
modelNetRetryTimes = 3
modelNetRetryBackoff = 500 * time.Millisecond
)
// isTransientNetError 判定是否可重试的瞬时网络错误。仅命中 DNS 解析失败(no such host):
// 模型域名解析抖动可重试恢复;连接拒绝/超时等其他网络错误可能反映真实配置问题,不纳入,避免掩盖错误。
func isTransientNetError(err error) bool {
if err == nil {
return false
}
return strings.Contains(err.Error(), "no such host") || strings.Contains(err.Error(), "timeout")
}
// modelDoRaw 模型 HTTP 请求(等价 commonHttp.doRequestRaw,但使用调大超时的 client)。
// DNS 解析失败等瞬时网络错误在请求层短退避重试(modelNetRetryTimes 次);
// 其余错误(含上游业务错误码)原样返回,由上层按错误码决定是否重试。
func modelDoRaw(ctx context.Context, method string, url string, headers map[string]string, data ...any) (*gclient.Response, error) {
client := modelHTTPClient()
if (method == http.MethodPost || method == http.MethodPut || method == http.MethodDelete) && len(data) > 0 {
client = client.ContentJson()
}
if len(headers) > 0 {
client.SetHeaderMap(headers)
} else if r := g.RequestFromCtx(ctx); r != nil {
client.SetHeader("Authorization", r.Request.Header.Get("Authorization"))
}
doOnce := func() (*gclient.Response, error) {
if method == http.MethodGet && len(data) > 0 && len(data)%2 == 0 {
queryParams := make(map[string]string)
for i := 0; i < len(data); i += 2 {
if key, ok := data[i].(string); ok && i+1 < len(data) {
queryParams[key] = gconv.String(data[i+1])
}
}
return client.DoRequest(ctx, method, url, queryParams)
}
if len(data) == 1 {
return client.DoRequest(ctx, method, url, data[0])
}
return client.DoRequest(ctx, method, url, data...)
}
var response *gclient.Response
var err error
for attempt := 0; ; attempt++ {
response, err = doOnce()
if err == nil || !isTransientNetError(err) {
return response, err
}
if attempt >= modelNetRetryTimes {
break
}
wait := time.Duration(1<<attempt) * modelNetRetryBackoff
g.Log().Warningf(ctx, "[HttpModel] 模型请求瞬时网络错误,第 %d/%d 次重试(等待 %v): %v", attempt+1, modelNetRetryTimes, wait, err)
select {
case <-ctx.Done():
return nil, ctx.Err()
case <-time.After(wait):
}
}
return response, err
}
// ModelHttpNormalRequest 同步/异步 普通HTTP全量请求
func ModelHttpNormalRequest(ctx context.Context, url string, headers map[string]string, httpMethod string, body map[string]any) (res []byte, err error) {
response, err := modelDoRaw(ctx, httpMethod, url, headers, body)
if err != nil {
g.Log().Errorf(ctx, "[HttpModel] 模型请求失败 [Error]: %v", err)
return nil, fmt.Errorf("模型请求失败: %w", err)
}
defer response.Close()
return response.ReadAll(), nil
}
// ModelHttpStreamRequest 通用流式请求
// stream=true 时设置 SSE 头并验证 Flusherstream=false 时只返回 Reader,不设置响应头
func ModelHttpStreamRequest(ctx context.Context, w http.ResponseWriter, url string, headers map[string]string, httpMethod string, body map[string]any) (io.Reader, error) {
// 1) 先发起上游请求(此时还没写任何 SSE 头,失败可以正常返回 error)
response, err := modelDoRaw(ctx, httpMethod, url, headers, body)
if err != nil {
g.Log().Errorf(ctx, "[HttpModel] 模型流式请求失败 [Error]: %v", err)
return nil, fmt.Errorf("模型流式请求失败: %w", err)
}
// 检查 HTTP 状态码
if response.StatusCode < 200 || response.StatusCode >= 300 {
bodyBytes, _ := io.ReadAll(response.Body)
response.Close()
return nil, fmt.Errorf("[HTTP][Stream] 状态码异常: %d, body=%s", response.StatusCode, string(bodyBytes))
}
if w != nil {
// 2) 上游连接成功,再设置 SSE 头
h := w.Header()
h.Set("Content-Type", "text/event-stream; charset=utf-8")
h.Set("Cache-Control", "no-cache")
h.Set("Connection", "keep-alive")
h.Set("X-Accel-Buffering", "no")
if _, ok := w.(http.Flusher); !ok {
response.Close()
return nil, errors.New("response writer not support flush")
}
}
// 下层统一托管关闭:用包装器保证流最终关闭
return &autoCloseReader{r: response.Body}, nil
}
// autoCloseReader 包装 io.ReadCloser,读取结束/销毁时自动 Close
type autoCloseReader struct {
r io.ReadCloser
}
func (a *autoCloseReader) Read(p []byte) (int, error) {
n, err := a.r.Read(p)
// 读取完毕 / 读出错,主动关闭流
if err != nil {
_ = a.r.Close()
}
return n, err
}
+85
View File
@@ -0,0 +1,85 @@
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] 流式读取正常结束")
}