From 6b28adfacf7c8d96ec4ae84be8035ab6a196597f Mon Sep 17 00:00:00 2001 From: ryan Date: Thu, 3 Sep 2026 08:50:37 +0800 Subject: [PATCH] refactor(admin): decouple system cleanup with event bus and enforce single owner principle --- backend/core/contracts/events.go | 5 + backend/plugins/domain/admin/plugin.go | 13 +- backend/plugins/domain/admin/plugin_test.go | 4 +- .../plugins/domain/admin/service/cleanup.go | 70 ++++++++++ .../domain/admin/service/cleanup_test.go | 78 +++++++++++ backend/plugins/domain/domain_test.go | 129 +++++++++++++++++- .../plugins/domain/msg_gateway/dao/push.go | 10 ++ .../domain/msg_gateway/dao/push_test.go | 49 +++++++ backend/plugins/domain/msg_gateway/plugin.go | 8 ++ .../domain/msg_gateway/service/push_worker.go | 13 ++ backend/plugins/domain/upload/exports.go | 3 + backend/plugins/domain/upload/plugin.go | 8 +- backend/plugins/domain/upload/task/cleanup.go | 72 ++++------ backend/plugins/domain/user/plugin.go | 7 + 14 files changed, 419 insertions(+), 50 deletions(-) create mode 100644 backend/plugins/domain/admin/service/cleanup.go create mode 100644 backend/plugins/domain/admin/service/cleanup_test.go diff --git a/backend/core/contracts/events.go b/backend/core/contracts/events.go index 8de13c0a..43f52740 100644 --- a/backend/core/contracts/events.go +++ b/backend/core/contracts/events.go @@ -151,3 +151,8 @@ type UserDeletedEvent struct { CurrentUserID uint64 `json:"current_user_id,string"` TargetUserID uint64 `json:"target_user_id,string"` } + +// SystemCleanupEvent fires when a periodic system cleanup is triggered. +type SystemCleanupEvent struct { + TriggeredAt string `json:"triggered_at"` +} diff --git a/backend/plugins/domain/admin/plugin.go b/backend/plugins/domain/admin/plugin.go index 2512a9e5..fcbf9607 100644 --- a/backend/plugins/domain/admin/plugin.go +++ b/backend/plugins/domain/admin/plugin.go @@ -21,6 +21,12 @@ import ( // SystemConfig aliases model.SystemConfig for external compatibility. type SystemConfig = model.SystemConfig +// TaskExecution aliases model.TaskExecution for external compatibility. +type TaskExecution = model.TaskExecution + +// Schedule aliases model.Schedule for external compatibility. +type Schedule = model.Schedule + //go:embed migrations/*/*.sql var adminMigrations embed.FS @@ -132,12 +138,17 @@ func (p *Plugin) Apply(ctx *core.Context) error { ctx.Router().RegisterWhitelist("/robots.txt") // 2. Register Background Tasks + const defaultCleanupRetry = 3 ctx.Task().Register(service.LogDBSwitchTask, &service.LogDBSwitchHandler{}, extpoints.WithTaskMeta(service.LogDBSwitchMeta)) + ctx.Task().Register(service.SystemCleanupTask, &service.SystemCleanupHandler{}, extpoints.WithTaskMeta(service.SystemCleanupMeta), extpoints.WithTaskRetry(defaultCleanupRetry)) + + // 2.1 Register Cron Schedule + ctx.Schedule().RegisterCron("0 3 * * *", service.SystemCleanupTask, nil) // 3. Register Settings Schemas ctx.Settings().Register(extpoints.SettingSchema{ Key: "admin.system_cleanup_cron", - Default: "0 4 * * *", + Default: "0 3 * * *", Description: "Cron expression for nightly system logs and expired tokens cleanup", Type: "string", Category: "maintenance", diff --git a/backend/plugins/domain/admin/plugin_test.go b/backend/plugins/domain/admin/plugin_test.go index 4aace149..dab839a3 100644 --- a/backend/plugins/domain/admin/plugin_test.go +++ b/backend/plugins/domain/admin/plugin_test.go @@ -36,11 +36,13 @@ func TestAdminPluginUnit(t *testing.T) { // Verify tasks _, ok := ctx.Tasks().Get("logs:db_switch") require.True(t, ok) + _, ok = ctx.Tasks().Get("system:cleanup") + require.True(t, ok) // Verify settings setting, ok := ctx.Settings().Get("admin.system_cleanup_cron") require.True(t, ok) - assert.Equal(t, "0 4 * * *", setting.Default) + assert.Equal(t, "0 3 * * *", setting.Default) provider, err := core.Inject[contracts.PublicConfigProvider](ctx) require.NoError(t, err) diff --git a/backend/plugins/domain/admin/service/cleanup.go b/backend/plugins/domain/admin/service/cleanup.go new file mode 100644 index 00000000..18de5ca1 --- /dev/null +++ b/backend/plugins/domain/admin/service/cleanup.go @@ -0,0 +1,70 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package service + +import ( + "Wavelet/core/contracts" + "Wavelet/pkg/logger" + "Wavelet/plugins/domain/admin/model" + "context" + "errors" + "fmt" + "time" +) + +const ( + // SystemCleanupTask 系统定期垃圾清理任务标识 + SystemCleanupTask = "system:cleanup" + // TaskTypeSystemCleanup 系统定期垃圾清理管理类型 + TaskTypeSystemCleanup = "system_cleanup" + taskQueueDefault = "default" +) + +// SystemCleanupMeta describes the system-wide cleanup task metadata. +var SystemCleanupMeta = contracts.TaskMetaDTO{ + Type: TaskTypeSystemCleanup, + AsynqTask: SystemCleanupTask, + Name: "系统垃圾清理", + DisplayName: "系统垃圾清理", + Description: "定期清理过期任务执行记录,并通过领域事件广播触发各业务域自治清理(临时文件、历史推送等)", + Category: "maintenance", + SupportsTime: false, + MaxRetry: 3, + Queue: taskQueueDefault, + Retryable: true, +} + +// SystemCleanupHandler handles the system-wide garbage cleanup task. +type SystemCleanupHandler struct{} + +// Execute executes system cleanup: clears old task executions and emits EventTopicSystemCleanup. +func (h *SystemCleanupHandler) Execute(ctx context.Context, _ []byte) (*contracts.TaskResultDTO, error) { + db := GetDB(ctx) + if db == nil { + return nil, errors.New("database service not available") + } + + // 1. 清理自身域(admin 域)的过期任务执行记录(7 天前) + var deletedExecutions int64 + sevenDaysAgo := time.Now().Add(-7 * 24 * time.Hour) + res := db.Where("created_at < ?", sevenDaysAgo).Delete(&model.TaskExecution{}) + if err := res.Error; err != nil { + logger.WarnF(ctx, "清理过期任务执行日志失败: %v", err) + } else { + deletedExecutions = res.RowsAffected + logger.InfoF(ctx, "已清理 7 天前任务执行日志,共 %d 条", deletedExecutions) + } + + // 2. 广播 EventTopicSystemCleanup 领域事件,由各业务域插件(upload, msg_gateway, user 等)自治执行各自的清理逻辑 + nowStr := time.Now().Format(time.RFC3339) + if err := EmitEvent(ctx, contracts.EventTopicSystemCleanup, contracts.SystemCleanupEvent{ + TriggeredAt: nowStr, + }); err != nil { + logger.WarnF(ctx, "广播系统清理领域事件失败: %v", err) + } + + msg := fmt.Sprintf("系统垃圾清理完成,已清理过期任务执行日志 %d 条,并已广播领域清理事件", deletedExecutions) + logger.InfoF(ctx, "%s", msg) + return &contracts.TaskResultDTO{Message: msg}, nil +} diff --git a/backend/plugins/domain/admin/service/cleanup_test.go b/backend/plugins/domain/admin/service/cleanup_test.go new file mode 100644 index 00000000..787536d9 --- /dev/null +++ b/backend/plugins/domain/admin/service/cleanup_test.go @@ -0,0 +1,78 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package service_test + +import ( + "Wavelet/core/contracts" + "Wavelet/plugins/domain/admin/model" + "Wavelet/plugins/domain/admin/service" + "context" + "sync/atomic" + "testing" + "time" + + "github.com/glebarez/sqlite" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + "gorm.io/gorm" +) + +func TestSystemCleanupHandler_Execute(t *testing.T) { + sqliteDB, err := gorm.Open(sqlite.Open(":memory:"), &gorm.Config{}) + require.NoError(t, err) + + require.NoError(t, sqliteDB.AutoMigrate(&model.TaskExecution{})) + + service.SetDBService(&testDBService{db: sqliteDB}) + defer service.ResetServices() + + now := time.Now() + oldTime := now.Add(-10 * 24 * time.Hour) + recentTime := now.Add(-1 * time.Hour) + + // Seed old execution (should be cleaned) + oldExec := model.TaskExecution{ + ID: 1, + TaskID: "task-old", + TaskType: "sample_task", + Status: "success", + CreatedAt: oldTime, + } + require.NoError(t, sqliteDB.Create(&oldExec).Error) + + // Seed recent execution (should remain) + recentExec := model.TaskExecution{ + ID: 2, + TaskID: "task-recent", + TaskType: "sample_task", + Status: "success", + CreatedAt: recentTime, + } + require.NoError(t, sqliteDB.Create(&recentExec).Error) + + var eventFired atomic.Bool + service.SetEventEmitter(func(ctx context.Context, topic string, payload any) error { + if topic == contracts.EventTopicSystemCleanup { + eventFired.Store(true) + } + return nil + }) + + handler := &service.SystemCleanupHandler{} + res, err := handler.Execute(context.Background(), nil) + require.NoError(t, err) + require.NotNil(t, res) + assert.Contains(t, res.Message, "系统垃圾清理完成") + + // Verify old execution is deleted and recent remains + var count int64 + sqliteDB.Model(&model.TaskExecution{}).Count(&count) + assert.Equal(t, int64(1), count) + + var remaining model.TaskExecution + sqliteDB.First(&remaining) + assert.Equal(t, uint64(2), remaining.ID) + + assert.True(t, eventFired.Load(), "EventTopicSystemCleanup must be emitted") +} diff --git a/backend/plugins/domain/domain_test.go b/backend/plugins/domain/domain_test.go index 8bc87a52..60c12476 100644 --- a/backend/plugins/domain/domain_test.go +++ b/backend/plugins/domain/domain_test.go @@ -12,6 +12,7 @@ import ( "Wavelet/plugins/domain/msg_gateway" "Wavelet/plugins/domain/risk_control" "Wavelet/plugins/domain/system" + "Wavelet/plugins/domain/upload" "Wavelet/plugins/domain/user" "Wavelet/plugins/infra/cache" "Wavelet/plugins/infra/logger" @@ -23,6 +24,7 @@ import ( "net/http/httptest" "path/filepath" "testing" + "time" "github.com/alicebob/miniredis/v2" "github.com/gin-gonic/gin" @@ -51,9 +53,12 @@ func setupTestDB(t *testing.T) *gorm.DB { &msg_gateway.MessageBinding{}, &msg_gateway.MessagePairingCode{}, &admin.SystemConfig{}, + &admin.TaskExecution{}, &msg_gateway.PushChannel{}, &msg_gateway.PushEvent{}, &msg_gateway.PushHistory{}, + &upload.Upload{}, + &upload.UploadStat{}, )) db.SetDB(testDB) @@ -392,7 +397,7 @@ func TestAdminPlugin(t *testing.T) { // 3. Settings schema, ok := ctx.Settings().Get("admin.system_cleanup_cron") require.True(t, ok) - assert.Equal(t, "0 4 * * *", schema.Default) + assert.Equal(t, "0 3 * * *", schema.Default) provider, err := core.Inject[contracts.PublicConfigProvider](ctx) require.NoError(t, err) @@ -532,3 +537,125 @@ func TestAllDomainPluginsCombined(t *testing.T) { // Clean shutdown require.NoError(t, ctx.Dispose()) } + +func TestSystemCleanupEventDrivenCoordination(t *testing.T) { + ctx := core.NewContext(context.Background()) + ctx.Config().SetSource(core.NewMapSource(nil)) + require.NoError(t, ctx.Config().Resolve()) + testDB := setupTestDB(t) + + // Apply Infra & Domain plugins + require.NoError(t, db.New(db.WithDB(testDB)).Apply(ctx)) + require.NoError(t, cache.New().Apply(ctx)) + require.NoError(t, logger.New().Apply(ctx)) + require.NoError(t, storage.New().Apply(ctx)) + + require.NoError(t, user.New().Apply(ctx)) + require.NoError(t, msg_gateway.New().Apply(ctx)) + require.NoError(t, upload.New().Apply(ctx)) + require.NoError(t, admin.New().Apply(ctx)) + + // 1. Seed old and recent task executions (admin domain) + now := time.Now() + oldTime := now.Add(-10 * 24 * time.Hour) + recentTime := now.Add(-1 * time.Hour) + + oldExec := admin.TaskExecution{ + ID: 101, + TaskID: "task-old-exec", + TaskType: "test_task", + Status: "success", + CreatedAt: oldTime, + } + recentExec := admin.TaskExecution{ + ID: 102, + TaskID: "task-recent-exec", + TaskType: "test_task", + Status: "success", + CreatedAt: recentTime, + } + require.NoError(t, testDB.Create(&oldExec).Error) + require.NoError(t, testDB.Model(&oldExec).UpdateColumn("created_at", oldTime).Error) + require.NoError(t, testDB.Create(&recentExec).Error) + + // 2. Seed old and recent push histories (msg_gateway domain, 30 days retention) + oldHistoryTime := now.Add(-40 * 24 * time.Hour) + oldHistory := msg_gateway.PushHistory{ + EventKey: "login", + Channel: "telegram", + Target: "123", + Title: "Old login", + Content: "Old content", + Level: "info", + Status: "success", + CreatedAt: oldHistoryTime, + } + recentHistory := msg_gateway.PushHistory{ + EventKey: "login", + Channel: "telegram", + Target: "123", + Title: "Recent login", + Content: "Recent content", + Level: "info", + Status: "success", + CreatedAt: recentTime, + } + require.NoError(t, testDB.Create(&oldHistory).Error) + require.NoError(t, testDB.Model(&oldHistory).UpdateColumn("created_at", oldHistoryTime).Error) + require.NoError(t, testDB.Create(&recentHistory).Error) + + // 3. Seed old pending upload and recent pending upload (upload domain) + oldUpload := upload.Upload{ + ID: 901, + UserID: 1, + FileName: "old.png", + FilePath: "uploads/old.png", + FileSize: 100, + Status: upload.UploadStatusPending, + CreatedAt: now.Add(-2 * time.Hour), + } + recentUpload := upload.Upload{ + ID: 902, + UserID: 1, + FileName: "recent.png", + FilePath: "uploads/recent.png", + FileSize: 100, + Status: upload.UploadStatusPending, + CreatedAt: now.Add(-10 * time.Minute), + } + require.NoError(t, testDB.Create(&oldUpload).Error) + require.NoError(t, testDB.Model(&oldUpload).UpdateColumn("created_at", now.Add(-2*time.Hour)).Error) + require.NoError(t, testDB.Create(&recentUpload).Error) + + // Dispatch admin system cleanup task handler + taskDef, ok := ctx.Tasks().Get("system:cleanup") + require.True(t, ok, "system:cleanup task must be registered in admin") + require.NotNil(t, taskDef.Handler) + + type resultExecutor interface { + Execute(ctx context.Context, payload []byte) (*contracts.TaskResultDTO, error) + } + handler, ok := taskDef.Handler.(resultExecutor) + require.True(t, ok) + + res, err := handler.Execute(context.Background(), nil) + require.NoError(t, err) + require.NotNil(t, res) + + // Verify Admin cleanup: old task execution deleted, recent remains + var execCount int64 + testDB.Model(&admin.TaskExecution{}).Count(&execCount) + assert.Equal(t, int64(1), execCount) + + // Verify MsgGateway cleanup: old push history deleted, recent remains + var historyCount int64 + testDB.Model(&msg_gateway.PushHistory{}).Count(&historyCount) + assert.Equal(t, int64(1), historyCount) + + // Verify Upload cleanup: old pending upload deleted, recent remains + var uploadCount int64 + testDB.Model(&upload.Upload{}).Count(&uploadCount) + assert.Equal(t, int64(1), uploadCount) + + require.NoError(t, ctx.Dispose()) +} diff --git a/backend/plugins/domain/msg_gateway/dao/push.go b/backend/plugins/domain/msg_gateway/dao/push.go index 4c3e24e2..f3a22d03 100644 --- a/backend/plugins/domain/msg_gateway/dao/push.go +++ b/backend/plugins/domain/msg_gateway/dao/push.go @@ -255,3 +255,13 @@ func CreatePushHistoryRecord(ctx context.Context, history *entity.PushHistory) e func PushHistoryQuery(ctx context.Context) *gorm.DB { return GetDB(ctx).Model(&entity.PushHistory{}) } + +// DeletePushHistoriesBeforeRecord deletes push history records created before cutoff time. +func DeletePushHistoriesBeforeRecord(ctx context.Context, cutoff time.Time) (int64, error) { + db := GetDB(ctx) + if db == nil { + return 0, nil + } + result := db.Where("created_at < ?", cutoff).Delete(&entity.PushHistory{}) + return result.RowsAffected, result.Error +} diff --git a/backend/plugins/domain/msg_gateway/dao/push_test.go b/backend/plugins/domain/msg_gateway/dao/push_test.go index d8d479eb..19712ba5 100644 --- a/backend/plugins/domain/msg_gateway/dao/push_test.go +++ b/backend/plugins/domain/msg_gateway/dao/push_test.go @@ -10,6 +10,7 @@ import ( "Wavelet/plugins/domain/msg_gateway/model/entity" "context" "testing" + "time" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" @@ -94,3 +95,51 @@ func TestPushEventDAO_CRUD(t *testing.T) { _, err = dao.GetPushEventByIDRecord(ctx, ev.ID) assert.Error(t, err) } + +func TestPushHistoryDAO_Cleanup(t *testing.T) { + _ = idgen.Init(1) + db, _, cleanup := testhelper.SetupTestEnvironment(t) + defer cleanup() + require.NoError(t, db.AutoMigrate(&entity.PushHistory{})) + + dao.SetDBServiceForTest(stubDBService{db: db}) + t.Cleanup(func() { dao.SetDBServiceForTest(nil) }) + + ctx := context.Background() + + now := time.Now() + oldTime := now.Add(-40 * 24 * time.Hour) + recentTime := now.Add(-5 * 24 * time.Hour) + + oldHistory := entity.PushHistory{ + EventKey: "login", + Channel: "telegram", + Target: "123", + Title: "Old login", + Content: "Old content", + Level: "info", + Status: "success", + CreatedAt: oldTime, + } + recentHistory := entity.PushHistory{ + EventKey: "login", + Channel: "telegram", + Target: "123", + Title: "Recent login", + Content: "Recent content", + Level: "info", + Status: "success", + CreatedAt: recentTime, + } + require.NoError(t, db.Create(&oldHistory).Error) + require.NoError(t, db.Create(&recentHistory).Error) + + cutoff := now.Add(-30 * 24 * time.Hour) + deleted, err := dao.DeletePushHistoriesBeforeRecord(ctx, cutoff) + require.NoError(t, err) + assert.Equal(t, int64(1), deleted) + + var count int64 + db.Model(&entity.PushHistory{}).Count(&count) + assert.Equal(t, int64(1), count) +} diff --git a/backend/plugins/domain/msg_gateway/plugin.go b/backend/plugins/domain/msg_gateway/plugin.go index b32c4481..85e1d58b 100644 --- a/backend/plugins/domain/msg_gateway/plugin.go +++ b/backend/plugins/domain/msg_gateway/plugin.go @@ -21,6 +21,7 @@ import ( "context" "embed" "reflect" + "time" "github.com/gin-gonic/gin" ) @@ -204,6 +205,13 @@ func (p *Plugin) Apply(ctx *core.Context) error { return nil }) + // 8.1 Register system cleanup event listener + ctx.Events().On(contracts.EventTopicSystemCleanup, func(c context.Context, _ contracts.SystemCleanupEvent) error { + const defaultPushHistoryRetention = 30 * 24 * time.Hour + _, err := service.CleanupPushHistories(c, defaultPushHistoryRetention) + return err + }) + // 9. Register built-in domain events and provide PushRegistry service.RegisterCustomEvents() core.Provide[contracts.PushRegistry](ctx, service.PushRegistryAdapter{}) diff --git a/backend/plugins/domain/msg_gateway/service/push_worker.go b/backend/plugins/domain/msg_gateway/service/push_worker.go index 7b5337bd..a29ec3f6 100644 --- a/backend/plugins/domain/msg_gateway/service/push_worker.go +++ b/backend/plugins/domain/msg_gateway/service/push_worker.go @@ -15,6 +15,7 @@ import ( "encoding/json" "errors" "fmt" + "time" ) const ( @@ -174,3 +175,15 @@ func RecordPushHistory(ctx context.Context, req do.SendPayload, status, errMsg s func ListPushHistories(ctx context.Context, filter do.PushHistoryListFilter) (int64, []entity.PushHistory, error) { return dao.ListPushHistoriesRecord(ctx, filter) } + +// CleanupPushHistories removes push delivery audit records older than the retention duration. +func CleanupPushHistories(ctx context.Context, retention time.Duration) (int64, error) { + cutoff := time.Now().Add(-retention) + deleted, err := dao.DeletePushHistoriesBeforeRecord(ctx, cutoff) + if err != nil { + logger.WarnF(ctx, "[Push] 清理历史推送日志失败: %v", err) + return 0, err + } + logger.InfoF(ctx, "[Push] 已清理 %s 前推送历史日志,共 %d 条", cutoff.Format(time.RFC3339), deleted) + return deleted, nil +} diff --git a/backend/plugins/domain/upload/exports.go b/backend/plugins/domain/upload/exports.go index dc789f5f..596f7600 100644 --- a/backend/plugins/domain/upload/exports.go +++ b/backend/plugins/domain/upload/exports.go @@ -112,3 +112,6 @@ type RebuildUploadStatsHandler = uploadtask.RebuildUploadStatsHandler // WarmImageCachePayload is the payload for image cache warmup tasks. type WarmImageCachePayload = uploadtask.WarmImageCachePayload + +// CleanupOrphanUploads removes unconfirmed upload files. +var CleanupOrphanUploads = uploadtask.CleanupOrphanUploads diff --git a/backend/plugins/domain/upload/plugin.go b/backend/plugins/domain/upload/plugin.go index 2a7b95eb..b24a660f 100644 --- a/backend/plugins/domain/upload/plugin.go +++ b/backend/plugins/domain/upload/plugin.go @@ -13,6 +13,7 @@ import ( "Wavelet/plugins/domain/upload/handler" "Wavelet/plugins/domain/upload/shared" "Wavelet/plugins/domain/upload/task" + "context" "embed" "reflect" @@ -126,8 +127,11 @@ func (p *Plugin) Apply(ctx *core.Context) error { ctx.Task().Register(task.StorageMigrationTask, &task.MigrationHandler{}, extpoints.WithTaskMeta(task.StorageMigrationMeta), extpoints.WithTaskRetry(defaultSingleRetry)) ctx.Task().Register(task.WarmImageCacheTask, &task.WarmImageCacheHandler{}, extpoints.WithTaskMeta(task.WarmImageCacheMeta), extpoints.WithTaskRetry(1)) - // 4. Register Cron Schedule - ctx.Schedule().RegisterCron("0 3 * * *", task.SystemCleanupTask, nil) + // 4. Register Event Listeners for domain events + ctx.Events().On(contracts.EventTopicSystemCleanup, func(c context.Context, _ contracts.SystemCleanupEvent) error { + _, _, err := task.CleanupOrphanUploads(c) + return err + }) // 5. Register Settings Schemas ctx.Settings().Register(extpoints.SettingSchema{ diff --git a/backend/plugins/domain/upload/task/cleanup.go b/backend/plugins/domain/upload/task/cleanup.go index 7b09ee2c..743ba74e 100644 --- a/backend/plugins/domain/upload/task/cleanup.go +++ b/backend/plugins/domain/upload/task/cleanup.go @@ -33,9 +33,9 @@ const ( var SystemCleanupMeta = contracts.TaskMetaDTO{ Type: TaskTypeSystemCleanup, AsynqTask: SystemCleanupTask, - Name: "系统垃圾清理", - DisplayName: "系统垃圾清理", - Description: "定期清理未使用上传文件、历史推送记录和过期任务执行日志", + Name: "清理未确认上传文件", + DisplayName: "清理未确认上传文件", + Description: "定期清理超过1小时未使用的上传临时文件与底层存储资源", Category: "maintenance", SupportsTime: false, MaxRetry: 3, @@ -43,13 +43,29 @@ var SystemCleanupMeta = contracts.TaskMetaDTO{ Retryable: true, } -// SystemCleanupHandler 系统定期垃圾清理异步任务处理器 +// SystemCleanupHandler 未确认上传文件清理异步任务处理器 type SystemCleanupHandler struct{} -// Execute 执行系统清理(包含文件清理、历史推送日志和任务执行日志清理) +// Execute 执行系统清理(清理未使用上传文件) func (h *SystemCleanupHandler) Execute(ctx context.Context, _ []byte) (*contracts.TaskResultDTO, error) { + totalProcessed, totalDeleted, err := CleanupOrphanUploads(ctx) + if err != nil { + return nil, err + } + + msg := fmt.Sprintf( + "系统垃圾清理完成,处理未确认文件: %d 个,物理删除: %d 个", + totalProcessed, + totalDeleted, + ) + logger.InfoF(ctx, "%s", msg) + return &contracts.TaskResultDTO{Message: msg}, nil +} + +// CleanupOrphanUploads 扫描并清理超过1小时未确认的 pending 状态上传文件及物理存储 +func CleanupOrphanUploads(ctx context.Context) (int, int, error) { if uploadstorage.ReadOnly(ctx) { - return nil, errors.New(shared.ErrStorageReadOnly) + return 0, 0, errors.New(shared.ErrStorageReadOnly) } const batchSize = 100 var lastID uint64 @@ -62,14 +78,14 @@ func (h *SystemCleanupHandler) Execute(ctx context.Context, _ []byte) (*contract db := shared.GetDB(ctx) if db == nil { - return nil, errors.New("database service not available") + return 0, 0, errors.New("database service not available") } storageSvc := shared.GetStorage(ctx) for { if err := ctx.Err(); err != nil { - return nil, fmt.Errorf("system cleanup canceled: %w", err) + return totalProcessed, totalDeleted, fmt.Errorf("system cleanup canceled: %w", err) } var pendingUploads []models.Upload @@ -79,7 +95,7 @@ func (h *SystemCleanupHandler) Execute(ctx context.Context, _ []byte) (*contract Limit(batchSize). Find(&pendingUploads).Error; err != nil { logger.ErrorF(ctx, "查询过期待使用上传文件失败: %v", err) - return nil, fmt.Errorf("failed to query pending uploads: %w", err) + return totalProcessed, totalDeleted, fmt.Errorf("failed to query pending uploads: %w", err) } if len(pendingUploads) == 0 { @@ -88,7 +104,7 @@ func (h *SystemCleanupHandler) Execute(ctx context.Context, _ []byte) (*contract for i := range pendingUploads { if err := ctx.Err(); err != nil { - return nil, fmt.Errorf("system cleanup canceled: %w", err) + return totalProcessed, totalDeleted, fmt.Errorf("system cleanup canceled: %w", err) } upload := &pendingUploads[i] @@ -117,39 +133,5 @@ func (h *SystemCleanupHandler) Execute(ctx context.Context, _ []byte) (*contract } } - // 清理过期任务执行记录 - var deletedExecutions int64 - sevenDaysAgo := time.Now().Add(-7 * 24 * time.Hour) - if err := db. - Table("w_task_executions"). - Where("created_at < ?", sevenDaysAgo). - Delete(&struct{}{}).Error; err != nil { - logger.WarnF(ctx, "清理过期任务执行日志失败: %v", err) - } else { - deletedExecutions = db.RowsAffected - logger.InfoF(ctx, "已清理 7 天前任务执行日志,共 %d 条", deletedExecutions) - } - - // 清理已过期推送日志 - var deletedPushLogs int64 - thirtyDaysAgo := time.Now().Add(-30 * 24 * time.Hour) - if err := db. - Table("w_push_logs"). - Where("created_at < ?", thirtyDaysAgo). - Delete(&struct{}{}).Error; err != nil { - logger.WarnF(ctx, "清理历史推送日志失败: %v", err) - } else { - deletedPushLogs = db.RowsAffected - logger.InfoF(ctx, "已清理 30 天前推送日志,共 %d 条", deletedPushLogs) - } - - msg := fmt.Sprintf( - "系统垃圾清理完成,处理未确认文件: %d 个,物理删除: %d 个,清理过期任务日志: %d 条,清理历史推送日志: %d 条", - totalProcessed, - totalDeleted, - deletedExecutions, - deletedPushLogs, - ) - logger.InfoF(ctx, "%s", msg) - return &contracts.TaskResultDTO{Message: msg}, nil + return totalProcessed, totalDeleted, nil } diff --git a/backend/plugins/domain/user/plugin.go b/backend/plugins/domain/user/plugin.go index d66864ae..6d41c35e 100644 --- a/backend/plugins/domain/user/plugin.go +++ b/backend/plugins/domain/user/plugin.go @@ -151,6 +151,13 @@ func (p *Plugin) Apply(ctx *core.Context) error { ctx.Task().Register(TaskCleanupInactive, &CleanupInactiveHandler{}, extpoints.WithTaskMeta(CleanupInactiveMeta)) + // 4.1 Register Event Listeners for domain events + ctx.Events().On(contracts.EventTopicSystemCleanup, func(c context.Context, _ contracts.SystemCleanupEvent) error { + handler := &CleanupInactiveHandler{} + _, err := handler.Execute(c, nil) + return err + }) + // 5. Register Settings Schemas ctx.Settings().Register(extpoints.SettingSchema{ Key: "user.registration_enabled",