653 lines
16 KiB
Go
653 lines
16 KiB
Go
package ossdrv
|
||
|
||
import (
|
||
"bytes"
|
||
"context"
|
||
"encoding/base64"
|
||
"fmt"
|
||
"os"
|
||
"path"
|
||
"strings"
|
||
"sync"
|
||
"time"
|
||
|
||
"u-desk/internal/filesystem"
|
||
"u-desk/internal/oss"
|
||
"u-desk/internal/oss/aliyun"
|
||
"u-desk/internal/oss/qiniu"
|
||
)
|
||
|
||
// accountCredentials 账户级凭据
|
||
type accountCredentials struct {
|
||
Provider string
|
||
AccessKey string
|
||
SecretKey string
|
||
Endpoint string
|
||
}
|
||
|
||
// Manager OSS 连接管理器(两级:账户 + 桶级客户端缓存)
|
||
type Manager struct {
|
||
accounts sync.Map // map[string]*accountCredentials key=provider
|
||
clients sync.Map // map[string]oss.OSSProvider key="provider:bucket"
|
||
bucketRegions sync.Map // map[string]string key="provider:bucket" → region
|
||
}
|
||
|
||
var globalManager = &Manager{}
|
||
|
||
func GetManager() *Manager { return globalManager }
|
||
|
||
// Connect 建立账户级连接(验证凭据通过 ListBuckets)
|
||
func (m *Manager) Connect(provider, accessKey, secretKey, endpoint string) error {
|
||
// 验证凭据
|
||
switch provider {
|
||
case "qiniu":
|
||
_, err := qiniu.ListBuckets(accessKey, secretKey)
|
||
if err != nil {
|
||
return fmt.Errorf("七牛云连接失败: %w", err)
|
||
}
|
||
case "aliyun":
|
||
_, err := aliyun.ListBuckets(accessKey, secretKey, endpoint)
|
||
if err != nil {
|
||
return fmt.Errorf("阿里云连接失败: %w", err)
|
||
}
|
||
default:
|
||
return fmt.Errorf("不支持的 OSS 提供商: %s", provider)
|
||
}
|
||
|
||
m.accounts.Store(provider, &accountCredentials{
|
||
Provider: provider,
|
||
AccessKey: accessKey,
|
||
SecretKey: secretKey,
|
||
Endpoint: endpoint,
|
||
})
|
||
return nil
|
||
}
|
||
|
||
// getOrCreateBucketClient 懒创建桶级 OSSProvider
|
||
func (m *Manager) getOrCreateBucketClient(provider, bucket, region string) (oss.OSSProvider, error) {
|
||
key := provider + ":" + bucket
|
||
if v, ok := m.clients.Load(key); ok {
|
||
return v.(oss.OSSProvider), nil
|
||
}
|
||
|
||
cred, ok := m.accounts.Load(provider)
|
||
if !ok {
|
||
return nil, fmt.Errorf("OSS 账户未连接: %s", provider)
|
||
}
|
||
c := cred.(*accountCredentials)
|
||
|
||
// 如果未传 region,从缓存取
|
||
if region == "" {
|
||
if v, ok := m.bucketRegions.Load(key); ok {
|
||
region = v.(string)
|
||
}
|
||
}
|
||
|
||
var client oss.OSSProvider
|
||
var err error
|
||
|
||
switch provider {
|
||
case "qiniu":
|
||
client, err = qiniu.NewClient(&qiniu.Config{
|
||
AccessKey: c.AccessKey,
|
||
SecretKey: c.SecretKey,
|
||
Bucket: bucket,
|
||
Region: region,
|
||
UseHTTPS: true,
|
||
})
|
||
case "aliyun":
|
||
client, err = aliyun.NewClient(&aliyun.Config{
|
||
AccessKeyID: c.AccessKey,
|
||
AccessKeySecret: c.SecretKey,
|
||
Bucket: bucket,
|
||
Region: region,
|
||
Endpoint: c.Endpoint,
|
||
UseHTTPS: true,
|
||
})
|
||
default:
|
||
return nil, fmt.Errorf("不支持的提供商: %s", provider)
|
||
}
|
||
|
||
if err != nil {
|
||
return nil, fmt.Errorf("创建桶客户端失败: %w", err)
|
||
}
|
||
|
||
m.clients.Store(key, client)
|
||
return client, nil
|
||
}
|
||
|
||
// GetClient 获取已有的桶级客户端
|
||
func (m *Manager) GetClient(provider, bucket string) oss.OSSProvider {
|
||
if v, ok := m.clients.Load(provider + ":" + bucket); ok {
|
||
return v.(oss.OSSProvider)
|
||
}
|
||
return nil
|
||
}
|
||
|
||
// Disconnect 关闭账户及所有桶级客户端
|
||
func (m *Manager) Disconnect(provider string) {
|
||
m.accounts.Delete(provider)
|
||
prefix := provider + ":"
|
||
m.clients.Range(func(key, value any) bool {
|
||
if strings.HasPrefix(key.(string), prefix) {
|
||
value.(oss.OSSProvider).Close()
|
||
m.clients.Delete(key)
|
||
}
|
||
return true
|
||
})
|
||
m.bucketRegions.Range(func(key, value any) bool {
|
||
if strings.HasPrefix(key.(string), prefix) {
|
||
m.bucketRegions.Delete(key)
|
||
}
|
||
return true
|
||
})
|
||
}
|
||
|
||
// Shutdown 关闭所有连接
|
||
func (m *Manager) Shutdown() {
|
||
m.clients.Range(func(key, value any) bool {
|
||
value.(oss.OSSProvider).Close()
|
||
m.clients.Delete(key)
|
||
return true
|
||
})
|
||
m.accounts.Range(func(key, value any) bool {
|
||
m.accounts.Delete(key)
|
||
return true
|
||
})
|
||
}
|
||
|
||
// Service OSS 文件操作服务
|
||
type Service struct {
|
||
manager *Manager
|
||
}
|
||
|
||
func NewService() *Service {
|
||
return &Service{manager: GetManager()}
|
||
}
|
||
|
||
func (s *Service) GetManager() *Manager { return s.manager }
|
||
|
||
// parseBucketPath 解析路径中的桶名和对象键
|
||
// "/my-bucket/photos/img.jpg" → bucket="my-bucket", key="photos/img.jpg"
|
||
func parseBucketPath(rawPath string) (bucket, key string) {
|
||
rawPath = strings.TrimPrefix(rawPath, "/")
|
||
if rawPath == "" {
|
||
return "", ""
|
||
}
|
||
parts := strings.SplitN(rawPath, "/", 2)
|
||
bucket = parts[0]
|
||
if len(parts) > 1 {
|
||
key = parts[1]
|
||
}
|
||
return
|
||
}
|
||
|
||
// listBuckets 列出所有桶
|
||
func (s *Service) listBuckets(provider string) ([]map[string]interface{}, error) {
|
||
cred, ok := s.manager.accounts.Load(provider)
|
||
if !ok {
|
||
return nil, fmt.Errorf("OSS 账户未连接: %s", provider)
|
||
}
|
||
c := cred.(*accountCredentials)
|
||
|
||
var entries []oss.BucketEntry
|
||
var err error
|
||
|
||
switch provider {
|
||
case "qiniu":
|
||
entries, err = qiniu.ListBuckets(c.AccessKey, c.SecretKey)
|
||
case "aliyun":
|
||
entries, err = aliyun.ListBuckets(c.AccessKey, c.SecretKey, c.Endpoint)
|
||
default:
|
||
return nil, fmt.Errorf("不支持的提供商: %s", provider)
|
||
}
|
||
if err != nil {
|
||
return nil, fmt.Errorf("列举存储桶失败: %w", err)
|
||
}
|
||
|
||
// 缓存桶区域信息
|
||
for _, e := range entries {
|
||
if e.Region != "" {
|
||
s.manager.bucketRegions.Store(provider+":"+e.Name, e.Region)
|
||
}
|
||
}
|
||
|
||
items := make([]map[string]interface{}, len(entries))
|
||
for i, e := range entries {
|
||
items[i] = map[string]interface{}{
|
||
"name": e.Name,
|
||
"path": "/" + e.Name,
|
||
"is_dir": true,
|
||
"is_bucket": true,
|
||
"size": int64(0),
|
||
}
|
||
}
|
||
return items, nil
|
||
}
|
||
|
||
// ListDir 列出目录内容
|
||
func (s *Service) ListDir(connID string, prefix string) ([]map[string]interface{}, error) {
|
||
prefix = strings.TrimPrefix(prefix, "/")
|
||
|
||
// 根目录 → 列出所有桶
|
||
if prefix == "" {
|
||
return s.listBuckets(connID)
|
||
}
|
||
|
||
// 解析桶名和对象前缀
|
||
bucket, objectPrefix := parseBucketPath(prefix)
|
||
if bucket == "" {
|
||
return s.listBuckets(connID)
|
||
}
|
||
|
||
if objectPrefix != "" && !strings.HasSuffix(objectPrefix, "/") {
|
||
objectPrefix += "/"
|
||
}
|
||
|
||
c, err := s.manager.getOrCreateBucketClient(connID, bucket, "")
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
|
||
ctx := context.Background()
|
||
result, err := c.ListFiles(ctx, &oss.ListOptions{
|
||
Prefix: objectPrefix,
|
||
Delimiter: "/",
|
||
MaxKeys: 1000,
|
||
})
|
||
if err != nil {
|
||
return nil, fmt.Errorf("列举文件失败: %w", err)
|
||
}
|
||
|
||
items := make([]map[string]interface{}, 0, len(result.Files)+len(result.Prefixes))
|
||
bucketPrefix := "/" + bucket + "/"
|
||
|
||
for _, p := range result.Prefixes {
|
||
name := strings.TrimSuffix(strings.TrimPrefix(p, objectPrefix), "/")
|
||
if name == "" {
|
||
continue
|
||
}
|
||
items = append(items, map[string]interface{}{
|
||
"name": name,
|
||
"path": bucketPrefix + p,
|
||
"is_dir": true,
|
||
"size": int64(0),
|
||
})
|
||
}
|
||
|
||
for _, f := range result.Files {
|
||
if strings.HasSuffix(f.Key, "/") && f.Size == 0 {
|
||
continue
|
||
}
|
||
items = append(items, map[string]interface{}{
|
||
"name": path.Base(f.Key),
|
||
"path": bucketPrefix + f.Key,
|
||
"is_dir": false,
|
||
"size": f.Size,
|
||
"mod_time": f.LastModified.Format("2006-01-02 15:04:05"),
|
||
})
|
||
}
|
||
|
||
return items, nil
|
||
}
|
||
|
||
// ReadFile 读取文件内容
|
||
func (s *Service) ReadFile(connID string, rawPath string) (string, error) {
|
||
bucket, key := parseBucketPath(rawPath)
|
||
if bucket == "" {
|
||
return "", fmt.Errorf("路径中缺少桶名")
|
||
}
|
||
|
||
c, err := s.manager.getOrCreateBucketClient(connID, bucket, "")
|
||
if err != nil {
|
||
return "", err
|
||
}
|
||
|
||
const maxSize int64 = 10 << 20
|
||
ctx := context.Background()
|
||
|
||
info, err := c.GetFileInfo(ctx, key)
|
||
if err != nil {
|
||
return "", fmt.Errorf("获取文件信息失败: %w", err)
|
||
}
|
||
if info.Size > maxSize {
|
||
return "", fmt.Errorf("文件过大 (%s),超过 %d 限制", filesystem.FormatBytes(info.Size), maxSize)
|
||
}
|
||
|
||
var buf bytes.Buffer
|
||
if err := c.Download(ctx, key, &buf); err != nil {
|
||
return "", fmt.Errorf("读取文件失败: %w", err)
|
||
}
|
||
return filesystem.BytesToString(buf.Bytes()), nil
|
||
}
|
||
|
||
// WriteFile 写入文件内容
|
||
func (s *Service) WriteFile(connID string, rawPath string, content string) error {
|
||
bucket, key := parseBucketPath(rawPath)
|
||
if bucket == "" {
|
||
return fmt.Errorf("路径中缺少桶名")
|
||
}
|
||
|
||
c, err := s.manager.getOrCreateBucketClient(connID, bucket, "")
|
||
if err != nil {
|
||
return err
|
||
}
|
||
|
||
_, err = c.Upload(context.Background(), key, strings.NewReader(content), nil)
|
||
if err != nil {
|
||
return fmt.Errorf("写入文件失败: %w", err)
|
||
}
|
||
return nil
|
||
}
|
||
|
||
// WriteBase64File 写入 base64 编码的二进制文件
|
||
func (s *Service) WriteBase64File(connID string, rawPath string, base64Content string) error {
|
||
bucket, key := parseBucketPath(rawPath)
|
||
if bucket == "" {
|
||
return fmt.Errorf("路径中缺少桶名")
|
||
}
|
||
|
||
data, err := base64.StdEncoding.DecodeString(base64Content)
|
||
if err != nil {
|
||
return fmt.Errorf("base64 解码失败: %w", err)
|
||
}
|
||
|
||
c, err := s.manager.getOrCreateBucketClient(connID, bucket, "")
|
||
if err != nil {
|
||
return err
|
||
}
|
||
|
||
_, err = c.Upload(context.Background(), key, bytes.NewReader(data), nil)
|
||
if err != nil {
|
||
return fmt.Errorf("写入文件失败: %w", err)
|
||
}
|
||
return nil
|
||
}
|
||
|
||
// GetFileInfo 获取文件信息
|
||
func (s *Service) GetFileInfo(connID string, rawPath string) (map[string]interface{}, error) {
|
||
bucket, key := parseBucketPath(rawPath)
|
||
if bucket == "" {
|
||
return nil, fmt.Errorf("路径中缺少桶名")
|
||
}
|
||
|
||
c, err := s.manager.getOrCreateBucketClient(connID, bucket, "")
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
|
||
info, err := c.GetFileInfo(context.Background(), key)
|
||
if err != nil {
|
||
return nil, fmt.Errorf("获取文件信息失败: %w", err)
|
||
}
|
||
|
||
bucketPrefix := "/" + bucket + "/"
|
||
return map[string]interface{}{
|
||
"name": path.Base(info.Key),
|
||
"path": bucketPrefix + info.Key,
|
||
"size": info.Size,
|
||
"size_str": filesystem.FormatBytes(info.Size),
|
||
"is_dir": strings.HasSuffix(info.Key, "/"),
|
||
"mod_time": info.LastModified.Format("2006-01-02 15:04:05"),
|
||
}, nil
|
||
}
|
||
|
||
// CreateDir 创建目录
|
||
func (s *Service) CreateDir(connID string, rawPath string) (*filesystem.FileOperationResult, error) {
|
||
bucket, key := parseBucketPath(rawPath)
|
||
if bucket == "" {
|
||
return nil, fmt.Errorf("路径中缺少桶名")
|
||
}
|
||
|
||
c, err := s.manager.getOrCreateBucketClient(connID, bucket, "")
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
|
||
if !strings.HasSuffix(key, "/") {
|
||
key += "/"
|
||
}
|
||
|
||
_, err = c.Upload(context.Background(), key, strings.NewReader(""), nil)
|
||
if err != nil {
|
||
return nil, fmt.Errorf("创建目录失败: %w", err)
|
||
}
|
||
|
||
name := path.Base(strings.TrimSuffix(key, "/"))
|
||
return &filesystem.FileOperationResult{
|
||
Path: "/" + bucket + "/" + key,
|
||
Name: name,
|
||
IsDir: true,
|
||
SizeStr: filesystem.FormatBytes(0),
|
||
}, nil
|
||
}
|
||
|
||
// CreateFile 创建空文件
|
||
func (s *Service) CreateFile(connID string, rawPath string) (*filesystem.FileOperationResult, error) {
|
||
bucket, key := parseBucketPath(rawPath)
|
||
if bucket == "" {
|
||
return nil, fmt.Errorf("路径中缺少桶名")
|
||
}
|
||
|
||
c, err := s.manager.getOrCreateBucketClient(connID, bucket, "")
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
|
||
_, err = c.Upload(context.Background(), key, strings.NewReader(""), nil)
|
||
if err != nil {
|
||
return nil, fmt.Errorf("创建文件失败: %w", err)
|
||
}
|
||
|
||
return &filesystem.FileOperationResult{
|
||
Path: "/" + bucket + "/" + key,
|
||
Name: path.Base(key),
|
||
IsDir: false,
|
||
SizeStr: filesystem.FormatBytes(0),
|
||
}, nil
|
||
}
|
||
|
||
// DeletePath 删除文件或目录
|
||
func (s *Service) DeletePath(connID string, rawPath string) (*filesystem.FileOperationResult, error) {
|
||
bucket, key := parseBucketPath(rawPath)
|
||
if bucket == "" {
|
||
return nil, fmt.Errorf("路径中缺少桶名")
|
||
}
|
||
|
||
c, err := s.manager.getOrCreateBucketClient(connID, bucket, "")
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
|
||
ctx := context.Background()
|
||
|
||
isDir := strings.HasSuffix(key, "/")
|
||
if !isDir {
|
||
prefix := key + "/"
|
||
listResult, listErr := c.ListFiles(ctx, &oss.ListOptions{Prefix: prefix, MaxKeys: 1})
|
||
if listErr == nil && len(listResult.Files) > 0 {
|
||
isDir = true
|
||
key = prefix
|
||
}
|
||
}
|
||
|
||
infoMap, _ := s.GetFileInfo(connID, "/"+bucket+"/"+key)
|
||
|
||
if isDir {
|
||
prefix := key
|
||
if !strings.HasSuffix(prefix, "/") {
|
||
prefix += "/"
|
||
}
|
||
for {
|
||
listResult, err := c.ListFiles(ctx, &oss.ListOptions{Prefix: prefix, MaxKeys: 1000})
|
||
if err != nil {
|
||
return nil, fmt.Errorf("列举目录文件失败: %w", err)
|
||
}
|
||
if len(listResult.Files) == 0 {
|
||
break
|
||
}
|
||
keys := make([]string, len(listResult.Files))
|
||
for i, f := range listResult.Files {
|
||
keys[i] = f.Key
|
||
}
|
||
if _, err := c.DeleteMultiple(ctx, keys); err != nil {
|
||
return nil, fmt.Errorf("批量删除失败: %w", err)
|
||
}
|
||
if !listResult.IsTruncated {
|
||
break
|
||
}
|
||
}
|
||
c.Delete(ctx, key) // marker 非关键,忽略错误
|
||
} else {
|
||
if err := c.Delete(ctx, key); err != nil {
|
||
return nil, fmt.Errorf("删除失败: %w", err)
|
||
}
|
||
}
|
||
|
||
result := toOssOperationResult(infoMap, isDir)
|
||
result.Deleted = true
|
||
return result, nil
|
||
}
|
||
|
||
// RenamePath 重命名(Copy + Delete)
|
||
func (s *Service) RenamePath(connID string, oldPath string, newPath string) (*filesystem.FileOperationResult, error) {
|
||
oldBucket, oldKey := parseBucketPath(oldPath)
|
||
newBucket, newKey := parseBucketPath(newPath)
|
||
if oldBucket == "" || newBucket == "" {
|
||
return nil, fmt.Errorf("路径中缺少桶名")
|
||
}
|
||
|
||
if oldBucket != newBucket {
|
||
return nil, fmt.Errorf("不支持跨桶重命名")
|
||
}
|
||
|
||
c, err := s.manager.getOrCreateBucketClient(connID, oldBucket, "")
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
|
||
ctx := context.Background()
|
||
|
||
isDir := strings.HasSuffix(oldKey, "/")
|
||
if !isDir {
|
||
prefix := oldKey + "/"
|
||
listResult, listErr := c.ListFiles(ctx, &oss.ListOptions{Prefix: prefix, MaxKeys: 1})
|
||
if listErr == nil && len(listResult.Files) > 0 {
|
||
isDir = true
|
||
oldKey = prefix
|
||
}
|
||
}
|
||
|
||
if isDir {
|
||
oldPrefix := oldKey
|
||
newPrefix := newKey
|
||
if !strings.HasSuffix(oldPrefix, "/") {
|
||
oldPrefix += "/"
|
||
}
|
||
if !strings.HasSuffix(newPrefix, "/") {
|
||
newPrefix += "/"
|
||
}
|
||
|
||
for {
|
||
listResult, err := c.ListFiles(ctx, &oss.ListOptions{Prefix: oldPrefix, MaxKeys: 1000})
|
||
if err != nil {
|
||
return nil, fmt.Errorf("列举目录文件失败: %w", err)
|
||
}
|
||
if len(listResult.Files) == 0 {
|
||
break
|
||
}
|
||
for _, f := range listResult.Files {
|
||
relativeKey := strings.TrimPrefix(f.Key, oldPrefix)
|
||
if err := c.Copy(ctx, f.Key, newPrefix+relativeKey); err != nil {
|
||
return nil, fmt.Errorf("复制失败: %w", err)
|
||
}
|
||
c.Delete(ctx, f.Key)
|
||
}
|
||
if !listResult.IsTruncated {
|
||
break
|
||
}
|
||
}
|
||
c.Delete(ctx, oldKey) // marker
|
||
} else {
|
||
if err := c.Copy(ctx, oldKey, newKey); err != nil {
|
||
return nil, fmt.Errorf("复制失败: %w", err)
|
||
}
|
||
if err := c.Delete(ctx, oldKey); err != nil {
|
||
return nil, fmt.Errorf("删除源文件失败: %w", err)
|
||
}
|
||
}
|
||
|
||
infoMap, _ := s.GetFileInfo(connID, newPath)
|
||
result := toOssOperationResult(infoMap, isDir)
|
||
result.OldPath = oldPath
|
||
return result, nil
|
||
}
|
||
|
||
// DownloadToTemp 下载文件到本地临时目录
|
||
func (s *Service) DownloadToTemp(connID string, rawPath string) (string, error) {
|
||
bucket, key := parseBucketPath(rawPath)
|
||
if bucket == "" {
|
||
return "", fmt.Errorf("路径中缺少桶名")
|
||
}
|
||
|
||
c, err := s.manager.getOrCreateBucketClient(connID, bucket, "")
|
||
if err != nil {
|
||
return "", err
|
||
}
|
||
|
||
f, err := os.CreateTemp("", "udesk-oss-*-"+path.Base(key))
|
||
if err != nil {
|
||
return "", fmt.Errorf("创建临时文件失败: %w", err)
|
||
}
|
||
localPath := f.Name()
|
||
defer f.Close()
|
||
|
||
if err := c.Download(context.Background(), key, f); err != nil {
|
||
os.Remove(localPath)
|
||
return "", fmt.Errorf("下载文件失败: %w", err)
|
||
}
|
||
return localPath, nil
|
||
}
|
||
|
||
// GetCommonPaths 返回常用路径
|
||
func (s *Service) GetCommonPaths(connID string) (map[string]string, error) {
|
||
return map[string]string{
|
||
"root": "/",
|
||
}, nil
|
||
}
|
||
|
||
// GetSignedURL 获取预签名 URL
|
||
func (s *Service) GetSignedURL(connID string, rawPath string) (string, error) {
|
||
bucket, key := parseBucketPath(rawPath)
|
||
if bucket == "" {
|
||
return "", fmt.Errorf("路径中缺少桶名")
|
||
}
|
||
|
||
c, err := s.manager.getOrCreateBucketClient(connID, bucket, "")
|
||
if err != nil {
|
||
return "", err
|
||
}
|
||
|
||
url, err := c.GetSignedURL(context.Background(), key, 1*time.Hour)
|
||
if err != nil {
|
||
return "", fmt.Errorf("获取签名 URL 失败: %w", err)
|
||
}
|
||
return url, nil
|
||
}
|
||
|
||
func toOssOperationResult(m map[string]interface{}, isDir bool) *filesystem.FileOperationResult {
|
||
name, _ := m["name"].(string)
|
||
p, _ := m["path"].(string)
|
||
size, _ := m["size"].(int64)
|
||
modTime, _ := m["mod_time"].(string)
|
||
|
||
return &filesystem.FileOperationResult{
|
||
Path: p,
|
||
Name: name,
|
||
Size: size,
|
||
SizeStr: filesystem.FormatBytes(size),
|
||
IsDir: isDir,
|
||
ModTime: modTime,
|
||
}
|
||
}
|