topfans/backend/services/assetService/util/ossutil/qr_uploader.go

110 lines
3.4 KiB
Go

// Package ossutil 提供 assetService 侧对 OSS 的轻量操作封装。
//
// 涉及 STS 凭证的操作建议在请求粒度调用即可(不要跨请求复用 STS,
// STS token 过期会导致后续调用 403)。每次调用都换一次凭证,延迟成本可接受。
package ossutil
import (
"bytes"
"context"
"fmt"
"github.com/aliyun/aliyun-oss-go-sdk/oss"
"github.com/aliyun/credentials-go/credentials"
"github.com/topfans/backend/pkg/logger"
"github.com/topfans/backend/services/assetService/util"
"go.uber.org/zap"
)
// OSSQRUploader 将内存字节流上传到 OSS 的简易封装。
//
// 复用 assetService/util.OSSConfig 避免与现有 OSS 工具并存两套配置。
// 每次 UploadBytes 都重新申请 STS,凭证不跨请求复用。
type OSSQRUploader struct {
config util.OSSConfig
}
// NewOSSQRUploader 创建上传器
func NewOSSQRUploader(cfg util.OSSConfig) *OSSQRUploader {
return &OSSQRUploader{config: cfg}
}
// UploadBytes 把 data 上传到 key,返回完整可访问的 CDN URL。
//
// contentType 例如 "image/png"。
// 返回的 URL 形如 https://<bucket>.oss-<region>.aliyuncs.com/<key>。
func (u *OSSQRUploader) UploadBytes(ctx context.Context, key string, data []byte, contentType string) (string, error) {
if u == nil {
return "", fmt.Errorf("ossutil: uploader is nil")
}
if key == "" {
return "", fmt.Errorf("ossutil: key is empty")
}
if len(data) == 0 {
return "", fmt.Errorf("ossutil: data is empty")
}
if u.config.BucketName == "" || u.config.Region == "" {
return "", fmt.Errorf("ossutil: bucket/region not configured")
}
// 提前检查 ctx 是否已取消,避免无谓的 STS 申请
if cerr := ctx.Err(); cerr != nil {
return "", fmt.Errorf("ossutil: context cancelled: %w", cerr)
}
// 1. 获取 STS 临时凭证
credConfig := new(credentials.Config).
SetType("ram_role_arn").
SetAccessKeyId(u.config.AccessKeyID).
SetAccessKeySecret(u.config.AccessKeySecret).
SetRoleArn(u.config.RoleArn).
SetRoleSessionName("topfans-asset-qrcode").
SetPolicy("").
SetRoleSessionExpiration(3600)
provider, err := credentials.NewCredential(credConfig)
if err != nil {
return "", fmt.Errorf("创建凭证提供器失败: %w", err)
}
cred, err := provider.GetCredential()
if err != nil {
return "", fmt.Errorf("获取临时凭证失败: %w", err)
}
// 2. 创建 OSS 客户端
endpoint := fmt.Sprintf("https://oss-%s.aliyuncs.com", u.config.Region)
client, err := oss.New(endpoint, *cred.AccessKeyId, *cred.AccessKeySecret,
oss.SecurityToken(*cred.SecurityToken))
if err != nil {
return "", fmt.Errorf("创建OSS客户端失败: %w", err)
}
// 3. 获取 Bucket
bucket, err := client.Bucket(u.config.BucketName)
if err != nil {
return "", fmt.Errorf("获取Bucket失败: %w", err)
}
// 4. 上传字节流
opts := []oss.Option{}
if contentType != "" {
opts = append(opts, oss.ContentType(contentType))
}
if err := bucket.PutObject(key, bytes.NewReader(data), opts...); err != nil {
return "", fmt.Errorf("上传OSS对象失败: %w", err)
}
// 5. 构造公开访问 URL
cdnURL := fmt.Sprintf("https://%s.oss-%s.aliyuncs.com/%s", u.config.BucketName, u.config.Region, key)
logger.Logger.Info("ossutil: upload bytes to OSS ok",
zap.String("oss_key", key),
zap.String("content_type", contentType),
zap.Int("bytes", len(data)),
zap.String("url", cdnURL),
)
_ = ctx // 已在函数入口处使用 ctx.Err() 校验取消
return cdnURL, nil
}