mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-10-08 00:26:37 +08:00
feat(upload): add rebuild stats async task
This commit is contained in:
@@ -39,9 +39,9 @@ var (
|
|||||||
|
|
||||||
// Ingest policy constants
|
// Ingest policy constants
|
||||||
const (
|
const (
|
||||||
PolicyCreate = ingest.PolicyCreate
|
PolicyCreate = ingest.PolicyCreate
|
||||||
PolicyDedupNewRecord = ingest.PolicyDedupNewRecord
|
PolicyDedupNewRecord = ingest.PolicyDedupNewRecord
|
||||||
PolicyResolveExisting = ingest.PolicyResolveExisting
|
PolicyResolveExisting = ingest.PolicyResolveExisting
|
||||||
)
|
)
|
||||||
|
|
||||||
type (
|
type (
|
||||||
@@ -55,7 +55,7 @@ type (
|
|||||||
|
|
||||||
// Ingest errors
|
// Ingest errors
|
||||||
var (
|
var (
|
||||||
ErrIngestForbidden = ingest.ErrForbidden
|
ErrIngestForbidden = ingest.ErrForbidden
|
||||||
ErrIngestStorageReadOnly = ingest.ErrStorageReadOnly
|
ErrIngestStorageReadOnly = ingest.ErrStorageReadOnly
|
||||||
)
|
)
|
||||||
|
|
||||||
@@ -82,9 +82,10 @@ var (
|
|||||||
|
|
||||||
// Task identifiers and metadata
|
// Task identifiers and metadata
|
||||||
const (
|
const (
|
||||||
StorageMigrationTask = uploadtask.StorageMigrationTask
|
StorageMigrationTask = uploadtask.StorageMigrationTask
|
||||||
SystemCleanupTask = uploadtask.SystemCleanupTask
|
SystemCleanupTask = uploadtask.SystemCleanupTask
|
||||||
WarmImageCacheTask = uploadtask.WarmImageCacheTask
|
WarmImageCacheTask = uploadtask.WarmImageCacheTask
|
||||||
|
RebuildUploadStatsTask = uploadtask.RebuildUploadStatsTask
|
||||||
)
|
)
|
||||||
|
|
||||||
var (
|
var (
|
||||||
@@ -94,6 +95,8 @@ var (
|
|||||||
SystemCleanupMeta = uploadtask.SystemCleanupMeta
|
SystemCleanupMeta = uploadtask.SystemCleanupMeta
|
||||||
// WarmImageCacheMeta describes the image compression cache warmup task.
|
// WarmImageCacheMeta describes the image compression cache warmup task.
|
||||||
WarmImageCacheMeta = uploadtask.WarmImageCacheMeta
|
WarmImageCacheMeta = uploadtask.WarmImageCacheMeta
|
||||||
|
// RebuildUploadStatsMeta describes the upload stats rebuild task.
|
||||||
|
RebuildUploadStatsMeta = uploadtask.RebuildUploadStatsMeta
|
||||||
)
|
)
|
||||||
|
|
||||||
// MigrationHandler executes storage migration tasks.
|
// MigrationHandler executes storage migration tasks.
|
||||||
@@ -105,6 +108,9 @@ type SystemCleanupHandler = uploadtask.SystemCleanupHandler
|
|||||||
// WarmImageCacheHandler pre-warms compressed image caches.
|
// WarmImageCacheHandler pre-warms compressed image caches.
|
||||||
type WarmImageCacheHandler = uploadtask.WarmImageCacheHandler
|
type WarmImageCacheHandler = uploadtask.WarmImageCacheHandler
|
||||||
|
|
||||||
|
// RebuildUploadStatsHandler rebuilds upload stats from active records.
|
||||||
|
type RebuildUploadStatsHandler = uploadtask.RebuildUploadStatsHandler
|
||||||
|
|
||||||
// WarmImageCachePayload is the payload for image cache warmup tasks.
|
// WarmImageCachePayload is the payload for image cache warmup tasks.
|
||||||
type WarmImageCachePayload = uploadtask.WarmImageCachePayload
|
type WarmImageCachePayload = uploadtask.WarmImageCachePayload
|
||||||
|
|
||||||
@@ -112,8 +118,9 @@ type WarmImageCachePayload = uploadtask.WarmImageCachePayload
|
|||||||
var (
|
var (
|
||||||
_ task.TaskHandler = (*MigrationHandler)(nil)
|
_ task.TaskHandler = (*MigrationHandler)(nil)
|
||||||
_ task.TaskHandler = (*SystemCleanupHandler)(nil)
|
_ task.TaskHandler = (*SystemCleanupHandler)(nil)
|
||||||
|
_ task.TaskHandler = (*RebuildUploadStatsHandler)(nil)
|
||||||
_ interface {
|
_ interface {
|
||||||
task.TaskHandler
|
task.TaskHandler
|
||||||
ValidatePayload([]byte) ([]byte, error)
|
ValidatePayload([]byte) ([]byte, error)
|
||||||
} = (*WarmImageCacheHandler)(nil)
|
} = (*WarmImageCacheHandler)(nil)
|
||||||
)
|
)
|
||||||
|
|||||||
@@ -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
|
||||||
|
}
|
||||||
@@ -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)
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -25,6 +25,9 @@ func Register() {
|
|||||||
task.RegisterHandler(upload.WarmImageCacheTask, &upload.WarmImageCacheHandler{})
|
task.RegisterHandler(upload.WarmImageCacheTask, &upload.WarmImageCacheHandler{})
|
||||||
task.RegisterTaskMeta(upload.WarmImageCacheMeta)
|
task.RegisterTaskMeta(upload.WarmImageCacheMeta)
|
||||||
|
|
||||||
|
task.RegisterHandler(upload.RebuildUploadStatsTask, &upload.RebuildUploadStatsHandler{})
|
||||||
|
task.RegisterTaskMeta(upload.RebuildUploadStatsMeta)
|
||||||
|
|
||||||
// user
|
// user
|
||||||
task.RegisterHandler(user.SendEmailTask, &user.SendEmailHandler{})
|
task.RegisterHandler(user.SendEmailTask, &user.SendEmailHandler{})
|
||||||
task.RegisterTaskMeta(user.SendEmailMeta)
|
task.RegisterTaskMeta(user.SendEmailMeta)
|
||||||
|
|||||||
Reference in New Issue
Block a user