From 0847d126a27879e850daef1f7e449ecdf5b3734e Mon Sep 17 00:00:00 2001 From: wonder Date: Mon, 25 May 2026 18:16:21 +0800 Subject: [PATCH] =?UTF-8?q?fix:=20=E4=BB=BB=E5=8A=A1=E5=88=9B=E5=BB=BA?= =?UTF-8?q?=E5=90=8E=E7=9B=B4=E6=8E=A5=E8=BF=9B=E5=85=A5=20running=20?= =?UTF-8?q?=E7=8A=B6=E6=80=81=EF=BC=8C=E9=81=BF=E5=85=8D=E5=89=8D=E7=AB=AF?= =?UTF-8?q?=E6=98=BE=E7=A4=BA=E6=8E=92=E9=98=9F=E4=B8=AD?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 移除任务初始状态 pending,改为直接创建为 running - 初始进度设为 5%,stage 设为 prompt_builder - 更新文档以反映实际执行流程(无队列机制) - 修复前端轮询时因竞态条件显示排队中的问题 --- backend/internal/handler/generate.go | 5 +- docs/api.md | 7 +- docs/async-tasks.md | 111 +++++++-------------------- 3 files changed, 32 insertions(+), 91 deletions(-) diff --git a/backend/internal/handler/generate.go b/backend/internal/handler/generate.go index 0c5864b..7994a52 100755 --- a/backend/internal/handler/generate.go +++ b/backend/internal/handler/generate.go @@ -236,8 +236,9 @@ func saveTaskToDB(ctx context.Context, projectID, taskID string, req GenerateReq ProjectID: uint(projectIDUint), Prompt: req.Prompt, AssetType: req.AssetType, - Status: "pending", - Progress: 0, + Status: "running", + Progress: 5, + Stage: "prompt_builder", RetryCount: 0, CreatedAt: time.Now(), UpdatedAt: time.Now(), diff --git a/docs/api.md b/docs/api.md index a78417f..9ebe4e2 100755 --- a/docs/api.md +++ b/docs/api.md @@ -423,8 +423,7 @@ POST /api/v1/prompt/optimize "code": 0, "message": "ok", "data": { - "taskId": "task_xyz789", - "status": "pending" + "taskId": "task_xyz789" } } ``` @@ -466,8 +465,8 @@ POST /api/v1/prompt/optimize | 字段 | 类型 | 说明 | |------|------|------| -| `status` | string | 任务状态:`pending` / `running` / `completed` / `failed` | -| `stage` | string | 当前管线阶段,`running` 时有值 | +| `status` | string | 任务状态:`running` / `completed` / `failed` / `saving` | +| `stage` | string | 当前管线阶段,`running` 时有值:`prompt_builder` / `asset_generator` / `quality_supervisor` / `format_adapter` | | `progress` | int | 0-100,整体进度 | | `retryCount` | int | 质检重试次数 | | `error` | string | 失败原因,仅 `failed` 时有值 | diff --git a/docs/async-tasks.md b/docs/async-tasks.md index 8f9ca9c..209a3a2 100755 --- a/docs/async-tasks.md +++ b/docs/async-tasks.md @@ -8,8 +8,7 @@ ```mermaid stateDiagram-v2 - [*] --> pending : POST /api/v1/generate - pending --> running : Worker 领取任务 + [*] --> running : POST /api/v1/generate running --> completed : 四阶段全部通过 running --> failed : 节点执行失败 / 超过重试次数 completed --> [*] @@ -18,106 +17,48 @@ stateDiagram-v2 | 状态 | 说明 | 持久化 | |------|------|--------| -| `pending` | 已入队,等待 Worker 领取 | `task.status = 'pending'` | -| `running` | 管线执行中,`task.stage` 记录当前阶段 | `task.status = 'running'` | +| `running` | 管线执行中,`task.stage` 记录当前阶段 | `task.status = 'running'`, `task.stage`, `task.progress` | | `completed` | 管线完成,素材已保存到本地 | `task.status = 'completed'`, 写入 `asset` 表 | | `failed` | 执行失败或超过重试上限 | `task.status = 'failed'`, `task.error` 记录原因 | +| `saving` | 保存素材中(内部状态) | `task.status = 'saving'` | -## 任务队列 +## 任务执行 -采用内存 Channel 队列(非 Redis),适合单实例部署场景: - -```go -// jobqueue.go -type JobQueue struct { - ch chan string // taskId channel - done chan struct{} -} - -func NewJobQueue(size int) *JobQueue { - return &JobQueue{ - ch: make(chan string, size), - done: make(chan struct{}), - } -} - -func (q *JobQueue) Push(taskId string) { - q.ch <- taskId -} - -func (q *JobQueue) Pop() (string, bool) { - select { - case taskId := <-q.ch: - return taskId, true - case <-q.done: - return "", false - } -} -``` +各任务独立通过 goroutine 异步执行,无队列机制: ### 提交流程 ``` POST /api/v1/generate │ - ├── 1. 去重检查(sync.Map,hash(prompt + assetType + params)) - ├── 2. 写入 task 表(status=pending) - ├── 3. 推入 JobQueue - └── 4. 同步返回 taskId + ├── 1. 创建任务记录(status=running, stage=prompt_builder, progress=5) + ├── 2. 同步返回 taskId + └── 3. 后台 goroutine 执行管线 ``` -### Worker 流程 +### 执行流程 ``` -JobQueue.Pop() 取出 taskId +后台 goroutine 执行 RunPipeline(ctx, input) │ - ├── 1. 更新 task.status = running - ├── 2. 查询 task + project_style - ├── 3. 构造 PipelineInput - ├── 4. RunPipeline(ctx, input) - │ ├── PromptOptimizer → stage 更新 - │ ├── AssetGenerator → stage 更新 - │ ├── QualitySupervisor → stage 更新(可能触发重优化分支) - │ └── FormatAdapter → stage 更新 - ├── 5. 保存素材到本地,写入 asset 表 - ├── 6. 更新 task.status = completed - └── 7. 异常时更新 task.status = failed, task.error = 错误信息 + ├── 1. 查询 task + project_style + ├── 2. 构造 PipelineInput + ├── 3. 执行管线,实时更新 stage & progress + │ ├── prompt_builder → stage 更新,progress 5-25% + │ ├── asset_generator → stage 更新,progress 25-60% + │ ├── quality_supervisor → stage 更新,progress 60-80%(可能触发重优化分支) + │ └── format_adapter → stage 更新,progress 80-90% + ├── 4. 上传素材至存储并写入 asset 表(status=saving, progress 90%) + ├── 5. 更新 task.status = completed(progress 100%) + └── 6. 异常时更新 task.status = failed, task.error = 错误信息 ``` ### 并发控制 -通过固定数量的 Worker goroutine 控制并发: - -```go -func (s *TaskService) StartWorkers(n int) { - for i := 0; i < n; i++ { - go func(workerID int) { - for { - taskId, ok := s.queue.Pop() - if !ok { - return - } - s.processTask(context.Background(), taskId) - } - }(i) - } -} -``` - -| 参数 | 默认值 | 说明 | -|------|--------|------| -| 队列容量 | 100 | 内存 channel 缓冲大小 | -| Worker 数 | 3 | 并发执行任务数 | - -### 超时控制 - -每个任务的 `processTask` 执行设置超时上下文(默认 10 分钟),超时后: - -1. 取消正在执行的管线节点(如 AI 推理调用) -2. 更新 `task.status = failed`,`task.error = "任务执行超时"` -3. Worker 释放,继续处理下一个任务 - -避免长时间卡住的任务永久占用 Worker。 +无固定 Worker 限制,每个请求启动一个独立 goroutine。实际并发受: +- 系统资源(CPU、内存) +- AI 推理 API 限流 +- 数据库连接池 ## 管线阶段与进度 @@ -125,8 +66,8 @@ func (s *TaskService) StartWorkers(n int) { | 阶段 | stage 值 | 进度范围 | 说明 | |------|----------|---------|------| -| 提示词优化 | `prompt_optimizer` | 0-20% | 调用 PromptAgent/LLM 优化提示词,合并风格与技术参数 | -| 素材生成 | `asset_generator` | 20-60% | 调用 AI 推理 API 出图 | +| 提示词优化 | `prompt_builder` | 5-25% | 调用 PromptAgent/LLM 优化提示词,合并风格与技术参数 | +| 素材生成 | `asset_generator` | 25-60% | 调用 AI 推理 API 出图 | | 质量检查 | `quality_supervisor` | 60-80% | 视觉模型质检 | | 格式适配 | `format_adapter` | 80-100% | 格式转换、spritesheet 打包、保存素材 |