199 lines
4.8 KiB
Go
199 lines
4.8 KiB
Go
package file
|
|
|
|
import (
|
|
"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
|
|
}
|
|
|
|
type Object struct {
|
|
Reader io.ReadCloser
|
|
ContentType string
|
|
Size int64
|
|
}
|
|
|
|
func NewStorage(cfg config.StorageConfig) (*Storage, error) {
|
|
endpoint, secure, err := normalizeEndpoint(cfg.Endpoint)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
client, err := minio.New(endpoint, &minio.Options{
|
|
Creds: credentials.NewStaticV4(cfg.AccessKeyID, cfg.SecretAccessKey, ""),
|
|
Secure: secure,
|
|
})
|
|
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
|
|
}
|
|
|
|
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 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
|
|
}
|