From f826807cf13be87b193504f9097fce4e442db1b9 Mon Sep 17 00:00:00 2001 From: ryan Date: Wed, 17 Jun 2026 14:29:38 +0800 Subject: [PATCH] feat(task): clean task log --- internal/apps/upload/cleanup.go | 25 ++++++++++-- internal/apps/upload/tasks_test.go | 21 +++++++++- internal/model/task_execution.go | 59 +++++++++++++++++++++++++++ internal/model/task_execution_test.go | 51 +++++++++++++++++++++++ 4 files changed, 151 insertions(+), 5 deletions(-) diff --git a/internal/apps/upload/cleanup.go b/internal/apps/upload/cleanup.go index e3121ad2..adfa80c1 100644 --- a/internal/apps/upload/cleanup.go +++ b/internal/apps/upload/cleanup.go @@ -35,7 +35,7 @@ var SystemCleanupMeta = task.TaskMeta{ Type: TaskTypeSystemCleanup, AsynqTask: SystemCleanupTask, Name: "系统垃圾清理", - Description: "定期清理超过1小时的未使用上传文件和超过7天的历史推送记录", + Description: "定期清理未使用上传文件、历史推送记录和过期任务执行日志", SupportsTime: false, MaxRetry: task.DefaultMaxRetry, Queue: task.QueueDefault, @@ -45,7 +45,7 @@ var SystemCleanupMeta = task.TaskMeta{ // SystemCleanupHandler 系统定期垃圾清理异步任务处理器 type SystemCleanupHandler struct{} -// Execute 执行系统清理(包含文件清理和历史消息推送日志清理) +// Execute 执行系统清理(包含文件清理、历史推送日志和任务执行日志清理) func (h *SystemCleanupHandler) Execute(ctx context.Context, _ []byte) (*task.TaskResult, error) { if storageReadOnly(ctx) { return nil, errors.New(errStorageReadOnly) @@ -132,7 +132,26 @@ func (h *SystemCleanupHandler) Execute(ctx context.Context, _ []byte) (*task.Tas task.AppendLog(ctx, "没有需要清理的历史推送记录 (截止时间: %s)", cutoff.Format("2006-01-02 15:04:05")) } - msg := fmt.Sprintf("系统清理完成。成功清理未使用的上传文件 %d/%d 个;清理历史推送审计日志 %d 条。", totalDeleted, totalProcessed, pushHistoryCount) + // 3. 清理任务执行日志:高频任务保留3天,低频任务保留30天。 + task.AppendLog(ctx, "开始清理任务执行日志:高频任务保留最近3天,低频任务保留最近30天...") + taskLogStats, err := model.CleanupTaskExecutionLogs(ctx, time.Now()) + if err != nil { + task.AppendLog(ctx, "清理任务执行日志失败: %v", err) + logger.ErrorF(ctx, "清理任务执行日志失败: %v", err) + } else { + task.AppendLog(ctx, "成功清理任务执行日志 %d 条(高频 %d 条,低频 %d 条)", + taskLogStats.HighFrequencyDeleted+taskLogStats.LowFrequencyDeleted, + taskLogStats.HighFrequencyDeleted, + taskLogStats.LowFrequencyDeleted, + ) + } + + msg := fmt.Sprintf("系统清理完成。成功清理未使用的上传文件 %d/%d 个;清理历史推送审计日志 %d 条;清理任务执行日志 %d 条。", + totalDeleted, + totalProcessed, + pushHistoryCount, + taskLogStats.HighFrequencyDeleted+taskLogStats.LowFrequencyDeleted, + ) task.AppendLog(ctx, "%s", msg) return &task.TaskResult{Message: msg}, nil } diff --git a/internal/apps/upload/tasks_test.go b/internal/apps/upload/tasks_test.go index 83476cbb..27bb5439 100644 --- a/internal/apps/upload/tasks_test.go +++ b/internal/apps/upload/tasks_test.go @@ -109,6 +109,18 @@ func TestSystemCleanupHandler_Execute(t *testing.T) { err = db.DB(ctx).Create(newPush).Error require.NoError(t, err) + oldTaskLog := &model.TaskExecution{ + TaskID: "old_low_frequency_task_log", + TaskType: "low:frequency", + TaskName: "低频任务", + Status: model.TaskExecutionStatusSucceeded, + CreatedAt: now.AddDate(0, 0, -31), + UpdatedAt: now.AddDate(0, 0, -31), + TriggeredBy: "system", + } + err = model.CreateTaskExecution(ctx, oldTaskLog) + require.NoError(t, err) + // 执行 handler handler := &SystemCleanupHandler{} result, err := handler.Execute(ctx, nil) @@ -116,7 +128,7 @@ func TestSystemCleanupHandler_Execute(t *testing.T) { // 验证结果 require.NoError(t, err) require.NotNil(t, result) - assert.Contains(t, result.Message, "系统清理完成。成功清理未使用的上传文件 2/2 个;清理历史推送审计日志 1 条。") + assert.Contains(t, result.Message, "系统清理完成。成功清理未使用的上传文件 2/2 个;清理历史推送审计日志 1 条;清理任务执行日志 1 条。") // 验证数据库状态:pending 且超过1小时的应被标记为 deleted var pendingCount int64 @@ -140,6 +152,11 @@ func TestSystemCleanupHandler_Execute(t *testing.T) { err = db.DB(ctx).First(&remainingPush).Error require.NoError(t, err) assert.Equal(t, "New Login", remainingPush.Title) + + var taskLogCount int64 + err = db.DB(ctx).Model(&model.TaskExecution{}).Where("task_id = ?", "old_low_frequency_task_log").Count(&taskLogCount).Error + require.NoError(t, err) + assert.Equal(t, int64(0), taskLogCount, "过期低频任务日志应被清理") } func TestSystemCleanupHandler_ExecuteNoFiles(t *testing.T) { @@ -166,7 +183,7 @@ func TestSystemCleanupHandler_ExecuteNoFiles(t *testing.T) { require.NoError(t, err) require.NotNil(t, result) - assert.Contains(t, result.Message, "系统清理完成。成功清理未使用的上传文件 0/0 个;清理历史推送审计日志 0 条。") + assert.Contains(t, result.Message, "系统清理完成。成功清理未使用的上传文件 0/0 个;清理历史推送审计日志 0 条;清理任务执行日志 0 条。") } func TestSystemCleanupHandler_ImplementsTaskHandler(t *testing.T) { diff --git a/internal/model/task_execution.go b/internal/model/task_execution.go index 81a1d193..4c0aac57 100644 --- a/internal/model/task_execution.go +++ b/internal/model/task_execution.go @@ -53,6 +53,12 @@ type TaskExecution struct { UpdatedAt time.Time `json:"updated_at" gorm:"autoUpdateTime"` } +// TaskExecutionCleanupStats describes task execution log cleanup results. +type TaskExecutionCleanupStats struct { + HighFrequencyDeleted int64 + LowFrequencyDeleted int64 +} + // TableName 表名 func (TaskExecution) TableName() string { return "w_task_executions" @@ -190,6 +196,59 @@ func ListTaskExecutions(ctx context.Context, req ListTaskExecutionsRequest) ([]T return executions, total, nil } +// CleanupTaskExecutionLogs removes finished task execution logs according to frequency-based retention. +func CleanupTaskExecutionLogs(ctx context.Context, now time.Time) (TaskExecutionCleanupStats, error) { + const ( + frequencyWindowDays = 30 + highFrequencyThreshold = frequencyWindowDays + ) + + frequencyWindowStart := now.AddDate(0, 0, -frequencyWindowDays) + highFrequencyCutoff := now.AddDate(0, 0, -3) + lowFrequencyCutoff := now.AddDate(0, 0, -30) + terminalStatuses := []TaskExecutionStatus{TaskExecutionStatusSucceeded, TaskExecutionStatusFailed} + + var highFrequencyTaskTypes []string + if err := db.DB(ctx). + Model(&TaskExecution{}). + Select("task_type"). + Where("created_at >= ?", frequencyWindowStart). + Group("task_type"). + Having("COUNT(*) > ?", highFrequencyThreshold). + Pluck("task_type", &highFrequencyTaskTypes).Error; err != nil { + return TaskExecutionCleanupStats{}, fmt.Errorf("query high-frequency task types: %w", err) + } + + var highFrequencyDeleted int64 + if len(highFrequencyTaskTypes) > 0 { + highFrequencyResult := db.DB(ctx). + Where("status IN ?", terminalStatuses). + Where("created_at < ?", highFrequencyCutoff). + Where("task_type IN ?", highFrequencyTaskTypes). + Delete(&TaskExecution{}) + if highFrequencyResult.Error != nil { + return TaskExecutionCleanupStats{}, fmt.Errorf("delete high-frequency task execution logs: %w", highFrequencyResult.Error) + } + highFrequencyDeleted = highFrequencyResult.RowsAffected + } + + lowFrequencyQuery := db.DB(ctx). + Where("status IN ?", terminalStatuses). + Where("created_at < ?", lowFrequencyCutoff) + if len(highFrequencyTaskTypes) > 0 { + lowFrequencyQuery = lowFrequencyQuery.Where("task_type NOT IN ?", highFrequencyTaskTypes) + } + lowFrequencyResult := lowFrequencyQuery.Delete(&TaskExecution{}) + if lowFrequencyResult.Error != nil { + return TaskExecutionCleanupStats{}, fmt.Errorf("delete low-frequency task execution logs: %w", lowFrequencyResult.Error) + } + + return TaskExecutionCleanupStats{ + HighFrequencyDeleted: highFrequencyDeleted, + LowFrequencyDeleted: lowFrequencyResult.RowsAffected, + }, nil +} + func taskExecutionLogRedisKey(taskID string) string { return db.PrefixedKey(taskExecutionLogRedisKeyPrefix + taskID) } diff --git a/internal/model/task_execution_test.go b/internal/model/task_execution_test.go index 3c99ba48..dc769462 100644 --- a/internal/model/task_execution_test.go +++ b/internal/model/task_execution_test.go @@ -427,7 +427,58 @@ func TestListTaskExecutionsDefaultPaging(t *testing.T) { assert.Len(t, items, 0) } +func TestCleanupTaskExecutionLogs(t *testing.T) { + cleanup := setupTaskExecutionTestEnvironment(t) + defer cleanup() + ctx := context.Background() + + now := time.Date(2026, 6, 17, 12, 0, 0, 0, time.UTC) + for i := 0; i < 31; i++ { + createTaskExecutionForCleanup(t, ctx, fmt.Sprintf("high_recent_%02d", i), "high:task", TaskExecutionStatusSucceeded, now.Add(-2*time.Hour)) + } + createTaskExecutionForCleanup(t, ctx, "high_old_4d", "high:task", TaskExecutionStatusSucceeded, now.AddDate(0, 0, -4)) + createTaskExecutionForCleanup(t, ctx, "high_old_40d", "high:task", TaskExecutionStatusFailed, now.AddDate(0, 0, -40)) + createTaskExecutionForCleanup(t, ctx, "high_running_old", "high:task", TaskExecutionStatusRunning, now.AddDate(0, 0, -10)) + createTaskExecutionForCleanup(t, ctx, "low_old_31d", "low:task", TaskExecutionStatusSucceeded, now.AddDate(0, 0, -31)) + createTaskExecutionForCleanup(t, ctx, "low_recent_29d", "low:task", TaskExecutionStatusSucceeded, now.AddDate(0, 0, -29)) + createTaskExecutionForCleanup(t, ctx, "low_pending_old", "low:task", TaskExecutionStatusPending, now.AddDate(0, 0, -45)) + + stats, err := CleanupTaskExecutionLogs(ctx, now) + require.NoError(t, err) + assert.Equal(t, int64(2), stats.HighFrequencyDeleted) + assert.Equal(t, int64(1), stats.LowFrequencyDeleted) + + for _, taskID := range []string{"high_old_4d", "high_old_40d", "low_old_31d"} { + var count int64 + err := db.DB(ctx).Model(&TaskExecution{}).Where("task_id = ?", taskID).Count(&count).Error + require.NoError(t, err) + assert.Equal(t, int64(0), count, "CleanupTaskExecutionLogs(%s) should delete expired log", taskID) + } + for _, taskID := range []string{"high_recent_00", "high_running_old", "low_recent_29d", "low_pending_old"} { + var count int64 + err := db.DB(ctx).Model(&TaskExecution{}).Where("task_id = ?", taskID).Count(&count).Error + require.NoError(t, err) + assert.Equal(t, int64(1), count, "CleanupTaskExecutionLogs(%s) should keep retained log", taskID) + } +} + func TestTaskExecutionTableName(t *testing.T) { execution := TaskExecution{} assert.Equal(t, "w_task_executions", execution.TableName()) } + +func createTaskExecutionForCleanup(t *testing.T, ctx context.Context, taskID string, taskType string, status TaskExecutionStatus, createdAt time.Time) { + t.Helper() + + execution := &TaskExecution{ + TaskID: taskID, + TaskType: taskType, + TaskName: taskType, + Status: status, + CreatedAt: createdAt, + UpdatedAt: createdAt, + TriggeredBy: "system", + } + err := CreateTaskExecution(ctx, execution) + require.NoError(t, err) +}