up
This commit is contained in:
+15
-7
@@ -9,13 +9,14 @@ import (
|
||||
|
||||
type OssApi struct {
|
||||
gorm.Model
|
||||
Name string `gorm:"size:100;not null" json:"name"` // Alias o nombre interno
|
||||
Endpoint string `gorm:"not null" json:"endpoint"` // Endpoint del access point
|
||||
AccessKeyID string `gorm:"not null" json:"access_key_id"` // Access Key ID
|
||||
AccessKeySecret string `gorm:"not null" json:"access_key_secret"` // Access Key Secret
|
||||
BucketName string `gorm:"not null" json:"bucket_name"` // Nombre del bucket o access point
|
||||
Region string `gorm:"size:50" json:"region"` // Región, opcional
|
||||
IsActive bool `gorm:"default:true" json:"is_active"` // Activar o desactivar config
|
||||
Name string `gorm:"size:100;not null" json:"name"` // Alias o nombre interno
|
||||
Provider string `gorm:"size:20;default:alibaba" json:"provider"` // alibaba | s3
|
||||
Endpoint string `gorm:"not null" json:"endpoint"` // Endpoint del access point
|
||||
AccessKeyID string `gorm:"not null" json:"access_key_id"` // Access Key ID
|
||||
AccessKeySecret string `gorm:"not null" json:"access_key_secret"` // Access Key Secret
|
||||
BucketName string `gorm:"not null" json:"bucket_name"` // Nombre del bucket o access point
|
||||
Region string `gorm:"size:50" json:"region"` // Región, opcional
|
||||
IsActive bool `gorm:"default:true" json:"is_active"` // Activar o desactivar config
|
||||
Notes string `gorm:"type:text" json:"notes"`
|
||||
}
|
||||
|
||||
@@ -84,6 +85,13 @@ func GetOssApiByID(id uint) (*OssApi, error) {
|
||||
return &item, nil
|
||||
}
|
||||
|
||||
// GetActiveOssApis obtiene todas las configuraciones activas
|
||||
func GetActiveOssApis() ([]OssApi, error) {
|
||||
var items []OssApi
|
||||
err := app.Http.Database.DB.Where("is_active = ?", true).Order("id DESC").Find(&items).Error
|
||||
return items, err
|
||||
}
|
||||
|
||||
// GetLastActiveOssApi obtiene el último registro activo
|
||||
func GetLastActiveOssApi() (*OssApi, error) {
|
||||
var ossConfig OssApi
|
||||
|
||||
+245
-84
@@ -1,12 +1,15 @@
|
||||
package services
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"io"
|
||||
"log"
|
||||
"time"
|
||||
|
||||
"github.com/aliyun/aliyun-oss-go-sdk/oss"
|
||||
"github.com/minio/minio-go/v7"
|
||||
"github.com/minio/minio-go/v7/pkg/credentials"
|
||||
"github.com/sujit-baniya/fiber-boilerplate/pkg/models"
|
||||
)
|
||||
|
||||
@@ -21,97 +24,73 @@ type OSSObject struct {
|
||||
// OSSListResult resultado paginado de listado de objetos
|
||||
type OSSListResult struct {
|
||||
Objects []OSSObject `json:"objects"`
|
||||
CommonPrefixes []string `json:"prefixes"` // "carpetas" virtuales
|
||||
CommonPrefixes []string `json:"prefixes"`
|
||||
IsTruncated bool `json:"is_truncated"`
|
||||
NextMarker string `json:"next_marker"`
|
||||
}
|
||||
|
||||
type OSSService struct {
|
||||
Client *oss.Client
|
||||
Bucket *oss.Bucket
|
||||
BucketName string
|
||||
Endpoint string
|
||||
// OSSProvider interface común para Alibaba OSS y S3-compatible (MinIO)
|
||||
type OSSProvider interface {
|
||||
UploadFile(objectKey, filePath string) error
|
||||
DeleteFile(objectKey string) error
|
||||
ListObjects(prefix, marker string, maxKeys int) (*OSSListResult, error)
|
||||
SignedURL(objectKey string, expireSeconds int) (string, error)
|
||||
UploadFromReader(objectKey, contentType string, r io.Reader) error
|
||||
DeleteObjects(keys []string) error
|
||||
BucketName() string
|
||||
Endpoint() string
|
||||
PublicURL(objectKey string) string
|
||||
}
|
||||
|
||||
// Configuración del OSS
|
||||
type OSSConfig struct {
|
||||
Endpoint string
|
||||
AccessKeyID string
|
||||
AccessKeySecret string
|
||||
BucketName string
|
||||
// ── Alibaba OSS ──────────────────────────────────────────────────────────────
|
||||
|
||||
type alibabaOSS struct {
|
||||
client *oss.Client
|
||||
bucket *oss.Bucket
|
||||
bucketName string
|
||||
endpoint string
|
||||
}
|
||||
|
||||
// Constructor del servicio OSS
|
||||
func NewOSSService(cfg OSSConfig) (*OSSService, error) {
|
||||
func newAlibabaOSS(cfg *models.OssApi) (OSSProvider, error) {
|
||||
client, err := oss.New(cfg.Endpoint, cfg.AccessKeyID, cfg.AccessKeySecret)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("error creando cliente OSS: %w", err)
|
||||
return nil, fmt.Errorf("error creando cliente Alibaba OSS: %w", err)
|
||||
}
|
||||
|
||||
bucket, err := client.Bucket(cfg.BucketName)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("error obteniendo bucket: %w", err)
|
||||
return nil, fmt.Errorf("error obteniendo bucket Alibaba: %w", err)
|
||||
}
|
||||
|
||||
return &OSSService{
|
||||
Client: client,
|
||||
Bucket: bucket,
|
||||
BucketName: cfg.BucketName,
|
||||
Endpoint: cfg.Endpoint,
|
||||
return &alibabaOSS{
|
||||
client: client,
|
||||
bucket: bucket,
|
||||
bucketName: cfg.BucketName,
|
||||
endpoint: cfg.Endpoint,
|
||||
}, nil
|
||||
}
|
||||
|
||||
// Subir un archivo
|
||||
func (s *OSSService) UploadFile(objectKey string, filePath string) error {
|
||||
err := s.Bucket.PutObjectFromFile(objectKey, filePath)
|
||||
if err != nil {
|
||||
return fmt.Errorf("error subiendo archivo: %w", err)
|
||||
func (s *alibabaOSS) BucketName() string { return s.bucketName }
|
||||
func (s *alibabaOSS) Endpoint() string { return s.endpoint }
|
||||
func (s *alibabaOSS) PublicURL(objectKey string) string {
|
||||
return fmt.Sprintf("https://%s.%s/%s", s.bucketName, s.endpoint, objectKey)
|
||||
}
|
||||
|
||||
func (s *alibabaOSS) UploadFile(objectKey, filePath string) error {
|
||||
if err := s.bucket.PutObjectFromFile(objectKey, filePath); err != nil {
|
||||
return fmt.Errorf("error subiendo archivo a Alibaba OSS: %w", err)
|
||||
}
|
||||
log.Printf("Archivo '%s' subido como '%s'", filePath, objectKey)
|
||||
log.Printf("Archivo '%s' subido como '%s' (Alibaba)", filePath, objectKey)
|
||||
return nil
|
||||
}
|
||||
|
||||
// Listar archivos en el bucket
|
||||
func (s *OSSService) ListFiles() ([]string, error) {
|
||||
objects, err := s.Bucket.ListObjects()
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("error listando archivos: %w", err)
|
||||
func (s *alibabaOSS) DeleteFile(objectKey string) error {
|
||||
if err := s.bucket.DeleteObject(objectKey); err != nil {
|
||||
return fmt.Errorf("error eliminando archivo de Alibaba OSS: %w", err)
|
||||
}
|
||||
|
||||
var keys []string
|
||||
for _, obj := range objects.Objects {
|
||||
keys = append(keys, obj.Key)
|
||||
}
|
||||
return keys, nil
|
||||
}
|
||||
|
||||
// NewOSSServiceFromDB construye OSSService desde el último registro activo
|
||||
func NewOSSServiceFromDB() (*OSSService, error) {
|
||||
cfg, err := models.GetLastActiveOssApi()
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("error obteniendo configuración OSS activa: %w", err)
|
||||
}
|
||||
|
||||
return NewOSSService(OSSConfig{
|
||||
Endpoint: cfg.Endpoint,
|
||||
AccessKeyID: cfg.AccessKeyID,
|
||||
AccessKeySecret: cfg.AccessKeySecret,
|
||||
BucketName: cfg.BucketName,
|
||||
})
|
||||
}
|
||||
|
||||
// Eliminar un archivo del bucket
|
||||
func (s *OSSService) DeleteFile(objectKey string) error {
|
||||
err := s.Bucket.DeleteObject(objectKey)
|
||||
if err != nil {
|
||||
return fmt.Errorf("error eliminando archivo '%s' de OSS: %w", objectKey, err)
|
||||
}
|
||||
log.Printf("Archivo eliminado de OSS: %s", objectKey)
|
||||
log.Printf("Archivo eliminado de Alibaba OSS: %s", objectKey)
|
||||
return nil
|
||||
}
|
||||
|
||||
// ListObjects lista objetos del bucket con soporte de prefijo (carpeta), delimitador y paginación.
|
||||
func (s *OSSService) ListObjects(prefix, marker string, maxKeys int) (*OSSListResult, error) {
|
||||
func (s *alibabaOSS) ListObjects(prefix, marker string, maxKeys int) (*OSSListResult, error) {
|
||||
if maxKeys <= 0 || maxKeys > 1000 {
|
||||
maxKeys = 100
|
||||
}
|
||||
@@ -125,12 +104,10 @@ func (s *OSSService) ListObjects(prefix, marker string, maxKeys int) (*OSSListRe
|
||||
if marker != "" {
|
||||
opts = append(opts, oss.Marker(marker))
|
||||
}
|
||||
|
||||
resp, err := s.Bucket.ListObjects(opts...)
|
||||
resp, err := s.bucket.ListObjects(opts...)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("error listando objetos OSS: %w", err)
|
||||
return nil, fmt.Errorf("error listando objetos Alibaba OSS: %w", err)
|
||||
}
|
||||
|
||||
result := &OSSListResult{
|
||||
IsTruncated: resp.IsTruncated,
|
||||
NextMarker: resp.NextMarker,
|
||||
@@ -149,37 +126,221 @@ func (s *OSSService) ListObjects(prefix, marker string, maxKeys int) (*OSSListRe
|
||||
return result, nil
|
||||
}
|
||||
|
||||
// SignedURL genera una URL firmada temporal para descargar/previsualizar un objeto.
|
||||
func (s *OSSService) SignedURL(objectKey string, expireSeconds int) (string, error) {
|
||||
func (s *alibabaOSS) SignedURL(objectKey string, expireSeconds int) (string, error) {
|
||||
if expireSeconds <= 0 {
|
||||
expireSeconds = 3600
|
||||
}
|
||||
url, err := s.Bucket.SignURL(objectKey, oss.HTTPGet, int64(expireSeconds))
|
||||
url, err := s.bucket.SignURL(objectKey, oss.HTTPGet, int64(expireSeconds))
|
||||
if err != nil {
|
||||
return "", fmt.Errorf("error generando URL firmada: %w", err)
|
||||
return "", fmt.Errorf("error generando URL firmada Alibaba: %w", err)
|
||||
}
|
||||
return url, nil
|
||||
}
|
||||
|
||||
// UploadFromReader sube un archivo desde un io.Reader con el content-type dado.
|
||||
func (s *OSSService) UploadFromReader(objectKey, contentType string, r io.Reader) error {
|
||||
func (s *alibabaOSS) UploadFromReader(objectKey, contentType string, r io.Reader) error {
|
||||
var opts []oss.Option
|
||||
if contentType != "" {
|
||||
opts = append(opts, oss.ContentType(contentType))
|
||||
}
|
||||
err := s.Bucket.PutObject(objectKey, r, opts...)
|
||||
if err != nil {
|
||||
return fmt.Errorf("error subiendo objeto '%s': %w", objectKey, err)
|
||||
if err := s.bucket.PutObject(objectKey, r, opts...); err != nil {
|
||||
return fmt.Errorf("error subiendo objeto a Alibaba OSS: %w", err)
|
||||
}
|
||||
log.Printf("Objeto subido a OSS: %s", objectKey)
|
||||
log.Printf("Objeto subido a Alibaba OSS: %s", objectKey)
|
||||
return nil
|
||||
}
|
||||
|
||||
// DeleteObjects elimina múltiples objetos en una sola llamada (batch).
|
||||
func (s *OSSService) DeleteObjects(keys []string) error {
|
||||
_, err := s.Bucket.DeleteObjects(keys)
|
||||
func (s *alibabaOSS) DeleteObjects(keys []string) error {
|
||||
_, err := s.bucket.DeleteObjects(keys)
|
||||
if err != nil {
|
||||
return fmt.Errorf("error eliminando objetos en batch: %w", err)
|
||||
return fmt.Errorf("error eliminando objetos en batch de Alibaba OSS: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// ── S3-compatible (MinIO) ────────────────────────────────────────────────────
|
||||
|
||||
type s3OSS struct {
|
||||
client *minio.Client
|
||||
bucketName string
|
||||
endpoint string
|
||||
}
|
||||
|
||||
func newS3OSS(cfg *models.OssApi) (OSSProvider, error) {
|
||||
useSSL := false
|
||||
endpoint := cfg.Endpoint
|
||||
// Detectar si el endpoint usa HTTPS
|
||||
if len(endpoint) > 8 && endpoint[:8] == "https://" {
|
||||
useSSL = true
|
||||
endpoint = endpoint[8:]
|
||||
} else if len(endpoint) > 7 && endpoint[:7] == "http://" {
|
||||
endpoint = endpoint[7:]
|
||||
}
|
||||
|
||||
client, err := minio.New(endpoint, &minio.Options{
|
||||
Creds: credentials.NewStaticV4(cfg.AccessKeyID, cfg.AccessKeySecret, ""),
|
||||
Secure: useSSL,
|
||||
Region: cfg.Region,
|
||||
})
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("error creando cliente S3/MinIO: %w", err)
|
||||
}
|
||||
|
||||
// Verificar que el bucket existe
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
|
||||
defer cancel()
|
||||
exists, err := client.BucketExists(ctx, cfg.BucketName)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("error verificando bucket S3: %w", err)
|
||||
}
|
||||
if !exists {
|
||||
return nil, fmt.Errorf("bucket '%s' no existe en S3/MinIO", cfg.BucketName)
|
||||
}
|
||||
|
||||
return &s3OSS{
|
||||
client: client,
|
||||
bucketName: cfg.BucketName,
|
||||
endpoint: cfg.Endpoint,
|
||||
}, nil
|
||||
}
|
||||
|
||||
func (s *s3OSS) BucketName() string { return s.bucketName }
|
||||
func (s *s3OSS) Endpoint() string { return s.endpoint }
|
||||
func (s *s3OSS) PublicURL(objectKey string) string {
|
||||
return fmt.Sprintf("%s/%s/%s", s.endpoint, s.bucketName, objectKey)
|
||||
}
|
||||
|
||||
func (s *s3OSS) UploadFile(objectKey, filePath string) error {
|
||||
ctx := context.Background()
|
||||
_, err := s.client.FPutObject(ctx, s.bucketName, objectKey, filePath, minio.PutObjectOptions{})
|
||||
if err != nil {
|
||||
return fmt.Errorf("error subiendo archivo a S3: %w", err)
|
||||
}
|
||||
log.Printf("Archivo '%s' subido como '%s' (S3)", filePath, objectKey)
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *s3OSS) DeleteFile(objectKey string) error {
|
||||
ctx := context.Background()
|
||||
if err := s.client.RemoveObject(ctx, s.bucketName, objectKey, minio.RemoveObjectOptions{}); err != nil {
|
||||
return fmt.Errorf("error eliminando archivo de S3: %w", err)
|
||||
}
|
||||
log.Printf("Archivo eliminado de S3: %s", objectKey)
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *s3OSS) ListObjects(prefix, marker string, maxKeys int) (*OSSListResult, error) {
|
||||
ctx := context.Background()
|
||||
if maxKeys <= 0 || maxKeys > 1000 {
|
||||
maxKeys = 100
|
||||
}
|
||||
|
||||
opts := minio.ListObjectsOptions{
|
||||
Prefix: prefix,
|
||||
Recursive: false, // false = usa '/' como delimiter
|
||||
MaxKeys: maxKeys,
|
||||
}
|
||||
|
||||
if marker != "" {
|
||||
opts.StartAfter = marker
|
||||
}
|
||||
|
||||
result := &OSSListResult{}
|
||||
for obj := range s.client.ListObjects(ctx, s.bucketName, opts) {
|
||||
if obj.Err != nil {
|
||||
return nil, fmt.Errorf("error listando objetos S3: %w", obj.Err)
|
||||
}
|
||||
if len(obj.Key) > 0 && obj.Key[len(obj.Key)-1] == '/' {
|
||||
result.CommonPrefixes = append(result.CommonPrefixes, obj.Key)
|
||||
} else {
|
||||
result.Objects = append(result.Objects, OSSObject{
|
||||
Key: obj.Key,
|
||||
Size: obj.Size,
|
||||
LastModified: obj.LastModified,
|
||||
ETag: obj.ETag,
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// Determinar si hay más resultados y generar next_marker
|
||||
if len(result.Objects) > 0 {
|
||||
last := result.Objects[len(result.Objects)-1]
|
||||
result.NextMarker = last.Key
|
||||
result.IsTruncated = len(result.Objects) >= maxKeys
|
||||
}
|
||||
|
||||
return result, nil
|
||||
}
|
||||
|
||||
func (s *s3OSS) SignedURL(objectKey string, expireSeconds int) (string, error) {
|
||||
if expireSeconds <= 0 {
|
||||
expireSeconds = 3600
|
||||
}
|
||||
ctx := context.Background()
|
||||
url, err := s.client.PresignedGetObject(ctx, s.bucketName, objectKey, time.Duration(expireSeconds)*time.Second, nil)
|
||||
if err != nil {
|
||||
return "", fmt.Errorf("error generando URL firmada S3: %w", err)
|
||||
}
|
||||
return url.String(), nil
|
||||
}
|
||||
|
||||
func (s *s3OSS) UploadFromReader(objectKey, contentType string, r io.Reader) error {
|
||||
ctx := context.Background()
|
||||
opts := minio.PutObjectOptions{}
|
||||
if contentType != "" {
|
||||
opts.ContentType = contentType
|
||||
}
|
||||
_, err := s.client.PutObject(ctx, s.bucketName, objectKey, r, -1, opts)
|
||||
if err != nil {
|
||||
return fmt.Errorf("error subiendo objeto a S3: %w", err)
|
||||
}
|
||||
log.Printf("Objeto subido a S3: %s", objectKey)
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *s3OSS) DeleteObjects(keys []string) error {
|
||||
ctx := context.Background()
|
||||
opts := minio.RemoveObjectsOptions{}
|
||||
ch := make(chan minio.ObjectInfo, len(keys))
|
||||
go func() {
|
||||
defer close(ch)
|
||||
for _, k := range keys {
|
||||
ch <- minio.ObjectInfo{Key: k}
|
||||
}
|
||||
}()
|
||||
for err := range s.client.RemoveObjects(ctx, s.bucketName, ch, opts) {
|
||||
if err.Err != nil {
|
||||
return fmt.Errorf("error eliminando objetos en batch de S3: %w", err.Err)
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// ── Factory ──────────────────────────────────────────────────────────────────
|
||||
|
||||
// NewOSSProvider crea el provider correcto según la configuración
|
||||
func NewOSSProvider(cfg *models.OssApi) (OSSProvider, error) {
|
||||
switch cfg.Provider {
|
||||
case "s3":
|
||||
return newS3OSS(cfg)
|
||||
default:
|
||||
return newAlibabaOSS(cfg)
|
||||
}
|
||||
}
|
||||
|
||||
// NewOSSProviderFromDB obtiene la configuración activa y crea el provider
|
||||
func NewOSSProviderFromDB() (OSSProvider, error) {
|
||||
cfg, err := models.GetLastActiveOssApi()
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("error obteniendo configuración OSS activa: %w", err)
|
||||
}
|
||||
return NewOSSProvider(cfg)
|
||||
}
|
||||
|
||||
// NewOSSProviderByID obtiene una configuración por ID y crea el provider
|
||||
func NewOSSProviderByID(id uint) (OSSProvider, error) {
|
||||
cfg, err := models.GetOssApiByID(id)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("error obteniendo configuración OSS #%d: %w", id, err)
|
||||
}
|
||||
return NewOSSProvider(cfg)
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user