Files
oci-portal/internal/service/objectstorage.go
T

350 lines
12 KiB
Go

package service
import (
"context"
"errors"
"fmt"
"log"
"strings"
"sync"
"time"
"oci-portal/internal/oci"
)
// ObjectStorageNamespace 查询租户 namespace。
func (s *OciConfigService) ObjectStorageNamespace(ctx context.Context, id uint, region string) (string, error) {
cred, err := s.credentialsByID(ctx, id)
if err != nil {
return "", err
}
return s.client.GetObjectStorageNamespace(ctx, cred, region)
}
// Buckets 列出指定区间的存储桶(compartmentID 空为生效 compartment)。
func (s *OciConfigService) Buckets(ctx context.Context, id uint, region, compartmentID string) ([]oci.Bucket, error) {
cred, err := s.credentialsByID(ctx, id)
if err != nil {
return nil, err
}
return s.client.ListBuckets(ctx, cred, region, compartmentID)
}
// CreateBucket 创建存储桶,名称做基本合法性校验。
func (s *OciConfigService) CreateBucket(ctx context.Context, id uint, region string, in oci.CreateBucketInput) (oci.Bucket, error) {
if err := validateBucketName(in.Name); err != nil {
return oci.Bucket{}, err
}
cred, err := s.credentialsByID(ctx, id)
if err != nil {
return oci.Bucket{}, err
}
return s.client.CreateBucket(ctx, cred, region, in)
}
// validateBucketName 桶名规则:1-256 位字母数字连字符下划线句点。
func validateBucketName(name string) error {
if name == "" || len(name) > 256 {
return fmt.Errorf("create bucket: name length must be 1-256")
}
for _, r := range name {
ok := r == '-' || r == '_' || r == '.' ||
(r >= '0' && r <= '9') || (r >= 'a' && r <= 'z') || (r >= 'A' && r <= 'Z')
if !ok {
return fmt.Errorf("create bucket: invalid character %q in name", r)
}
}
return nil
}
// UpdateBucket 更新桶可见性 / 版本控制。
func (s *OciConfigService) UpdateBucket(ctx context.Context, id uint, region, name string, in oci.UpdateBucketInput) (oci.Bucket, error) {
cred, err := s.credentialsByID(ctx, id)
if err != nil {
return oci.Bucket{}, err
}
return s.client.UpdateBucket(ctx, cred, region, name, in)
}
// DeleteBucket 删除桶:空桶同步秒删;非空桶转后台清空(对象全部版本 + PAR)后删除,
// 返回 queued=true 表示已排队。
func (s *OciConfigService) DeleteBucket(ctx context.Context, id uint, region, name string) (bool, error) {
cred, err := s.credentialsByID(ctx, id)
if err != nil {
return false, err
}
err = s.client.DeleteBucket(ctx, cred, region, name)
if err == nil {
return false, nil
}
if !errors.Is(err, oci.ErrBucketNotEmpty) {
return false, err
}
go s.purgeAndDeleteBucket(cred, region, name)
return true, nil
}
// purgeBucketTimeout 是后台清空删除桶的总时限;超大桶超时后可重试删除续跑。
const purgeBucketTimeout = 30 * time.Minute
// purgeAndDeleteBucket 后台清空并删除桶;失败只记日志(删除请求已受理)。
func (s *OciConfigService) purgeAndDeleteBucket(cred oci.Credentials, region, name string) {
ctx, cancel := context.WithTimeout(context.Background(), purgeBucketTimeout)
defer cancel()
start := time.Now()
if err := s.purgeBucketVersions(ctx, cred, region, name); err != nil {
log.Printf("[bucket-purge] %s purge objects failed: %v", name, err)
return
}
if err := s.purgeBucketPARs(ctx, cred, region, name); err != nil {
log.Printf("[bucket-purge] %s purge pars failed: %v", name, err)
}
if err := s.client.DeleteBucket(ctx, cred, region, name); err != nil {
log.Printf("[bucket-purge] %s delete failed: %v", name, err)
return
}
log.Printf("[bucket-purge] %s purged and deleted in %s", name, time.Since(start).Round(time.Second))
}
// purgeBucketVersions 逐页并发删除全部对象版本(未开版本控制的桶即当前对象);
// 结束时记总量与耗时,便于判断慢在对象清空还是 PAR 清空。
func (s *OciConfigService) purgeBucketVersions(ctx context.Context, cred oci.Credentials, region, name string) error {
start, total := time.Now(), 0
for {
versions, next, err := s.client.ListObjectVersions(ctx, cred, region, name, "")
if err != nil {
return err
}
if len(versions) == 0 {
break
}
err = forEachConcurrently(bulkDeleteWorkers, versions, func(v oci.ObjectVersion) error {
return s.client.DeleteObjectVersion(ctx, cred, region, name, v.Name, v.VersionID)
})
if err != nil {
return err
}
total += len(versions)
// 每轮都从首页重新列:删除后游标失效,直到列表为空
if next == "" && len(versions) < 1000 {
break
}
}
if total > 0 {
log.Printf("[bucket-purge] %s purged %d object versions in %s", name, total, time.Since(start).Round(time.Second))
}
return nil
}
// purgeBucketPARs 并发删除桶内全部预签名请求(与手动「全部删除」同一路径);
// 记录成功数/总数与耗时,速率长期顶在 ~10/s 即为 OCI 服务端限流。
func (s *OciConfigService) purgeBucketPARs(ctx context.Context, cred oci.Credentials, region, name string) error {
start := time.Now()
pars, err := s.client.ListPARs(ctx, cred, region, name)
if err != nil {
return err
}
if len(pars) == 0 {
return nil
}
deleted, err := s.deletePARsConcurrently(ctx, cred, region, name, pars)
log.Printf("[bucket-purge] %s purged %d/%d pars in %s", name, deleted, len(pars), time.Since(start).Round(time.Second))
return err
}
// Objects 分页列出对象(delimiter 前缀模式)。
func (s *OciConfigService) Objects(ctx context.Context, id uint, region, bucket, prefix, startWith string, limit int) (oci.ListObjectsResult, error) {
cred, err := s.credentialsByID(ctx, id)
if err != nil {
return oci.ListObjectsResult{}, err
}
return s.client.ListObjects(ctx, cred, region, bucket, prefix, startWith, limit)
}
// DeleteObject 删除对象。
func (s *OciConfigService) DeleteObject(ctx context.Context, id uint, region, bucket, object string) error {
if object == "" {
return fmt.Errorf("delete object: object name is required")
}
cred, err := s.credentialsByID(ctx, id)
if err != nil {
return err
}
return s.client.DeleteObject(ctx, cred, region, bucket, object)
}
// RenameObject 重命名对象。
func (s *OciConfigService) RenameObject(ctx context.Context, id uint, region, bucket, src, dst string) error {
if src == "" || dst == "" || src == dst {
return fmt.Errorf("rename object: source and distinct new name are required")
}
cred, err := s.credentialsByID(ctx, id)
if err != nil {
return err
}
return s.client.RenameObject(ctx, cred, region, bucket, src, dst)
}
// RestoreObject 取回 Archive 对象。
func (s *OciConfigService) RestoreObject(ctx context.Context, id uint, region, bucket, object string, hours int) error {
cred, err := s.credentialsByID(ctx, id)
if err != nil {
return err
}
return s.client.RestoreObject(ctx, cred, region, bucket, object, hours)
}
// ObjectDetail 查询对象元数据。
func (s *OciConfigService) ObjectDetail(ctx context.Context, id uint, region, bucket, object string) (oci.ObjectDetail, error) {
cred, err := s.credentialsByID(ctx, id)
if err != nil {
return oci.ObjectDetail{}, err
}
return s.client.HeadObject(ctx, cred, region, bucket, object)
}
// 对象内容中转上限:GET 服务预览(office 文档可达数 MB),PUT 服务文本编辑保存。
const (
maxObjectGetBytes = 20 << 20
maxObjectPutBytes = 5 << 20
)
// ObjectContent 读取对象内容(面板中转,预览/编辑用,不签发 PAR)。
func (s *OciConfigService) ObjectContent(ctx context.Context, id uint, region, bucket, object string) (oci.ObjectContent, error) {
if object == "" {
return oci.ObjectContent{}, fmt.Errorf("get object content: object name is required")
}
cred, err := s.credentialsByID(ctx, id)
if err != nil {
return oci.ObjectContent{}, err
}
return s.client.GetObject(ctx, cred, region, bucket, object, maxObjectGetBytes)
}
// PutObjectContent 保存对象内容(面板中转;ifMatch 冲突时 OCI 返回 412 透出)。
func (s *OciConfigService) PutObjectContent(ctx context.Context, id uint, region, bucket, object string, data []byte, contentType, ifMatch string) (string, error) {
if object == "" {
return "", fmt.Errorf("put object content: object name is required")
}
if int64(len(data)) > maxObjectPutBytes {
return "", fmt.Errorf("put object content %s: %w", object, oci.ErrObjectTooLarge)
}
cred, err := s.credentialsByID(ctx, id)
if err != nil {
return "", err
}
return s.client.PutObject(ctx, cred, region, bucket, object, data, contentType, ifMatch)
}
// parAccessTypes 是允许签发的 PAR 访问类型。
var parAccessTypes = map[string]bool{
"ObjectRead": true, "ObjectWrite": true, "ObjectReadWrite": true, "AnyObjectReadWrite": true,
}
// CreatePAR 签发预签名请求;过期时长限制在 1 小时到 30 天。
func (s *OciConfigService) CreatePAR(ctx context.Context, id uint, region, bucket string, in oci.CreatePARInput) (oci.PAR, error) {
if !parAccessTypes[in.AccessType] {
return oci.PAR{}, fmt.Errorf("create par: unsupported access type %q", in.AccessType)
}
if in.ExpiresHours < 1 || in.ExpiresHours > 30*24 {
return oci.PAR{}, fmt.Errorf("create par: expiresHours must be within 1-720")
}
if strings.TrimSpace(in.Name) == "" {
in.Name = "par-" + strings.ReplaceAll(in.ObjectName, "/", "-")
}
cred, err := s.credentialsByID(ctx, id)
if err != nil {
return oci.PAR{}, err
}
return s.client.CreatePAR(ctx, cred, region, bucket, in)
}
// PARsPage 分页列出桶内预签名请求;limit 归一到 1-1000,默认 100。
func (s *OciConfigService) PARsPage(ctx context.Context, id uint, region, bucket, page string, limit int) ([]oci.PAR, string, error) {
if limit <= 0 || limit > 1000 {
limit = 100
}
cred, err := s.credentialsByID(ctx, id)
if err != nil {
return nil, "", err
}
return s.client.ListPARsPage(ctx, cred, region, bucket, page, limit)
}
// DeletePAR 撤销预签名请求。
func (s *OciConfigService) DeletePAR(ctx context.Context, id uint, region, bucket, parID string) error {
if parID == "" {
return fmt.Errorf("delete par: parId is required")
}
cred, err := s.credentialsByID(ctx, id)
if err != nil {
return err
}
return s.client.DeletePAR(ctx, cred, region, bucket, parID)
}
// DeleteAllPARs 删除桶的全部预签名请求,返回成功删除数量;
// 个别失败不中断其余删除,返回首个错误供上层提示。
func (s *OciConfigService) DeleteAllPARs(ctx context.Context, id uint, region, bucket string) (int, error) {
cred, err := s.credentialsByID(ctx, id)
if err != nil {
return 0, err
}
pars, err := s.client.ListPARs(ctx, cred, region, bucket)
if err != nil {
return 0, err
}
return s.deletePARsConcurrently(ctx, cred, region, bucket, pars)
}
// bulkDeleteWorkers 批量删除并发度;OCI 无批量接口,串行逐条会到分钟级。
// 16 配合请求级 429/5xx 退避重试(internal/oci)在限流与吞吐间取平衡,不再往上提。
const bulkDeleteWorkers = 16
// forEachConcurrently 以固定并发度对 items 逐个执行 fn;
// 个别失败不中断其余执行,返回首个错误。
func forEachConcurrently[T any](workers int, items []T, fn func(T) error) error {
// 信号量按并发度限流,容量=workers(规范允许的注释说明场景)
sem := make(chan struct{}, workers)
var (
wg sync.WaitGroup
mu sync.Mutex
firstErr error
)
for _, it := range items {
wg.Add(1)
sem <- struct{}{}
go func(it T) {
defer wg.Done()
defer func() { <-sem }()
if err := fn(it); err != nil {
mu.Lock()
if firstErr == nil {
firstErr = err
}
mu.Unlock()
}
}(it)
}
wg.Wait()
return firstErr
}
// deletePARsConcurrently 并发删除给定 PAR,返回成功删除数量与首个错误。
func (s *OciConfigService) deletePARsConcurrently(ctx context.Context, cred oci.Credentials, region, bucket string, pars []oci.PAR) (int, error) {
var (
mu sync.Mutex
deleted int
)
err := forEachConcurrently(bulkDeleteWorkers, pars, func(p oci.PAR) error {
if err := s.client.DeletePAR(ctx, cred, region, bucket, p.ID); err != nil {
return fmt.Errorf("delete par %s: %w", p.Name, err)
}
mu.Lock()
deleted++
mu.Unlock()
return nil
})
return deleted, err
}