mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-09-29 22:06:38 +08:00
feat(task): clean task log
This commit is contained in:
@@ -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
|
||||
}
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user