Files
hfb_sys/backend/internal/modules/file/storage.go
T

254 lines
6.6 KiB
Go

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
}