Files
OpenList/server/s3/utils.go
T
Nostalgia 893457cd50 fix(aliyundrive): limit callback concurrency (#3071)
* fix(aliyundrive): limit callback concurrency

- Share proxy callback admission by Aliyun user identity and hold permits for complete response-body lifetimes.
- Retry only verified callback-capacity rejections while preserving direct redirects and server download limiting.
- Map exhausted temporary capacity to S3 SlowDown through the merged OpenListTeam gofakes3 module.
- Cover shared limits, lifecycle release, cancellation, retry classification, and the S3 HTTP response.

Co-authored-by: Codex <267193182+codex@users.noreply.github.com>

# Conflicts:
#	go.mod
#	go.sum
#	server/s3/pager.go

* fix(op): separate redirect and proxy link cache entries

- Include redirect mode in the link cache key for all drivers.

- Cover both redirect-to-proxy and proxy-to-redirect cache reuse.

Co-authored-by: Codex <267193182+codex@users.noreply.github.com>

* fix(proxy): close range bodies before opening next

- make ServeHTTP own each range body and preserve cleanup failures
- remove the aggregate range closer and pass range readers directly
- replace the obsolete callback transport test with focused lifecycle coverage

Co-authored-by: Codex <267193182+codex@users.noreply.github.com>

---------

Co-authored-by: nostalume <nostalucent@gmail.com>
Co-authored-by: Codex <267193182+codex@users.noreply.github.com>
2026-09-24 12:01:50 +08:00

176 lines
4.0 KiB
Go

// Credits: https://pkg.go.dev/github.com/rclone/rclone@v1.65.2/cmd/serve/s3
// Package s3 implements a fake s3 server for openlist
package s3
import (
"context"
"encoding/json"
stderrors "errors"
"strings"
"github.com/OpenListTeam/OpenList/v4/internal/conf"
"github.com/OpenListTeam/OpenList/v4/internal/errs"
"github.com/OpenListTeam/OpenList/v4/internal/fs"
"github.com/OpenListTeam/OpenList/v4/internal/model"
"github.com/OpenListTeam/OpenList/v4/internal/op"
"github.com/OpenListTeam/OpenList/v4/internal/setting"
"github.com/OpenListTeam/gofakes3"
)
type Bucket struct {
Name string `json:"name"`
Path string `json:"path"`
}
func mapBackendError(err error) error {
if stderrors.Is(err, errs.TemporaryCapacity) {
return gofakes3.ErrSlowDown
}
return err
}
const emptyObjectName = "ThisIsAnEmptyFolderInTheS3Bucket"
func getAndParseBuckets() ([]Bucket, error) {
var res []Bucket
err := json.Unmarshal([]byte(setting.GetStr(conf.S3Buckets)), &res)
return res, err
}
func getBucketByName(name string) (Bucket, error) {
buckets, err := getAndParseBuckets()
if err != nil {
return Bucket{}, err
}
for _, b := range buckets {
if b.Name == name {
return b, nil
}
}
return Bucket{}, gofakes3.BucketNotFound(name)
}
func getDirEntries(ctx context.Context, path string) ([]model.Obj, error) {
meta, _ := op.GetNearestMeta(path)
fi, err := fs.Get(context.WithValue(ctx, conf.MetaKey, meta), path, &fs.GetArgs{})
if ctxErr := ctx.Err(); ctxErr != nil {
return nil, ctxErr
}
if errs.IsNotFoundError(err) {
return nil, gofakes3.ErrNoSuchKey
} else if err != nil {
return nil, gofakes3.ErrNoSuchKey
}
if !fi.IsDir() {
return nil, gofakes3.ErrNoSuchKey
}
dirEntries, err := fs.List(context.WithValue(ctx, conf.MetaKey, meta), path, &fs.ListArgs{})
if ctxErr := ctx.Err(); ctxErr != nil {
return nil, ctxErr
}
if err != nil {
return nil, err
}
return dirEntries, nil
}
// func getFileHashByte(node interface{}) []byte {
// b, err := hex.DecodeString(getFileHash(node))
// if err != nil {
// return nil
// }
// return b
// }
func getFileHash(node interface{}) string {
// var o fs.Object
// switch b := node.(type) {
// case vfs.Node:
// fsObj, ok := b.DirEntry().(fs.Object)
// if !ok {
// fs.Debugf("serve s3", "File uploading - reading hash from VFS cache")
// in, err := b.Open(os.O_RDONLY)
// if err != nil {
// return ""
// }
// defer func() {
// _ = in.Close()
// }()
// h, err := hash.NewMultiHasherTypes(hash.NewHashSet(hash.MD5))
// if err != nil {
// return ""
// }
// _, err = io.Copy(h, in)
// if err != nil {
// return ""
// }
// return h.Sums()[hash.MD5]
// }
// o = fsObj
// case fs.Object:
// o = b
// }
// hash, err := o.Hash(context.Background(), hash.MD5)
// if err != nil {
// return ""
// }
// return hash
return ""
}
func prefixParser(p *gofakes3.Prefix) (path, remaining string) {
idx := strings.LastIndexByte(p.Prefix, '/')
if idx < 0 {
return "", p.Prefix
}
return p.Prefix[:idx], p.Prefix[idx+1:]
}
// // FIXME this could be implemented by VFS.MkdirAll()
// func mkdirRecursive(path string, VFS *vfs.VFS) error {
// path = strings.Trim(path, "/")
// dirs := strings.Split(path, "/")
// dir := ""
// for _, d := range dirs {
// dir += "/" + d
// if _, err := VFS.Stat(dir); err != nil {
// err := VFS.Mkdir(dir, 0777)
// if err != nil {
// return err
// }
// }
// }
// return nil
// }
// func rmdirRecursive(p string, VFS *vfs.VFS) {
// dir := path.Dir(p)
// if !strings.ContainsAny(dir, "/\\") {
// // might be bucket(root)
// return
// }
// if _, err := VFS.Stat(dir); err == nil {
// err := VFS.Remove(dir)
// if err != nil {
// return
// }
// rmdirRecursive(dir, VFS)
// }
// }
func authlistResolver() map[string]string {
s3accesskeyid := setting.GetStr(conf.S3AccessKeyId)
s3secretaccesskey := setting.GetStr(conf.S3SecretAccessKey)
if s3accesskeyid == "" && s3secretaccesskey == "" {
return nil
}
authList := make(map[string]string)
authList[s3accesskeyid] = s3secretaccesskey
return authList
}