图片预热任务

This commit is contained in:
ryan
2026-06-11 17:49:08 +08:00
parent 7da4b72d24
commit 5915b31519
6 changed files with 429 additions and 44 deletions
+7 -1
View File
@@ -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) {
+30 -25
View File
@@ -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"
)
+32 -18
View File
@@ -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 {
+165
View File
@@ -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
}
+193
View File
@@ -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)
}
}
+2
View File
@@ -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{})