package file import ( "bytes" "context" "crypto/rand" "encoding/hex" "fmt" "io" "mime/multipart" "net/url" "path" "strings" "sync" "time" "hfb_sys/backend/internal/config" "github.com/minio/minio-go/v7" "github.com/minio/minio-go/v7/pkg/credentials" ) type Storage struct { client *minio.Client bucket string bucketMu sync.Mutex bucketReady bool mirror *Storage mirrorErr error } type Object struct { Reader io.ReadCloser ContentType string Size int64 } func NewStorage(cfg config.StorageConfig) (*Storage, error) { client, err := NewStorageClient(cfg) if err != nil { return nil, err } storage := &Storage{client: client, bucket: cfg.Bucket} if err := storage.ensureBucket(context.Background()); err != nil { return storage, err } 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 { return "", err } if err := s.PutObject(ctx, key, reader, header.Size, contentType, map[string]string{ "original-filename": header.Filename, }); err != nil { return "", err } return key, nil } 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 } _, err := s.client.PutObject(ctx, s.bucket, key, reader, size, minio.PutObjectOptions{ ContentType: contentType, UserMetadata: sanitizeObjectMetadata(metadata), }) if isNoSuchBucketError(err) && rewindReader(reader) == nil { s.markBucketNotReady() if readyErr := s.ensureBucketReady(ctx); readyErr != nil { return readyErr } _, err = s.client.PutObject(ctx, s.bucket, key, reader, size, minio.PutObjectOptions{ ContentType: contentType, UserMetadata: sanitizeObjectMetadata(metadata), }) } return err } func (s *Storage) Get(ctx context.Context, key string) (*Object, error) { object, err := s.client.GetObject(ctx, s.bucket, key, minio.GetObjectOptions{}) if err != nil { return nil, err } info, err := object.Stat() if err != nil { _ = object.Close() return nil, err } return &Object{ Reader: object, ContentType: info.ContentType, Size: info.Size, }, nil } func (s *Storage) ensureBucket(ctx context.Context) error { exists, err := s.client.BucketExists(ctx, s.bucket) if err != nil { errResp := minio.ToErrorResponse(err) if errResp.Code != "NoSuchBucket" && !strings.Contains(strings.ToLower(err.Error()), "bucket does not exist") { return err } } if exists { s.bucketReady = true return nil } if err := s.client.MakeBucket(ctx, s.bucket, minio.MakeBucketOptions{}); err != nil { return err } s.bucketReady = true return nil } func (s *Storage) ensureBucketReady(ctx context.Context) error { s.bucketMu.Lock() defer s.bucketMu.Unlock() if s.bucketReady { return nil } return s.ensureBucket(ctx) } func (s *Storage) markBucketNotReady() { s.bucketMu.Lock() defer s.bucketMu.Unlock() s.bucketReady = false } func isNoSuchBucketError(err error) bool { if err == nil { return false } errResp := minio.ToErrorResponse(err) if errResp.Code == "NoSuchBucket" { return true } return strings.Contains(strings.ToLower(err.Error()), "bucket does not exist") } func rewindReader(reader io.Reader) error { seeker, ok := reader.(io.Seeker) if !ok { return fmt.Errorf("reader cannot rewind") } _, err := seeker.Seek(0, io.SeekStart) return err } func sanitizeObjectMetadata(metadata map[string]string) map[string]string { if len(metadata) == 0 { return nil } sanitized := make(map[string]string, len(metadata)) for key, value := range metadata { key = strings.TrimSpace(key) if key == "" { continue } // MinIO/S3 用户元数据会作为 HTTP header 发送,值统一转义为 ASCII,避免中文文件名导致上传失败。 sanitized[key] = url.QueryEscape(value) } return sanitized } func normalizeEndpoint(raw string) (string, bool, error) { parsed, err := url.Parse(raw) if err != nil { return "", false, err } if parsed.Scheme == "" { return raw, false, nil } return parsed.Host, parsed.Scheme == "https", nil } func newObjectKey(scene string, filename string) (string, error) { if scene == "" { scene = "misc" } scene = strings.ToLower(scene) now := time.Now() token := make([]byte, 12) if _, err := rand.Read(token); err != nil { return "", err } ext := strings.ToLower(path.Ext(filename)) return fmt.Sprintf("%s/%04d/%02d/%02d/%s%s", scene, now.Year(), now.Month(), now.Day(), hex.EncodeToString(token), ext), nil }