解耦任务框架与业务任务

This commit is contained in:
ryan
2026-06-11 08:50:38 +08:00
parent 5dedff1324
commit 3023d47eec
11 changed files with 214 additions and 177 deletions
+1 -1
View File
@@ -37,7 +37,7 @@ func init() {
// @Failure 403 {object} util.ResponseAny "无管理员权限"
// @Router /api/v1/admin/tasks/types [get]
func ListTaskTypes(c *gin.Context) {
c.JSON(http.StatusOK, util.OK(task.DispatchableTasks))
c.JSON(http.StatusOK, util.OK(task.GetDispatchableTasks()))
}
// DispatchTaskRequest 下发任务请求
+8 -6
View File
@@ -15,6 +15,8 @@ import (
"time"
"github.com/Rain-kl/Wavelet/internal/apps/oauth"
"github.com/Rain-kl/Wavelet/internal/apps/upload"
"github.com/Rain-kl/Wavelet/internal/apps/user"
"github.com/Rain-kl/Wavelet/internal/model"
"github.com/Rain-kl/Wavelet/internal/task"
"github.com/Rain-kl/Wavelet/internal/testhelper"
@@ -88,13 +90,13 @@ func TestListTaskTypes(t *testing.T) {
foundCleanup := false
for _, m := range taskMetas {
if m.Type == task.TaskTypeCleanupUploads {
if m.Type == upload.TaskTypeCleanupUploads {
foundCleanup = true
break
}
}
if !foundCleanup {
t.Errorf("expected task type %s to be listed", task.TaskTypeCleanupUploads)
t.Errorf("expected task type %s to be listed", upload.TaskTypeCleanupUploads)
}
}
@@ -107,7 +109,7 @@ func TestDispatchTask(t *testing.T) {
t.Run("dispatch valid task successfully", func(t *testing.T) {
payload := DispatchTaskRequest{
TaskType: task.TaskTypeCleanupUploads,
TaskType: upload.TaskTypeCleanupUploads,
}
body, _ := json.Marshal(payload)
req, _ := http.NewRequest("POST", "/api/v1/admin/tasks/dispatch", bytes.NewBuffer(body))
@@ -130,7 +132,7 @@ func TestDispatchTask(t *testing.T) {
t.Run("dispatch send_email task successfully with valid payload", func(t *testing.T) {
payload := DispatchTaskRequest{
TaskType: task.TaskTypeSendEmail,
TaskType: user.TaskTypeSendEmail,
Payload: `{"to":"receiver@example.com","subject":"Test Subject","body":"Test Body"}`,
}
body, _ := json.Marshal(payload)
@@ -149,7 +151,7 @@ func TestDispatchTask(t *testing.T) {
t.Run("dispatch send_email task failure with invalid payload json", func(t *testing.T) {
payload := DispatchTaskRequest{
TaskType: task.TaskTypeSendEmail,
TaskType: user.TaskTypeSendEmail,
Payload: `{"to":`,
}
body, _ := json.Marshal(payload)
@@ -166,7 +168,7 @@ func TestDispatchTask(t *testing.T) {
t.Run("dispatch send_email task failure with missing fields", func(t *testing.T) {
payload := DispatchTaskRequest{
TaskType: task.TaskTypeSendEmail,
TaskType: user.TaskTypeSendEmail,
Payload: `{"to":"","subject":"Test","body":"Test"}`,
}
body, _ := json.Marshal(payload)
+20
View File
@@ -16,6 +16,26 @@ import (
"gorm.io/gorm"
)
// 异步任务名称与管理类型定义
const (
// CleanupUnusedUploadsTask 清理未使用上传任务标识
CleanupUnusedUploadsTask = "upload:cleanup_unused"
// TaskTypeCleanupUploads 清理未使用上传管理类型
TaskTypeCleanupUploads = "cleanup_unused_uploads"
)
// CleanupUnusedUploadsMeta represents the task metadata.
var CleanupUnusedUploadsMeta = task.TaskMeta{
Type: TaskTypeCleanupUploads,
AsynqTask: CleanupUnusedUploadsTask,
Name: "清理未使用上传",
Description: "清理超过1小时未使用的上传文件",
SupportsTime: false,
MaxRetry: task.DefaultMaxRetry,
Queue: task.QueueDefault,
Retryable: true,
}
// CleanupUnusedUploadsHandler 清理未使用上传文件的异步任务处理器
type CleanupUnusedUploadsHandler struct{}
+1 -1
View File
@@ -109,7 +109,7 @@ func sendEmailVerificationCode(ctx context.Context, email, scene, templateName s
Body: emailBody,
}
payloadBytes, _ := json.Marshal(payload)
_, err = task.DispatchTask(ctx, task.TaskTypeSendEmail, payloadBytes, "system")
_, err = task.DispatchTask(ctx, TaskTypeSendEmail, payloadBytes, "system")
if err != nil {
return errors.New(errDispatchEmailTaskFailed)
}
+46 -1
View File
@@ -1,4 +1,3 @@
// Copyright 2025 linux.do
// Copyright 2026 Arctel.net
// SPDX-License-Identifier: Apache-2.0
@@ -17,6 +16,52 @@ import (
"github.com/Rain-kl/Wavelet/internal/util/mail"
)
// 异步任务名称与管理类型定义
const (
// SendEmailTask 发送邮件任务标识
SendEmailTask = "mail:send"
// TaskTypeSendEmail 发送邮件管理类型
TaskTypeSendEmail = "send_email"
)
// SendEmailMeta represents the task metadata.
var SendEmailMeta = task.TaskMeta{
Type: TaskTypeSendEmail,
AsynqTask: SendEmailTask,
Name: "发送邮件",
Description: "异步发送系统邮件",
SupportsTime: false,
MaxRetry: task.DefaultMaxRetry,
Queue: task.QueueDefault,
Retryable: true,
Params: []task.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: "请输入邮件内容(支持 HTML 格式)",
Description: "发送邮件的内容主体",
},
},
}
// SendEmailPayload 邮件发送任务载荷
type SendEmailPayload struct {
To string `json:"to"`
+2 -111
View File
@@ -5,119 +5,10 @@
// Package task 定义异步任务类型与调度常量
package task
// 异步任务类型标识
const (
CleanupUnusedUploadsTask = "upload:cleanup_unused"
SendEmailTask = "mail:send"
)
// 任务队列名称
const (
QueueDefault = "default"
)
// 管理员可下发的任务类型标识
const (
TaskTypeCleanupUploads = "cleanup_unused_uploads"
TaskTypeSendEmail = "send_email"
)
// defaultMaxRetry 任务默认最大重试次数
const defaultMaxRetry = 3
// TaskParam 任务参数定义
//
//nolint:revive // TaskParam 保留完整名称以避免与通用 Param 混淆
type TaskParam struct {
Name string `json:"Name"` // 参数键名
Label string `json:"Label"` // 显示名称
Type string `json:"Type"` // 类型:string, text, number, boolean
Required bool `json:"Required"` // 是否必填
Placeholder string `json:"Placeholder"` // 占位符
Description string `json:"Description"` // 描述
}
// TaskMeta 任务元数据
//
//nolint:revive // TaskMeta 保留完整名称以避免与通用 Meta 混淆
type TaskMeta struct {
Type string
AsynqTask string
Name string
Description string
SupportsTime bool
MaxRetry int
Queue string
Retryable bool // 是否支持手动重试
Params []TaskParam
}
// DispatchableTasks 可下发的任务列表
var DispatchableTasks = []TaskMeta{
{
Type: TaskTypeCleanupUploads,
AsynqTask: CleanupUnusedUploadsTask,
Name: "清理未使用上传",
Description: "清理超过1小时未使用的上传文件",
SupportsTime: false,
MaxRetry: defaultMaxRetry,
Queue: QueueDefault,
Retryable: true,
},
{
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: "请输入邮件内容(支持 HTML 格式)",
Description: "发送邮件的内容主体",
},
},
},
}
// GetTaskMeta 根据任务类型获取元数据
func GetTaskMeta(taskType string) *TaskMeta {
for _, t := range DispatchableTasks {
if t.Type == taskType {
return &t
}
}
return nil
}
// GetTaskMetaByAsynqTask 根据 Asynq 任务名称获取元数据
func GetTaskMetaByAsynqTask(asynqTask string) *TaskMeta {
for _, t := range DispatchableTasks {
if t.AsynqTask == asynqTask {
return &t
}
}
return nil
}
// DefaultMaxRetry 任务默认最大重试次数
const DefaultMaxRetry = 3
+8 -3
View File
@@ -11,8 +11,13 @@ import (
"github.com/Rain-kl/Wavelet/internal/task"
)
// Register registers all built-in task handlers.
// Register registers all built-in task handlers and their metadata.
func Register() {
task.RegisterHandler(task.CleanupUnusedUploadsTask, &upload.CleanupUnusedUploadsHandler{})
task.RegisterHandler(task.SendEmailTask, &user.SendEmailHandler{})
// upload
task.RegisterHandler(upload.CleanupUnusedUploadsTask, &upload.CleanupUnusedUploadsHandler{})
task.RegisterTaskMeta(upload.CleanupUnusedUploadsMeta)
// user
task.RegisterHandler(user.SendEmailTask, &user.SendEmailHandler{})
task.RegisterTaskMeta(user.SendEmailMeta)
}
+92
View File
@@ -0,0 +1,92 @@
// Copyright 2026 Arctel.net
// SPDX-License-Identifier: Apache-2.0
package task
import (
"sync"
)
// TaskParam 任务参数定义
//
//nolint:revive // TaskParam 保留完整名称以避免与通用 Param 混淆
type TaskParam struct {
Name string `json:"Name"` // 参数键名
Label string `json:"Label"` // 显示名称
Type string `json:"Type"` // 类型:string, text, number, boolean
Required bool `json:"Required"` // 是否必填
Placeholder string `json:"Placeholder"` // 占位符
Description string `json:"Description"` // 描述
}
// TaskMeta 任务元数据
//
//nolint:revive // TaskMeta 保留完整名称以避免与通用 Meta 混淆
type TaskMeta struct {
Type string
AsynqTask string
Name string
Description string
SupportsTime bool
MaxRetry int
Queue string
Retryable bool // 是否支持手动重试
Params []TaskParam
}
var (
dispatchableTasksMutex sync.RWMutex
dispatchableTasks []TaskMeta
)
// RegisterTaskMeta 注册任务元数据到全局列表
func RegisterTaskMeta(meta TaskMeta) {
dispatchableTasksMutex.Lock()
defer dispatchableTasksMutex.Unlock()
dispatchableTasks = append(dispatchableTasks, meta)
}
// GetDispatchableTasks 获取所有已注册的元数据列表(返回副本以避免并发并发读写冲突)
func GetDispatchableTasks() []TaskMeta {
dispatchableTasksMutex.RLock()
defer dispatchableTasksMutex.RUnlock()
metas := make([]TaskMeta, len(dispatchableTasks))
copy(metas, dispatchableTasks)
return metas
}
// GetTaskMeta 根据任务类型获取元数据
func GetTaskMeta(taskType string) *TaskMeta {
dispatchableTasksMutex.RLock()
defer dispatchableTasksMutex.RUnlock()
for _, t := range dispatchableTasks {
if t.Type == taskType {
copied := t
return &copied
}
}
return nil
}
// GetTaskMetaByAsynqTask 根据 Asynq 任务名称获取元数据
func GetTaskMetaByAsynqTask(asynqTask string) *TaskMeta {
dispatchableTasksMutex.RLock()
defer dispatchableTasksMutex.RUnlock()
for _, t := range dispatchableTasks {
if t.AsynqTask == asynqTask {
copied := t
return &copied
}
}
return nil
}
// GetRegisteredAsynqTasks 返回所有已注册的 Asynq 任务名称,以便动态注册路由
func GetRegisteredAsynqTasks() []string {
keys := make([]string, 0, len(handlerRegistry))
for k := range handlerRegistry {
keys = append(keys, k)
}
return keys
}
+4 -2
View File
@@ -39,8 +39,10 @@ func StartWorker() error {
// 统一使用 task.ProcessTask 处理所有任务类型
// 框架内部自动分发到对应的 TaskHandler 实现
mux.HandleFunc(task.CleanupUnusedUploadsTask, task.ProcessTask)
mux.HandleFunc(task.SendEmailTask, task.ProcessTask)
// 动态注册所有已注册的任务处理器路由,框架内部自动分发到对应的 TaskHandler 实现
for _, taskName := range task.GetRegisteredAsynqTasks() {
mux.HandleFunc(taskName, task.ProcessTask)
}
// 启动服务器
return asynqServer.Run(mux)