229 lines
7.3 KiB
Go
229 lines
7.3 KiB
Go
package main
|
||
|
||
import (
|
||
"context"
|
||
"flag"
|
||
"fmt"
|
||
"log"
|
||
"os"
|
||
"strings"
|
||
|
||
"hfb_sys/backend/internal/config"
|
||
filemodule "hfb_sys/backend/internal/modules/file"
|
||
|
||
"github.com/minio/minio-go/v7"
|
||
)
|
||
|
||
type migrationConfig struct {
|
||
source config.StorageConfig
|
||
target config.StorageConfig
|
||
}
|
||
|
||
type summary struct {
|
||
objectsCopied int64
|
||
bytesCopied int64
|
||
objectsSkipped int64
|
||
objectsFailed int64
|
||
}
|
||
|
||
func main() {
|
||
verifyOnly := flag.Bool("verify", false, "仅校验目标对象,不上传文件")
|
||
dryRun := flag.Bool("dry-run", false, "只输出待迁移对象,不上传文件")
|
||
prefix := flag.String("prefix", "", "仅迁移指定前缀")
|
||
flag.Parse()
|
||
if *verifyOnly && *dryRun {
|
||
log.Fatal("--verify 与 --dry-run 不能同时使用")
|
||
}
|
||
|
||
cfg, err := loadMigrationConfig()
|
||
if err != nil {
|
||
log.Fatal(err)
|
||
}
|
||
if sameStorage(cfg.source, cfg.target) {
|
||
log.Fatal("源存储和目标存储相同,已拒绝执行")
|
||
}
|
||
|
||
ctx := context.Background()
|
||
source, err := filemodule.NewStorageClient(cfg.source)
|
||
if err != nil {
|
||
log.Fatalf("创建 MinIO 源客户端失败: %v", err)
|
||
}
|
||
target, err := filemodule.NewStorageClient(cfg.target)
|
||
if err != nil {
|
||
log.Fatalf("创建 OSS 目标客户端失败: %v", err)
|
||
}
|
||
if err := requireBucket(ctx, source, cfg.source.Bucket, "源"); err != nil {
|
||
log.Fatal(err)
|
||
}
|
||
if err := requireBucket(ctx, target, cfg.target.Bucket, "目标"); err != nil {
|
||
log.Fatal(err)
|
||
}
|
||
|
||
mode := "同步"
|
||
if *verifyOnly {
|
||
mode = "校验"
|
||
} else if *dryRun {
|
||
mode = "预演"
|
||
}
|
||
log.Printf("开始%s:%s/%s -> %s/%s,前缀=%q", mode, cfg.source.Endpoint, cfg.source.Bucket, cfg.target.Endpoint, cfg.target.Bucket, *prefix)
|
||
|
||
result, err := migrate(ctx, source, cfg.source.Bucket, target, cfg.target.Bucket, *prefix, *verifyOnly, *dryRun)
|
||
log.Printf("结束:已复制 %d 个对象(%d 字节),已跳过 %d 个对象,失败 %d 个对象", result.objectsCopied, result.bytesCopied, result.objectsSkipped, result.objectsFailed)
|
||
if err != nil {
|
||
log.Fatal(err)
|
||
}
|
||
}
|
||
|
||
func loadMigrationConfig() (migrationConfig, error) {
|
||
source := config.StorageConfig{
|
||
Endpoint: os.Getenv("STORAGE_ENDPOINT"),
|
||
Bucket: os.Getenv("STORAGE_BUCKET"),
|
||
AccessKeyID: os.Getenv("STORAGE_ACCESS_KEY_ID"),
|
||
SecretAccessKey: os.Getenv("STORAGE_SECRET_ACCESS_KEY"),
|
||
Region: os.Getenv("STORAGE_REGION"),
|
||
BucketLookup: valueOrDefault("STORAGE_BUCKET_LOOKUP", "auto"),
|
||
}
|
||
target := config.StorageConfig{
|
||
Endpoint: os.Getenv("STORAGE_MIRROR_ENDPOINT"),
|
||
Bucket: os.Getenv("STORAGE_MIRROR_BUCKET"),
|
||
AccessKeyID: os.Getenv("STORAGE_MIRROR_ACCESS_KEY_ID"),
|
||
SecretAccessKey: os.Getenv("STORAGE_MIRROR_SECRET_ACCESS_KEY"),
|
||
Region: os.Getenv("STORAGE_MIRROR_REGION"),
|
||
BucketLookup: valueOrDefault("STORAGE_MIRROR_BUCKET_LOOKUP", "auto"),
|
||
}
|
||
if err := validateStorageConfig("源 MinIO", source); err != nil {
|
||
return migrationConfig{}, err
|
||
}
|
||
if err := validateStorageConfig("目标 OSS", target); err != nil {
|
||
return migrationConfig{}, err
|
||
}
|
||
return migrationConfig{source: source, target: target}, nil
|
||
}
|
||
|
||
func valueOrDefault(key, fallback string) string {
|
||
if value := strings.TrimSpace(os.Getenv(key)); value != "" {
|
||
return value
|
||
}
|
||
return fallback
|
||
}
|
||
|
||
func validateStorageConfig(name string, cfg config.StorageConfig) error {
|
||
if strings.TrimSpace(cfg.Endpoint) == "" || strings.TrimSpace(cfg.Bucket) == "" || strings.TrimSpace(cfg.AccessKeyID) == "" || strings.TrimSpace(cfg.SecretAccessKey) == "" {
|
||
return fmt.Errorf("%s 配置不完整,请检查 endpoint、bucket 和访问密钥", name)
|
||
}
|
||
return nil
|
||
}
|
||
|
||
func sameStorage(source, target config.StorageConfig) bool {
|
||
return strings.TrimRight(strings.ToLower(source.Endpoint), "/") == strings.TrimRight(strings.ToLower(target.Endpoint), "/") && source.Bucket == target.Bucket
|
||
}
|
||
|
||
func requireBucket(ctx context.Context, client *minio.Client, bucket, label string) error {
|
||
exists, err := client.BucketExists(ctx, bucket)
|
||
if err != nil {
|
||
return fmt.Errorf("检查%s Bucket 失败: %w", label, err)
|
||
}
|
||
if !exists {
|
||
return fmt.Errorf("%s Bucket 不存在: %s", label, bucket)
|
||
}
|
||
return nil
|
||
}
|
||
|
||
func migrate(ctx context.Context, source *minio.Client, sourceBucket string, target *minio.Client, targetBucket, prefix string, verifyOnly, dryRun bool) (summary, error) {
|
||
var result summary
|
||
for item := range source.ListObjects(ctx, sourceBucket, minio.ListObjectsOptions{Prefix: prefix, Recursive: true}) {
|
||
if item.Err != nil {
|
||
return result, fmt.Errorf("列出源对象失败: %w", item.Err)
|
||
}
|
||
sourceInfo, err := source.StatObject(ctx, sourceBucket, item.Key, minio.StatObjectOptions{})
|
||
if err != nil {
|
||
result.objectsFailed++
|
||
log.Printf("跳过对象 %s:读取源对象元数据失败: %v", item.Key, err)
|
||
continue
|
||
}
|
||
matches, err := targetMatches(ctx, target, targetBucket, item.Key, sourceInfo)
|
||
if err != nil {
|
||
result.objectsFailed++
|
||
log.Printf("跳过对象 %s:检查目标对象失败: %v", item.Key, err)
|
||
continue
|
||
}
|
||
if matches {
|
||
result.objectsSkipped++
|
||
continue
|
||
}
|
||
if verifyOnly {
|
||
result.objectsFailed++
|
||
continue
|
||
}
|
||
if dryRun {
|
||
log.Printf("待复制:%s(%d 字节)", item.Key, sourceInfo.Size)
|
||
continue
|
||
}
|
||
if err := copyObject(ctx, source, sourceBucket, target, targetBucket, item.Key, sourceInfo); err != nil {
|
||
result.objectsFailed++
|
||
log.Printf("复制对象 %s 失败: %v", item.Key, err)
|
||
continue
|
||
}
|
||
result.objectsCopied++
|
||
result.bytesCopied += sourceInfo.Size
|
||
}
|
||
if result.objectsFailed > 0 {
|
||
if verifyOnly {
|
||
return result, fmt.Errorf("校验失败:%d 个对象缺失、内容不匹配或无法检查", result.objectsFailed)
|
||
}
|
||
return result, fmt.Errorf("同步未完成:%d 个对象处理失败", result.objectsFailed)
|
||
}
|
||
return result, nil
|
||
}
|
||
|
||
func targetMatches(ctx context.Context, target *minio.Client, bucket, key string, source minio.ObjectInfo) (bool, error) {
|
||
targetInfo, err := target.StatObject(ctx, bucket, key, minio.StatObjectOptions{})
|
||
if err != nil {
|
||
if isNotFound(err) {
|
||
return false, nil
|
||
}
|
||
return false, fmt.Errorf("读取目标对象 %s 元数据失败: %w", key, err)
|
||
}
|
||
if targetInfo.Size != source.Size {
|
||
return false, nil
|
||
}
|
||
return sameETag(source.ETag, targetInfo.ETag), nil
|
||
}
|
||
|
||
func copyObject(ctx context.Context, source *minio.Client, sourceBucket string, target *minio.Client, targetBucket, key string, info minio.ObjectInfo) error {
|
||
object, err := source.GetObject(ctx, sourceBucket, key, minio.GetObjectOptions{})
|
||
if err != nil {
|
||
return fmt.Errorf("读取源对象 %s 失败: %w", key, err)
|
||
}
|
||
defer func() { _ = object.Close() }()
|
||
_, err = target.PutObject(ctx, targetBucket, key, object, info.Size, minio.PutObjectOptions{
|
||
ContentType: info.ContentType,
|
||
UserMetadata: info.UserMetadata,
|
||
})
|
||
if err != nil {
|
||
return fmt.Errorf("复制对象 %s 到目标 OSS 失败: %w", key, err)
|
||
}
|
||
verified, err := targetMatches(ctx, target, targetBucket, key, info)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
if !verified {
|
||
return fmt.Errorf("复制对象 %s 后校验失败", key)
|
||
}
|
||
return nil
|
||
}
|
||
|
||
func sameETag(left, right string) bool {
|
||
return strings.Trim(strings.ToLower(left), "\"") == strings.Trim(strings.ToLower(right), "\"")
|
||
}
|
||
|
||
func isNotFound(err error) bool {
|
||
response := minio.ToErrorResponse(err)
|
||
switch response.Code {
|
||
case "NoSuchKey", "NoSuchObject", "NoSuchBucket", "NotFound":
|
||
return true
|
||
default:
|
||
return false
|
||
}
|
||
}
|