diff --git a/backend/core/contracts/task.go b/backend/core/contracts/task.go index 838d6afa..60a051c1 100644 --- a/backend/core/contracts/task.go +++ b/backend/core/contracts/task.go @@ -71,6 +71,14 @@ type TaskExecutionDTO struct { UpdatedAt time.Time `json:"updated_at"` } +// Canonical triggered_by values persisted on task executions and shown in admin UI. +const ( + TaskTriggerSystem = "system" + TaskTriggerManual = "manual" + TaskTriggerRetry = "retry" + TaskTriggerSchedule = "schedule" +) + // TaskService defines the unified contract for dispatching and tracking background tasks. type TaskService interface { Dispatch(ctx context.Context, taskType string, payload []byte, triggeredBy string) (string, error) diff --git a/backend/plugins/domain/admin/service/task.go b/backend/plugins/domain/admin/service/task.go index 3226c31f..d092d4dd 100644 --- a/backend/plugins/domain/admin/service/task.go +++ b/backend/plugins/domain/admin/service/task.go @@ -41,7 +41,7 @@ func DispatchTask(ctx context.Context, req model.DispatchTaskRequest) (string, e return "", err } - taskID, err := taskSvc.Dispatch(ctx, req.TaskType, validated, "manual") + taskID, err := taskSvc.Dispatch(ctx, req.TaskType, validated, contracts.TaskTriggerManual) if err != nil { return "", fmt.Errorf("%s: %w", errs.TaskDispatchFailed, err) } @@ -110,6 +110,21 @@ func TaskExecution(ctx context.Context, id uint64) (*model.TaskExecution, error) return &row, nil } +func normalizeTaskTrigger(v string) string { + switch v { + case contracts.TaskTriggerManual, contracts.TaskTriggerSystem, contracts.TaskTriggerRetry, contracts.TaskTriggerSchedule: + return v + case "inproc_cron", "cron": + return contracts.TaskTriggerSchedule + case "http": + return contracts.TaskTriggerSystem + case "": + return contracts.TaskTriggerSystem + default: + return v + } +} + func executionFromDTO(dto contracts.TaskExecutionDTO) model.TaskExecution { return model.TaskExecution{ ID: dto.ID, @@ -127,7 +142,7 @@ func executionFromDTO(dto contracts.TaskExecutionDTO) model.TaskExecution { FinishedAt: dto.FinishedAt, Duration: dto.Duration, Payload: dto.Payload, - TriggeredBy: dto.TriggeredBy, + TriggeredBy: normalizeTaskTrigger(dto.TriggeredBy), CreatedAt: dto.CreatedAt, UpdatedAt: dto.UpdatedAt, } diff --git a/backend/plugins/domain/admin/service/task_trigger_test.go b/backend/plugins/domain/admin/service/task_trigger_test.go new file mode 100644 index 00000000..4667c9be --- /dev/null +++ b/backend/plugins/domain/admin/service/task_trigger_test.go @@ -0,0 +1,32 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package service + +import ( + "Wavelet/core/contracts" + "testing" +) + +func TestNormalizeTaskTrigger(t *testing.T) { + tests := []struct { + in string + want string + }{ + {in: contracts.TaskTriggerManual, want: contracts.TaskTriggerManual}, + {in: contracts.TaskTriggerSystem, want: contracts.TaskTriggerSystem}, + {in: contracts.TaskTriggerRetry, want: contracts.TaskTriggerRetry}, + {in: contracts.TaskTriggerSchedule, want: contracts.TaskTriggerSchedule}, + {in: "http", want: contracts.TaskTriggerSystem}, + {in: "inproc_cron", want: contracts.TaskTriggerSchedule}, + {in: "cron", want: contracts.TaskTriggerSchedule}, + {in: "", want: contracts.TaskTriggerSystem}, + {in: "custom", want: "custom"}, + } + for _, tt := range tests { + got := normalizeTaskTrigger(tt.in) + if got != tt.want { + t.Errorf("normalizeTaskTrigger(%q) = %q, want %q", tt.in, got, tt.want) + } + } +} diff --git a/backend/plugins/domain/message_gateway/service/push.go b/backend/plugins/domain/message_gateway/service/push.go index e7d6c646..38e903b5 100644 --- a/backend/plugins/domain/message_gateway/service/push.go +++ b/backend/plugins/domain/message_gateway/service/push.go @@ -697,7 +697,7 @@ func EnqueuePushTask(ctx context.Context, payload model.SendPayload) error { return err } if taskSvc := GetTaskService(ctx); taskSvc != nil { - _, err = taskSvc.Dispatch(ctx, "send_notification", payloadBytes, "system") + _, err = taskSvc.Dispatch(ctx, "send_notification", payloadBytes, contracts.TaskTriggerSystem) return err } return errors.New(errs.ErrTaskServiceUnavailable) diff --git a/backend/plugins/domain/user/handlers.go b/backend/plugins/domain/user/handlers.go index 4df8b824..75b9ace0 100644 --- a/backend/plugins/domain/user/handlers.go +++ b/backend/plugins/domain/user/handlers.go @@ -209,7 +209,7 @@ func SendEmailCode(c *gin.Context) { } ctx := c.Request.Context() if taskSvc := getTaskService(ctx); taskSvc != nil { - if _, err := taskSvc.Dispatch(ctx, TaskTypeSendEmailCode, payload, "http"); err != nil { + if _, err := taskSvc.Dispatch(ctx, TaskTypeSendEmailCode, payload, contracts.TaskTriggerSystem); err != nil { logger.ErrorF(ctx, "dispatch send_email_code failed: %v", err) response.AbortInternal(c, errSendEmailFailed) return diff --git a/backend/plugins/drivers/driver_asynq_worker/executor.go b/backend/plugins/drivers/driver_asynq_worker/executor.go index 7da4c364..6d214bf1 100644 --- a/backend/plugins/drivers/driver_asynq_worker/executor.go +++ b/backend/plugins/drivers/driver_asynq_worker/executor.go @@ -4,6 +4,7 @@ package driver_asynq_worker import ( + "Wavelet/core/contracts" "Wavelet/pkg/idgen" "Wavelet/pkg/logger" "Wavelet/pkg/util" @@ -208,7 +209,7 @@ func RetryTask(ctx context.Context, id uint64) (string, error) { MaxRetry: execution.MaxRetry, RetryCount: execution.RetryCount + 1, Payload: execution.Payload, - TriggeredBy: "retry", + TriggeredBy: contracts.TaskTriggerRetry, } if err := createTaskExecution(ctx, newExecution); err != nil { @@ -393,7 +394,7 @@ func getOrCreateTaskExecution(ctx context.Context, taskID string, t *asynq.Task, MaxRetry: meta.MaxRetry, RetryCount: 0, Payload: string(payload), - TriggeredBy: "schedule", + TriggeredBy: contracts.TaskTriggerSchedule, StartedAt: &now, } diff --git a/backend/plugins/drivers/driver_inproc_cron/scheduler.go b/backend/plugins/drivers/driver_inproc_cron/scheduler.go index 8b3b33f0..b0760c7e 100644 --- a/backend/plugins/drivers/driver_inproc_cron/scheduler.go +++ b/backend/plugins/drivers/driver_inproc_cron/scheduler.go @@ -99,7 +99,7 @@ func (s *inprocScheduler) registerJob(ctx context.Context, def extpoints.Schedul _, err := s.cronRunner.AddFunc(cronSpec, func() { if s.taskSvc != nil { - if _, dispatchErr := s.taskSvc.Dispatch(ctx, taskType, payloadBytes, "inproc_cron"); dispatchErr != nil { + if _, dispatchErr := s.taskSvc.Dispatch(ctx, taskType, payloadBytes, contracts.TaskTriggerSchedule); dispatchErr != nil { logger.ErrorF(ctx, "driver_inproc_cron: dispatch task %q failed: %v", taskType, dispatchErr) } return diff --git a/backend/plugins/drivers/driver_inproc_worker/executor.go b/backend/plugins/drivers/driver_inproc_worker/executor.go index d22b2363..1a57688b 100644 --- a/backend/plugins/drivers/driver_inproc_worker/executor.go +++ b/backend/plugins/drivers/driver_inproc_worker/executor.go @@ -76,7 +76,7 @@ func (q *InprocQueue) Enqueue(ctx context.Context, taskType string, payload []by } if source == "" { - source = "manual" + source = contracts.TaskTriggerManual } idType := td.Type if idType == "" { diff --git a/frontend/app/(main)/admin/tasks/components/task-executions.tsx b/frontend/app/(main)/admin/tasks/components/task-executions.tsx index 9c294eac..49609354 100644 --- a/frontend/app/(main)/admin/tasks/components/task-executions.tsx +++ b/frontend/app/(main)/admin/tasks/components/task-executions.tsx @@ -85,6 +85,16 @@ function statusVariant(status: TaskExecutionStatus) { return 'outline'; } +function mappedLabel( + t: (key: string) => string, + keys: Record, + value: string, +): string { + const key = keys[value]; + if (!key) return value || '-'; + return t(key); +} + export function TaskExecutionsManager() { const t = useTranslations('admin.tasks'); const queryClient = useQueryClient(); @@ -281,14 +291,16 @@ export function TaskExecutionsManager() { - {t(STATUS_LABELS_KEYS[execution.status]) || - execution.status} + {mappedLabel(t, STATUS_LABELS_KEYS, execution.status)} - {t(TRIGGER_LABELS_KEYS[execution.triggered_by]) || - execution.triggered_by} + {mappedLabel( + t, + TRIGGER_LABELS_KEYS, + execution.triggered_by, + )} @@ -368,8 +380,11 @@ export function TaskExecutionsManager() {
- {t(STATUS_LABELS_KEYS[selectedExecution.status]) || - selectedExecution.status} + {mappedLabel( + t, + STATUS_LABELS_KEYS, + selectedExecution.status, + )}
@@ -378,8 +393,11 @@ export function TaskExecutionsManager() { {t('detailTrigger')}
- {t(TRIGGER_LABELS_KEYS[selectedExecution.triggered_by]) || - selectedExecution.triggered_by} + {mappedLabel( + t, + TRIGGER_LABELS_KEYS, + selectedExecution.triggered_by, + )}