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 } }