From 990e3a6f51dd77d3aa0c9f4f9611d7deb99080d2 Mon Sep 17 00:00:00 2001 From: ryan Date: Thu, 18 Jun 2026 15:05:52 +0800 Subject: [PATCH] feat(upload): add rebuild stats async task --- internal/apps/upload/exports.go | 23 +++--- internal/apps/upload/task/rebuild_stats.go | 72 +++++++++++++++++++ .../apps/upload/task/rebuild_stats_test.go | 69 ++++++++++++++++++ internal/task/handlers/register.go | 3 + 4 files changed, 159 insertions(+), 8 deletions(-) create mode 100644 internal/apps/upload/task/rebuild_stats.go create mode 100644 internal/apps/upload/task/rebuild_stats_test.go diff --git a/internal/apps/upload/exports.go b/internal/apps/upload/exports.go index 56eb4a60..981f0c4c 100644 --- a/internal/apps/upload/exports.go +++ b/internal/apps/upload/exports.go @@ -39,9 +39,9 @@ var ( // Ingest policy constants const ( - PolicyCreate = ingest.PolicyCreate - PolicyDedupNewRecord = ingest.PolicyDedupNewRecord - PolicyResolveExisting = ingest.PolicyResolveExisting + PolicyCreate = ingest.PolicyCreate + PolicyDedupNewRecord = ingest.PolicyDedupNewRecord + PolicyResolveExisting = ingest.PolicyResolveExisting ) type ( @@ -55,7 +55,7 @@ type ( // Ingest errors var ( - ErrIngestForbidden = ingest.ErrForbidden + ErrIngestForbidden = ingest.ErrForbidden ErrIngestStorageReadOnly = ingest.ErrStorageReadOnly ) @@ -82,9 +82,10 @@ var ( // Task identifiers and metadata const ( - StorageMigrationTask = uploadtask.StorageMigrationTask - SystemCleanupTask = uploadtask.SystemCleanupTask - WarmImageCacheTask = uploadtask.WarmImageCacheTask + StorageMigrationTask = uploadtask.StorageMigrationTask + SystemCleanupTask = uploadtask.SystemCleanupTask + WarmImageCacheTask = uploadtask.WarmImageCacheTask + RebuildUploadStatsTask = uploadtask.RebuildUploadStatsTask ) var ( @@ -94,6 +95,8 @@ var ( SystemCleanupMeta = uploadtask.SystemCleanupMeta // WarmImageCacheMeta describes the image compression cache warmup task. WarmImageCacheMeta = uploadtask.WarmImageCacheMeta + // RebuildUploadStatsMeta describes the upload stats rebuild task. + RebuildUploadStatsMeta = uploadtask.RebuildUploadStatsMeta ) // MigrationHandler executes storage migration tasks. @@ -105,6 +108,9 @@ type SystemCleanupHandler = uploadtask.SystemCleanupHandler // WarmImageCacheHandler pre-warms compressed image caches. type WarmImageCacheHandler = uploadtask.WarmImageCacheHandler +// RebuildUploadStatsHandler rebuilds upload stats from active records. +type RebuildUploadStatsHandler = uploadtask.RebuildUploadStatsHandler + // WarmImageCachePayload is the payload for image cache warmup tasks. type WarmImageCachePayload = uploadtask.WarmImageCachePayload @@ -112,8 +118,9 @@ type WarmImageCachePayload = uploadtask.WarmImageCachePayload var ( _ task.TaskHandler = (*MigrationHandler)(nil) _ task.TaskHandler = (*SystemCleanupHandler)(nil) + _ task.TaskHandler = (*RebuildUploadStatsHandler)(nil) _ interface { task.TaskHandler ValidatePayload([]byte) ([]byte, error) } = (*WarmImageCacheHandler)(nil) -) \ No newline at end of file +) diff --git a/internal/apps/upload/task/rebuild_stats.go b/internal/apps/upload/task/rebuild_stats.go new file mode 100644 index 00000000..018def05 --- /dev/null +++ b/internal/apps/upload/task/rebuild_stats.go @@ -0,0 +1,72 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package task + +import ( + "context" + "fmt" + + uploadstats "github.com/Rain-kl/Wavelet/internal/apps/upload/stats" + "github.com/Rain-kl/Wavelet/internal/db" + "github.com/Rain-kl/Wavelet/internal/model" + "github.com/Rain-kl/Wavelet/internal/task" +) + +const ( + // RebuildUploadStatsTask is the Asynq task name for rebuilding upload stats. + RebuildUploadStatsTask = "upload:rebuild_stats" + // TaskTypeRebuildUploadStats is the admin-dispatchable task type. + TaskTypeRebuildUploadStats = "rebuild_upload_stats" +) + +// RebuildUploadStatsMeta describes the upload stats rebuild task. +var RebuildUploadStatsMeta = task.TaskMeta{ + Type: TaskTypeRebuildUploadStats, + AsynqTask: RebuildUploadStatsTask, + Name: "重算文件存储统计", + Description: "根据当前 w_uploads 活跃记录全量重建 w_upload_stats(总量、类型、分类、趋势)", + SupportsTime: false, + MaxRetry: task.DefaultMaxRetry, + Queue: task.QueueDefault, + Retryable: true, +} + +// RebuildUploadStatsHandler rebuilds incremental upload stats from active upload records. +type RebuildUploadStatsHandler struct{} + +// Execute scans active uploads and rebuilds all upload stat dimensions. +func (h *RebuildUploadStatsHandler) Execute(ctx context.Context, _ []byte) (*task.TaskResult, error) { + var activeCount int64 + if err := db.DB(ctx). + Model(&model.Upload{}). + Where("status != ?", model.UploadStatusDeleted). + Count(&activeCount).Error; err != nil { + task.AppendLog(ctx, "统计活跃上传记录失败: %v", err) + return nil, fmt.Errorf("count active uploads: %w", err) + } + + task.AppendLog(ctx, "开始重算文件存储统计,活跃记录数: %d", activeCount) + + if err := uploadstats.RebuildUploadStats(ctx); err != nil { + task.AppendLog(ctx, "重算文件存储统计失败: %v", err) + return nil, fmt.Errorf("rebuild upload stats: %w", err) + } + + var totalStat model.UploadStat + if err := db.DB(ctx). + Where("dimension = ? AND stat_key = ?", model.UploadStatDimensionTotal, ""). + First(&totalStat).Error; err != nil { + task.AppendLog(ctx, "读取总量统计失败: %v", err) + return nil, fmt.Errorf("load total upload stats: %w", err) + } + + msg := fmt.Sprintf( + "文件存储统计重算完成,活跃记录 %d 条,统计文件数 %d,总大小 %d 字节", + activeCount, + totalStat.FileCount, + totalStat.FileSize, + ) + task.AppendLog(ctx, "%s", msg) + return &task.TaskResult{Message: msg}, nil +} diff --git a/internal/apps/upload/task/rebuild_stats_test.go b/internal/apps/upload/task/rebuild_stats_test.go new file mode 100644 index 00000000..e0e185be --- /dev/null +++ b/internal/apps/upload/task/rebuild_stats_test.go @@ -0,0 +1,69 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package task + +import ( + "context" + "testing" + "time" + + "github.com/Rain-kl/Wavelet/internal/db" + "github.com/Rain-kl/Wavelet/internal/model" + "github.com/Rain-kl/Wavelet/internal/testhelper" +) + +func TestRebuildUploadStatsHandler_Execute(t *testing.T) { + _, _, cleanup := testhelper.SetupTestEnvironment(t) + defer cleanup() + + ctx := context.Background() + now := time.Now() + + uploads := []model.Upload{ + { + UserID: 1001, FileName: "a.jpg", FilePath: "uploads/a.jpg", + FileSize: 100, MimeType: "image/jpeg", Extension: "jpg", Hash: "hash-a", + Type: "pixez_mirror", Status: model.UploadStatusUsed, CreatedAt: now, + }, + { + UserID: 1001, FileName: "b.png", FilePath: "uploads/b.png", + FileSize: 200, MimeType: "image/png", Extension: "png", Hash: "hash-b", + Type: "attachment", Status: model.UploadStatusUsed, CreatedAt: now, + }, + } + for i := range uploads { + if err := db.DB(ctx).Create(&uploads[i]).Error; err != nil { + t.Fatalf("seed upload failed: %v", err) + } + } + + // Corrupt stats to ensure rebuild recalculates from uploads. + if err := db.DB(ctx).Create(&model.UploadStat{ + Dimension: model.UploadStatDimensionTotal, + StatKey: "", + FileCount: 0, + FileSize: 0, + }).Error; err != nil { + t.Fatalf("seed broken total stat failed: %v", err) + } + + handler := &RebuildUploadStatsHandler{} + result, err := handler.Execute(ctx, nil) + if err != nil { + t.Fatalf("Execute() error = %v", err) + } + if result == nil || result.Message == "" { + t.Fatalf("Execute() returned empty result: %+v", result) + } + + var totalStat model.UploadStat + if err := db.DB(ctx). + Where("dimension = ? AND stat_key = ?", model.UploadStatDimensionTotal, ""). + First(&totalStat).Error; err != nil { + t.Fatalf("load total stat failed: %v", err) + } + if totalStat.FileCount != 2 || totalStat.FileSize != 300 { + t.Fatalf("total stat = count %d size %d, want 2 / 300", totalStat.FileCount, totalStat.FileSize) + } +} diff --git a/internal/task/handlers/register.go b/internal/task/handlers/register.go index b9ef86b5..f032ed7f 100644 --- a/internal/task/handlers/register.go +++ b/internal/task/handlers/register.go @@ -25,6 +25,9 @@ func Register() { task.RegisterHandler(upload.WarmImageCacheTask, &upload.WarmImageCacheHandler{}) task.RegisterTaskMeta(upload.WarmImageCacheMeta) + task.RegisterHandler(upload.RebuildUploadStatsTask, &upload.RebuildUploadStatsHandler{}) + task.RegisterTaskMeta(upload.RebuildUploadStatsMeta) + // user task.RegisterHandler(user.SendEmailTask, &user.SendEmailHandler{}) task.RegisterTaskMeta(user.SendEmailMeta)