mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-10-11 01:36:37 +08:00
fix(persistence): migrate all pkg/persistence imports to plugins/infra/database and plugins/infra/cache
- Replace db.DB(ctx) with database.DB(ctx) from plugins/infra/database
- Replace db.Redis/db.PrefixedKey/db.GetJSON/db.SetJSON with cachepkg.* from plugins/infra/cache
- Replace pkg/persistence/idgen with pkg/idgen (already exists)
- Replace pkg/persistence/batchwriter with pkg/batchwriter (already exists)
- Replace pkg/persistence/migrator with pkg/migrator (already exists)
- Replace pkg/persistence/logstore with plugins/domain/risk_control/logstore
- Delete defunct pkg/{persistence,cap,message_gateway,push,shared,task}
- Fix vet issues: db alias in domain_test.go, driver_asynq_worker.TaskHandler reference
- Update Makefile architecture guard
- Update docs and skill references
- Update go.mod: gorilla/sessions promotion to direct dependency
This commit is contained in:
@@ -15,14 +15,16 @@ import (
|
||||
"sync/atomic"
|
||||
"time"
|
||||
|
||||
"github.com/Rain-kl/Wavelet/pkg/persistence"
|
||||
"github.com/Rain-kl/Wavelet/pkg/task"
|
||||
"golang.org/x/sync/errgroup"
|
||||
|
||||
cache "github.com/Rain-kl/Wavelet/plugins/infra/cache"
|
||||
database "github.com/Rain-kl/Wavelet/plugins/infra/database"
|
||||
"github.com/Rain-kl/Wavelet/pkg/util"
|
||||
"github.com/Rain-kl/Wavelet/plugins/domain/upload/models"
|
||||
uploadstats "github.com/Rain-kl/Wavelet/plugins/domain/upload/stats"
|
||||
uploadstorage "github.com/Rain-kl/Wavelet/plugins/domain/upload/storage"
|
||||
"github.com/Rain-kl/Wavelet/plugins/drivers/driver_asynq_worker"
|
||||
"github.com/Rain-kl/Wavelet/plugins/infra/storage/objectstore"
|
||||
"golang.org/x/sync/errgroup"
|
||||
)
|
||||
|
||||
const (
|
||||
@@ -33,16 +35,16 @@ const (
|
||||
)
|
||||
|
||||
// StorageMigrationMeta describes the manually dispatchable migration task.
|
||||
var StorageMigrationMeta = task.TaskMeta{
|
||||
var StorageMigrationMeta = driver_asynq_worker.TaskMeta{
|
||||
Type: TaskTypeStorageMigration,
|
||||
AsynqTask: StorageMigrationTask,
|
||||
Name: "迁移文件存储",
|
||||
Description: "将活动存储中的文件迁移到待切换的目标存储,迁移期间文件系统保持只读",
|
||||
SupportsTime: false,
|
||||
MaxRetry: task.DefaultMaxRetry,
|
||||
Queue: task.QueueDefault,
|
||||
MaxRetry: driver_asynq_worker.DefaultMaxRetry,
|
||||
Queue: driver_asynq_worker.QueueDefault,
|
||||
Retryable: true,
|
||||
Params: []task.TaskParam{
|
||||
Params: []driver_asynq_worker.TaskParam{
|
||||
{
|
||||
Name: "target",
|
||||
Label: "目标存储配置 (JSON)",
|
||||
@@ -74,15 +76,15 @@ func (h *MigrationHandler) ValidatePayload(payload []byte) ([]byte, error) {
|
||||
}
|
||||
|
||||
// Execute migrates all unique active-storage objects to the pending backend.
|
||||
func (h *MigrationHandler) Execute(ctx context.Context, payload []byte) (*task.TaskResult, error) {
|
||||
if db.Redis != nil {
|
||||
func (h *MigrationHandler) Execute(ctx context.Context, payload []byte) (*driver_asynq_worker.TaskResult, error) {
|
||||
if cache.Redis != nil {
|
||||
const (
|
||||
cleanupTimeout = 5 * time.Second
|
||||
renewalInterval = 10 * time.Minute
|
||||
)
|
||||
|
||||
lockKey := db.PrefixedKey("lock:storage:migrate")
|
||||
ok, err := db.Redis.SetNX(ctx, lockKey, "locked", time.Hour).Result()
|
||||
lockKey := cache.PrefixedKey("lock:storage:migrate")
|
||||
ok, err := cache.Redis.SetNX(ctx, lockKey, "locked", time.Hour).Result()
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("acquire migration lock: %w", err)
|
||||
}
|
||||
@@ -96,7 +98,7 @@ func (h *MigrationHandler) Execute(ctx context.Context, payload []byte) (*task.T
|
||||
close(stopRenewal)
|
||||
cleanupCtx, cancel := context.WithTimeout(context.Background(), cleanupTimeout)
|
||||
defer cancel()
|
||||
_ = db.Redis.Del(cleanupCtx, lockKey)
|
||||
_ = cache.Redis.Del(cleanupCtx, lockKey)
|
||||
}()
|
||||
|
||||
//nolint:contextcheck,gosec
|
||||
@@ -107,7 +109,7 @@ func (h *MigrationHandler) Execute(ctx context.Context, payload []byte) (*task.T
|
||||
select {
|
||||
case <-ticker.C:
|
||||
renewCtx, cancel := context.WithTimeout(context.Background(), cleanupTimeout)
|
||||
_ = db.Redis.Expire(renewCtx, lockKey, time.Hour).Err()
|
||||
_ = cache.Redis.Expire(renewCtx, lockKey, time.Hour).Err()
|
||||
cancel()
|
||||
case <-stopRenewal:
|
||||
return
|
||||
@@ -131,8 +133,8 @@ func (h *MigrationHandler) Execute(ctx context.Context, payload []byte) (*task.T
|
||||
return nil, fmt.Errorf("activate same-driver storage config: %w", err)
|
||||
}
|
||||
message := fmt.Sprintf("存储配置已更新,活动存储保持为 %s", target.Driver)
|
||||
task.AppendLog(ctx, "%s", message)
|
||||
return &task.TaskResult{Message: message}, nil
|
||||
driver_asynq_worker.AppendLog(ctx, "%s", message)
|
||||
return &driver_asynq_worker.TaskResult{Message: message}, nil
|
||||
}
|
||||
|
||||
total, err := countStorageObjects(ctx)
|
||||
@@ -144,8 +146,8 @@ func (h *MigrationHandler) Execute(ctx context.Context, payload []byte) (*task.T
|
||||
return nil, fmt.Errorf("activate empty storage config: %w", err)
|
||||
}
|
||||
message := fmt.Sprintf("当前存储没有需要迁移的对象,活动存储已切换为 %s", target.Driver)
|
||||
task.AppendLog(ctx, "%s", message)
|
||||
return &task.TaskResult{Message: message}, nil
|
||||
driver_asynq_worker.AppendLog(ctx, "%s", message)
|
||||
return &driver_asynq_worker.TaskResult{Message: message}, nil
|
||||
}
|
||||
|
||||
sourceBackend, err := objectstore.NewBackend(ctx, active, active.Driver)
|
||||
@@ -157,7 +159,7 @@ func (h *MigrationHandler) Execute(ctx context.Context, payload []byte) (*task.T
|
||||
return nil, fmt.Errorf("create target storage: %w", err)
|
||||
}
|
||||
|
||||
task.AppendLog(ctx, "开始存储迁移: %s -> %s,总对象数: %d", active.Driver, target.Driver, total)
|
||||
driver_asynq_worker.AppendLog(ctx, "开始存储迁移: %s -> %s,总对象数: %d", active.Driver, target.Driver, total)
|
||||
migrated, err := migrateObjects(ctx, sourceBackend, targetBackend, total)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
@@ -167,13 +169,13 @@ func (h *MigrationHandler) Execute(ctx context.Context, payload []byte) (*task.T
|
||||
return nil, fmt.Errorf("activate target storage: %w", err)
|
||||
}
|
||||
message := fmt.Sprintf("存储迁移完成,共迁移 %d 个对象,活动存储已切换为 %s", migrated, target.Driver)
|
||||
task.AppendLog(ctx, "%s", message)
|
||||
return &task.TaskResult{Message: message}, nil
|
||||
driver_asynq_worker.AppendLog(ctx, "%s", message)
|
||||
return &driver_asynq_worker.TaskResult{Message: message}, nil
|
||||
}
|
||||
|
||||
func countStorageObjects(ctx context.Context) (int64, error) {
|
||||
var count int64
|
||||
err := db.DB(ctx).Model(&models.Upload{}).
|
||||
err := database.DB(ctx).Model(&models.Upload{}).
|
||||
Where("status != ?", models.UploadStatusDeleted).
|
||||
Distinct("file_path").
|
||||
Count(&count).Error
|
||||
@@ -185,7 +187,7 @@ func hasUnresolvedMigrationTask(ctx context.Context) (bool, error) {
|
||||
if err != nil || !ok {
|
||||
return false, err
|
||||
}
|
||||
return execution.Status == task.TaskExecutionStatusPending || execution.Status == task.TaskExecutionStatusRunning, nil
|
||||
return execution.Status == driver_asynq_worker.TaskExecutionStatusPending || execution.Status == driver_asynq_worker.TaskExecutionStatusRunning, nil
|
||||
}
|
||||
|
||||
type migrationObject struct {
|
||||
@@ -211,10 +213,10 @@ func migrateObjects(
|
||||
return atomic.LoadInt64(&migrated), fmt.Errorf("storage migration canceled: %w", err)
|
||||
}
|
||||
|
||||
task.AppendLog(ctx, "正在查询待迁移对象批次,当前已完成迁移: %d/%d", atomic.LoadInt64(&migrated), total)
|
||||
driver_asynq_worker.AppendLog(ctx, "正在查询待迁移对象批次,当前已完成迁移: %d/%d", atomic.LoadInt64(&migrated), total)
|
||||
|
||||
var objects []migrationObject
|
||||
query := db.DB(ctx).Model(&models.Upload{}).
|
||||
query := database.DB(ctx).Model(&models.Upload{}).
|
||||
Select("file_path, MAX(file_size) AS file_size, MAX(mime_type) AS mime_type, MAX(hash) AS hash").
|
||||
Where("status != ?", models.UploadStatusDeleted)
|
||||
if lastFilePath != "" {
|
||||
@@ -227,12 +229,12 @@ func migrateObjects(
|
||||
return atomic.LoadInt64(&migrated), fmt.Errorf("query source objects: %w", err)
|
||||
}
|
||||
if len(objects) == 0 {
|
||||
task.AppendLog(ctx, "所有对象迁移完毕")
|
||||
driver_asynq_worker.AppendLog(ctx, "所有对象迁移完毕")
|
||||
break
|
||||
}
|
||||
|
||||
lastFilePath = objects[len(objects)-1].FilePath
|
||||
task.AppendLog(ctx, "获取当前批次迁移对象,批次大小: %d,实际获取对象数: %d", batchSize, len(objects))
|
||||
driver_asynq_worker.AppendLog(ctx, "获取当前批次迁移对象,批次大小: %d,实际获取对象数: %d", batchSize, len(objects))
|
||||
|
||||
var g errgroup.Group
|
||||
g.SetLimit(migrationConcurrency)
|
||||
@@ -252,7 +254,7 @@ func migrateObjects(
|
||||
return atomic.LoadInt64(&migrated), err
|
||||
}
|
||||
|
||||
task.AppendLog(ctx, "当前批次迁移完成。迁移进度: %d/%d", atomic.LoadInt64(&migrated), total)
|
||||
driver_asynq_worker.AppendLog(ctx, "当前批次迁移完成。迁移进度: %d/%d", atomic.LoadInt64(&migrated), total)
|
||||
}
|
||||
return atomic.LoadInt64(&migrated), nil
|
||||
}
|
||||
@@ -265,11 +267,11 @@ func migrateSingleObject(
|
||||
sha256HexLength int,
|
||||
) error {
|
||||
if shouldSkipMigration(ctx, targetBackend, obj) {
|
||||
task.AppendLog(ctx, "[跳过迁移] 目标存储已存在相同文件: %s", obj.FilePath)
|
||||
driver_asynq_worker.AppendLog(ctx, "[跳过迁移] 目标存储已存在相同文件: %s", obj.FilePath)
|
||||
return nil
|
||||
}
|
||||
|
||||
task.AppendLog(ctx, "[迁移开始] 正在从源存储读取文件: %s", obj.FilePath)
|
||||
driver_asynq_worker.AppendLog(ctx, "[迁移开始] 正在从源存储读取文件: %s", obj.FilePath)
|
||||
source, err := sourceBackend.Get(ctx, obj.FilePath)
|
||||
if err != nil {
|
||||
if isNotFoundError(err) {
|
||||
@@ -277,7 +279,7 @@ func migrateSingleObject(
|
||||
}
|
||||
return fmt.Errorf("open source object %q: %w", obj.FilePath, err)
|
||||
}
|
||||
task.AppendLog(ctx, "[传输中] 正在向目标存储上传文件: %s (大小: %d 字节, 类型: %s)", obj.FilePath, obj.FileSize, obj.MimeType)
|
||||
driver_asynq_worker.AppendLog(ctx, "[传输中] 正在向目标存储上传文件: %s (大小: %d 字节, 类型: %s)", obj.FilePath, obj.FileSize, obj.MimeType)
|
||||
targetResult, putErr := targetBackend.Put(ctx, obj.FilePath, source.Body, obj.FileSize, obj.MimeType)
|
||||
closeErr := source.Body.Close()
|
||||
if putErr != nil {
|
||||
@@ -288,7 +290,7 @@ func migrateSingleObject(
|
||||
}
|
||||
|
||||
if len(obj.Hash) == sha256HexLength {
|
||||
task.AppendLog(ctx, "[校验中] 正在对目标文件进行数据一致性校验 (SHA-256): %s", targetResult.Key)
|
||||
driver_asynq_worker.AppendLog(ctx, "[校验中] 正在对目标文件进行数据一致性校验 (SHA-256): %s", targetResult.Key)
|
||||
targetObj, getErr := targetBackend.Get(ctx, targetResult.Key)
|
||||
if getErr != nil {
|
||||
return fmt.Errorf("retrieve target object for verification %q: %w", obj.FilePath, getErr)
|
||||
@@ -306,18 +308,18 @@ func migrateSingleObject(
|
||||
if computedHash != obj.Hash {
|
||||
return fmt.Errorf("integrity check failed for %q: got hash %s, want %s", obj.FilePath, computedHash, obj.Hash)
|
||||
}
|
||||
task.AppendLog(ctx, "[校验通过] 文件一致性校验成功: %s", targetResult.Key)
|
||||
driver_asynq_worker.AppendLog(ctx, "[校验通过] 文件一致性校验成功: %s", targetResult.Key)
|
||||
}
|
||||
|
||||
if targetResult.Key != obj.FilePath {
|
||||
task.AppendLog(ctx, "[更新数据库] 正在更新文件路径: %s -> %s", obj.FilePath, targetResult.Key)
|
||||
if err := db.DB(ctx).Model(&models.Upload{}).
|
||||
driver_asynq_worker.AppendLog(ctx, "[更新数据库] 正在更新文件路径: %s -> %s", obj.FilePath, targetResult.Key)
|
||||
if err := database.DB(ctx).Model(&models.Upload{}).
|
||||
Where("file_path = ? AND status != ?", obj.FilePath, models.UploadStatusDeleted).
|
||||
Update("file_path", targetResult.Key).Error; err != nil {
|
||||
return fmt.Errorf("update migrated object %q: %w", obj.FilePath, err)
|
||||
}
|
||||
}
|
||||
task.AppendLog(ctx, "[迁移成功] 文件已完成迁移: %s", targetResult.Key)
|
||||
driver_asynq_worker.AppendLog(ctx, "[迁移成功] 文件已完成迁移: %s", targetResult.Key)
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -342,15 +344,15 @@ func markMissingMigrationObjectDeleted(
|
||||
filePath string,
|
||||
sourceErr error,
|
||||
) error {
|
||||
task.AppendLog(ctx, "警告: 源存储中物理文件不存在,标记为已删除并跳过: %s (错误: %v)", filePath, sourceErr)
|
||||
driver_asynq_worker.AppendLog(ctx, "警告: 源存储中物理文件不存在,标记为已删除并跳过: %s (错误: %v)", filePath, sourceErr)
|
||||
|
||||
var affectedUploads []models.Upload
|
||||
if err := db.DB(ctx).
|
||||
if err := database.DB(ctx).
|
||||
Where("file_path = ? AND status != ?", filePath, models.UploadStatusDeleted).
|
||||
Find(&affectedUploads).Error; err != nil {
|
||||
return fmt.Errorf("load missing object uploads %q: %w", filePath, err)
|
||||
}
|
||||
if err := db.DB(ctx).Model(&models.Upload{}).
|
||||
if err := database.DB(ctx).Model(&models.Upload{}).
|
||||
Where("file_path = ?", filePath).
|
||||
Update("status", models.UploadStatusDeleted).Error; err != nil {
|
||||
return fmt.Errorf("update missing object %q: %w", filePath, err)
|
||||
|
||||
Reference in New Issue
Block a user