mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-09-29 05:56:38 +08:00
fix(task): restore task metadata contract and type fields in task types api
This commit is contained in:
@@ -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.
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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"
|
||||
},
|
||||
|
||||
@@ -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"
|
||||
},
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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")
|
||||
}
|
||||
}
|
||||
@@ -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"})
|
||||
|
||||
@@ -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",
|
||||
},
|
||||
},
|
||||
}
|
||||
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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()
|
||||
|
||||
@@ -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"})
|
||||
|
||||
@@ -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: "目标接收者",
|
||||
},
|
||||
},
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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 系统定期垃圾清理异步任务处理器
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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)",
|
||||
},
|
||||
},
|
||||
}
|
||||
|
||||
|
||||
@@ -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",
|
||||
},
|
||||
},
|
||||
|
||||
@@ -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{
|
||||
|
||||
@@ -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 任务名称,以便动态注册路由
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
|
||||
@@ -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) {
|
||||
|
||||
Reference in New Issue
Block a user