From b5d2332ca93944d8c50fcb69a2f8597a7275edb3 Mon Sep 17 00:00:00 2001 From: qhd <1766646056@qq.com> Date: Thu, 9 Jul 2026 13:48:37 +0800 Subject: [PATCH] =?UTF-8?q?fix:=20=E4=BF=AE=E5=A4=8D=E7=BD=91=E5=85=B3?= =?UTF-8?q?=E4=BB=BB=E5=8A=A1=E5=88=9B=E5=BB=BA=E8=B6=85=E6=97=B6=E5=8F=8A?= =?UTF-8?q?=E5=93=8D=E5=BA=94=E8=A7=A3=E6=9E=90?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 1. 使用 context.WithoutCancel 防止上游断连导致请求取消,并设置30分钟超时 2. 显式设置 HTTP 客户端超时与 ResponseHeaderTimeout 3. 手动解析内部 API 响应格式并增加错误校验 4. 调整任务创建流程:先本地生成 UUID 并落库,再携带 taskId 调用网关 5. 提升 uuid 为直接依赖 --- go.mod | 2 +- service/gateway/gateway_http_service.go | 64 ++++++++++++++++++++---- service/prompt/prompt_compose_service.go | 27 ++++++---- 3 files changed, 70 insertions(+), 23 deletions(-) diff --git a/go.mod b/go.mod index 517dda7..3a36c4c 100644 --- a/go.mod +++ b/go.mod @@ -7,6 +7,7 @@ require ( github.com/gogf/gf/contrib/drivers/pgsql/v2 v2.10.2 github.com/gogf/gf/contrib/nosql/redis/v2 v2.10.2 github.com/gogf/gf/v2 v2.10.2 + github.com/google/uuid v1.6.0 ) require ( @@ -39,7 +40,6 @@ require ( github.com/golang/snappy v1.0.0 // indirect github.com/google/btree v1.1.3 // indirect github.com/google/flatbuffers v25.12.19+incompatible // indirect - github.com/google/uuid v1.6.0 // indirect github.com/gorilla/websocket v1.5.4-0.20250319132907-e064f32e3674 // indirect github.com/grokify/html-strip-tags-go v0.1.0 // indirect github.com/grpc-ecosystem/grpc-gateway/v2 v2.28.0 // indirect diff --git a/service/gateway/gateway_http_service.go b/service/gateway/gateway_http_service.go index d4c67e5..8ca6780 100644 --- a/service/gateway/gateway_http_service.go +++ b/service/gateway/gateway_http_service.go @@ -3,17 +3,21 @@ package gateway import ( "context" "encoding/json" + "errors" "fmt" "io" "net/http" "prompts-core/model/entity" "strings" + "time" "gitea.redpowerfuture.com/red-future/common/beans" commonHttp "gitea.redpowerfuture.com/red-future/common/http" "github.com/gogf/gf/v2/errors/gerror" "github.com/gogf/gf/v2/frame/g" + "github.com/gogf/gf/v2/net/ghttp" "github.com/gogf/gf/v2/os/gtime" + "github.com/gogf/gf/v2/util/gconv" ) // CreateTaskReq 创建任务请求 @@ -29,23 +33,61 @@ type CreateTaskReq struct { // CreateGatewayTask 创建网关异步任务 func CreateGatewayTask(ctx context.Context, payload map[string]any) (string, error) { fullURL := "model-gateway/task/createTask" - //headers := util.ForwardHeaders(ctx) - headers := make(map[string]string) - if r := g.RequestFromCtx(ctx); r != nil { - for k, v := range r.Request.Header { - if len(v) > 0 { - headers[k] = v[0] - } - } - } var req CreateTaskReq body, err := json.Marshal(payload) if err != nil { return "", err } - if err := commonHttp.Post(ctx, fullURL, headers, &req, body); err != nil { - return "", err + // 1. 隔离上游取消(防止 HTTP 客户端断连导致下游请求被 cancel)+ 设置独立超时 + baseCtx := context.WithoutCancel(ctx) + postCtx, cancel := context.WithTimeout(baseCtx, 30*time.Minute) + defer cancel() // 必须释放,防止上下文泄露 + + // 2. 克隆 commonHttp 客户端(保留 Consul 服务发现),显式设置超时和 ResponseHeaderTimeout + client := commonHttp.Httpclient.Clone() + client.SetTimeout(30 * time.Minute) + if tr, ok := client.Transport.(*http.Transport); ok { + tr.ResponseHeaderTimeout = 30 * time.Minute } + // 透传请求头 + if r := g.RequestFromCtx(ctx); r != nil { + for k, v := range r.Request.Header { + if len(v) > 0 { + client.SetHeader(k, v[0]) + } + } + } + resp, err := client.ContentJson().Post(postCtx, fullURL, body) + if err != nil { + return "", fmt.Errorf("请求创建任务失败: %w", err) + } + defer resp.Close() + result, err := io.ReadAll(resp.Body) + if err != nil { + return "", fmt.Errorf("读取响应失败: %w", err) + } + + // 统一处理内部API响应格式:{code:200,message:"",data:{...}} + resultStrut := &ghttp.DefaultHandlerResponse{} + + if err = gconv.Struct(result, &resultStrut); err != nil { // 修复:增加err检查 + return "", fmt.Errorf("响应解析失败: " + err.Error()) + } + + // 添加调试日志:打印解析后的结构 + g.Log().Debugf(ctx, "[HTTP] 解析后结构: Code=%d, Message=%s, Data类型=%T, Data值=%+v", + resultStrut.Code, resultStrut.Message, resultStrut.Data, resultStrut.Data) + + if resultStrut.Code == 200 || resultStrut.Code == 0 { + if err = gconv.Struct(resultStrut.Data, &req); err != nil { // 修复:增加err检查 + return "", fmt.Errorf("数据解析失败: " + err.Error()) + } + // 添加调试日志:打印最终的target + g.Log().Debugf(ctx, "[HTTP] 最终target: %+v", &req) + } else { + err = errors.New(resultStrut.Message) + } + return req.TaskId, nil } diff --git a/service/prompt/prompt_compose_service.go b/service/prompt/prompt_compose_service.go index 9f8ee21..7405d1d 100644 --- a/service/prompt/prompt_compose_service.go +++ b/service/prompt/prompt_compose_service.go @@ -17,6 +17,7 @@ import ( "gitea.redpowerfuture.com/red-future/common/utils" "github.com/gogf/gf/v2/frame/g" "github.com/gogf/gf/v2/util/gconv" + "github.com/google/uuid" ) // ComposeMessages 核心拼接提示词主流程 @@ -102,19 +103,11 @@ func handleBuild(ctx context.Context, req *dto.ComposeMessagesReq, chatModel, ai if err != nil { return nil, fmt.Errorf("构建推理请求失败: %w", err) } - - // 3) 调用网关创建任务 - taskID, err := gateway.CreateGatewayTask(ctx, taskReq) - if err != nil { - return nil, fmt.Errorf("创建网关任务失败: %w", err) - } - if taskID == "" { - return nil, errors.New("网关未返回taskId") - } - + taskId := uuid.NewString() + fmt.Println("taskId:", taskId) // 4) 保存任务记录 if _, err = dao.ComposeTask.Insert(ctx, &entity.ComposeTask{ - TaskId: taskID, + TaskId: taskId, ModelName: req.ModelName, SkillName: req.SkillName, BuildType: req.BuildType, @@ -124,6 +117,18 @@ func handleBuild(ctx context.Context, req *dto.ComposeMessagesReq, chatModel, ai }); err != nil { return nil, err } + + // 3) 调用网关创建任务 + taskReq["taskId"] = taskId + taskID, err := gateway.CreateGatewayTask(ctx, taskReq) + if err != nil { + return nil, fmt.Errorf("创建网关任务失败: %w", err) + } + if taskID == "" { + return nil, errors.New("网关未返回taskId") + } + fmt.Println("taskID:", taskID) + return &dto.ComposeMessagesRes{TaskId: taskID}, nil }