From ea97b64407f121e6d723753c880b959b8d4df0af Mon Sep 17 00:00:00 2001 From: ryan Date: Sat, 29 Aug 2026 12:22:51 +0800 Subject: [PATCH] fix(task): restore task metadata contract and type fields in task types api --- backend/core/contracts/task.go | 26 +-- backend/core/extpoints/extpoints_test.go | 29 +++ backend/core/extpoints/task.go | 167 +++++++++++++++++- backend/docs/docs.go | 18 ++ backend/docs/swagger.json | 18 ++ backend/docs/swagger.yaml | 12 ++ .../domain/admin/handler/tasks_test.go | 97 ++++++++++ backend/plugins/domain/admin/plugin.go | 16 +- .../domain/admin/service/log_switch.go | 24 ++- backend/plugins/domain/admin/service/task.go | 5 +- backend/plugins/domain/domain_test.go | 6 + .../plugins/domain/message_gateway/plugin.go | 30 +++- .../domain/message_gateway/service/push.go | 22 ++- backend/plugins/domain/upload/plugin.go | 8 +- backend/plugins/domain/upload/task/cleanup.go | 16 +- .../domain/upload/task/rebuild_stats.go | 16 +- .../domain/upload/task/storage_migration.go | 30 ++-- backend/plugins/domain/upload/task/tasks.go | 18 +- backend/plugins/domain/user/plugin.go | 80 ++++++++- .../drivers/driver_asynq_worker/meta.go | 110 +++++++++--- .../drivers/driver_asynq_worker/plugin.go | 68 +++++-- .../driver_inproc_worker/task_service.go | 21 +-- 22 files changed, 704 insertions(+), 133 deletions(-) create mode 100644 backend/plugins/domain/admin/handler/tasks_test.go diff --git a/backend/core/contracts/task.go b/backend/core/contracts/task.go index b61b8003..3fe1c01c 100644 --- a/backend/core/contracts/task.go +++ b/backend/core/contracts/task.go @@ -12,23 +12,29 @@ import ( // TaskParamDTO describes a parameter accepted by a background task. type TaskParamDTO struct { Name string `json:"name"` + Label string `json:"label"` Type string `json:"type"` - Description string `json:"description"` Required bool `json:"required"` + Placeholder string `json:"placeholder,omitempty"` + Description string `json:"description,omitempty"` Default any `json:"default,omitempty"` } // TaskMetaDTO describes the metadata and configuration of a registered background task. type TaskMetaDTO struct { - Name string `json:"name"` - DisplayName string `json:"display_name"` - Description string `json:"description"` - Category string `json:"category"` - Params []TaskParamDTO `json:"params,omitempty"` - MaxRetry int `json:"max_retry"` - Timeout time.Duration `json:"timeout"` - Queue string `json:"queue"` - Schedule string `json:"schedule,omitempty"` + Type string `json:"type"` + AsynqTask string `json:"asynq_task"` + Name string `json:"name"` + DisplayName string `json:"display_name,omitempty"` + Description string `json:"description"` + Category string `json:"category,omitempty"` + SupportsTime bool `json:"supports_time"` + Params []TaskParamDTO `json:"params,omitempty"` + MaxRetry int `json:"max_retry"` + Timeout time.Duration `json:"timeout,omitempty"` + Queue string `json:"queue"` + Retryable bool `json:"retryable"` + Schedule string `json:"schedule,omitempty"` } // TaskResultDTO represents the outcome of a background task execution. diff --git a/backend/core/extpoints/extpoints_test.go b/backend/core/extpoints/extpoints_test.go index adfc3e85..6d51b999 100644 --- a/backend/core/extpoints/extpoints_test.go +++ b/backend/core/extpoints/extpoints_test.go @@ -173,13 +173,29 @@ func TestTaskExtension(t *testing.T) { extpoints.WithTaskRetry(3), extpoints.WithTaskTimeout(10*time.Second), extpoints.WithTaskMetadata("queue", "critical"), + extpoints.WithTaskType("cancel_timeout"), + extpoints.WithTaskName("取消超时订单"), + extpoints.WithTaskDescription("自动关单"), + extpoints.WithTaskCategory("order"), + extpoints.WithTaskSupportsTime(true), + extpoints.WithTaskQueue("orders"), + extpoints.WithTaskRetryable(true), nil, // test nil option ) // Re-register to test update tr.Register("order:cancel_timeout", handler, extpoints.WithTaskConcurrency(10), + extpoints.WithTaskRetry(3), + extpoints.WithTaskTimeout(10*time.Second), extpoints.WithTaskMetadata("queue", "high"), + extpoints.WithTaskType("cancel_timeout"), + extpoints.WithTaskName("取消超时订单"), + extpoints.WithTaskDescription("自动关单"), + extpoints.WithTaskCategory("order"), + extpoints.WithTaskSupportsTime(true), + extpoints.WithTaskQueue("orders"), + extpoints.WithTaskRetryable(true), ) tasks := tr.Tasks() @@ -187,6 +203,19 @@ func TestTaskExtension(t *testing.T) { assert.Equal(t, "order:cancel_timeout", tasks[0].Pattern) assert.Equal(t, 10, tasks[0].Concurrency) assert.Equal(t, "high", tasks[0].Metadata["queue"]) + assert.Equal(t, "cancel_timeout", tasks[0].Type) + assert.Equal(t, "取消超时订单", tasks[0].Name) + + dto := tasks[0].ToDTO() + assert.Equal(t, "cancel_timeout", dto.Type) + assert.Equal(t, "order:cancel_timeout", dto.AsynqTask) + assert.Equal(t, "取消超时订单", dto.Name) + assert.Equal(t, "取消超时订单", dto.DisplayName) + assert.Equal(t, "自动关单", dto.Description) + assert.Equal(t, "order", dto.Category) + assert.True(t, dto.SupportsTime) + assert.Equal(t, "orders", dto.Queue) + assert.True(t, dto.Retryable) task, ok := tr.Get("order:cancel_timeout") assert.True(t, ok) diff --git a/backend/core/extpoints/task.go b/backend/core/extpoints/task.go index 7dafe3c9..1217ed64 100644 --- a/backend/core/extpoints/task.go +++ b/backend/core/extpoints/task.go @@ -4,23 +4,137 @@ package extpoints import ( + "Wavelet/core/contracts" "sync" "time" ) // TaskDefinition holds the definition and runtime options for an asynchronous background task. type TaskDefinition struct { - Pattern string - Handler any - Concurrency int - Retry int - Timeout time.Duration - Metadata map[string]any + Pattern string + Type string + Name string + DisplayName string + Description string + Category string + SupportsTime bool + Retryable bool + Queue string + Params []contracts.TaskParamDTO + Handler any + Concurrency int + Retry int + Timeout time.Duration + Metadata map[string]any } // TaskOption configures a TaskDefinition. type TaskOption func(*TaskDefinition) +// WithTaskType sets the admin task type identifier. +func WithTaskType(taskType string) TaskOption { + return func(td *TaskDefinition) { + td.Type = taskType + } +} + +// WithTaskName sets the task human-readable display name. +func WithTaskName(name string) TaskOption { + return func(td *TaskDefinition) { + td.Name = name + if td.DisplayName == "" { + td.DisplayName = name + } + } +} + +// WithTaskDisplayName sets the task display name. +func WithTaskDisplayName(displayName string) TaskOption { + return func(td *TaskDefinition) { + td.DisplayName = displayName + if td.Name == "" { + td.Name = displayName + } + } +} + +// WithTaskDescription sets the task description. +func WithTaskDescription(desc string) TaskOption { + return func(td *TaskDefinition) { + td.Description = desc + } +} + +// WithTaskCategory sets the task category grouping. +func WithTaskCategory(category string) TaskOption { + return func(td *TaskDefinition) { + td.Category = category + } +} + +// WithTaskSupportsTime sets whether the task supports time range filtering. +func WithTaskSupportsTime(supports bool) TaskOption { + return func(td *TaskDefinition) { + td.SupportsTime = supports + } +} + +// WithTaskQueue sets the task queue. +func WithTaskQueue(queue string) TaskOption { + return func(td *TaskDefinition) { + td.Queue = queue + } +} + +// WithTaskRetryable sets whether the task is retryable. +func WithTaskRetryable(retryable bool) TaskOption { + return func(td *TaskDefinition) { + td.Retryable = retryable + } +} + +// WithTaskParams sets the task parameter definitions. +func WithTaskParams(params ...contracts.TaskParamDTO) TaskOption { + return func(td *TaskDefinition) { + td.Params = append(td.Params, params...) + } +} + +// WithTaskMeta sets all task metadata from a TaskMetaDTO. +func WithTaskMeta(meta contracts.TaskMetaDTO) TaskOption { + return func(td *TaskDefinition) { + if meta.Type != "" { + td.Type = meta.Type + } + if meta.Name != "" { + td.Name = meta.Name + } + if meta.DisplayName != "" { + td.DisplayName = meta.DisplayName + } + if meta.Description != "" { + td.Description = meta.Description + } + if meta.Category != "" { + td.Category = meta.Category + } + td.SupportsTime = meta.SupportsTime + if meta.Queue != "" { + td.Queue = meta.Queue + } + td.Retryable = meta.Retryable + if meta.MaxRetry > 0 { + td.Retry = meta.MaxRetry + } + if meta.Timeout > 0 { + td.Timeout = meta.Timeout + } + if len(meta.Params) > 0 { + td.Params = append([]contracts.TaskParamDTO(nil), meta.Params...) + } + } +} + // WithTaskConcurrency sets the concurrency limit for the task. func WithTaskConcurrency(concurrency int) TaskOption { return func(td *TaskDefinition) { @@ -52,6 +166,47 @@ func WithTaskMetadata(key string, val any) TaskOption { } } +// ToDTO converts TaskDefinition to contracts.TaskMetaDTO. +func (td TaskDefinition) ToDTO() contracts.TaskMetaDTO { + taskType := td.Type + if taskType == "" { + taskType = td.Pattern + } + name := td.Name + if name == "" { + name = td.DisplayName + } + if name == "" { + name = td.Pattern + } + displayName := td.DisplayName + if displayName == "" { + displayName = name + } + queue := td.Queue + if queue == "" { + queue = "default" + } + retryable := td.Retryable + if !retryable && td.Retry > 0 { + retryable = true + } + return contracts.TaskMetaDTO{ + Type: taskType, + AsynqTask: td.Pattern, + Name: name, + DisplayName: displayName, + Description: td.Description, + Category: td.Category, + SupportsTime: td.SupportsTime, + Params: td.Params, + MaxRetry: td.Retry, + Timeout: td.Timeout, + Queue: queue, + Retryable: retryable, + } +} + // TaskExtension defines the interface for registering and querying background task handlers. type TaskExtension interface { Register(pattern string, handler any, opts ...TaskOption) diff --git a/backend/docs/docs.go b/backend/docs/docs.go index 7535ec1c..e3b023ed 100644 --- a/backend/docs/docs.go +++ b/backend/docs/docs.go @@ -4103,6 +4103,9 @@ const docTemplate = `{ "contracts.TaskMetaDTO": { "type": "object", "properties": { + "asynq_task": { + "type": "string" + }, "category": { "type": "string" }, @@ -4127,11 +4130,20 @@ const docTemplate = `{ "queue": { "type": "string" }, + "retryable": { + "type": "boolean" + }, "schedule": { "type": "string" }, + "supports_time": { + "type": "boolean" + }, "timeout": { "$ref": "#/definitions/time.Duration" + }, + "type": { + "type": "string" } } }, @@ -4142,9 +4154,15 @@ const docTemplate = `{ "description": { "type": "string" }, + "label": { + "type": "string" + }, "name": { "type": "string" }, + "placeholder": { + "type": "string" + }, "required": { "type": "boolean" }, diff --git a/backend/docs/swagger.json b/backend/docs/swagger.json index d997397b..e87db4e5 100644 --- a/backend/docs/swagger.json +++ b/backend/docs/swagger.json @@ -4096,6 +4096,9 @@ "contracts.TaskMetaDTO": { "type": "object", "properties": { + "asynq_task": { + "type": "string" + }, "category": { "type": "string" }, @@ -4120,11 +4123,20 @@ "queue": { "type": "string" }, + "retryable": { + "type": "boolean" + }, "schedule": { "type": "string" }, + "supports_time": { + "type": "boolean" + }, "timeout": { "$ref": "#/definitions/time.Duration" + }, + "type": { + "type": "string" } } }, @@ -4135,9 +4147,15 @@ "description": { "type": "string" }, + "label": { + "type": "string" + }, "name": { "type": "string" }, + "placeholder": { + "type": "string" + }, "required": { "type": "boolean" }, diff --git a/backend/docs/swagger.yaml b/backend/docs/swagger.yaml index b83f68f1..119fe1c1 100644 --- a/backend/docs/swagger.yaml +++ b/backend/docs/swagger.yaml @@ -49,6 +49,8 @@ definitions: type: object contracts.TaskMetaDTO: properties: + asynq_task: + type: string category: type: string description: @@ -65,18 +67,28 @@ definitions: type: array queue: type: string + retryable: + type: boolean schedule: type: string + supports_time: + type: boolean timeout: $ref: '#/definitions/time.Duration' + type: + type: string type: object contracts.TaskParamDTO: properties: default: {} description: type: string + label: + type: string name: type: string + placeholder: + type: string required: type: boolean type: diff --git a/backend/plugins/domain/admin/handler/tasks_test.go b/backend/plugins/domain/admin/handler/tasks_test.go new file mode 100644 index 00000000..b513ece6 --- /dev/null +++ b/backend/plugins/domain/admin/handler/tasks_test.go @@ -0,0 +1,97 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package handler_test + +import ( + "Wavelet/core" + "Wavelet/core/contracts" + "Wavelet/core/extpoints" + "Wavelet/plugins/domain/admin/handler" + "Wavelet/plugins/domain/admin/service" + "Wavelet/plugins/drivers/driver_asynq_worker" + "encoding/json" + "net/http" + "net/http/httptest" + "testing" + + "github.com/gin-gonic/gin" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +type listTaskTypesResponse struct { + ErrorMsg string `json:"error_msg"` + Data []contracts.TaskMetaDTO `json:"data"` +} + +func TestListTaskTypesHandler(t *testing.T) { + gin.SetMode(gin.TestMode) + + ctx := core.NewContext(t.Context()) + worker := driver_asynq_worker.New() + require.NoError(t, worker.Apply(ctx)) + + ctx.Task().Register("logs:db_switch", func(_ any) error { return nil }, + extpoints.WithTaskType("logs_db_switch"), + extpoints.WithTaskName("切换日志数据库"), + extpoints.WithTaskDescription("复制迁移用户访问日志并在成功后切换日志主库"), + extpoints.WithTaskCategory("system"), + extpoints.WithTaskRetry(3), + extpoints.WithTaskQueue("default"), + extpoints.WithTaskRetryable(true), + extpoints.WithTaskParams(contracts.TaskParamDTO{ + Name: "target", + Label: "目标日志库", + Type: "string", + Required: true, + Placeholder: "postgres|sqlite|clickhouse", + Description: "迁移目标", + }), + ) + + taskSvc, err := core.Inject[contracts.TaskService](ctx) + require.NoError(t, err) + service.SetTaskService(taskSvc) + + r := gin.New() + r.GET("/api/v1/admin/tasks/types", handler.ListTaskTypes) + + req := httptest.NewRequest(http.MethodGet, "/api/v1/admin/tasks/types", nil) + w := httptest.NewRecorder() + r.ServeHTTP(w, req) + + assert.Equal(t, http.StatusOK, w.Code) + + var resp listTaskTypesResponse + require.NoError(t, json.Unmarshal(w.Body.Bytes(), &resp)) + assert.Empty(t, resp.ErrorMsg) + require.NotEmpty(t, resp.Data) + + var found bool + for _, task := range resp.Data { + if task.Type == "logs_db_switch" { + found = true + assert.Equal(t, "logs:db_switch", task.AsynqTask) + assert.Equal(t, "切换日志数据库", task.Name) + assert.Equal(t, "复制迁移用户访问日志并在成功后切换日志主库", task.Description) + assert.Equal(t, "system", task.Category) + assert.Equal(t, 3, task.MaxRetry) + assert.Equal(t, "default", task.Queue) + assert.True(t, task.Retryable) + require.Len(t, task.Params, 1) + assert.Equal(t, "target", task.Params[0].Name) + assert.Equal(t, "目标日志库", task.Params[0].Label) + assert.Equal(t, "string", task.Params[0].Type) + assert.True(t, task.Params[0].Required) + assert.Equal(t, "postgres|sqlite|clickhouse", task.Params[0].Placeholder) + } + } + assert.True(t, found, "expected logs_db_switch in task types") + + // Ensure NO task in the list has an empty Type or Name, which breaks SelectItem key/value in frontend + for _, task := range resp.Data { + assert.NotEmpty(t, task.Type, "task.type must never be empty") + assert.NotEmpty(t, task.Name, "task.name must never be empty") + } +} diff --git a/backend/plugins/domain/admin/plugin.go b/backend/plugins/domain/admin/plugin.go index af9b122d..ecbb30b4 100644 --- a/backend/plugins/domain/admin/plugin.go +++ b/backend/plugins/domain/admin/plugin.go @@ -171,9 +171,23 @@ func (p *Plugin) Apply(ctx *core.Context) error { handler.RegisterRoutes(adminRouter) // 2. Register Background Tasks + logSwitchHandler := &service.LogDBSwitchHandler{} + ctx.Task().Register(service.LogDBSwitchTask, func(c context.Context, payload []byte) error { + _, err := logSwitchHandler.Execute(c, payload) + return err + }, extpoints.WithTaskMeta(service.LogDBSwitchMeta)) + ctx.Task().Register("admin:system_cleanup", func(_ context.Context, _ []byte) error { return nil - }, extpoints.WithTaskRetry(1)) + }, + extpoints.WithTaskType("system_cleanup"), + extpoints.WithTaskName("系统垃圾清理"), + extpoints.WithTaskDescription("定期清理未使用上传文件、历史推送记录和过期任务执行日志"), + extpoints.WithTaskCategory("maintenance"), + extpoints.WithTaskRetry(1), + extpoints.WithTaskQueue("default"), + extpoints.WithTaskRetryable(true), + ) // 3. Register Cron Schedules ctx.Schedule().RegisterCron("0 4 * * *", "admin:system_cleanup", map[string]string{"type": "daily"}) diff --git a/backend/plugins/domain/admin/service/log_switch.go b/backend/plugins/domain/admin/service/log_switch.go index 6145191b..792eb95f 100644 --- a/backend/plugins/domain/admin/service/log_switch.go +++ b/backend/plugins/domain/admin/service/log_switch.go @@ -31,13 +31,25 @@ const ( // LogDBSwitchMeta 描述切换日志数据库任务。 var LogDBSwitchMeta = contracts.TaskMetaDTO{ - Name: LogDBSwitchTask, - DisplayName: "切换日志数据库", - Description: "复制迁移用户访问日志并在成功后切换日志主库(期间禁止日志写入)", - MaxRetry: 3, - Queue: "default", + Type: TaskTypeLogDBSwitch, + AsynqTask: LogDBSwitchTask, + Name: "切换日志数据库", + DisplayName: "切换日志数据库", + Description: "复制迁移用户访问日志并在成功后切换日志主库(期间禁止日志写入)", + Category: "system", + SupportsTime: false, + MaxRetry: 3, + Queue: "default", + Retryable: true, Params: []contracts.TaskParamDTO{ - {Name: "target", Description: "迁移目标:postgres(主库为 PG 时)、sqlite(主库为 SQLite 时)或 clickhouse", Type: "string", Required: true}, + { + Name: "target", + Label: "目标日志库", + Type: "string", + Required: true, + Placeholder: "postgres|sqlite|clickhouse", + Description: "迁移目标:postgres(主库为 PG 时)、sqlite(主库为 SQLite 时)或 clickhouse", + }, }, } diff --git a/backend/plugins/domain/admin/service/task.go b/backend/plugins/domain/admin/service/task.go index 02150fe7..2543ee85 100644 --- a/backend/plugins/domain/admin/service/task.go +++ b/backend/plugins/domain/admin/service/task.go @@ -32,12 +32,11 @@ func DispatchTask(ctx context.Context, req model.DispatchTaskRequest) (string, e return "", err } - meta, ok := taskSvc.GetTaskMeta(req.TaskType) - if !ok { + if _, ok := taskSvc.GetTaskMeta(req.TaskType); !ok { return "", errs.ErrInvalidTaskType } - validated, err := validateTaskPayload(taskSvc, meta.Name, req.Payload) + validated, err := validateTaskPayload(taskSvc, req.TaskType, req.Payload) if err != nil { return "", err } diff --git a/backend/plugins/domain/domain_test.go b/backend/plugins/domain/domain_test.go index 93a57740..9f3f1b6a 100644 --- a/backend/plugins/domain/domain_test.go +++ b/backend/plugins/domain/domain_test.go @@ -446,6 +446,12 @@ func TestAllDomainPluginsCombined(t *testing.T) { // Verify total tasks registered allTasks := ctx.Tasks().Tasks() assert.GreaterOrEqual(t, len(allTasks), 4) + for _, task := range allTasks { + dto := task.ToDTO() + assert.NotEmpty(t, dto.Type, "task %s should have type", task.Pattern) + assert.NotEmpty(t, dto.AsynqTask, "task %s should have asynq_task", task.Pattern) + assert.NotEmpty(t, dto.Name, "task %s should have name", task.Pattern) + } // Verify total schedules registered allSchedules := ctx.Schedules().Schedules() diff --git a/backend/plugins/domain/message_gateway/plugin.go b/backend/plugins/domain/message_gateway/plugin.go index fb8e14e0..16c1b463 100644 --- a/backend/plugins/domain/message_gateway/plugin.go +++ b/backend/plugins/domain/message_gateway/plugin.go @@ -165,19 +165,41 @@ func (p *Plugin) Apply(ctx *core.Context) error { // 5. Register background tasks ctx.Task().Register("message_gateway:push_notification", func(c context.Context, payload []byte) error { return pushHandler.Execute(c, payload) - }, extpoints.WithTaskRetry(defaultTaskRetry)) + }, + extpoints.WithTaskType("push_notification"), + extpoints.WithTaskName("消息网关推送通知"), + extpoints.WithTaskDescription("异步执行系统通知的多渠道派发与推送"), + extpoints.WithTaskCategory("push"), + extpoints.WithTaskRetry(defaultTaskRetry), + extpoints.WithTaskQueue("default"), + extpoints.WithTaskRetryable(true), + ) ctx.Task().Register(service.SendNotificationTask, func(c context.Context, payload []byte) error { return pushHandler.Execute(c, payload) - }, extpoints.WithTaskRetry(defaultTaskRetry)) + }, extpoints.WithTaskMeta(service.SendNotificationMeta), extpoints.WithTaskRetry(defaultTaskRetry)) ctx.Task().Register("message_gateway:dispatch_bot_msg", func(_ context.Context, _ []byte) error { return nil - }) + }, + extpoints.WithTaskType("dispatch_bot_msg"), + extpoints.WithTaskName("分发 Bot 消息"), + extpoints.WithTaskDescription("异步处理与分发 Bot 下行消息"), + extpoints.WithTaskCategory("messaging"), + extpoints.WithTaskQueue("default"), + ) ctx.Task().Register("message_gateway:cleanup_pairing_codes", func(c context.Context, _ []byte) error { return repository.DeleteExpiredPairingCodes(c) - }, extpoints.WithTaskRetry(defaultTaskRetry)) + }, + extpoints.WithTaskType("cleanup_pairing_codes"), + extpoints.WithTaskName("清理过期配对码"), + extpoints.WithTaskDescription("定时清理已过期的平台 Bot 配对码"), + extpoints.WithTaskCategory("messaging"), + extpoints.WithTaskRetry(defaultTaskRetry), + extpoints.WithTaskQueue("default"), + extpoints.WithTaskRetryable(true), + ) // 6. Register Cron Schedules ctx.Schedule().RegisterCron("*/10 * * * *", "message_gateway:cleanup_pairing_codes", map[string]any{"action": "cleanup"}) diff --git a/backend/plugins/domain/message_gateway/service/push.go b/backend/plugins/domain/message_gateway/service/push.go index c7e60b7f..6f675b21 100644 --- a/backend/plugins/domain/message_gateway/service/push.go +++ b/backend/plugins/domain/message_gateway/service/push.go @@ -708,23 +708,31 @@ const ( // SendNotificationMeta represents the task metadata. var SendNotificationMeta = contracts.TaskMetaDTO{ - Name: TaskTypeSendNotification, - DisplayName: "推送通知", - Description: "异步执行系统通知的多渠道派发与推送", - MaxRetry: 3, - Queue: "default", + Type: TaskTypeSendNotification, + AsynqTask: SendNotificationTask, + Name: "推送通知", + DisplayName: "推送通知", + Description: "异步执行系统通知的多渠道派发与推送", + Category: "push", + SupportsTime: false, + MaxRetry: 3, + Queue: "default", + Retryable: true, Params: []contracts.TaskParamDTO{ { Name: "event_key", + Label: "事件标识", Type: "string", - Description: "事件标识 (如 admin_login)", Required: true, + Placeholder: "admin_login", + Description: "事件标识 (如 admin_login)", }, { Name: "target", + Label: "目标接收者", Type: "string", - Description: "目标接收者", Required: false, + Description: "目标接收者", }, }, } diff --git a/backend/plugins/domain/upload/plugin.go b/backend/plugins/domain/upload/plugin.go index a4c1f39c..e6b998bb 100644 --- a/backend/plugins/domain/upload/plugin.go +++ b/backend/plugins/domain/upload/plugin.go @@ -150,25 +150,25 @@ func (p *Plugin) Apply(ctx *core.Context) error { ctx.Task().Register(task.SystemCleanupTask, func(c context.Context, payload []byte) error { _, err := cleanupHandler.Execute(c, payload) return err - }, extpoints.WithTaskRetry(defaultCleanupRetry)) + }, extpoints.WithTaskMeta(task.SystemCleanupMeta), extpoints.WithTaskRetry(defaultCleanupRetry)) rebuildStatsHandler := &task.RebuildUploadStatsHandler{} ctx.Task().Register(task.RebuildUploadStatsTask, func(c context.Context, payload []byte) error { _, err := rebuildStatsHandler.Execute(c, payload) return err - }, extpoints.WithTaskRetry(defaultStatsRetry)) + }, extpoints.WithTaskMeta(task.RebuildUploadStatsMeta), extpoints.WithTaskRetry(defaultStatsRetry)) migrationHandler := &task.MigrationHandler{} ctx.Task().Register(task.StorageMigrationTask, func(c context.Context, payload []byte) error { _, err := migrationHandler.Execute(c, payload) return err - }, extpoints.WithTaskRetry(defaultSingleRetry)) + }, extpoints.WithTaskMeta(task.StorageMigrationMeta), extpoints.WithTaskRetry(defaultSingleRetry)) warmHandler := &task.WarmImageCacheHandler{} ctx.Task().Register(task.WarmImageCacheTask, func(c context.Context, payload []byte) error { _, err := warmHandler.Execute(c, payload) return err - }, extpoints.WithTaskRetry(1)) + }, extpoints.WithTaskMeta(task.WarmImageCacheMeta), extpoints.WithTaskRetry(1)) // 4. Register Cron Schedule ctx.Schedule().RegisterCron("0 3 * * *", task.SystemCleanupTask, nil) diff --git a/backend/plugins/domain/upload/task/cleanup.go b/backend/plugins/domain/upload/task/cleanup.go index 4bc3cb15..7b09ee2c 100644 --- a/backend/plugins/domain/upload/task/cleanup.go +++ b/backend/plugins/domain/upload/task/cleanup.go @@ -31,12 +31,16 @@ const ( // SystemCleanupMeta represents the task metadata. var SystemCleanupMeta = contracts.TaskMetaDTO{ - Name: SystemCleanupTask, - DisplayName: "系统垃圾清理", - Description: "定期清理未使用上传文件、历史推送记录和过期任务执行日志", - Category: "maintenance", - MaxRetry: 3, - Queue: taskQueueDefault, + Type: TaskTypeSystemCleanup, + AsynqTask: SystemCleanupTask, + Name: "系统垃圾清理", + DisplayName: "系统垃圾清理", + Description: "定期清理未使用上传文件、历史推送记录和过期任务执行日志", + Category: "maintenance", + SupportsTime: false, + MaxRetry: 3, + Queue: taskQueueDefault, + Retryable: true, } // SystemCleanupHandler 系统定期垃圾清理异步任务处理器 diff --git a/backend/plugins/domain/upload/task/rebuild_stats.go b/backend/plugins/domain/upload/task/rebuild_stats.go index 68f3fe8c..2b192fba 100644 --- a/backend/plugins/domain/upload/task/rebuild_stats.go +++ b/backend/plugins/domain/upload/task/rebuild_stats.go @@ -24,12 +24,16 @@ const ( // RebuildUploadStatsMeta describes the upload stats rebuild task. var RebuildUploadStatsMeta = contracts.TaskMetaDTO{ - Name: RebuildUploadStatsTask, - DisplayName: "重算文件存储统计", - Description: "根据当前 w_uploads 活跃记录全量重建 w_upload_stats(总量、类型、分类、趋势)", - Category: taskCategoryUpload, - MaxRetry: 3, - Queue: taskQueueDefault, + Type: TaskTypeRebuildUploadStats, + AsynqTask: RebuildUploadStatsTask, + Name: "重算文件存储统计", + DisplayName: "重算文件存储统计", + Description: "根据当前 w_uploads 活跃记录全量重建 w_upload_stats(总量、类型、分类、趋势)", + Category: taskCategoryUpload, + SupportsTime: false, + MaxRetry: 2, + Queue: taskQueueDefault, + Retryable: true, } // RebuildUploadStatsHandler rebuilds incremental upload stats from active upload records. diff --git a/backend/plugins/domain/upload/task/storage_migration.go b/backend/plugins/domain/upload/task/storage_migration.go index 1697305a..e0d927e1 100644 --- a/backend/plugins/domain/upload/task/storage_migration.go +++ b/backend/plugins/domain/upload/task/storage_migration.go @@ -36,31 +36,25 @@ const ( // StorageMigrationMeta describes the manually dispatchable migration task. var StorageMigrationMeta = contracts.TaskMetaDTO{ - Name: StorageMigrationTask, - DisplayName: "迁移文件存储", - Description: "将活动存储中的文件迁移到待切换的目标存储,迁移期间文件系统保持只读", - Category: taskCategoryUpload, - MaxRetry: 3, - Queue: taskQueueDefault, + Type: TaskTypeStorageMigration, + AsynqTask: StorageMigrationTask, + Name: "迁移文件存储", + DisplayName: "迁移文件存储", + Description: "将活动存储中的文件迁移到待切换的目标存储,迁移期间文件系统保持只读", + Category: taskCategoryUpload, + SupportsTime: false, + MaxRetry: 1, + Queue: taskQueueDefault, + Retryable: true, Params: []contracts.TaskParamDTO{ { Name: "target", + Label: "目标存储配置 (JSON)", Type: "text", Required: true, + Placeholder: `{"driver": "s3", "local": {"root": "."}, "s3": {"bucket": "my-bucket"}}`, Description: "待迁移到的目标存储引擎完整配置 JSON 字符串", }, - { - Name: "batch_size", - Type: "number", - Required: false, - Description: "每批扫描的文件数量(默认 100)", - }, - { - Name: "concurrency", - Type: "number", - Required: false, - Description: "并发迁移 worker 数量(默认 4)", - }, }, } diff --git a/backend/plugins/domain/upload/task/tasks.go b/backend/plugins/domain/upload/task/tasks.go index 43ec9ae6..dacb08a5 100644 --- a/backend/plugins/domain/upload/task/tasks.go +++ b/backend/plugins/domain/upload/task/tasks.go @@ -34,17 +34,23 @@ var warmImageCacheMu sync.Mutex // WarmImageCacheMeta represents the image cache warmup task metadata. var WarmImageCacheMeta = contracts.TaskMetaDTO{ - Name: WarmImageCacheTask, - DisplayName: "预热图片压缩缓存", - Description: "串行将文件管理中的图片转换为指定质量的 WebP 并写入永久缓存", - Category: taskCategoryUpload, - MaxRetry: 3, - Queue: taskQueueDefault, + Type: TaskTypeWarmImageCache, + AsynqTask: WarmImageCacheTask, + Name: "预热图片压缩缓存", + DisplayName: "预热图片压缩缓存", + Description: "串行将文件管理中的图片转换为指定质量的 WebP 并写入永久缓存", + Category: taskCategoryUpload, + SupportsTime: false, + MaxRetry: 1, + Queue: taskQueueDefault, + Retryable: true, Params: []contracts.TaskParamDTO{ { Name: "quality", + Label: "图片质量", Type: "string", Required: true, + Placeholder: "low / medium / high", Description: "WebP 压缩质量,仅支持 low、medium、high", }, }, diff --git a/backend/plugins/domain/user/plugin.go b/backend/plugins/domain/user/plugin.go index b0a17af5..b9818d54 100644 --- a/backend/plugins/domain/user/plugin.go +++ b/backend/plugins/domain/user/plugin.go @@ -140,16 +140,90 @@ func (p *Plugin) Apply(ctx *core.Context) error { } } - const defaultUserTaskRetry = 3 + const ( + defaultUserTaskRetry = 3 + paramTypeString = "string" + paramNameEmail = "email" + ) // 4. Register background tasks ctx.Task().Register("user:send_email_code", func(_ context.Context, _ []byte) error { return nil - }, extpoints.WithTaskRetry(defaultUserTaskRetry)) + }, + extpoints.WithTaskType("send_email_code"), + extpoints.WithTaskName("发送邮箱验证码"), + extpoints.WithTaskDescription("异步发送用户注册与验证邮箱验证码"), + extpoints.WithTaskCategory("user"), + extpoints.WithTaskRetry(defaultUserTaskRetry), + extpoints.WithTaskQueue("default"), + extpoints.WithTaskRetryable(true), + extpoints.WithTaskParams( + contracts.TaskParamDTO{ + Name: paramNameEmail, + Label: "目标邮箱", + Type: paramTypeString, + Required: true, + Placeholder: "user@example.com", + Description: "接收验证码的目标邮箱", + }, + contracts.TaskParamDTO{ + Name: "code", + Label: "验证码", + Type: paramTypeString, + Required: true, + Placeholder: "123456", + Description: "6 位数字验证码", + }, + ), + ) + + ctx.Task().Register("mail:send", func(_ context.Context, _ []byte) error { + return nil + }, + extpoints.WithTaskType("send_email"), + extpoints.WithTaskName("发送邮件"), + extpoints.WithTaskDescription("异步发送系统邮件"), + extpoints.WithTaskCategory("mail"), + extpoints.WithTaskRetry(defaultUserTaskRetry), + extpoints.WithTaskQueue("default"), + extpoints.WithTaskRetryable(true), + extpoints.WithTaskParams( + contracts.TaskParamDTO{ + Name: "to", + Label: "接收邮箱 (To)", + Type: paramTypeString, + Required: true, + Placeholder: "receiver@example.com", + Description: "接收邮件的目标邮箱地址", + }, + contracts.TaskParamDTO{ + Name: "subject", + Label: "邮件主题 (Subject)", + Type: paramTypeString, + Required: true, + Placeholder: "请输入邮件主题", + Description: "发送邮件的主题标题", + }, + contracts.TaskParamDTO{ + Name: "body", + Label: "邮件内容 (Body)", + Type: "text", + Required: true, + Placeholder: "请输入邮件内容(支持 HTML格式)", + Description: "发送邮件的内容主体", + }, + ), + ) ctx.Task().Register("user:cleanup_inactive", func(_ context.Context, _ []byte) error { return nil - }) + }, + extpoints.WithTaskType("cleanup_inactive_users"), + extpoints.WithTaskName("清理未激活用户"), + extpoints.WithTaskDescription("清理长期未激活的注册用户与临时凭据"), + extpoints.WithTaskCategory("user"), + extpoints.WithTaskQueue("default"), + ) // 5. Register Settings Schemas ctx.Settings().Register(extpoints.SettingSchema{ diff --git a/backend/plugins/drivers/driver_asynq_worker/meta.go b/backend/plugins/drivers/driver_asynq_worker/meta.go index 7dfbfdf7..0504a3d2 100644 --- a/backend/plugins/drivers/driver_asynq_worker/meta.go +++ b/backend/plugins/drivers/driver_asynq_worker/meta.go @@ -4,6 +4,7 @@ package driver_asynq_worker import ( + "Wavelet/core/contracts" "Wavelet/core/extpoints" "sync" ) @@ -28,6 +29,7 @@ type TaskMeta struct { AsynqTask string `json:"asynq_task"` Name string `json:"name"` Description string `json:"description"` + Category string `json:"category,omitempty"` SupportsTime bool `json:"supports_time"` MaxRetry int `json:"max_retry"` Queue string `json:"queue"` @@ -52,16 +54,82 @@ func RegisterTaskMeta(meta TaskMeta) { dispatchableTasks = append(dispatchableTasks, meta) } -// GetDispatchableTasks 获取所有已注册的元数据列表(返回副本以避免并发读写冲突) +// GetDispatchableTasks 获取所有已注册的元数据列表(优先结合 activeTaskReg 和 dispatchableTasks) func GetDispatchableTasks() []TaskMeta { - dispatchableTasksMutex.RLock() - defer dispatchableTasksMutex.RUnlock() + activeTaskRegMutex.RLock() + reg := activeTaskReg + activeTaskRegMutex.RUnlock() + + var metas []TaskMeta + seen := make(map[string]bool) + + if reg != nil { + for _, td := range reg.Tasks() { + dto := td.ToDTO() + m := toInternalTaskMeta(dto) + metas = append(metas, m) + seen[m.Type] = true + seen[m.AsynqTask] = true + } + } + + dispatchableTasksMutex.RLock() + for _, t := range dispatchableTasks { + if !seen[t.Type] && !seen[t.AsynqTask] { + metas = append(metas, t) + seen[t.Type] = true + seen[t.AsynqTask] = true + } + } + dispatchableTasksMutex.RUnlock() - metas := make([]TaskMeta, len(dispatchableTasks)) - copy(metas, dispatchableTasks) return metas } +func toInternalTaskMeta(dto contracts.TaskMetaDTO) TaskMeta { + params := make([]TaskParam, 0, len(dto.Params)) + for _, p := range dto.Params { + params = append(params, TaskParam{ + Name: p.Name, + Label: p.Label, + Type: p.Type, + Required: p.Required, + Placeholder: p.Placeholder, + Description: p.Description, + }) + } + taskType := dto.Type + if taskType == "" { + taskType = dto.AsynqTask + } + if taskType == "" { + taskType = dto.Name + } + asynqTask := dto.AsynqTask + if asynqTask == "" { + asynqTask = taskType + } + name := dto.Name + if name == "" { + name = dto.DisplayName + } + if name == "" { + name = taskType + } + return TaskMeta{ + Type: taskType, + AsynqTask: asynqTask, + Name: name, + Description: dto.Description, + Category: dto.Category, + SupportsTime: dto.SupportsTime, + MaxRetry: dto.MaxRetry, + Queue: dto.Queue, + Retryable: dto.Retryable, + Params: params, + } +} + var ( activeTaskRegMutex sync.RWMutex activeTaskReg extpoints.TaskExtension @@ -80,14 +148,10 @@ func getFromActiveTaskExtension(taskType string) *TaskMeta { if activeTaskReg == nil { return nil } - if td, ok := activeTaskReg.Get(taskType); ok { - return &TaskMeta{ - Type: td.Pattern, - Name: td.Pattern, - AsynqTask: td.Pattern, - Queue: "default", - Retryable: td.Retry > 0, - MaxRetry: td.Retry, + for _, td := range activeTaskReg.Tasks() { + if td.Type == taskType || td.Pattern == taskType { + m := toInternalTaskMeta(td.ToDTO()) + return &m } } return nil @@ -95,9 +159,13 @@ func getFromActiveTaskExtension(taskType string) *TaskMeta { // GetTaskMeta 根据任务类型获取元数据 func GetTaskMeta(taskType string) *TaskMeta { + if m := getFromActiveTaskExtension(taskType); m != nil { + return m + } + dispatchableTasksMutex.RLock() for _, t := range dispatchableTasks { - if t.Type == taskType { + if t.Type == taskType || t.AsynqTask == taskType { copied := t dispatchableTasksMutex.RUnlock() return &copied @@ -105,22 +173,12 @@ func GetTaskMeta(taskType string) *TaskMeta { } dispatchableTasksMutex.RUnlock() - return getFromActiveTaskExtension(taskType) + return nil } // GetTaskMetaByAsynqTask 根据 Asynq 任务名称获取元数据 func GetTaskMetaByAsynqTask(asynqTask string) *TaskMeta { - dispatchableTasksMutex.RLock() - for _, t := range dispatchableTasks { - if t.AsynqTask == asynqTask { - copied := t - dispatchableTasksMutex.RUnlock() - return &copied - } - } - dispatchableTasksMutex.RUnlock() - - return getFromActiveTaskExtension(asynqTask) + return GetTaskMeta(asynqTask) } // GetRegisteredAsynqTasks 返回所有已注册的 Asynq 任务名称,以便动态注册路由 diff --git a/backend/plugins/drivers/driver_asynq_worker/plugin.go b/backend/plugins/drivers/driver_asynq_worker/plugin.go index 8f31d4d9..05905182 100644 --- a/backend/plugins/drivers/driver_asynq_worker/plugin.go +++ b/backend/plugins/drivers/driver_asynq_worker/plugin.go @@ -330,18 +330,40 @@ func (s *taskServiceImpl) ListTasks() []contracts.TaskMetaDTO { for _, param := range m.Params { params = append(params, contracts.TaskParamDTO{ Name: param.Name, + Label: param.Label, Type: param.Type, Description: param.Description, + Placeholder: param.Placeholder, Required: param.Required, }) } + taskType := m.Type + if taskType == "" { + taskType = m.AsynqTask + } + if taskType == "" { + taskType = m.Name + } + asynqTask := m.AsynqTask + if asynqTask == "" { + asynqTask = taskType + } + name := m.Name + if name == "" { + name = taskType + } res = append(res, contracts.TaskMetaDTO{ - Name: m.Type, - DisplayName: m.Name, - Description: m.Description, - Params: params, - MaxRetry: m.MaxRetry, - Queue: m.Queue, + Type: taskType, + AsynqTask: asynqTask, + Name: name, + DisplayName: name, + Description: m.Description, + Category: m.Category, + SupportsTime: m.SupportsTime, + Params: params, + MaxRetry: m.MaxRetry, + Queue: m.Queue, + Retryable: m.Retryable, }) } return res @@ -356,18 +378,40 @@ func (s *taskServiceImpl) GetTaskMeta(taskType string) (contracts.TaskMetaDTO, b for _, param := range m.Params { params = append(params, contracts.TaskParamDTO{ Name: param.Name, + Label: param.Label, Type: param.Type, Description: param.Description, + Placeholder: param.Placeholder, Required: param.Required, }) } + tType := m.Type + if tType == "" { + tType = m.AsynqTask + } + if tType == "" { + tType = m.Name + } + asynqTask := m.AsynqTask + if asynqTask == "" { + asynqTask = tType + } + name := m.Name + if name == "" { + name = tType + } return contracts.TaskMetaDTO{ - Name: m.Type, - DisplayName: m.Name, - Description: m.Description, - Params: params, - MaxRetry: m.MaxRetry, - Queue: m.Queue, + Type: tType, + AsynqTask: asynqTask, + Name: name, + DisplayName: name, + Description: m.Description, + Category: m.Category, + SupportsTime: m.SupportsTime, + Params: params, + MaxRetry: m.MaxRetry, + Queue: m.Queue, + Retryable: m.Retryable, }, true } diff --git a/backend/plugins/drivers/driver_inproc_worker/task_service.go b/backend/plugins/drivers/driver_inproc_worker/task_service.go index bbeb0beb..0f3e7d83 100644 --- a/backend/plugins/drivers/driver_inproc_worker/task_service.go +++ b/backend/plugins/drivers/driver_inproc_worker/task_service.go @@ -36,12 +36,7 @@ func (s *inprocTaskService) ListTasks() []contracts.TaskMetaDTO { tasks := s.taskReg.Tasks() res := make([]contracts.TaskMetaDTO, 0, len(tasks)) for _, td := range tasks { - res = append(res, contracts.TaskMetaDTO{ - Name: td.Pattern, - DisplayName: td.Pattern, - MaxRetry: td.Retry, - Timeout: td.Timeout, - }) + res = append(res, td.ToDTO()) } return res } @@ -50,16 +45,12 @@ func (s *inprocTaskService) GetTaskMeta(taskType string) (contracts.TaskMetaDTO, if s.taskReg == nil { return contracts.TaskMetaDTO{}, false } - td, ok := s.taskReg.Get(taskType) - if !ok { - return contracts.TaskMetaDTO{}, false + for _, td := range s.taskReg.Tasks() { + if td.Type == taskType || td.Pattern == taskType { + return td.ToDTO(), true + } } - return contracts.TaskMetaDTO{ - Name: td.Pattern, - DisplayName: td.Pattern, - MaxRetry: td.Retry, - Timeout: td.Timeout, - }, true + return contracts.TaskMetaDTO{}, false } func (s *inprocTaskService) ValidatePayload(_ string, payload []byte) ([]byte, error) {