Files
gen2d/docs/async-tasks.md
T
Gmarker689 ecb74a2aed feat: 异步生成管线 + 前端轮询 + 图片编辑 + JWT 中间件
后端:
- 异步生成: POST /api/v1/generate 立即返回 taskId,后台执行管线
- 任务轮询: GET /api/v1/tasks/:id + GET /api/v1/tasks/:id/assets
- 图片保存: 生成图片写入 ../generation/{projectId}/{taskId}/,静态服务
- 图片编辑: POST /api/v1/images/edit (multipart/form-data)
- JWT 中间件: mildware/auth.go 保护生成/编辑端点
- config.yml 清空敏感默认值,交由 .env 控制
- ImageGenConfig 新增 Quality 字段

前端:
- api/generate.ts: 对接真实 API (submitGenerate + poll getTask/getAssets)
- api/types.ts: 新增 GenerateResponse, AssetsResponse, Task 类型
- stores/generation.ts: 异步提交→轮询进度→获取素材→完成
- stores/task.ts: 默认分辨率 256→1024
- GenerateForm: 分辨率范围 1024-1536
- GeneratePage: 显示状态文本,完成后可查看结果/继续生成
- ResultPage: 从 store 读取,下载功能实现
2026-05-25 14:08:08 +08:00

191 lines
5.7 KiB
Markdown
Executable File
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
# 异步任务
## 概述
生成任务采用「提交-异步执行-轮询/推送」模式:API 同步返回 taskId,后端异步执行四阶段管线,前端通过轮询或 WebSocket 获取进度与结果。
## 任务生命周期
```mermaid
stateDiagram-v2
[*] --> pending : POST /api/v1/generate
pending --> running : Worker 领取任务
running --> completed : 四阶段全部通过
running --> failed : 节点执行失败 / 超过重试次数
completed --> [*]
failed --> [*]
```
| 状态 | 说明 | 持久化 |
|------|------|--------|
| `pending` | 已入队,等待 Worker 领取 | `task.status = 'pending'` |
| `running` | 管线执行中,`task.stage` 记录当前阶段 | `task.status = 'running'` |
| `completed` | 管线完成,素材已保存到本地 | `task.status = 'completed'`, 写入 `asset` 表 |
| `failed` | 执行失败或超过重试上限 | `task.status = 'failed'`, `task.error` 记录原因 |
## 任务队列
采用内存 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
}
}
```
### 提交流程
```
POST /api/v1/generate
│
├── 1. 去重检查(sync.Map,hash(prompt + assetType + params))
├── 2. 写入 task 表(status=pending)
├── 3. 推入 JobQueue
└── 4. 同步返回 taskId
```
### Worker 流程
```
JobQueue.Pop() 取出 taskId
│
├── 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 = 错误信息
```
### 并发控制
通过固定数量的 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。
## 管线阶段与进度
每个阶段对应 Eino Graph 的一个节点,执行过程中更新 `task.stage` 和 `task.progress`:
| 阶段 | stage 值 | 进度范围 | 说明 |
|------|----------|---------|------|
| 提示词优化 | `prompt_optimizer` | 0-20% | 调用 PromptAgent/LLM 优化提示词,合并风格与技术参数 |
| 素材生成 | `asset_generator` | 20-60% | 调用 AI 推理 API 出图 |
| 质量检查 | `quality_supervisor` | 60-80% | 视觉模型质检 |
| 格式适配 | `format_adapter` | 80-100% | 格式转换、spritesheet 打包、保存素材 |
进度更新通过 WebSocket 实时推送给前端(参见 [API 设计 - WebSocket 消息格式](api.md))。
## 重试策略
### 质检重试(管线内重试)
QualitySupervisor 质检不通过时,Eino Graph 分支回到 PromptOptimizer 重新优化提示词:
```
QualitySupervisor -- fail --> PromptOptimizer --> AssetGenerator --> QualitySupervisor
QualitySupervisor -- pass --> FormatAdapter
```
| 参数 | 默认值 | 说明 |
|------|--------|------|
| 最大重试次数 | 3 | `task.retry_count` 达到上限后降级输出 |
| 重试触发条件 | 质检不通过 | 视觉模型判定风格不一致 |
超过重试次数后,跳过质检直接进入 FormatAdapter 输出(降级策略,保证任务不会无限循环)。
### 任务级重试(管线外重试)
管线执行过程中发生不可恢复的错误(如 AI API 超时、文件写入失败):
| 参数 | 默认值 | 说明 |
|------|--------|------|
| 最大重试次数 | 1 | 仅重试一次 |
| 重试间隔 | 5 秒 | 固定间隔 |
重试时重新执行完整管线,不保留上次中间状态。
## 进度推送
前端可通过两种方式获取任务进度:
### 轮询
```
GET /api/v1/tasks/:taskId
```
前端定时调用(建议间隔 2-3 秒),简单可靠,适合不需要实时性的场景。
### WebSocket
```
ws://host/api/v1/tasks/:taskId/ws
```
长连接实时推送,通过 httpOnly Cookie 自动认证,每个阶段的状态变更立即通知前端。消息格式参见 [API 设计](api.md)。
## 去重
去重策略详见 [数据存储 — 去重策略](database.md#去重策略)。相同输入的并发请求直接返回已有 taskId,任务完成后从内存清除。
## 文件存储
任务完成后,FormatAdapter 将素材保存到本地文件系统。素材文件的存储结构与访问方式详见 [数据存储 — 文件存储](database.md#文件存储)。