Files

229 lines
7.3 KiB
Go
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
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
}
}