mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-10-07 08:06:37 +08:00
refactor(structure): group platform, infra, and shared packages
Reorganize internal packages into platform/infra/shared layers and update imports, docs, and seed-count tests to match current system configs.
This commit is contained in:
+3
-3
@@ -13,10 +13,10 @@ import (
|
||||
|
||||
"github.com/Rain-kl/Wavelet/internal/apps/upload/shared"
|
||||
uploadstorage "github.com/Rain-kl/Wavelet/internal/apps/upload/storage"
|
||||
"github.com/Rain-kl/Wavelet/internal/db"
|
||||
"github.com/Rain-kl/Wavelet/internal/infra/objectstore"
|
||||
db "github.com/Rain-kl/Wavelet/internal/infra/persistence"
|
||||
"github.com/Rain-kl/Wavelet/internal/model"
|
||||
"github.com/Rain-kl/Wavelet/internal/repository"
|
||||
"github.com/Rain-kl/Wavelet/internal/storage"
|
||||
)
|
||||
|
||||
const fileAccessInvalidationChannel = "upload:file_access_invalidation"
|
||||
@@ -59,7 +59,7 @@ func startAccessCacheInvalidationListener() {
|
||||
go func() {
|
||||
pubsub := db.Redis.Subscribe(
|
||||
context.Background(),
|
||||
storage.ConfigInvalidationChannel,
|
||||
objectstore.ConfigInvalidationChannel,
|
||||
fileAccessInvalidationChannel,
|
||||
)
|
||||
defer func() {
|
||||
|
||||
+1
-1
@@ -10,7 +10,7 @@ import (
|
||||
|
||||
"github.com/Rain-kl/Wavelet/internal/apps/upload/shared"
|
||||
uploadstorage "github.com/Rain-kl/Wavelet/internal/apps/upload/storage"
|
||||
"github.com/Rain-kl/Wavelet/internal/db"
|
||||
db "github.com/Rain-kl/Wavelet/internal/infra/persistence"
|
||||
"github.com/Rain-kl/Wavelet/internal/model"
|
||||
"github.com/Rain-kl/Wavelet/internal/repository"
|
||||
"github.com/Rain-kl/Wavelet/internal/testhelper"
|
||||
|
||||
+1
-1
@@ -9,7 +9,7 @@ import (
|
||||
"fmt"
|
||||
"sync"
|
||||
|
||||
"github.com/Rain-kl/Wavelet/internal/db"
|
||||
db "github.com/Rain-kl/Wavelet/internal/infra/persistence"
|
||||
"github.com/Rain-kl/Wavelet/internal/model"
|
||||
"github.com/Rain-kl/Wavelet/pkg/cache/ram"
|
||||
)
|
||||
|
||||
+1
-1
@@ -9,7 +9,7 @@ import (
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/Rain-kl/Wavelet/internal/db"
|
||||
db "github.com/Rain-kl/Wavelet/internal/infra/persistence"
|
||||
"github.com/Rain-kl/Wavelet/internal/model"
|
||||
"github.com/Rain-kl/Wavelet/internal/testhelper"
|
||||
"gorm.io/gorm"
|
||||
|
||||
@@ -12,7 +12,7 @@ import (
|
||||
uploadstats "github.com/Rain-kl/Wavelet/internal/apps/upload/stats"
|
||||
uploadtask "github.com/Rain-kl/Wavelet/internal/apps/upload/task"
|
||||
"github.com/Rain-kl/Wavelet/internal/apps/upload/util"
|
||||
"github.com/Rain-kl/Wavelet/internal/task"
|
||||
"github.com/Rain-kl/Wavelet/internal/infra/task"
|
||||
)
|
||||
|
||||
// HTTP handlers
|
||||
|
||||
@@ -20,10 +20,10 @@ import (
|
||||
"github.com/Rain-kl/Wavelet/internal/apps/upload/shared"
|
||||
uploadstorage "github.com/Rain-kl/Wavelet/internal/apps/upload/storage"
|
||||
"github.com/Rain-kl/Wavelet/internal/apps/upload/util"
|
||||
"github.com/Rain-kl/Wavelet/internal/common"
|
||||
"github.com/Rain-kl/Wavelet/internal/common/response"
|
||||
"github.com/Rain-kl/Wavelet/internal/diskcache"
|
||||
"github.com/Rain-kl/Wavelet/internal/infra/diskcache"
|
||||
"github.com/Rain-kl/Wavelet/internal/model"
|
||||
appshared "github.com/Rain-kl/Wavelet/internal/shared"
|
||||
"github.com/Rain-kl/Wavelet/internal/shared/response"
|
||||
|
||||
"github.com/Rain-kl/Wavelet/pkg/logger"
|
||||
"github.com/gin-gonic/gin"
|
||||
@@ -77,7 +77,7 @@ func ServeFileByID(c *gin.Context) {
|
||||
}
|
||||
|
||||
if err := CheckFileAccessPermission(c, upload); err != nil {
|
||||
response.AbortUnauthorized(c, common.UnAuthorized)
|
||||
response.AbortUnauthorized(c, appshared.UnAuthorized)
|
||||
return
|
||||
}
|
||||
|
||||
|
||||
@@ -20,13 +20,13 @@ import (
|
||||
"github.com/Rain-kl/Wavelet/internal/apps/upload/cache"
|
||||
"github.com/Rain-kl/Wavelet/internal/apps/upload/shared"
|
||||
"github.com/Rain-kl/Wavelet/internal/apps/upload/util"
|
||||
"github.com/Rain-kl/Wavelet/internal/common"
|
||||
"github.com/Rain-kl/Wavelet/internal/common/response"
|
||||
"github.com/Rain-kl/Wavelet/internal/db"
|
||||
"github.com/Rain-kl/Wavelet/internal/diskcache"
|
||||
"github.com/Rain-kl/Wavelet/internal/infra/diskcache"
|
||||
"github.com/Rain-kl/Wavelet/internal/infra/objectstore"
|
||||
db "github.com/Rain-kl/Wavelet/internal/infra/persistence"
|
||||
"github.com/Rain-kl/Wavelet/internal/model"
|
||||
"github.com/Rain-kl/Wavelet/internal/repository"
|
||||
"github.com/Rain-kl/Wavelet/internal/storage"
|
||||
appshared "github.com/Rain-kl/Wavelet/internal/shared"
|
||||
"github.com/Rain-kl/Wavelet/internal/shared/response"
|
||||
"github.com/Rain-kl/Wavelet/internal/testhelper"
|
||||
"github.com/gin-contrib/sessions"
|
||||
"github.com/gin-contrib/sessions/cookie"
|
||||
@@ -138,8 +138,8 @@ func TestServeFileByIDAccessControl(t *testing.T) {
|
||||
if err := json.Unmarshal(w.Body.Bytes(), &body); err != nil {
|
||||
t.Fatalf("failed to parse JSON: %v", err)
|
||||
}
|
||||
if body["error_msg"] != common.UnAuthorized {
|
||||
t.Errorf("expected error_msg %q, got %v", common.UnAuthorized, body["error_msg"])
|
||||
if body["error_msg"] != appshared.UnAuthorized {
|
||||
t.Errorf("expected error_msg %q, got %v", appshared.UnAuthorized, body["error_msg"])
|
||||
}
|
||||
})
|
||||
|
||||
@@ -361,7 +361,7 @@ func configureLocalStorageRoot(t *testing.T, dbConn *gorm.DB, tempDir string) {
|
||||
if err := dbConn.Where("key = ?", model.ConfigKeyStorageConfig).First(&sc).Error; err != nil {
|
||||
t.Fatalf("failed to find storage config: %v", err)
|
||||
}
|
||||
var cfg storage.Config
|
||||
var cfg objectstore.Config
|
||||
if err := json.Unmarshal([]byte(sc.Value), &cfg); err != nil {
|
||||
t.Fatalf("failed to unmarshal storage config: %v", err)
|
||||
}
|
||||
@@ -376,5 +376,5 @@ func configureLocalStorageRoot(t *testing.T, dbConn *gorm.DB, tempDir string) {
|
||||
}
|
||||
_ = db.HSetJSON(context.Background(), repository.SystemConfigRedisHashKey, sc.Key, &sc)
|
||||
repository.ResetSystemConfigRAMCacheForTest()
|
||||
storage.ResetCache()
|
||||
objectstore.ResetCache()
|
||||
}
|
||||
|
||||
@@ -12,9 +12,9 @@ import (
|
||||
"github.com/Rain-kl/Wavelet/internal/apps/upload/ingest"
|
||||
"github.com/Rain-kl/Wavelet/internal/apps/upload/shared"
|
||||
uploadstorage "github.com/Rain-kl/Wavelet/internal/apps/upload/storage"
|
||||
"github.com/Rain-kl/Wavelet/internal/common/response"
|
||||
"github.com/Rain-kl/Wavelet/internal/model"
|
||||
"github.com/Rain-kl/Wavelet/internal/repository"
|
||||
"github.com/Rain-kl/Wavelet/internal/shared/response"
|
||||
|
||||
"github.com/gin-gonic/gin"
|
||||
)
|
||||
|
||||
@@ -28,9 +28,9 @@ import (
|
||||
"github.com/Rain-kl/Wavelet/internal/apps/upload/shared"
|
||||
uploadstorage "github.com/Rain-kl/Wavelet/internal/apps/upload/storage"
|
||||
"github.com/Rain-kl/Wavelet/internal/apps/upload/util"
|
||||
"github.com/Rain-kl/Wavelet/internal/common"
|
||||
"github.com/Rain-kl/Wavelet/internal/common/response"
|
||||
"github.com/Rain-kl/Wavelet/internal/model"
|
||||
appshared "github.com/Rain-kl/Wavelet/internal/shared"
|
||||
"github.com/Rain-kl/Wavelet/internal/shared/response"
|
||||
"github.com/Rain-kl/Wavelet/pkg/logger"
|
||||
"github.com/gin-gonic/gin"
|
||||
"gorm.io/gorm"
|
||||
@@ -187,7 +187,7 @@ func DownloadFile(c *gin.Context) {
|
||||
}
|
||||
|
||||
if err := filesrv.CheckFileAccessPermission(c, upload); err != nil {
|
||||
response.AbortUnauthorized(c, common.UnAuthorized)
|
||||
response.AbortUnauthorized(c, appshared.UnAuthorized)
|
||||
return
|
||||
}
|
||||
|
||||
|
||||
@@ -23,11 +23,11 @@ import (
|
||||
"github.com/Rain-kl/Wavelet/internal/apps/oauth"
|
||||
"github.com/Rain-kl/Wavelet/internal/apps/upload/shared"
|
||||
uploadstats "github.com/Rain-kl/Wavelet/internal/apps/upload/stats"
|
||||
"github.com/Rain-kl/Wavelet/internal/common/response"
|
||||
"github.com/Rain-kl/Wavelet/internal/db"
|
||||
"github.com/Rain-kl/Wavelet/internal/infra/objectstore"
|
||||
db "github.com/Rain-kl/Wavelet/internal/infra/persistence"
|
||||
"github.com/Rain-kl/Wavelet/internal/model"
|
||||
"github.com/Rain-kl/Wavelet/internal/repository"
|
||||
"github.com/Rain-kl/Wavelet/internal/storage"
|
||||
"github.com/Rain-kl/Wavelet/internal/shared/response"
|
||||
"github.com/Rain-kl/Wavelet/internal/testhelper"
|
||||
"github.com/gin-gonic/gin"
|
||||
"gorm.io/gorm"
|
||||
@@ -117,7 +117,7 @@ func TestUploadFile(t *testing.T) {
|
||||
mockFiles := make(map[string][]byte)
|
||||
var putCount int
|
||||
|
||||
restoreStorage := storage.MockStorage(
|
||||
restoreStorage := objectstore.MockStorage(
|
||||
func(ctx context.Context, key string, body io.Reader, size int64, contentType string) error {
|
||||
data, err := io.ReadAll(body)
|
||||
if err != nil {
|
||||
@@ -127,12 +127,12 @@ func TestUploadFile(t *testing.T) {
|
||||
putCount++
|
||||
return nil
|
||||
},
|
||||
func(ctx context.Context, key string) (*storage.Object, error) {
|
||||
func(ctx context.Context, key string) (*objectstore.Object, error) {
|
||||
data, ok := mockFiles[key]
|
||||
if !ok {
|
||||
return nil, os.ErrNotExist
|
||||
}
|
||||
return &storage.Object{
|
||||
return &objectstore.Object{
|
||||
Body: io.NopCloser(bytes.NewReader(data)),
|
||||
ContentLength: int64(len(data)),
|
||||
ContentType: "application/octet-stream",
|
||||
@@ -146,9 +146,9 @@ func TestUploadFile(t *testing.T) {
|
||||
defer restoreStorage()
|
||||
|
||||
// 开启 S3 Storage
|
||||
storage.IsEnabledFunc = func() bool { return true }
|
||||
objectstore.IsEnabledFunc = func() bool { return true }
|
||||
defer func() {
|
||||
storage.IsEnabledFunc = func() bool { return false }
|
||||
objectstore.IsEnabledFunc = func() bool { return false }
|
||||
}()
|
||||
|
||||
t.Run("upload allowed image file successfully", func(t *testing.T) {
|
||||
@@ -326,7 +326,7 @@ func TestUploadFile(t *testing.T) {
|
||||
|
||||
t.Run("upload in local storage fallback mode", func(t *testing.T) {
|
||||
// Turn off S3
|
||||
storage.IsEnabledFunc = func() bool { return false }
|
||||
objectstore.IsEnabledFunc = func() bool { return false }
|
||||
|
||||
// Seed allowed extensions configuration to allow txt files
|
||||
var sc model.SystemConfig
|
||||
@@ -1090,7 +1090,7 @@ func configureLocalStorageRoot(t *testing.T, dbConn *gorm.DB, tempDir string) {
|
||||
if err := dbConn.Where("key = ?", model.ConfigKeyStorageConfig).First(&sc).Error; err != nil {
|
||||
t.Fatalf("failed to find storage config: %v", err)
|
||||
}
|
||||
var cfg storage.Config
|
||||
var cfg objectstore.Config
|
||||
if err := json.Unmarshal([]byte(sc.Value), &cfg); err != nil {
|
||||
t.Fatalf("failed to unmarshal storage config: %v", err)
|
||||
}
|
||||
@@ -1105,5 +1105,5 @@ func configureLocalStorageRoot(t *testing.T, dbConn *gorm.DB, tempDir string) {
|
||||
}
|
||||
_ = db.HSetJSON(context.Background(), repository.SystemConfigRedisHashKey, sc.Key, &sc)
|
||||
repository.ResetSystemConfigRAMCacheForTest()
|
||||
storage.ResetCache()
|
||||
objectstore.ResetCache()
|
||||
}
|
||||
|
||||
@@ -8,8 +8,8 @@ import (
|
||||
"time"
|
||||
|
||||
"github.com/Rain-kl/Wavelet/internal/apps/upload/shared"
|
||||
"github.com/Rain-kl/Wavelet/internal/common/response"
|
||||
"github.com/Rain-kl/Wavelet/internal/model"
|
||||
"github.com/Rain-kl/Wavelet/internal/shared/response"
|
||||
"github.com/gin-gonic/gin"
|
||||
)
|
||||
|
||||
|
||||
@@ -15,11 +15,11 @@ import (
|
||||
"github.com/Rain-kl/Wavelet/internal/apps/upload/shared"
|
||||
uploadstats "github.com/Rain-kl/Wavelet/internal/apps/upload/stats"
|
||||
uploadstorage "github.com/Rain-kl/Wavelet/internal/apps/upload/storage"
|
||||
"github.com/Rain-kl/Wavelet/internal/db"
|
||||
"github.com/Rain-kl/Wavelet/internal/db/idgen"
|
||||
"github.com/Rain-kl/Wavelet/internal/infra/objectstore"
|
||||
db "github.com/Rain-kl/Wavelet/internal/infra/persistence"
|
||||
"github.com/Rain-kl/Wavelet/internal/infra/persistence/idgen"
|
||||
"github.com/Rain-kl/Wavelet/internal/model"
|
||||
"github.com/Rain-kl/Wavelet/internal/repository"
|
||||
"github.com/Rain-kl/Wavelet/internal/storage"
|
||||
"github.com/Rain-kl/Wavelet/pkg/logger"
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
@@ -84,7 +84,7 @@ func storeObject(ctx context.Context, objectKey string, reader io.Reader, size i
|
||||
return "", ErrStorageReadOnly
|
||||
}
|
||||
|
||||
driver, backend, err := storage.Active(ctx)
|
||||
driver, backend, err := objectstore.Active(ctx)
|
||||
if err != nil {
|
||||
logger.ErrorF(ctx, "初始化活动存储失败: %v", err)
|
||||
return "", errors.New(shared.ErrSaveFileFailed)
|
||||
@@ -112,7 +112,7 @@ func persistUploadRecord(ctx context.Context, upload *model.Upload, objectKey st
|
||||
}
|
||||
|
||||
func cleanupUnpersistedObject(ctx context.Context, objectKey string) {
|
||||
_, backend, err := storage.Active(ctx)
|
||||
_, backend, err := objectstore.Active(ctx)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
|
||||
@@ -17,9 +17,9 @@ import (
|
||||
|
||||
uploadcache "github.com/Rain-kl/Wavelet/internal/apps/upload/cache"
|
||||
"github.com/Rain-kl/Wavelet/internal/apps/upload/shared"
|
||||
"github.com/Rain-kl/Wavelet/internal/db"
|
||||
"github.com/Rain-kl/Wavelet/internal/infra/objectstore"
|
||||
db "github.com/Rain-kl/Wavelet/internal/infra/persistence"
|
||||
"github.com/Rain-kl/Wavelet/internal/model"
|
||||
"github.com/Rain-kl/Wavelet/internal/storage"
|
||||
"github.com/Rain-kl/Wavelet/internal/testhelper"
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
@@ -270,7 +270,7 @@ func TestDedupRecordFailureDoesNotDeleteSharedObject(t *testing.T) {
|
||||
t.Fatalf("shared object delete count = %d, want 0", deleteCount)
|
||||
}
|
||||
|
||||
_, backend, err := storage.Active(ctx)
|
||||
_, backend, err := objectstore.Active(ctx)
|
||||
if err != nil {
|
||||
t.Fatalf("load active storage: %v", err)
|
||||
}
|
||||
@@ -519,7 +519,7 @@ func setupMockStorage(t *testing.T, putCount *int) (restore func(), disable func
|
||||
func setupMockStorageWithDeleteCount(t *testing.T, putCount, deleteCount *int) (restore func(), disable func()) {
|
||||
t.Helper()
|
||||
mockFiles := make(map[string][]byte)
|
||||
restore = storage.MockStorage(
|
||||
restore = objectstore.MockStorage(
|
||||
func(ctx context.Context, key string, body io.Reader, size int64, contentType string) error {
|
||||
data, err := io.ReadAll(body)
|
||||
if err != nil {
|
||||
@@ -531,12 +531,12 @@ func setupMockStorageWithDeleteCount(t *testing.T, putCount, deleteCount *int) (
|
||||
}
|
||||
return nil
|
||||
},
|
||||
func(ctx context.Context, key string) (*storage.Object, error) {
|
||||
func(ctx context.Context, key string) (*objectstore.Object, error) {
|
||||
data, ok := mockFiles[key]
|
||||
if !ok {
|
||||
return nil, os.ErrNotExist
|
||||
}
|
||||
return &storage.Object{
|
||||
return &objectstore.Object{
|
||||
Body: io.NopCloser(bytes.NewReader(data)),
|
||||
ContentLength: int64(len(data)),
|
||||
ContentType: "application/octet-stream",
|
||||
@@ -550,11 +550,11 @@ func setupMockStorageWithDeleteCount(t *testing.T, putCount, deleteCount *int) (
|
||||
return nil
|
||||
},
|
||||
)
|
||||
storage.IsEnabledFunc = func() bool { return true }
|
||||
storage.ResetCache()
|
||||
objectstore.IsEnabledFunc = func() bool { return true }
|
||||
objectstore.ResetCache()
|
||||
disable = func() {
|
||||
storage.IsEnabledFunc = func() bool { return false }
|
||||
storage.ResetCache()
|
||||
objectstore.IsEnabledFunc = func() bool { return false }
|
||||
objectstore.ResetCache()
|
||||
}
|
||||
return restore, disable
|
||||
}
|
||||
|
||||
@@ -10,7 +10,7 @@ import (
|
||||
"path/filepath"
|
||||
"strings"
|
||||
|
||||
"github.com/Rain-kl/Wavelet/internal/storage"
|
||||
"github.com/Rain-kl/Wavelet/internal/infra/objectstore"
|
||||
)
|
||||
|
||||
// LocalFileCandidateRequest describes filesystem locations that may host a legacy blob.
|
||||
@@ -103,7 +103,7 @@ func buildLocalFileCandidates(ctx context.Context, req LocalFileCandidateRequest
|
||||
}
|
||||
|
||||
func localStorageRoots(ctx context.Context) []string {
|
||||
cfg, err := storage.LoadConfig(ctx)
|
||||
cfg, err := objectstore.LoadConfig(ctx)
|
||||
if err != nil {
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -9,7 +9,7 @@ import (
|
||||
uploadcache "github.com/Rain-kl/Wavelet/internal/apps/upload/cache"
|
||||
"github.com/Rain-kl/Wavelet/internal/apps/upload/shared"
|
||||
uploadstats "github.com/Rain-kl/Wavelet/internal/apps/upload/stats"
|
||||
"github.com/Rain-kl/Wavelet/internal/db"
|
||||
db "github.com/Rain-kl/Wavelet/internal/infra/persistence"
|
||||
"github.com/Rain-kl/Wavelet/internal/model"
|
||||
"github.com/Rain-kl/Wavelet/internal/repository"
|
||||
"gorm.io/gorm"
|
||||
|
||||
@@ -7,7 +7,7 @@ import (
|
||||
"context"
|
||||
"time"
|
||||
|
||||
"github.com/Rain-kl/Wavelet/internal/db"
|
||||
db "github.com/Rain-kl/Wavelet/internal/infra/persistence"
|
||||
"github.com/Rain-kl/Wavelet/internal/model"
|
||||
"github.com/Rain-kl/Wavelet/pkg/logger"
|
||||
"gorm.io/gorm"
|
||||
|
||||
@@ -8,7 +8,7 @@ import (
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/Rain-kl/Wavelet/internal/db"
|
||||
db "github.com/Rain-kl/Wavelet/internal/infra/persistence"
|
||||
"github.com/Rain-kl/Wavelet/internal/model"
|
||||
"github.com/Rain-kl/Wavelet/internal/testhelper"
|
||||
"gorm.io/gorm"
|
||||
|
||||
@@ -10,14 +10,14 @@ import (
|
||||
"time"
|
||||
|
||||
"github.com/Rain-kl/Wavelet/internal/apps/upload/shared"
|
||||
"github.com/Rain-kl/Wavelet/internal/infra/objectstore"
|
||||
"github.com/Rain-kl/Wavelet/internal/model"
|
||||
"github.com/Rain-kl/Wavelet/internal/storage"
|
||||
)
|
||||
|
||||
// MigrationAccessState captures cached migration maintenance state.
|
||||
type MigrationAccessState struct {
|
||||
ReadOnly bool
|
||||
Target storage.Config
|
||||
Target objectstore.Config
|
||||
HasTarget bool
|
||||
TargetErr error
|
||||
LoadErr error
|
||||
|
||||
@@ -10,9 +10,9 @@ import (
|
||||
"fmt"
|
||||
"strings"
|
||||
|
||||
"github.com/Rain-kl/Wavelet/internal/db"
|
||||
"github.com/Rain-kl/Wavelet/internal/infra/objectstore"
|
||||
db "github.com/Rain-kl/Wavelet/internal/infra/persistence"
|
||||
"github.com/Rain-kl/Wavelet/internal/model"
|
||||
"github.com/Rain-kl/Wavelet/internal/storage"
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
|
||||
@@ -36,20 +36,20 @@ func LatestMigrationExecution(ctx context.Context) (*model.TaskExecution, bool,
|
||||
}
|
||||
|
||||
// ParseMigrationTargetConfig parses and validates a storage migration target payload.
|
||||
func ParseMigrationTargetConfig(ctx context.Context, payload []byte) (storage.Config, error) {
|
||||
func ParseMigrationTargetConfig(ctx context.Context, payload []byte) (objectstore.Config, error) {
|
||||
if strings.TrimSpace(string(payload)) == "" {
|
||||
return storage.Config{}, errors.New("storage migration target payload is required")
|
||||
return objectstore.Config{}, errors.New("storage migration target payload is required")
|
||||
}
|
||||
|
||||
var raw struct {
|
||||
Target json.RawMessage `json:"target"`
|
||||
}
|
||||
if err := json.Unmarshal(payload, &raw); err != nil {
|
||||
return storage.Config{}, fmt.Errorf("parse storage migration payload envelope: %w", err)
|
||||
return objectstore.Config{}, fmt.Errorf("parse storage migration payload envelope: %w", err)
|
||||
}
|
||||
|
||||
if len(raw.Target) == 0 {
|
||||
return storage.Config{}, errors.New("storage migration target payload is required")
|
||||
return objectstore.Config{}, errors.New("storage migration target payload is required")
|
||||
}
|
||||
|
||||
var targetBytes []byte
|
||||
@@ -60,34 +60,34 @@ func ParseMigrationTargetConfig(ctx context.Context, payload []byte) (storage.Co
|
||||
targetBytes = raw.Target
|
||||
}
|
||||
|
||||
var target storage.Config
|
||||
var target objectstore.Config
|
||||
if err := json.Unmarshal(targetBytes, &target); err != nil {
|
||||
return storage.Config{}, fmt.Errorf("parse target storage config: %w", err)
|
||||
return objectstore.Config{}, fmt.Errorf("parse target storage config: %w", err)
|
||||
}
|
||||
|
||||
current, err := storage.LoadConfig(ctx)
|
||||
current, err := objectstore.LoadConfig(ctx)
|
||||
if err != nil {
|
||||
return storage.Config{}, fmt.Errorf("load active storage config: %w", err)
|
||||
return objectstore.Config{}, fmt.Errorf("load active storage config: %w", err)
|
||||
}
|
||||
target = storage.MergeMaskedSecrets(target, current)
|
||||
if err := storage.ValidateConfig(target); err != nil {
|
||||
return storage.Config{}, fmt.Errorf("validate target storage config: %w", err)
|
||||
target = objectstore.MergeMaskedSecrets(target, current)
|
||||
if err := objectstore.ValidateConfig(target); err != nil {
|
||||
return objectstore.Config{}, fmt.Errorf("validate target storage config: %w", err)
|
||||
}
|
||||
return target, nil
|
||||
}
|
||||
|
||||
// NormalizeMigrationPayload validates and normalizes a storage migration payload.
|
||||
func NormalizeMigrationPayload(ctx context.Context, payload []byte) ([]byte, storage.Config, error) {
|
||||
func NormalizeMigrationPayload(ctx context.Context, payload []byte) ([]byte, objectstore.Config, error) {
|
||||
target, err := ParseMigrationTargetConfig(ctx, payload)
|
||||
if err != nil {
|
||||
return nil, storage.Config{}, err
|
||||
return nil, objectstore.Config{}, err
|
||||
}
|
||||
type storageMigrationPayload struct {
|
||||
Target storage.Config `json:"target"`
|
||||
Target objectstore.Config `json:"target"`
|
||||
}
|
||||
normalized, err := json.Marshal(storageMigrationPayload{Target: target})
|
||||
if err != nil {
|
||||
return nil, storage.Config{}, fmt.Errorf("marshal storage migration payload: %w", err)
|
||||
return nil, objectstore.Config{}, fmt.Errorf("marshal storage migration payload: %w", err)
|
||||
}
|
||||
return normalized, target, nil
|
||||
}
|
||||
|
||||
@@ -6,8 +6,8 @@ package storage
|
||||
import (
|
||||
"context"
|
||||
|
||||
"github.com/Rain-kl/Wavelet/internal/infra/objectstore"
|
||||
"github.com/Rain-kl/Wavelet/internal/model"
|
||||
"github.com/Rain-kl/Wavelet/internal/storage"
|
||||
"github.com/Rain-kl/Wavelet/pkg/logger"
|
||||
)
|
||||
|
||||
@@ -22,8 +22,8 @@ func ReadOnly(ctx context.Context) bool {
|
||||
}
|
||||
|
||||
// OpenStoredObject opens a stored upload object from the active storage backend.
|
||||
func OpenStoredObject(ctx context.Context, upload *model.Upload) (*storage.Object, error) {
|
||||
_, backend, err := storage.Active(ctx)
|
||||
func OpenStoredObject(ctx context.Context, upload *model.Upload) (*objectstore.Object, error) {
|
||||
_, backend, err := objectstore.Active(ctx)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
@@ -13,9 +13,9 @@ import (
|
||||
"github.com/Rain-kl/Wavelet/internal/apps/upload/ingest"
|
||||
"github.com/Rain-kl/Wavelet/internal/apps/upload/shared"
|
||||
uploadstorage "github.com/Rain-kl/Wavelet/internal/apps/upload/storage"
|
||||
"github.com/Rain-kl/Wavelet/internal/db"
|
||||
db "github.com/Rain-kl/Wavelet/internal/infra/persistence"
|
||||
"github.com/Rain-kl/Wavelet/internal/infra/task"
|
||||
"github.com/Rain-kl/Wavelet/internal/model"
|
||||
"github.com/Rain-kl/Wavelet/internal/task"
|
||||
"github.com/Rain-kl/Wavelet/pkg/logger"
|
||||
"gorm.io/gorm"
|
||||
"gorm.io/gorm/clause"
|
||||
|
||||
@@ -8,9 +8,9 @@ import (
|
||||
"fmt"
|
||||
|
||||
uploadstats "github.com/Rain-kl/Wavelet/internal/apps/upload/stats"
|
||||
"github.com/Rain-kl/Wavelet/internal/db"
|
||||
db "github.com/Rain-kl/Wavelet/internal/infra/persistence"
|
||||
"github.com/Rain-kl/Wavelet/internal/infra/task"
|
||||
"github.com/Rain-kl/Wavelet/internal/model"
|
||||
"github.com/Rain-kl/Wavelet/internal/task"
|
||||
)
|
||||
|
||||
const (
|
||||
|
||||
@@ -8,7 +8,7 @@ import (
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/Rain-kl/Wavelet/internal/db"
|
||||
db "github.com/Rain-kl/Wavelet/internal/infra/persistence"
|
||||
"github.com/Rain-kl/Wavelet/internal/model"
|
||||
"github.com/Rain-kl/Wavelet/internal/testhelper"
|
||||
)
|
||||
|
||||
@@ -17,10 +17,10 @@ import (
|
||||
|
||||
uploadstats "github.com/Rain-kl/Wavelet/internal/apps/upload/stats"
|
||||
uploadstorage "github.com/Rain-kl/Wavelet/internal/apps/upload/storage"
|
||||
"github.com/Rain-kl/Wavelet/internal/db"
|
||||
"github.com/Rain-kl/Wavelet/internal/infra/objectstore"
|
||||
db "github.com/Rain-kl/Wavelet/internal/infra/persistence"
|
||||
"github.com/Rain-kl/Wavelet/internal/infra/task"
|
||||
"github.com/Rain-kl/Wavelet/internal/model"
|
||||
"github.com/Rain-kl/Wavelet/internal/storage"
|
||||
"github.com/Rain-kl/Wavelet/internal/task"
|
||||
"golang.org/x/sync/errgroup"
|
||||
)
|
||||
|
||||
@@ -117,7 +117,7 @@ func (h *MigrationHandler) Execute(ctx context.Context, payload []byte) (*task.T
|
||||
}()
|
||||
}
|
||||
|
||||
active, err := storage.LoadConfig(ctx)
|
||||
active, err := objectstore.LoadConfig(ctx)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("load active storage config: %w", err)
|
||||
}
|
||||
@@ -126,7 +126,7 @@ func (h *MigrationHandler) Execute(ctx context.Context, payload []byte) (*task.T
|
||||
return nil, err
|
||||
}
|
||||
if target.Driver == active.Driver {
|
||||
if err := storage.SaveActiveConfig(ctx, target); err != nil {
|
||||
if err := objectstore.SaveActiveConfig(ctx, target); err != nil {
|
||||
return nil, fmt.Errorf("activate same-driver storage config: %w", err)
|
||||
}
|
||||
message := fmt.Sprintf("存储配置已更新,活动存储保持为 %s", target.Driver)
|
||||
@@ -139,7 +139,7 @@ func (h *MigrationHandler) Execute(ctx context.Context, payload []byte) (*task.T
|
||||
return nil, fmt.Errorf("count source objects: %w", err)
|
||||
}
|
||||
if total == 0 {
|
||||
if err := storage.SaveActiveConfig(ctx, target); err != nil {
|
||||
if err := objectstore.SaveActiveConfig(ctx, target); err != nil {
|
||||
return nil, fmt.Errorf("activate empty storage config: %w", err)
|
||||
}
|
||||
message := fmt.Sprintf("当前存储没有需要迁移的对象,活动存储已切换为 %s", target.Driver)
|
||||
@@ -147,11 +147,11 @@ func (h *MigrationHandler) Execute(ctx context.Context, payload []byte) (*task.T
|
||||
return &task.TaskResult{Message: message}, nil
|
||||
}
|
||||
|
||||
sourceBackend, err := storage.NewBackend(ctx, active, active.Driver)
|
||||
sourceBackend, err := objectstore.NewBackend(ctx, active, active.Driver)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("create source storage: %w", err)
|
||||
}
|
||||
targetBackend, err := storage.NewBackend(ctx, target, target.Driver)
|
||||
targetBackend, err := objectstore.NewBackend(ctx, target, target.Driver)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("create target storage: %w", err)
|
||||
}
|
||||
@@ -162,7 +162,7 @@ func (h *MigrationHandler) Execute(ctx context.Context, payload []byte) (*task.T
|
||||
return nil, err
|
||||
}
|
||||
|
||||
if err := storage.SaveActiveConfig(ctx, target); err != nil {
|
||||
if err := objectstore.SaveActiveConfig(ctx, target); err != nil {
|
||||
return nil, fmt.Errorf("activate target storage: %w", err)
|
||||
}
|
||||
message := fmt.Sprintf("存储迁移完成,共迁移 %d 个对象,活动存储已切换为 %s", migrated, target.Driver)
|
||||
@@ -196,8 +196,8 @@ type migrationObject struct {
|
||||
|
||||
func migrateObjects(
|
||||
ctx context.Context,
|
||||
sourceBackend storage.Backend,
|
||||
targetBackend storage.Backend,
|
||||
sourceBackend objectstore.Backend,
|
||||
targetBackend objectstore.Backend,
|
||||
total int64,
|
||||
) (int64, error) {
|
||||
const batchSize = 50
|
||||
@@ -258,8 +258,8 @@ func migrateObjects(
|
||||
|
||||
func migrateSingleObject(
|
||||
ctx context.Context,
|
||||
sourceBackend storage.Backend,
|
||||
targetBackend storage.Backend,
|
||||
sourceBackend objectstore.Backend,
|
||||
targetBackend objectstore.Backend,
|
||||
obj migrationObject,
|
||||
sha256HexLength int,
|
||||
) error {
|
||||
@@ -322,7 +322,7 @@ func migrateSingleObject(
|
||||
|
||||
func shouldSkipMigration(
|
||||
ctx context.Context,
|
||||
targetBackend storage.Backend,
|
||||
targetBackend objectstore.Backend,
|
||||
obj migrationObject,
|
||||
) bool {
|
||||
targetObj, err := targetBackend.Get(ctx, obj.FilePath)
|
||||
|
||||
@@ -16,9 +16,9 @@ import (
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/Rain-kl/Wavelet/internal/db"
|
||||
"github.com/Rain-kl/Wavelet/internal/infra/objectstore"
|
||||
db "github.com/Rain-kl/Wavelet/internal/infra/persistence"
|
||||
"github.com/Rain-kl/Wavelet/internal/model"
|
||||
"github.com/Rain-kl/Wavelet/internal/storage"
|
||||
"github.com/Rain-kl/Wavelet/internal/testhelper"
|
||||
"github.com/alicebob/miniredis/v2"
|
||||
"github.com/redis/go-redis/v9"
|
||||
@@ -39,21 +39,21 @@ func TestMigrationHandlerExecute(t *testing.T) {
|
||||
}
|
||||
|
||||
ctx := context.Background()
|
||||
active := storage.DefaultConfig()
|
||||
active := objectstore.DefaultConfig()
|
||||
active.Local.Root = sourceRoot
|
||||
if err := storage.SaveActiveConfig(ctx, active); err != nil {
|
||||
if err := objectstore.SaveActiveConfig(ctx, active); err != nil {
|
||||
t.Fatalf("SaveActiveConfig() returned error: %v", err)
|
||||
}
|
||||
target := storage.DefaultConfig()
|
||||
target.Driver = storage.DriverS3
|
||||
target.S3 = storage.ObjectConfig{
|
||||
target := objectstore.DefaultConfig()
|
||||
target.Driver = objectstore.DriverS3
|
||||
target.S3 = objectstore.ObjectConfig{
|
||||
Region: "us-east-1",
|
||||
Bucket: "target",
|
||||
AccessKeyID: "key",
|
||||
SecretAccessKey: "secret",
|
||||
}
|
||||
payload, err := json.Marshal(struct {
|
||||
Target storage.Config `json:"target"`
|
||||
Target objectstore.Config `json:"target"`
|
||||
}{Target: target})
|
||||
if err != nil {
|
||||
t.Fatalf("Marshal(storageMigrationPayload) returned error: %v", err)
|
||||
@@ -76,12 +76,12 @@ func TestMigrationHandlerExecute(t *testing.T) {
|
||||
}
|
||||
|
||||
var copied bytes.Buffer
|
||||
restore := storage.MockStorage(
|
||||
restore := objectstore.MockStorage(
|
||||
func(_ context.Context, _ string, body io.Reader, _ int64, _ string) error {
|
||||
_, err := io.Copy(&copied, body)
|
||||
return err
|
||||
},
|
||||
func(context.Context, string) (*storage.Object, error) {
|
||||
func(context.Context, string) (*objectstore.Object, error) {
|
||||
return nil, nil
|
||||
},
|
||||
func(context.Context, string) error {
|
||||
@@ -105,12 +105,12 @@ func TestMigrationHandlerExecute(t *testing.T) {
|
||||
if err := dbConn.First(&migrated, upload.ID).Error; err != nil {
|
||||
t.Fatalf("First(upload) returned error: %v", err)
|
||||
}
|
||||
current, err := storage.LoadConfig(ctx)
|
||||
current, err := objectstore.LoadConfig(ctx)
|
||||
if err != nil {
|
||||
t.Fatalf("LoadConfig() returned error: %v", err)
|
||||
}
|
||||
if current.Driver != storage.DriverS3 {
|
||||
t.Errorf("active driver = %q, want %q", current.Driver, storage.DriverS3)
|
||||
if current.Driver != objectstore.DriverS3 {
|
||||
t.Errorf("active driver = %q, want %q", current.Driver, objectstore.DriverS3)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -134,22 +134,22 @@ func TestMigrationHandlerExecuteWithHashValidation(t *testing.T) {
|
||||
correctHash := hex.EncodeToString(h.Sum(nil))
|
||||
|
||||
ctx := context.Background()
|
||||
active := storage.DefaultConfig()
|
||||
active := objectstore.DefaultConfig()
|
||||
active.Local.Root = sourceRoot
|
||||
if err := storage.SaveActiveConfig(ctx, active); err != nil {
|
||||
if err := objectstore.SaveActiveConfig(ctx, active); err != nil {
|
||||
t.Fatalf("SaveActiveConfig() returned error: %v", err)
|
||||
}
|
||||
|
||||
target := storage.DefaultConfig()
|
||||
target.Driver = storage.DriverS3
|
||||
target.S3 = storage.ObjectConfig{
|
||||
target := objectstore.DefaultConfig()
|
||||
target.Driver = objectstore.DriverS3
|
||||
target.S3 = objectstore.ObjectConfig{
|
||||
Region: "us-east-1",
|
||||
Bucket: "target",
|
||||
AccessKeyID: "key",
|
||||
SecretAccessKey: "secret",
|
||||
}
|
||||
payload, err := json.Marshal(struct {
|
||||
Target storage.Config `json:"target"`
|
||||
Target objectstore.Config `json:"target"`
|
||||
}{Target: target})
|
||||
if err != nil {
|
||||
t.Fatalf("Marshal(storageMigrationPayload) returned error: %v", err)
|
||||
@@ -173,14 +173,14 @@ func TestMigrationHandlerExecuteWithHashValidation(t *testing.T) {
|
||||
}
|
||||
|
||||
var copied bytes.Buffer
|
||||
restore := storage.MockStorage(
|
||||
restore := objectstore.MockStorage(
|
||||
func(_ context.Context, _ string, body io.Reader, _ int64, _ string) error {
|
||||
copied.Reset()
|
||||
_, err := io.Copy(&copied, body)
|
||||
return err
|
||||
},
|
||||
func(context.Context, string) (*storage.Object, error) {
|
||||
return &storage.Object{
|
||||
func(context.Context, string) (*objectstore.Object, error) {
|
||||
return &objectstore.Object{
|
||||
Body: io.NopCloser(bytes.NewBuffer(copied.Bytes())),
|
||||
ContentLength: int64(copied.Len()),
|
||||
ContentType: "text/plain",
|
||||
@@ -253,13 +253,13 @@ func TestMigrationHandlerExecuteWithRedisLock(t *testing.T) {
|
||||
t.Fatalf("Failed to set manual lock in Redis: %v", err)
|
||||
}
|
||||
|
||||
active := storage.DefaultConfig()
|
||||
if err := storage.SaveActiveConfig(ctx, active); err != nil {
|
||||
active := objectstore.DefaultConfig()
|
||||
if err := objectstore.SaveActiveConfig(ctx, active); err != nil {
|
||||
t.Fatalf("SaveActiveConfig() returned error: %v", err)
|
||||
}
|
||||
|
||||
payload, err := json.Marshal(struct {
|
||||
Target storage.Config `json:"target"`
|
||||
Target objectstore.Config `json:"target"`
|
||||
}{Target: active})
|
||||
if err != nil {
|
||||
t.Fatalf("Marshal payload failed: %v", err)
|
||||
|
||||
@@ -14,9 +14,9 @@ import (
|
||||
|
||||
"github.com/Rain-kl/Wavelet/internal/apps/upload/filesrv"
|
||||
"github.com/Rain-kl/Wavelet/internal/apps/upload/shared"
|
||||
"github.com/Rain-kl/Wavelet/internal/db"
|
||||
db "github.com/Rain-kl/Wavelet/internal/infra/persistence"
|
||||
"github.com/Rain-kl/Wavelet/internal/infra/task"
|
||||
"github.com/Rain-kl/Wavelet/internal/model"
|
||||
"github.com/Rain-kl/Wavelet/internal/task"
|
||||
)
|
||||
|
||||
const (
|
||||
|
||||
@@ -20,11 +20,11 @@ import (
|
||||
"github.com/Rain-kl/Wavelet/internal/apps/upload/filesrv"
|
||||
"github.com/Rain-kl/Wavelet/internal/apps/upload/shared"
|
||||
uploadstats "github.com/Rain-kl/Wavelet/internal/apps/upload/stats"
|
||||
"github.com/Rain-kl/Wavelet/internal/db"
|
||||
"github.com/Rain-kl/Wavelet/internal/diskcache"
|
||||
"github.com/Rain-kl/Wavelet/internal/infra/diskcache"
|
||||
"github.com/Rain-kl/Wavelet/internal/infra/objectstore"
|
||||
db "github.com/Rain-kl/Wavelet/internal/infra/persistence"
|
||||
"github.com/Rain-kl/Wavelet/internal/infra/task"
|
||||
"github.com/Rain-kl/Wavelet/internal/model"
|
||||
"github.com/Rain-kl/Wavelet/internal/storage"
|
||||
"github.com/Rain-kl/Wavelet/internal/task"
|
||||
"github.com/Rain-kl/Wavelet/internal/testhelper"
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
@@ -36,20 +36,20 @@ func TestSystemCleanupHandler_Execute(t *testing.T) {
|
||||
|
||||
deleteCount := 0
|
||||
// Mock S3 存储并记录 Delete,cleanup 不应物理删除共享对象。
|
||||
storageMock := storage.MockStorage(
|
||||
storageMock := objectstore.MockStorage(
|
||||
func(ctx context.Context, key string, body io.Reader, size int64, contentType string) error {
|
||||
return nil
|
||||
},
|
||||
func(ctx context.Context, key string) (*storage.Object, error) { return nil, nil },
|
||||
func(ctx context.Context, key string) (*objectstore.Object, error) { return nil, nil },
|
||||
func(ctx context.Context, key string) error {
|
||||
deleteCount++
|
||||
return nil
|
||||
},
|
||||
)
|
||||
defer storageMock()
|
||||
storage.IsEnabledFunc = func() bool { return true }
|
||||
defer func() { storage.IsEnabledFunc = func() bool { return false } }()
|
||||
storage.ResetCache()
|
||||
objectstore.IsEnabledFunc = func() bool { return true }
|
||||
defer func() { objectstore.IsEnabledFunc = func() bool { return false } }()
|
||||
objectstore.ResetCache()
|
||||
|
||||
ctx := context.Background()
|
||||
err := db.DB(ctx).AutoMigrate(&model.PushHistory{})
|
||||
@@ -195,11 +195,11 @@ func TestSystemCleanupHandler_ExecuteNoFiles(t *testing.T) {
|
||||
defer cleanup()
|
||||
|
||||
// Mock S3 存储
|
||||
storageMock := storage.MockStorage(
|
||||
storageMock := objectstore.MockStorage(
|
||||
func(ctx context.Context, key string, body io.Reader, size int64, contentType string) error {
|
||||
return nil
|
||||
},
|
||||
func(ctx context.Context, key string) (*storage.Object, error) { return nil, nil },
|
||||
func(ctx context.Context, key string) (*objectstore.Object, error) { return nil, nil },
|
||||
func(ctx context.Context, key string) error { return nil },
|
||||
)
|
||||
defer storageMock()
|
||||
@@ -293,9 +293,9 @@ func TestWarmImageCacheHandlerExecute(t *testing.T) {
|
||||
|
||||
testDir := t.TempDir()
|
||||
ctx := context.Background()
|
||||
active := storage.DefaultConfig()
|
||||
active := objectstore.DefaultConfig()
|
||||
active.Local.Root = testDir
|
||||
if err := storage.SaveActiveConfig(ctx, active); err != nil {
|
||||
if err := objectstore.SaveActiveConfig(ctx, active); err != nil {
|
||||
t.Fatalf("SaveActiveConfig() returned error: %v", err)
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user