diff --git a/backend/plugins/domain/upload/task/storage_migration.go b/backend/plugins/domain/upload/task/storage_migration.go index e0d927e1..bf91c49a 100644 --- a/backend/plugins/domain/upload/task/storage_migration.go +++ b/backend/plugins/domain/upload/task/storage_migration.go @@ -12,7 +12,6 @@ import ( "context" "crypto/sha256" "encoding/hex" - "encoding/json" "errors" "fmt" "io" @@ -121,7 +120,7 @@ func (h *MigrationHandler) Execute(ctx context.Context, payload []byte) (*contra }) } - active, err := loadActiveStorageConfig(ctx) + active, err := uploadstorage.LoadStorageConfig(ctx) if err != nil { return nil, fmt.Errorf("load active storage config: %w", err) } @@ -130,7 +129,7 @@ func (h *MigrationHandler) Execute(ctx context.Context, payload []byte) (*contra return nil, err } if target.Driver == active.Driver { - if err := saveActiveStorageConfig(ctx, target); err != nil { + if err := uploadstorage.SaveActiveConfig(ctx, target); err != nil { return nil, fmt.Errorf("activate same-driver storage config: %w", err) } message := fmt.Sprintf("存储配置已更新,活动存储保持为 %s", target.Driver) @@ -143,7 +142,7 @@ func (h *MigrationHandler) Execute(ctx context.Context, payload []byte) (*contra return nil, fmt.Errorf("count source objects: %w", err) } if total == 0 { - if err := saveActiveStorageConfig(ctx, target); err != nil { + if err := uploadstorage.SaveActiveConfig(ctx, target); err != nil { return nil, fmt.Errorf("activate empty storage config: %w", err) } message := fmt.Sprintf("当前存储没有需要迁移的对象,活动存储已切换为 %s", target.Driver) @@ -162,7 +161,7 @@ func (h *MigrationHandler) Execute(ctx context.Context, payload []byte) (*contra return nil, err } - if err := saveActiveStorageConfig(ctx, target); err != nil { + if err := uploadstorage.SaveActiveConfig(ctx, target); err != nil { return nil, fmt.Errorf("activate target storage: %w", err) } message := fmt.Sprintf("存储迁移完成,共迁移 %d 个对象,活动存储已切换为 %s", migrated, target.Driver) @@ -170,33 +169,6 @@ func (h *MigrationHandler) Execute(ctx context.Context, payload []byte) (*contra return &contracts.TaskResultDTO{Message: message}, nil } -func loadActiveStorageConfig(ctx context.Context) (contracts.StorageConfigDTO, error) { - var val string - db := shared.GetDB(ctx) - if db != nil { - _ = db.Table("w_system_configs").Where("key = ?", "storage_config").Pluck("value", &val).Error - } - var cfg contracts.StorageConfigDTO - if val != "" { - _ = json.Unmarshal([]byte(val), &cfg) - } - return cfg, nil -} - -func saveActiveStorageConfig(ctx context.Context, cfg contracts.StorageConfigDTO) error { - data, err := json.Marshal(cfg) - if err != nil { - return err - } - db := shared.GetDB(ctx) - if db == nil { - return errors.New("database not available") - } - return db.Table("w_system_configs"). - Where("key = ?", "storage_config"). - Update("value", string(data)).Error -} - func countStorageObjects(ctx context.Context) (int64, error) { var count int64 db := shared.GetDB(ctx) diff --git a/backend/plugins/domain/upload/task/storage_migration_config_test.go b/backend/plugins/domain/upload/task/storage_migration_config_test.go new file mode 100644 index 00000000..f56bba10 --- /dev/null +++ b/backend/plugins/domain/upload/task/storage_migration_config_test.go @@ -0,0 +1,37 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package task + +import ( + "Wavelet/plugins/domain/upload/shared" + + "context" + "strings" + "testing" +) + +// TestMigrationHandlerRejectsCorruptActiveConfig pins the failure path the task +// package used to swallow. Its local loader discarded the read error and the +// JSON parse error and returned a zero config with a nil error, so a migration +// could proceed believing the active storage driver was simply unknown. +func TestMigrationHandlerRejectsCorruptActiveConfig(t *testing.T) { + dbConn, cleanup := shared.SetupTestEnv(t) + defer cleanup() + + if err := dbConn.Table("w_system_configs"). + Where("key = ?", "storage_config"). + Update("value", "{ this is not json").Error; err != nil { + t.Fatalf("corrupt stored config: %v", err) + } + + payload := []byte(`{"target":{"driver":"s3","s3":{"region":"us-east-1","bucket":"dst"}}}`) + + result, err := (&MigrationHandler{}).Execute(context.Background(), payload) + if err == nil { + t.Fatalf("Execute() = (%v, nil), want an error for an unparsable active config", result) + } + if !strings.Contains(err.Error(), "load active storage config") { + t.Errorf("Execute() error = %q, want it to name the failed active-config read", err) + } +}