支持迁移对象存储到 OSS
This commit is contained in:
@@ -24,6 +24,7 @@ type Config struct {
|
||||
BootstrapAdminPassword string
|
||||
BootstrapAdminNickname string
|
||||
Storage StorageConfig
|
||||
StorageMirror StorageConfig
|
||||
SMS SMSConfig
|
||||
Realname RealnameConfig
|
||||
Log LogConfig
|
||||
@@ -35,6 +36,8 @@ type StorageConfig struct {
|
||||
Bucket string
|
||||
AccessKeyID string
|
||||
SecretAccessKey string
|
||||
Region string
|
||||
BucketLookup string
|
||||
}
|
||||
|
||||
type SMSConfig struct {
|
||||
@@ -91,6 +94,17 @@ func Load() Config {
|
||||
Bucket: getEnv("STORAGE_BUCKET", "hfb-sys"),
|
||||
AccessKeyID: getEnv("STORAGE_ACCESS_KEY_ID", "minioadmin"),
|
||||
SecretAccessKey: getEnv("STORAGE_SECRET_ACCESS_KEY", "minioadmin"),
|
||||
Region: getEnv("STORAGE_REGION", ""),
|
||||
BucketLookup: getEnv("STORAGE_BUCKET_LOOKUP", "auto"),
|
||||
},
|
||||
// StorageMirror 仅在对象存储迁移期间使用。镜像端配置完整时,新上传文件会同步写入两个存储端。
|
||||
StorageMirror: StorageConfig{
|
||||
Endpoint: getEnv("STORAGE_MIRROR_ENDPOINT", ""),
|
||||
Bucket: getEnv("STORAGE_MIRROR_BUCKET", ""),
|
||||
AccessKeyID: getEnv("STORAGE_MIRROR_ACCESS_KEY_ID", ""),
|
||||
SecretAccessKey: getEnv("STORAGE_MIRROR_SECRET_ACCESS_KEY", ""),
|
||||
Region: getEnv("STORAGE_MIRROR_REGION", ""),
|
||||
BucketLookup: getEnv("STORAGE_MIRROR_BUCKET_LOOKUP", "auto"),
|
||||
},
|
||||
SMS: SMSConfig{
|
||||
Provider: getEnv("SMS_PROVIDER", "mock"),
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
package file
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"crypto/rand"
|
||||
"encoding/hex"
|
||||
@@ -24,6 +25,8 @@ type Storage struct {
|
||||
bucket string
|
||||
bucketMu sync.Mutex
|
||||
bucketReady bool
|
||||
mirror *Storage
|
||||
mirrorErr error
|
||||
}
|
||||
|
||||
type Object struct {
|
||||
@@ -33,14 +36,7 @@ type Object struct {
|
||||
}
|
||||
|
||||
func NewStorage(cfg config.StorageConfig) (*Storage, error) {
|
||||
endpoint, secure, err := normalizeEndpoint(cfg.Endpoint)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
client, err := minio.New(endpoint, &minio.Options{
|
||||
Creds: credentials.NewStaticV4(cfg.AccessKeyID, cfg.SecretAccessKey, ""),
|
||||
Secure: secure,
|
||||
})
|
||||
client, err := NewStorageClient(cfg)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
@@ -51,6 +47,43 @@ func NewStorage(cfg config.StorageConfig) (*Storage, error) {
|
||||
return storage, nil
|
||||
}
|
||||
|
||||
// NewStorageClient 根据项目存储配置创建 S3 兼容客户端,供迁移工具复用。
|
||||
func NewStorageClient(cfg config.StorageConfig) (*minio.Client, error) {
|
||||
endpoint, secure, err := normalizeEndpoint(cfg.Endpoint)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
lookup, err := parseBucketLookup(cfg.BucketLookup)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return minio.New(endpoint, &minio.Options{
|
||||
Creds: credentials.NewStaticV4(cfg.AccessKeyID, cfg.SecretAccessKey, ""),
|
||||
Secure: secure,
|
||||
Region: cfg.Region,
|
||||
BucketLookup: lookup,
|
||||
})
|
||||
}
|
||||
|
||||
func parseBucketLookup(raw string) (minio.BucketLookupType, error) {
|
||||
switch strings.ToLower(strings.TrimSpace(raw)) {
|
||||
case "", "auto":
|
||||
return minio.BucketLookupAuto, nil
|
||||
case "dns":
|
||||
return minio.BucketLookupDNS, nil
|
||||
case "path":
|
||||
return minio.BucketLookupPath, nil
|
||||
default:
|
||||
return minio.BucketLookupAuto, fmt.Errorf("unsupported storage bucket lookup: %s", raw)
|
||||
}
|
||||
}
|
||||
|
||||
// SetMirror 启用迁移期间的同步双写。镜像端不可用时上传会失败,避免出现未同步的业务文件。
|
||||
func (s *Storage) SetMirror(mirror *Storage, err error) {
|
||||
s.mirror = mirror
|
||||
s.mirrorErr = err
|
||||
}
|
||||
|
||||
func (s *Storage) Put(ctx context.Context, scene string, header *multipart.FileHeader, reader io.Reader, contentType string) (string, error) {
|
||||
key, err := newObjectKey(scene, header.Filename)
|
||||
if err != nil {
|
||||
@@ -65,6 +98,28 @@ func (s *Storage) Put(ctx context.Context, scene string, header *multipart.FileH
|
||||
}
|
||||
|
||||
func (s *Storage) PutObject(ctx context.Context, key string, reader io.Reader, size int64, contentType string, metadata map[string]string) error {
|
||||
if s.mirrorErr != nil {
|
||||
return fmt.Errorf("storage mirror unavailable: %w", s.mirrorErr)
|
||||
}
|
||||
if s.mirror == nil {
|
||||
return s.putObject(ctx, key, reader, size, contentType, metadata)
|
||||
}
|
||||
|
||||
// 当前文件上传大小被限制在 10MB;镜像写入前缓存在内存中,确保两个端点得到完全相同的字节流。
|
||||
data, err := io.ReadAll(reader)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if err := s.putObject(ctx, key, bytes.NewReader(data), int64(len(data)), contentType, metadata); err != nil {
|
||||
return err
|
||||
}
|
||||
if err := s.mirror.putObject(ctx, key, bytes.NewReader(data), int64(len(data)), contentType, metadata); err != nil {
|
||||
return fmt.Errorf("storage mirror upload failed: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *Storage) putObject(ctx context.Context, key string, reader io.Reader, size int64, contentType string, metadata map[string]string) error {
|
||||
if err := s.ensureBucketReady(ctx); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
@@ -3,6 +3,8 @@ package file
|
||||
import (
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"github.com/minio/minio-go/v7"
|
||||
)
|
||||
|
||||
func TestSanitizeObjectMetadataEscapesNonASCII(t *testing.T) {
|
||||
@@ -15,6 +17,29 @@ func TestSanitizeObjectMetadataEscapesNonASCII(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestParseBucketLookup(t *testing.T) {
|
||||
tests := []struct {
|
||||
name string
|
||||
raw string
|
||||
want minio.BucketLookupType
|
||||
}{
|
||||
{name: "auto", raw: "auto", want: minio.BucketLookupAuto},
|
||||
{name: "dns", raw: "dns", want: minio.BucketLookupDNS},
|
||||
{name: "path", raw: "path", want: minio.BucketLookupPath},
|
||||
}
|
||||
for _, tt := range tests {
|
||||
t.Run(tt.name, func(t *testing.T) {
|
||||
got, err := parseBucketLookup(tt.raw)
|
||||
if err != nil || got != tt.want {
|
||||
t.Fatalf("parseBucketLookup(%q) = %v, %v", tt.raw, got, err)
|
||||
}
|
||||
})
|
||||
}
|
||||
if _, err := parseBucketLookup("invalid"); err == nil {
|
||||
t.Fatal("invalid bucket lookup should fail")
|
||||
}
|
||||
}
|
||||
|
||||
func TestAnnouncementSceneUsesPublicFileURL(t *testing.T) {
|
||||
key := "announcement/2026/06/15/example.webp"
|
||||
got := fileURLForScene("announcement", key)
|
||||
|
||||
@@ -333,6 +333,26 @@ func New(cfg config.Config, deps Dependencies, logger *zap.Logger) *gin.Engine {
|
||||
if err != nil {
|
||||
logger.Warn("文件存储桶尚未就绪,接口请求时将重试", zap.Error(err))
|
||||
}
|
||||
if cfg.StorageMirror.Endpoint != "" {
|
||||
mirrorStorage, mirrorErr := filemodule.NewStorage(cfg.StorageMirror)
|
||||
if fileStorage != nil {
|
||||
// 客户端已创建时,即使启动检查失败也由上传请求重新探测 Bucket。
|
||||
if mirrorStorage != nil {
|
||||
fileStorage.SetMirror(mirrorStorage, nil)
|
||||
} else {
|
||||
fileStorage.SetMirror(nil, mirrorErr)
|
||||
}
|
||||
}
|
||||
if mirrorErr != nil {
|
||||
if mirrorStorage != nil {
|
||||
logger.Warn("文件存储镜像启动检查失败,上传时将重试", zap.Error(mirrorErr))
|
||||
} else {
|
||||
logger.Error("文件存储镜像客户端创建失败,已阻止新文件上传", zap.Error(mirrorErr))
|
||||
}
|
||||
} else {
|
||||
logger.Info("文件存储镜像已启用")
|
||||
}
|
||||
}
|
||||
}
|
||||
fileService := filemodule.NewService(fileStorage)
|
||||
fileHandler := filemodule.NewHandler(fileService, fileStorage)
|
||||
|
||||
Reference in New Issue
Block a user