From 5915b31519f25e28e43ca5c48a736e14f93690aa Mon Sep 17 00:00:00 2001 From: ryan Date: Thu, 11 Jun 2026 17:49:08 +0800 Subject: [PATCH] =?UTF-8?q?=E5=9B=BE=E7=89=87=E9=A2=84=E7=83=AD=E4=BB=BB?= =?UTF-8?q?=E5=8A=A1?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- internal/apps/admin/task/routers_test.go | 8 +- internal/apps/upload/errs.go | 55 ++++--- internal/apps/upload/file_server.go | 50 +++--- internal/apps/upload/tasks.go | 165 +++++++++++++++++++ internal/apps/upload/tasks_test.go | 193 +++++++++++++++++++++++ internal/task/handlers/register.go | 2 + 6 files changed, 429 insertions(+), 44 deletions(-) diff --git a/internal/apps/admin/task/routers_test.go b/internal/apps/admin/task/routers_test.go index 5a0a60bc..8664c101 100644 --- a/internal/apps/admin/task/routers_test.go +++ b/internal/apps/admin/task/routers_test.go @@ -89,15 +89,21 @@ func TestListTaskTypes(t *testing.T) { } foundCleanup := false + foundWarmImageCache := false for _, m := range taskMetas { if m.Type == upload.TaskTypeCleanupUploads { foundCleanup = true - break + } + if m.Type == upload.TaskTypeWarmImageCache { + foundWarmImageCache = true } } if !foundCleanup { t.Errorf("expected task type %s to be listed", upload.TaskTypeCleanupUploads) } + if !foundWarmImageCache { + t.Errorf("expected task type %s to be listed", upload.TaskTypeWarmImageCache) + } } func TestDispatchTask(t *testing.T) { diff --git a/internal/apps/upload/errs.go b/internal/apps/upload/errs.go index 7bd57f2e..1660dc74 100644 --- a/internal/apps/upload/errs.go +++ b/internal/apps/upload/errs.go @@ -7,29 +7,34 @@ package upload // 文件管理常量 const ( - ErrNoFileSelected = "请选择要上传的文件" - ErrUnsupportedFormat = "只支持 JPG、PNG、WEBP 格式的图片" - ErrProcessFileFailed = "处理文件失败" - ErrSaveFileFailed = "保存文件失败" - ErrOpenFileFailed = "打开文件失败" - ErrSaveUploadRecordFailed = "保存上传记录失败" - ErrGenericFileTooLarge = "文件大小不能超过 32MB" - ErrFileContentExtensionMismatch = "文件内容与扩展名不匹配,可能包含安全风险" - ErrFileValidationFailed = "文件校验失败" - ErrInvalidMetadataJSON = "元数据 JSON 格式不合法" - ErrInvalidFileID = "无效的文件 ID" - ErrQueryUploadRecordFailed = "查询文件记录失败" - ErrInvalidBatchDownloadRequest = "参数绑定失败,请传入有效的文件 ID 数组" - ErrInvalidIDValueFormat = "无效的 ID 值: %s" - ErrRetrieveUploadRecordsFailed = "检索文件记录失败" - ErrNoValidFilesForArchive = "没有找到任何有效的文件记录进行打包" - ErrInvalidParams = "参数错误" - ErrQueryFileCountFailed = "查询文件数量失败" - ErrQueryFileListFailed = "查询文件列表失败" - ErrDeleteFileFailed = "删除文件失败" - ErrS3KeyRequired = "s3 key must not be empty" - ErrS3KeyTooLongFormat = "s3 key exceeds maximum length of %d" - ErrS3KeyStartsWithSlash = "s3 key must not start with /" - ErrS3KeyContainsNullBytes = "s3 key must not contain null bytes" - ErrQueryUnusedUploadsFailed = "查询未使用的上传文件失败: %w" + ErrNoFileSelected = "请选择要上传的文件" + ErrUnsupportedFormat = "只支持 JPG、PNG、WEBP 格式的图片" + ErrProcessFileFailed = "处理文件失败" + ErrSaveFileFailed = "保存文件失败" + ErrOpenFileFailed = "打开文件失败" + ErrSaveUploadRecordFailed = "保存上传记录失败" + ErrGenericFileTooLarge = "文件大小不能超过 32MB" + ErrFileContentExtensionMismatch = "文件内容与扩展名不匹配,可能包含安全风险" + ErrFileValidationFailed = "文件校验失败" + ErrInvalidMetadataJSON = "元数据 JSON 格式不合法" + ErrInvalidFileID = "无效的文件 ID" + ErrQueryUploadRecordFailed = "查询文件记录失败" + ErrInvalidBatchDownloadRequest = "参数绑定失败,请传入有效的文件 ID 数组" + ErrInvalidIDValueFormat = "无效的 ID 值: %s" + ErrRetrieveUploadRecordsFailed = "检索文件记录失败" + ErrNoValidFilesForArchive = "没有找到任何有效的文件记录进行打包" + ErrInvalidParams = "参数错误" + ErrQueryFileCountFailed = "查询文件数量失败" + ErrQueryFileListFailed = "查询文件列表失败" + ErrDeleteFileFailed = "删除文件失败" + ErrS3KeyRequired = "s3 key must not be empty" + ErrS3KeyTooLongFormat = "s3 key exceeds maximum length of %d" + ErrS3KeyStartsWithSlash = "s3 key must not start with /" + ErrS3KeyContainsNullBytes = "s3 key must not contain null bytes" + ErrQueryUnusedUploadsFailed = "查询未使用的上传文件失败: %w" + errImageCacheWarmupPayloadRequired = "图片缓存预热参数不能为空" + errInvalidImageCacheWarmupPayload = "图片缓存预热参数格式无效: %w" + errInvalidImageCacheWarmupQuality = "图片质量仅支持 low、medium、high" + errParseImageCacheWarmupPayload = "解析图片缓存预热参数失败: %w" + errQueryImagesForCacheWarmup = "查询待预热图片失败: %w" ) diff --git a/internal/apps/upload/file_server.go b/internal/apps/upload/file_server.go index 607d0c82..2df9febb 100644 --- a/internal/apps/upload/file_server.go +++ b/internal/apps/upload/file_server.go @@ -96,37 +96,51 @@ func ServeUpload(c *gin.Context, upload *model.Upload) { return } + webpBytes, _, err := ensureCompressedImageCache(c.Request.Context(), upload, quality) + if err != nil { + if len(webpBytes) > 0 { + logger.WarnF(c.Request.Context(), "failed to cache compressed image: %v", err) + c.Data(http.StatusOK, "image/webp", webpBytes) + return + } + logger.ErrorF(c.Request.Context(), "failed to prepare compressed image cache: %v", err) + serveOriginal(c, upload) + return + } + + c.Data(http.StatusOK, "image/webp", webpBytes) +} + +func ensureCompressedImageCache( + ctx context.Context, + upload *model.Upload, + quality string, +) ([]byte, bool, error) { cache := diskcache.GetGlobalCache() cacheKey := imageCompressionCacheKey(upload, quality) - if webpBytes, err := cache.Get(cacheKey); err == nil { - c.Data(http.StatusOK, "image/webp", webpBytes) - return - } else if !errors.Is(err, diskcache.ErrCacheMiss) { - logger.WarnF(c.Request.Context(), "failed to read compressed image cache: %v", err) + webpBytes, err := cache.Get(cacheKey) + if err == nil { + return webpBytes, true, nil + } + if !errors.Is(err, diskcache.ErrCacheMiss) { + return nil, false, fmt.Errorf("read compressed image cache: %w", err) } - // Cache miss: retrieve original file content - origBytes, err := getOriginalFileBytes(c.Request.Context(), upload) + origBytes, err := getOriginalFileBytes(ctx, upload) if err != nil { - logger.ErrorF(c.Request.Context(), "failed to retrieve original file bytes for compression: %v", err) - serveOriginal(c, upload) - return + return nil, false, fmt.Errorf("read original image: %w", err) } - // Compress to WebP - webpBytes, err := CompressImageToWebP(bytes.NewReader(origBytes), quality) + webpBytes, err = CompressImageToWebP(bytes.NewReader(origBytes), quality) if err != nil { - logger.ErrorF(c.Request.Context(), "failed to compress image to WebP: %v", err) - serveOriginal(c, upload) - return + return nil, false, fmt.Errorf("compress image to WebP: %w", err) } if err := cache.Set(cacheKey, webpBytes, diskcache.NoExpiration); err != nil { - logger.WarnF(c.Request.Context(), "failed to cache compressed image: %v", err) + return webpBytes, false, fmt.Errorf("write compressed image cache: %w", err) } - // Serve compressed WebP - c.Data(http.StatusOK, "image/webp", webpBytes) + return webpBytes, false, nil } func imageCompressionCacheKey(upload *model.Upload, quality string) string { diff --git a/internal/apps/upload/tasks.go b/internal/apps/upload/tasks.go index 7cf7da4b..3d5e7dd2 100644 --- a/internal/apps/upload/tasks.go +++ b/internal/apps/upload/tasks.go @@ -6,7 +6,11 @@ package upload import ( "context" + "encoding/json" + "errors" "fmt" + "strings" + "sync" "time" "github.com/Rain-kl/Wavelet/internal/db" @@ -22,8 +26,14 @@ const ( CleanupUnusedUploadsTask = "upload:cleanup_unused" // TaskTypeCleanupUploads 清理未使用上传管理类型 TaskTypeCleanupUploads = "cleanup_unused_uploads" + // WarmImageCacheTask 图片压缩缓存预热任务标识 + WarmImageCacheTask = "upload:warm_image_cache" + // TaskTypeWarmImageCache 图片压缩缓存预热管理类型 + TaskTypeWarmImageCache = "warm_image_cache" ) +var warmImageCacheMu sync.Mutex + // CleanupUnusedUploadsMeta represents the task metadata. var CleanupUnusedUploadsMeta = task.TaskMeta{ Type: TaskTypeCleanupUploads, @@ -36,9 +46,39 @@ var CleanupUnusedUploadsMeta = task.TaskMeta{ Retryable: true, } +// WarmImageCacheMeta represents the image cache warmup task metadata. +var WarmImageCacheMeta = task.TaskMeta{ + Type: TaskTypeWarmImageCache, + AsynqTask: WarmImageCacheTask, + Name: "预热图片压缩缓存", + Description: "串行将文件管理中的图片转换为指定质量的 WebP 并写入永久缓存", + SupportsTime: false, + MaxRetry: task.DefaultMaxRetry, + Queue: task.QueueDefault, + Retryable: true, + Params: []task.TaskParam{ + { + Name: "quality", + Label: "图片质量", + Type: "string", + Required: true, + Placeholder: "low / medium / high", + Description: "WebP 压缩质量,仅支持 low、medium、high", + }, + }, +} + +// WarmImageCachePayload is the image cache warmup task payload. +type WarmImageCachePayload struct { + Quality string `json:"quality"` +} + // CleanupUnusedUploadsHandler 清理未使用上传文件的异步任务处理器 type CleanupUnusedUploadsHandler struct{} +// WarmImageCacheHandler serially warms compressed image cache entries. +type WarmImageCacheHandler struct{} + // Execute 执行清理未使用上传文件的业务逻辑 func (h *CleanupUnusedUploadsHandler) Execute(ctx context.Context, _ []byte) (*task.TaskResult, error) { const batchSize = 100 // 每批处理100个文件 @@ -103,3 +143,128 @@ func (h *CleanupUnusedUploadsHandler) Execute(ctx context.Context, _ []byte) (*t task.AppendLog(ctx, "%s", msg) return &task.TaskResult{Message: msg}, nil } + +// ValidatePayload validates and normalizes image cache warmup parameters. +func (h *WarmImageCacheHandler) ValidatePayload(payload []byte) ([]byte, error) { + if len(payload) == 0 { + return nil, errors.New(errImageCacheWarmupPayloadRequired) + } + + var req WarmImageCachePayload + if err := json.Unmarshal(payload, &req); err != nil { + return nil, fmt.Errorf(errInvalidImageCacheWarmupPayload, err) + } + + req.Quality = strings.ToLower(strings.TrimSpace(req.Quality)) + if req.Quality != imageQualityLow && + req.Quality != imageQualityMedium && + req.Quality != imageQualityHigh { + return nil, errors.New(errInvalidImageCacheWarmupQuality) + } + + return json.Marshal(req) +} + +// Execute serially converts all managed images to WebP cache entries. +func (h *WarmImageCacheHandler) Execute(ctx context.Context, payload []byte) (*task.TaskResult, error) { + normalizedPayload, err := h.ValidatePayload(payload) + if err != nil { + task.AppendLog(ctx, "图片缓存预热参数无效: %v", err) + return nil, err + } + + var req WarmImageCachePayload + if err := json.Unmarshal(normalizedPayload, &req); err != nil { + return nil, fmt.Errorf(errParseImageCacheWarmupPayload, err) + } + + task.AppendLog(ctx, "等待获取图片缓存预热执行锁,质量: %s", req.Quality) + warmImageCacheMu.Lock() + defer warmImageCacheMu.Unlock() + + const ( + batchSize = 50 + maxFailureLogs = 5 + ) + var lastID uint64 + var totalProcessed int + var totalCached int + var totalGenerated int + var totalFailed int + + task.AppendLog(ctx, "开始串行预热图片压缩缓存,质量: %s,每批: %d", req.Quality, batchSize) + + for { + if err := ctx.Err(); err != nil { + return nil, fmt.Errorf("image cache warmup canceled: %w", err) + } + + var uploads []model.Upload + if err := db.DB(ctx). + Where("id > ? AND status != ? AND (LOWER(mime_type) LIKE ? OR LOWER(extension) IN ?)", + lastID, + model.UploadStatusDeleted, + "image/%", + []string{"jpg", "jpeg", "png", "webp", "gif"}, + ). + Order("id ASC"). + Limit(batchSize). + Find(&uploads).Error; err != nil { + task.AppendLog(ctx, "查询图片上传记录失败: %v", err) + return nil, fmt.Errorf(errQueryImagesForCacheWarmup, err) + } + + if len(uploads) == 0 { + break + } + + batchGenerated := 0 + batchCached := 0 + batchFailed := 0 + for i := range uploads { + if err := ctx.Err(); err != nil { + return nil, fmt.Errorf("image cache warmup canceled: %w", err) + } + + upload := &uploads[i] + totalProcessed++ + lastID = upload.ID + + _, cacheHit, err := ensureCompressedImageCache(ctx, upload, req.Quality) + if err != nil { + totalFailed++ + batchFailed++ + if totalFailed <= maxFailureLogs { + task.AppendLog(ctx, "图片处理失败 [ID:%d]: %v", upload.ID, err) + } + continue + } + if cacheHit { + totalCached++ + batchCached++ + continue + } + totalGenerated++ + batchGenerated++ + } + + task.AppendLog( + ctx, + "批次完成,末尾 ID: %d,生成: %d,命中: %d,失败: %d", + lastID, + batchGenerated, + batchCached, + batchFailed, + ) + } + + msg := fmt.Sprintf( + "图片缓存预热完成,共处理 %d 张,生成 %d 张,命中 %d 张,失败 %d 张", + totalProcessed, + totalGenerated, + totalCached, + totalFailed, + ) + 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 bb3c8cd9..71cf3555 100644 --- a/internal/apps/upload/tasks_test.go +++ b/internal/apps/upload/tasks_test.go @@ -5,12 +5,20 @@ package upload import ( + "bytes" "context" + "encoding/json" + "image" + "image/color" + "image/png" "io" + "os" + "path/filepath" "testing" "time" "github.com/Rain-kl/Wavelet/internal/db" + "github.com/Rain-kl/Wavelet/internal/diskcache" "github.com/Rain-kl/Wavelet/internal/model" "github.com/Rain-kl/Wavelet/internal/storage" "github.com/Rain-kl/Wavelet/internal/task" @@ -125,3 +133,188 @@ func TestCleanupUnusedUploadsHandler_ImplementsTaskHandler(t *testing.T) { // 编译期验证 CleanupUnusedUploadsHandler 实现了 TaskHandler 接口 var _ task.TaskHandler = (*CleanupUnusedUploadsHandler)(nil) } + +func TestWarmImageCacheHandlerValidatePayload(t *testing.T) { + tests := []struct { + name string + payload []byte + wantQuality string + wantErr bool + }{ + { + name: "normalizes quality", + payload: []byte(`{"quality":" HIGH "}`), + wantQuality: imageQualityHigh, + }, + { + name: "empty payload", + wantErr: true, + }, + { + name: "invalid json", + payload: []byte(`{`), + wantErr: true, + }, + { + name: "origin is not a compressed quality", + payload: []byte(`{"quality":"origin"}`), + wantErr: true, + }, + { + name: "unsupported quality", + payload: []byte(`{"quality":"maximum"}`), + wantErr: true, + }, + } + + handler := &WarmImageCacheHandler{} + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + gotPayload, err := handler.ValidatePayload(tt.payload) + if gotErr := err != nil; gotErr != tt.wantErr { + t.Fatalf("ValidatePayload(%s) error = %v, want error presence = %t", tt.payload, err, tt.wantErr) + } + if tt.wantErr { + return + } + + var got WarmImageCachePayload + if err := json.Unmarshal(gotPayload, &got); err != nil { + t.Fatalf("json.Unmarshal(%s) returned error: %v", gotPayload, err) + } + if got.Quality != tt.wantQuality { + t.Errorf("ValidatePayload(%s).Quality = %q, want %q", tt.payload, got.Quality, tt.wantQuality) + } + }) + } +} + +func TestWarmImageCacheHandlerExecute(t *testing.T) { + dbConn, _, cleanup := testhelper.SetupTestEnvironment(t) + defer cleanup() + + cache := diskcache.GetGlobalCache() + if err := cache.Clear(); err != nil { + t.Fatalf("Clear() before test returned error: %v", err) + } + t.Cleanup(func() { + if err := cache.Clear(); err != nil { + t.Errorf("Clear() after test returned error: %v", err) + } + }) + + testDir := t.TempDir() + firstPath := filepath.Join(testDir, "first.png") + secondPath := filepath.Join(testDir, "second.jpg") + writeTaskTestPNG(t, firstPath, color.RGBA{R: 255, A: 255}) + writeTaskTestPNG(t, secondPath, color.RGBA{G: 255, A: 255}) + + records := []model.Upload{ + { + ID: 4101, + UserID: 1001, + FileName: "first.png", + FilePath: firstPath, + MimeType: "image/png", + Extension: "png", + StorageDriver: storageDriverLocal, + Status: model.UploadStatusUsed, + }, + { + ID: 4102, + UserID: 1001, + FileName: "second.jpg", + FilePath: secondPath, + MimeType: "application/octet-stream", + Extension: "jpg", + StorageDriver: storageDriverLocal, + Status: model.UploadStatusPending, + }, + { + ID: 4103, + UserID: 1001, + FileName: "notes.txt", + FilePath: filepath.Join(testDir, "notes.txt"), + MimeType: "text/plain", + Extension: "txt", + StorageDriver: storageDriverLocal, + Status: model.UploadStatusUsed, + }, + { + ID: 4104, + UserID: 1001, + FileName: "deleted.png", + FilePath: firstPath, + MimeType: "image/png", + Extension: "png", + StorageDriver: storageDriverLocal, + Status: model.UploadStatusDeleted, + }, + } + for i := range records { + if info, err := os.Stat(records[i].FilePath); err == nil { + records[i].FileSize = info.Size() + } + if err := dbConn.Create(&records[i]).Error; err != nil { + t.Fatalf("failed to create upload %d: %v", records[i].ID, err) + } + } + + handler := &WarmImageCacheHandler{} + payload := []byte(`{"quality":"low"}`) + + result, err := handler.Execute(context.Background(), payload) + if err != nil { + t.Fatalf("Execute(%s) returned error: %v", payload, err) + } + if result == nil { + t.Fatal("Execute() result = nil, want non-nil") + } + if result.Message != "图片缓存预热完成,共处理 2 张,生成 2 张,命中 0 张,失败 0 张" { + t.Errorf("Execute() message = %q, want generated summary", result.Message) + } + + for i := range records[:2] { + key := imageCompressionCacheKey(&records[i], imageQualityLow) + got, err := cache.Get(key) + if err != nil { + t.Errorf("cache.Get(%q) returned error: %v", key, err) + continue + } + if len(got) == 0 { + t.Errorf("cache.Get(%q) returned empty WebP data", key) + } + } + + secondResult, err := handler.Execute(context.Background(), payload) + if err != nil { + t.Fatalf("second Execute(%s) returned error: %v", payload, err) + } + if secondResult.Message != "图片缓存预热完成,共处理 2 张,生成 0 张,命中 2 张,失败 0 张" { + t.Errorf("second Execute() message = %q, want cache-hit summary", secondResult.Message) + } +} + +func TestWarmImageCacheHandlerImplementsTaskInterfaces(t *testing.T) { + var _ task.TaskHandler = (*WarmImageCacheHandler)(nil) + var _ task.PayloadValidator = (*WarmImageCacheHandler)(nil) +} + +func writeTaskTestPNG(t *testing.T, path string, fill color.RGBA) { + t.Helper() + + img := image.NewRGBA(image.Rect(0, 0, 2, 2)) + for y := 0; y < 2; y++ { + for x := 0; x < 2; x++ { + img.Set(x, y, fill) + } + } + + var buf bytes.Buffer + if err := png.Encode(&buf, img); err != nil { + t.Fatalf("png.Encode() returned error: %v", err) + } + if err := os.WriteFile(path, buf.Bytes(), 0o600); err != nil { + t.Fatalf("os.WriteFile(%q) returned error: %v", path, err) + } +} diff --git a/internal/task/handlers/register.go b/internal/task/handlers/register.go index f980e485..9d7304e7 100644 --- a/internal/task/handlers/register.go +++ b/internal/task/handlers/register.go @@ -16,6 +16,8 @@ func Register() { // upload task.RegisterHandler(upload.CleanupUnusedUploadsTask, &upload.CleanupUnusedUploadsHandler{}) task.RegisterTaskMeta(upload.CleanupUnusedUploadsMeta) + task.RegisterHandler(upload.WarmImageCacheTask, &upload.WarmImageCacheHandler{}) + task.RegisterTaskMeta(upload.WarmImageCacheMeta) // user task.RegisterHandler(user.SendEmailTask, &user.SendEmailHandler{})