mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-10-11 09:46:37 +08:00
async task skill
This commit is contained in:
@@ -32,6 +32,8 @@ Admin dispatch -> ValidateAndNormalizePayload -> DispatchTask
|
|||||||
|
|
||||||
按任务影响面选择对应步骤。不要只改其中一条链路。
|
按任务影响面选择对应步骤。不要只改其中一条链路。
|
||||||
|
|
||||||
|
> 需要可复制的代码模板时,阅读 [references/CODE-EXAMPLES.md](references/CODE-EXAMPLES.md)。那里包含任务常量、无参数 handler、带参数 `PayloadValidator`、统一注册、Worker 路由、Cron 配置和测试示例。
|
||||||
|
|
||||||
1. 定义任务元数据。
|
1. 定义任务元数据。
|
||||||
- 在 `internal/task/constants.go` 添加 Asynq task type 常量,例如 `upload:cleanup_unused`。
|
- 在 `internal/task/constants.go` 添加 Asynq task type 常量,例如 `upload:cleanup_unused`。
|
||||||
- 添加 Admin 可下发 task type 常量,例如 `cleanup_unused_uploads`。
|
- 添加 Admin 可下发 task type 常量,例如 `cleanup_unused_uploads`。
|
||||||
@@ -69,67 +71,6 @@ Admin dispatch -> ValidateAndNormalizePayload -> DispatchTask
|
|||||||
|
|
||||||
## Handler 模式
|
## Handler 模式
|
||||||
|
|
||||||
无参数任务:
|
|
||||||
|
|
||||||
```go
|
|
||||||
type CleanupUnusedUploadsHandler struct{}
|
|
||||||
|
|
||||||
func (h *CleanupUnusedUploadsHandler) Execute(ctx context.Context, payload []byte) (*task.TaskResult, error) {
|
|
||||||
task.AppendLog(ctx, "开始扫描未使用上传")
|
|
||||||
|
|
||||||
// 调用 model/service 完成业务逻辑。
|
|
||||||
|
|
||||||
msg := "清理完成"
|
|
||||||
task.AppendLog(ctx, "%s", msg)
|
|
||||||
return &task.TaskResult{Message: msg}, nil
|
|
||||||
}
|
|
||||||
```
|
|
||||||
|
|
||||||
带参数任务:
|
|
||||||
|
|
||||||
```go
|
|
||||||
type SendEmailPayload struct {
|
|
||||||
To string `json:"to"`
|
|
||||||
Subject string `json:"subject"`
|
|
||||||
Body string `json:"body"`
|
|
||||||
}
|
|
||||||
|
|
||||||
type SendEmailHandler struct{}
|
|
||||||
|
|
||||||
func (h *SendEmailHandler) ValidatePayload(payload []byte) ([]byte, error) {
|
|
||||||
if len(payload) == 0 {
|
|
||||||
return nil, errors.New("任务参数不能为空")
|
|
||||||
}
|
|
||||||
|
|
||||||
var req SendEmailPayload
|
|
||||||
if err := json.Unmarshal(payload, &req); err != nil {
|
|
||||||
return nil, fmt.Errorf("无效的 JSON 格式: %w", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
req.To = strings.TrimSpace(req.To)
|
|
||||||
req.Subject = strings.TrimSpace(req.Subject)
|
|
||||||
req.Body = strings.TrimSpace(req.Body)
|
|
||||||
if req.To == "" || req.Subject == "" || req.Body == "" {
|
|
||||||
return nil, errors.New("to、subject、body 不能为空")
|
|
||||||
}
|
|
||||||
|
|
||||||
return json.Marshal(req)
|
|
||||||
}
|
|
||||||
|
|
||||||
func (h *SendEmailHandler) Execute(ctx context.Context, payload []byte) (*task.TaskResult, error) {
|
|
||||||
var req SendEmailPayload
|
|
||||||
if err := json.Unmarshal(payload, &req); err != nil {
|
|
||||||
return nil, fmt.Errorf("解析参数失败: %w", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
task.AppendLog(ctx, "开始发送邮件到: %s", req.To)
|
|
||||||
|
|
||||||
msg := fmt.Sprintf("邮件成功发送至: %s", req.To)
|
|
||||||
task.AppendLog(ctx, "%s", msg)
|
|
||||||
return &task.TaskResult{Message: msg}, nil
|
|
||||||
}
|
|
||||||
```
|
|
||||||
|
|
||||||
约定:
|
约定:
|
||||||
|
|
||||||
- `ValidatePayload` 是 Admin 下发时的服务端校验入口,返回值会作为标准化 payload 存库和入队。
|
- `ValidatePayload` 是 Admin 下发时的服务端校验入口,返回值会作为标准化 payload 存库和入队。
|
||||||
@@ -137,6 +78,8 @@ func (h *SendEmailHandler) Execute(ctx context.Context, payload []byte) (*task.T
|
|||||||
- 成功返回 `&task.TaskResult{Message: "...", Detail: "..."}`;失败返回 `nil, fmt.Errorf("...")`,由 `ProcessTask` 标记失败并交给 Asynq 重试。
|
- 成功返回 `&task.TaskResult{Message: "...", Detail: "..."}`;失败返回 `nil, fmt.Errorf("...")`,由 `ProcessTask` 标记失败并交给 Asynq 重试。
|
||||||
- 错误要返回给框架,不在 handler 内吞掉;可继续的单条失败可用 `AppendLog` 记录后继续处理。
|
- 错误要返回给框架,不在 handler 内吞掉;可继续的单条失败可用 `AppendLog` 记录后继续处理。
|
||||||
|
|
||||||
|
> 在新增无参数任务、带参数任务或 `PayloadValidator` 时,阅读 [references/CODE-EXAMPLES.md](references/CODE-EXAMPLES.md) 的 Handler 示例。
|
||||||
|
|
||||||
## AppendLog 规则
|
## AppendLog 规则
|
||||||
|
|
||||||
在 `TaskHandler.Execute` 内用 `task.AppendLog(ctx, format, args...)` 写任务日志。
|
在 `TaskHandler.Execute` 内用 `task.AppendLog(ctx, format, args...)` 写任务日志。
|
||||||
@@ -192,3 +135,11 @@ make code-check
|
|||||||
```
|
```
|
||||||
|
|
||||||
如涉及前端或整体构建,补跑 `make build-test`。测试中需要 Redis/Asynq 时,可用现有测试模式或 `miniredis`;不要为了测试便利把 `internal/task` 反向塞进通用 util/testhelper,避免 import cycle。
|
如涉及前端或整体构建,补跑 `make build-test`。测试中需要 Redis/Asynq 时,可用现有测试模式或 `miniredis`;不要为了测试便利把 `internal/task` 反向塞进通用 util/testhelper,避免 import cycle。
|
||||||
|
|
||||||
|
## 相关 Skills
|
||||||
|
|
||||||
|
- Go 错误处理:在设计 handler 返回错误、包装底层错误或避免“记录并返回”时,参见 [go-error-handling](../go-error-handling/SKILL.md)。
|
||||||
|
- Go 测试:在编写 handler、dispatch 或 retry 单测时,参见 [go-testing](../go-testing/SKILL.md)。
|
||||||
|
- Go context:在任务业务逻辑传播取消、超时或 request scoped 值时,参见 [go-context](../go-context/SKILL.md)。
|
||||||
|
- Go logging:在决定任务日志、应用日志和日志级别边界时,参见 [go-logging](../go-logging/SKILL.md)。
|
||||||
|
- shadcn:在修改 Admin 任务 UI 时,参见 [shadcn](../shadcn/SKILL.md)。
|
||||||
|
|||||||
@@ -0,0 +1,307 @@
|
|||||||
|
# Wavelet 异步任务代码示例
|
||||||
|
|
||||||
|
这些示例用于新增或修改 Wavelet Asynq 任务时快速套用。复制前先对照当前代码,因为任务框架可能随项目演进。
|
||||||
|
|
||||||
|
## 任务元数据
|
||||||
|
|
||||||
|
在 `internal/task/constants.go` 添加 Asynq task type、Admin task type 和 `TaskMeta`。
|
||||||
|
|
||||||
|
```go
|
||||||
|
// 异步任务类型标识。格式建议为 "{module}:{action}"。
|
||||||
|
const CleanupUnusedUploadsTask = "upload:cleanup_unused"
|
||||||
|
|
||||||
|
// 管理员可下发的任务类型标识。用于 Admin API 的 task_type。
|
||||||
|
const TaskTypeCleanupUploads = "cleanup_unused_uploads"
|
||||||
|
|
||||||
|
var DispatchableTasks = []TaskMeta{
|
||||||
|
{
|
||||||
|
Type: TaskTypeCleanupUploads,
|
||||||
|
AsynqTask: CleanupUnusedUploadsTask,
|
||||||
|
Name: "清理未使用上传",
|
||||||
|
Description: "清理超过1小时未使用的上传文件",
|
||||||
|
SupportsTime: false,
|
||||||
|
MaxRetry: defaultMaxRetry,
|
||||||
|
Queue: QueueDefault,
|
||||||
|
Retryable: true,
|
||||||
|
},
|
||||||
|
}
|
||||||
|
```
|
||||||
|
|
||||||
|
带参数任务把前端表单元数据放在 `Params`。`Name` 必须和 payload JSON tag 对齐。
|
||||||
|
|
||||||
|
```go
|
||||||
|
{
|
||||||
|
Type: TaskTypeSendEmail,
|
||||||
|
AsynqTask: SendEmailTask,
|
||||||
|
Name: "发送邮件",
|
||||||
|
Description: "异步发送系统邮件",
|
||||||
|
SupportsTime: false,
|
||||||
|
MaxRetry: defaultMaxRetry,
|
||||||
|
Queue: QueueDefault,
|
||||||
|
Retryable: true,
|
||||||
|
Params: []TaskParam{
|
||||||
|
{
|
||||||
|
Name: "to",
|
||||||
|
Label: "接收邮箱 (To)",
|
||||||
|
Type: "string",
|
||||||
|
Required: true,
|
||||||
|
Placeholder: "receiver@example.com",
|
||||||
|
Description: "接收邮件的目标邮箱地址",
|
||||||
|
},
|
||||||
|
{
|
||||||
|
Name: "subject",
|
||||||
|
Label: "邮件主题 (Subject)",
|
||||||
|
Type: "string",
|
||||||
|
Required: true,
|
||||||
|
Placeholder: "请输入邮件主题",
|
||||||
|
Description: "发送邮件的主题标题",
|
||||||
|
},
|
||||||
|
{
|
||||||
|
Name: "body",
|
||||||
|
Label: "邮件内容 (Body)",
|
||||||
|
Type: "text",
|
||||||
|
Required: true,
|
||||||
|
Placeholder: "请输入邮件内容",
|
||||||
|
Description: "发送邮件的内容主体",
|
||||||
|
},
|
||||||
|
},
|
||||||
|
}
|
||||||
|
```
|
||||||
|
|
||||||
|
## 无参数 Handler
|
||||||
|
|
||||||
|
放在对应业务模块,例如 `internal/apps/upload/tasks.go`。
|
||||||
|
|
||||||
|
```go
|
||||||
|
package upload
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
|
||||||
|
"github.com/Rain-kl/Wavelet/internal/task"
|
||||||
|
)
|
||||||
|
|
||||||
|
type CleanupUnusedUploadsHandler struct{}
|
||||||
|
|
||||||
|
func (h *CleanupUnusedUploadsHandler) Execute(ctx context.Context, payload []byte) (*task.TaskResult, error) {
|
||||||
|
task.AppendLog(ctx, "开始扫描未使用上传")
|
||||||
|
|
||||||
|
// 调用 model/service 完成业务逻辑。
|
||||||
|
// 批量处理时按批次记录日志,不要每条记录都 AppendLog。
|
||||||
|
|
||||||
|
msg := "清理完成"
|
||||||
|
task.AppendLog(ctx, "%s", msg)
|
||||||
|
return &task.TaskResult{Message: msg}, nil
|
||||||
|
}
|
||||||
|
```
|
||||||
|
|
||||||
|
## 带参数 Handler
|
||||||
|
|
||||||
|
实现 `PayloadValidator` 做 Admin 下发时的服务端校验和标准化。`Execute` 仍然解析 payload,因为 Scheduler 和 Retry 不一定经过 Admin 校验路径。
|
||||||
|
|
||||||
|
```go
|
||||||
|
package user
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"encoding/json"
|
||||||
|
"errors"
|
||||||
|
"fmt"
|
||||||
|
"strings"
|
||||||
|
|
||||||
|
"github.com/Rain-kl/Wavelet/internal/task"
|
||||||
|
)
|
||||||
|
|
||||||
|
type SendEmailPayload struct {
|
||||||
|
To string `json:"to"`
|
||||||
|
Subject string `json:"subject"`
|
||||||
|
Body string `json:"body"`
|
||||||
|
}
|
||||||
|
|
||||||
|
type SendEmailHandler struct{}
|
||||||
|
|
||||||
|
func (h *SendEmailHandler) ValidatePayload(payload []byte) ([]byte, error) {
|
||||||
|
if len(payload) == 0 {
|
||||||
|
return nil, errors.New("任务参数不能为空")
|
||||||
|
}
|
||||||
|
|
||||||
|
var req SendEmailPayload
|
||||||
|
if err := json.Unmarshal(payload, &req); err != nil {
|
||||||
|
return nil, fmt.Errorf("无效的 JSON 格式: %w", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
req.To = strings.TrimSpace(req.To)
|
||||||
|
req.Subject = strings.TrimSpace(req.Subject)
|
||||||
|
req.Body = strings.TrimSpace(req.Body)
|
||||||
|
if req.To == "" || req.Subject == "" || req.Body == "" {
|
||||||
|
return nil, errors.New("to、subject、body 不能为空")
|
||||||
|
}
|
||||||
|
|
||||||
|
return json.Marshal(req)
|
||||||
|
}
|
||||||
|
|
||||||
|
func (h *SendEmailHandler) Execute(ctx context.Context, payload []byte) (*task.TaskResult, error) {
|
||||||
|
var req SendEmailPayload
|
||||||
|
if err := json.Unmarshal(payload, &req); err != nil {
|
||||||
|
return nil, fmt.Errorf("解析任务参数: %w", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
task.AppendLog(ctx, "开始发送邮件到: %s", req.To)
|
||||||
|
|
||||||
|
// 调用业务服务发送邮件。
|
||||||
|
|
||||||
|
msg := fmt.Sprintf("邮件成功发送至: %s", req.To)
|
||||||
|
task.AppendLog(ctx, "%s", msg)
|
||||||
|
return &task.TaskResult{Message: msg}, nil
|
||||||
|
}
|
||||||
|
```
|
||||||
|
|
||||||
|
## 统一注册
|
||||||
|
|
||||||
|
在 `internal/task/handlers/register.go` 注册。Admin dispatch 的 `ValidateAndNormalizePayload` 和 Worker 执行都依赖这里。
|
||||||
|
|
||||||
|
```go
|
||||||
|
package handlers
|
||||||
|
|
||||||
|
import (
|
||||||
|
"github.com/Rain-kl/Wavelet/internal/apps/upload"
|
||||||
|
"github.com/Rain-kl/Wavelet/internal/apps/user"
|
||||||
|
"github.com/Rain-kl/Wavelet/internal/task"
|
||||||
|
)
|
||||||
|
|
||||||
|
func Register() {
|
||||||
|
task.RegisterHandler(task.CleanupUnusedUploadsTask, &upload.CleanupUnusedUploadsHandler{})
|
||||||
|
task.RegisterHandler(task.SendEmailTask, &user.SendEmailHandler{})
|
||||||
|
}
|
||||||
|
```
|
||||||
|
|
||||||
|
## Worker 路由
|
||||||
|
|
||||||
|
在 `internal/task/worker/worker.go` 的 mux 上添加任务类型。所有业务任务都指向 `task.ProcessTask`。
|
||||||
|
|
||||||
|
```go
|
||||||
|
func StartWorker() error {
|
||||||
|
asynqServer := asynq.NewServer(task.RedisOpt, asynq.Config{
|
||||||
|
Concurrency: config.Config.Worker.Concurrency,
|
||||||
|
ShutdownTimeout: workerShutdownTimeout,
|
||||||
|
Queues: buildQueuesFromConfig(),
|
||||||
|
StrictPriority: config.Config.Worker.StrictPriority,
|
||||||
|
})
|
||||||
|
|
||||||
|
mux := asynq.NewServeMux()
|
||||||
|
mux.Use(taskLoggingMiddleware)
|
||||||
|
mux.HandleFunc(task.CleanupUnusedUploadsTask, task.ProcessTask)
|
||||||
|
mux.HandleFunc(task.SendEmailTask, task.ProcessTask)
|
||||||
|
|
||||||
|
return asynqServer.Run(mux)
|
||||||
|
}
|
||||||
|
```
|
||||||
|
|
||||||
|
## Cron 调度和配置
|
||||||
|
|
||||||
|
如果任务需要定时运行,补齐 scheduler、config struct 和 `config.example.yaml`。
|
||||||
|
|
||||||
|
```go
|
||||||
|
const (
|
||||||
|
cleanupDedupWindow = 23 * time.Hour
|
||||||
|
cleanupMaxRetry = 3
|
||||||
|
)
|
||||||
|
|
||||||
|
if _, err = scheduler.Register(
|
||||||
|
config.Config.Scheduler.CleanupUnusedUploadsTaskCron,
|
||||||
|
asynq.NewTask(task.CleanupUnusedUploadsTask, nil),
|
||||||
|
asynq.Unique(cleanupDedupWindow),
|
||||||
|
asynq.MaxRetry(cleanupMaxRetry),
|
||||||
|
); err != nil {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
```
|
||||||
|
|
||||||
|
```go
|
||||||
|
type schedulerConfig struct {
|
||||||
|
CleanupUnusedUploadsTaskCron string `mapstructure:"cleanup_unused_uploads_task_cron"`
|
||||||
|
}
|
||||||
|
```
|
||||||
|
|
||||||
|
```yaml
|
||||||
|
scheduler:
|
||||||
|
cleanup_unused_uploads_task_cron: "0 */2 * * *"
|
||||||
|
```
|
||||||
|
|
||||||
|
## Handler 测试
|
||||||
|
|
||||||
|
带参数任务至少覆盖合法 payload、空 payload、非法 JSON、缺失必填和标准化。
|
||||||
|
|
||||||
|
```go
|
||||||
|
func TestSendEmailHandlerValidatePayload(t *testing.T) {
|
||||||
|
tests := []struct {
|
||||||
|
name string
|
||||||
|
payload []byte
|
||||||
|
want SendEmailPayload
|
||||||
|
wantErr bool
|
||||||
|
}{
|
||||||
|
{
|
||||||
|
name: "valid payload is normalized",
|
||||||
|
payload: []byte(`{"to":" user@example.com ","subject":" hi ","body":" body "}`),
|
||||||
|
want: SendEmailPayload{
|
||||||
|
To: "user@example.com",
|
||||||
|
Subject: "hi",
|
||||||
|
Body: "body",
|
||||||
|
},
|
||||||
|
},
|
||||||
|
{
|
||||||
|
name: "empty payload",
|
||||||
|
payload: nil,
|
||||||
|
wantErr: true,
|
||||||
|
},
|
||||||
|
{
|
||||||
|
name: "invalid json",
|
||||||
|
payload: []byte(`{`),
|
||||||
|
wantErr: true,
|
||||||
|
},
|
||||||
|
{
|
||||||
|
name: "missing required field",
|
||||||
|
payload: []byte(`{"to":"user@example.com","subject":"","body":"body"}`),
|
||||||
|
wantErr: true,
|
||||||
|
},
|
||||||
|
}
|
||||||
|
|
||||||
|
h := &SendEmailHandler{}
|
||||||
|
for _, tt := range tests {
|
||||||
|
t.Run(tt.name, func(t *testing.T) {
|
||||||
|
gotPayload, err := h.ValidatePayload(tt.payload)
|
||||||
|
if gotErr := err != nil; gotErr != tt.wantErr {
|
||||||
|
t.Fatalf("ValidatePayload(%s) error = %v, want error presence = %t", tt.payload, err, tt.wantErr)
|
||||||
|
}
|
||||||
|
if tt.wantErr {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
var got SendEmailPayload
|
||||||
|
if err := json.Unmarshal(gotPayload, &got); err != nil {
|
||||||
|
t.Fatalf("json.Unmarshal(%s) error = %v", gotPayload, err)
|
||||||
|
}
|
||||||
|
if diff := cmp.Diff(tt.want, got); diff != "" {
|
||||||
|
t.Errorf("ValidatePayload(%s) mismatch (-want +got):\n%s", tt.payload, diff)
|
||||||
|
}
|
||||||
|
})
|
||||||
|
}
|
||||||
|
}
|
||||||
|
```
|
||||||
|
|
||||||
|
`Execute` 测试优先验证业务服务调用、错误返回和结果摘要;日志可只验证关键路径,避免把精确日志文本写成脆弱断言。
|
||||||
|
|
||||||
|
## Admin Dispatch 测试形状
|
||||||
|
|
||||||
|
Admin dispatch 测试关注通用链路是否调用了 `PayloadValidator`,不要为每种任务在 handler 里写 if 分支。
|
||||||
|
|
||||||
|
```go
|
||||||
|
func TestDispatchTaskValidatesPayload(t *testing.T) {
|
||||||
|
// 1. 初始化测试 DB 和 task.AsynqClient。
|
||||||
|
// 2. 注册测试 handler: task.RegisterHandler(task.SendEmailTask, &user.SendEmailHandler{})
|
||||||
|
// 3. POST /api/v1/admin/tasks/dispatch,传入非法 payload。
|
||||||
|
// 4. 断言响应为 400,错误信息清晰,且没有创建可执行任务。
|
||||||
|
}
|
||||||
|
```
|
||||||
|
|
||||||
|
需要 Redis/Asynq 时优先复用项目现有测试模式;没有现成依赖时可用 `miniredis` 初始化 `task.AsynqClient`。不要把 `internal/task` 依赖塞进通用 testhelper 造成 import cycle。
|
||||||
Reference in New Issue
Block a user