From 96ac56cb4e04e9b10b2c856ee0693be2f38bbc15 Mon Sep 17 00:00:00 2001 From: Gmaker689 <1711322114@qq.com> Date: Mon, 25 May 2026 19:10:46 +0800 Subject: [PATCH] =?UTF-8?q?fix:=20=E4=BF=AE=E5=A4=8D=E7=AE=A1=E7=BA=BFcont?= =?UTF-8?q?ext=E5=8F=96=E6=B6=88=E3=80=81=E6=B7=BB=E5=8A=A0FIFO=E4=BB=BB?= =?UTF-8?q?=E5=8A=A1=E9=98=9F=E5=88=97=E3=80=81=E7=A7=BB=E9=99=A4=E9=87=8D?= =?UTF-8?q?=E5=A4=8D=E6=8F=90=E7=A4=BA=E8=AF=8D=E4=BC=98=E5=8C=96?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 三个修复: 1. runPipelineBg 改用 context.Background(),避免 HTTP 响应返回后 Gin 取消 request context 导致后台管线静默失败 2. 新增 TaskQueue FIFO 串行队列,任务提交后进入 pending 状态排队, 按提交顺序逐个执行,前端轮询显示排队中 3. promptOptimizerNode 移除 RunPromptAgent 调用,提示词优化仅由 前端在提交前通过 /api/v1/prompt/optimize 执行一次,管线内只做 风格合并和技术参数追加 --- backend/cmd/main.go | 5 ++ backend/internal/handler/generate.go | 38 +++++++--- backend/internal/service/nodes.go | 27 ++----- backend/internal/service/queue.go | 104 +++++++++++++++++++++++++++ frontend/src/stores/generation.ts | 2 +- 5 files changed, 145 insertions(+), 31 deletions(-) create mode 100644 backend/internal/service/queue.go diff --git a/backend/cmd/main.go b/backend/cmd/main.go index 967a175..5d06703 100755 --- a/backend/cmd/main.go +++ b/backend/cmd/main.go @@ -44,6 +44,11 @@ func main() { // 初始化项目服务 service.InitProjectService(db.GetDB()) + // 初始化生成任务队列(FIFO 串行执行) + generateQueue := service.NewTaskQueue() + handler.SetGenerateQueue(generateQueue) + defer generateQueue.Stop() + r := gin.New() r.Use(mildware.Logger()) r.Use(mildware.Recovery()) diff --git a/backend/internal/handler/generate.go b/backend/internal/handler/generate.go index b5fb767..e9f928f 100755 --- a/backend/internal/handler/generate.go +++ b/backend/internal/handler/generate.go @@ -16,6 +16,14 @@ import ( "github.com/gin-gonic/gin" ) +// generateQueue 全局生成任务队列,由 main 通过 SetGenerateQueue 注入。 +var generateQueue *service.TaskQueue + +// SetGenerateQueue 设置生成任务队列。 +func SetGenerateQueue(q *service.TaskQueue) { + generateQueue = q +} + // GenerateRequest 素材生成请求。 type GenerateRequest struct { ProjectID string `json:"projectId"` @@ -43,7 +51,7 @@ type AssetsResponse struct { } // Generate 素材生成接口(异步)。 -// 立即返回 taskId,后台执行管线,前端通过 GET /tasks/:taskId 轮询进度。 +// 立即返回 taskId,任务进入 FIFO 队列串行执行,前端通过 GET /tasks/:taskId 轮询进度。 func Generate(c *gin.Context) { var req GenerateRequest if err := c.ShouldBindJSON(&req); err != nil { @@ -57,18 +65,29 @@ func Generate(c *gin.Context) { } taskID := fmt.Sprintf("task-%d", time.Now().UnixMilli()) - // 保存任务到数据库 + // 保存任务到数据库,初始状态为 pending if err := saveTaskToDB(c.Request.Context(), projectID, taskID, req); err != nil { logger.FromCtx(c.Request.Context()).Error("failed to save task", "error", err) c.JSON(http.StatusInternalServerError, model.Fail(http.StatusInternalServerError, "创建任务失败")) return } - // 返回 taskId - c.JSON(http.StatusOK, model.OK(GenerateResponse{TaskID: taskID})) + // 加入 FIFO 队列 + queuePos := 1 + if generateQueue != nil { + queuePos = generateQueue.Enqueue(&service.TaskJob{ + Ctx: context.Background(), + ProjectID: projectID, + TaskID: taskID, + Execute: func(ctx context.Context) error { + runPipelineBg(ctx, projectID, taskID, req) + return nil + }, + }) + } - // 后台执行管线 - go runPipelineBg(c.Request.Context(), projectID, taskID, req) + c.JSON(http.StatusOK, model.OK(GenerateResponse{TaskID: taskID})) + _ = queuePos } // runPipelineBg 后台执行生成管线,更新任务状态。 @@ -82,6 +101,7 @@ func runPipelineBg(ctx context.Context, projectID, taskID string, req GenerateRe updateTaskInDB(ctx, taskID, "running", stage, "", progress) }) + // 队列调度后才标记为 running,初始写入时是 pending updateTaskInDB(ctx, taskID, "running", "prompt_builder", "", 5) in := service.PipelineInput{ @@ -240,9 +260,9 @@ func saveTaskToDB(ctx context.Context, projectID, taskID string, req GenerateReq ProjectID: uint(projectIDUint), Prompt: req.Prompt, AssetType: req.AssetType, - Status: "running", - Progress: 5, - Stage: "prompt_builder", + Status: "pending", + Progress: 0, + Stage: "", RetryCount: 0, CreatedAt: time.Now(), UpdatedAt: time.Now(), diff --git a/backend/internal/service/nodes.go b/backend/internal/service/nodes.go index 56c35b3..ce7c847 100755 --- a/backend/internal/service/nodes.go +++ b/backend/internal/service/nodes.go @@ -14,10 +14,10 @@ import ( "github.com/cloudwego/eino/compose" ) -// promptOptimizerNode 节点:调用 PromptAgent 生成规范提示词,合并风格与重试信息。 -// 输入 PipelineInput,输出最终提示词字符串(直接供 AssetGenerator 消费)。 +// promptOptimizerNode 节点:合并风格描述、重试信息与技术参数,输出最终提示词。 +// 提示词优化已由前端在提交前完成,管线内不再重复调用 PromptAgent。 var promptOptimizerNode = compose.InvokableLambda(func(ctx context.Context, in PipelineInput) (string, error) { - // 合并风格描述,注入原始 Prompt 中 + // 合并风格描述 styleDesc := buildStyleDescription(in.ProjectStyle, in.TaskStyle) if styleDesc != "" { if in.Prompt != "" { @@ -36,26 +36,11 @@ var promptOptimizerNode = compose.InvokableLambda(func(ctx context.Context, in P } } - if len(in.Tags) == 0 && in.Prompt == "" { - return "", fmt.Errorf("pipeline: Prompt and Tags are both empty") + if in.Prompt == "" { + return "", fmt.Errorf("pipeline: prompt is empty") } - // 有标签时调用 PromptAgent 优化提示词 - if len(in.Tags) > 0 { - agentIn := PromptAgentInput{ - Tags: in.Tags, - AssetType: in.AssetType, - Prompt: in.Prompt, - UserNote: in.UserNote, - } - output, err := RunPromptAgent(ctx, agentIn) - if err != nil { - return "", fmt.Errorf("prompt agent: %w", err) - } - return output.Prompt, nil - } - - // 无标签时直接使用原始 Prompt,补上技术参数段 + // 追加技术参数段,不再调用 PromptAgent return appendTechNotes(in.Prompt, in.AssetType, in.Params), nil }) diff --git a/backend/internal/service/queue.go b/backend/internal/service/queue.go new file mode 100644 index 0000000..ea89ce0 --- /dev/null +++ b/backend/internal/service/queue.go @@ -0,0 +1,104 @@ +package service + +import ( + "context" + "sync" + + "gen2d/internal/logger" +) + +// TaskJob 队列中的任务。 +type TaskJob struct { + Ctx context.Context + ProjectID string + TaskID string + Execute func(ctx context.Context) error +} + +// TaskQueue 串行 FIFO 任务队列。 +type TaskQueue struct { + mu sync.Mutex + jobs []*TaskJob + ready chan struct{} + stop chan struct{} + stopped bool +} + +// NewTaskQueue 创建任务队列并启动调度协程。 +func NewTaskQueue() *TaskQueue { + q := &TaskQueue{ + jobs: make([]*TaskJob, 0), + ready: make(chan struct{}, 1), + stop: make(chan struct{}), + } + go q.run() + return q +} + +// Enqueue 将任务加入队尾,返回队列中的位置(1-based)。 +func (q *TaskQueue) Enqueue(job *TaskJob) int { + q.mu.Lock() + defer q.mu.Unlock() + q.jobs = append(q.jobs, job) + pos := len(q.jobs) + select { + case q.ready <- struct{}{}: + default: + } + return pos +} + +// QueueLen 返回当前队列长度。 +func (q *TaskQueue) QueueLen() int { + q.mu.Lock() + defer q.mu.Unlock() + return len(q.jobs) +} + +// Stop 优雅关闭队列。 +func (q *TaskQueue) Stop() { + q.mu.Lock() + defer q.mu.Unlock() + if !q.stopped { + q.stopped = true + close(q.stop) + } +} + +func (q *TaskQueue) run() { + for { + select { + case <-q.stop: + return + case <-q.ready: + q.processNext() + } + } +} + +func (q *TaskQueue) processNext() { + q.mu.Lock() + if len(q.jobs) == 0 { + q.mu.Unlock() + return + } + job := q.jobs[0] + q.jobs = q.jobs[1:] + // 如果队列还有任务,重新发信号 + if len(q.jobs) > 0 { + select { + case q.ready <- struct{}{}: + default: + } + } + q.mu.Unlock() + + l := logger.With("task_id", job.TaskID, "project_id", job.ProjectID) + l.Info("task queue executing job", "queue_remaining", len(q.jobs)) + + if err := job.Execute(job.Ctx); err != nil { + l.Error("task job failed", "error", err) + } else { + l.Info("task job completed") + } +} diff --git a/frontend/src/stores/generation.ts b/frontend/src/stores/generation.ts index 6686a2f..9f740ea 100755 --- a/frontend/src/stores/generation.ts +++ b/frontend/src/stores/generation.ts @@ -120,7 +120,7 @@ export const useGenerationStore = create((set) => ({ function stageLabel(stage?: string): string { switch (stage) { - case 'prompt_builder': return '优化提示词...' + case 'prompt_builder': return '构建提示词...' case 'asset_generator': return '生成素材中...' case 'quality_supervisor': return '质检中...' case 'format_adapter': return '格式转换中...'