350 lines
12 KiB
Go
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
|
|
}
|