fix: 任务创建后直接进入 running 状态,避免前端显示排队中
- 移除任务初始状态 pending,改为直接创建为 running - 初始进度设为 5%,stage 设为 prompt_builder - 更新文档以反映实际执行流程(无队列机制) - 修复前端轮询时因竞态条件显示排队中的问题
This commit is contained in:
+26
-85
@@ -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 打包、保存素材 |
|
||||
|
||||
|
||||
Reference in New Issue
Block a user