// 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<= 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 }