feat(server/s3): support multipart upload

- Implement gofakes3 MultipartBackend (Create/UploadPart/Complete/Abort)
  on s3Backend so multipart parts stream to local temp files instead of
  being buffered in memory
- Track in-progress uploads via a new uploads sync.Map on s3Backend
- Refactor the PutObject body into a reusable putStream helper shared with
  the multipart Complete path
- Validate part ordering, etags and existence, reject short reads, and
  return S3-style multipart etags
- Add end-to-end multipart tests against a real Local driver

Co-authored-by: Codex <267193182+codex@users.noreply.github.com>
This commit is contained in:
MadDogOwner
2026-07-21 00:22:24 +08:00
parent 7d91598f17
commit dd8247e578
3 changed files with 541 additions and 15 deletions
+24 -15
View File
@@ -33,17 +33,21 @@ var (
timeFormat = "Mon, 2 Jan 2006 15:04:05 GMT"
)
// s3Backend implements the gofacess3.Backend interface to make an S3
// backend for gofakes3
// s3Backend implements the gofakes3.Backend interface to make an S3
// backend for gofakes3. It also implements gofakes3.MultipartBackend so that
// multipart uploads are streamed to local temp files part-by-part and
// assembled into storage on completion, instead of being buffered in memory.
type s3Backend struct {
meta *sync.Map
listDir func(context.Context, string) ([]model.Obj, error)
uploads *sync.Map // map[gofakes3.UploadID]*multipartState
}
// newBackend creates a new SimpleBucketBackend.
func newBackend() gofakes3.Backend {
return &s3Backend{
meta: new(sync.Map),
uploads: new(sync.Map),
listDir: getDirEntries,
}
}
@@ -233,9 +237,20 @@ func (b *s3Backend) PutObject(
meta map[string]string,
input io.Reader, size int64,
) (result gofakes3.PutObjectResult, err error) {
return result, b.putStream(ctx, bucketName, objectName, meta, input, size)
}
// putStream stores the given object into the underlying storage. It is shared
// by PutObject and the multipart-upload Complete step so both paths apply the
// same directory creation, metadata and ignore rules.
func (b *s3Backend) putStream(
ctx context.Context, bucketName, objectName string,
meta map[string]string,
input io.Reader, size int64,
) error {
bucket, err := getBucketByName(bucketName)
if err != nil {
return result, err
return err
}
bucketPath := bucket.Path
@@ -261,15 +276,15 @@ func (b *s3Backend) PutObject(
log.Debugf("reqPath: %s not found and objectName contains /, need to makeDir", reqPath)
err = fs.MakeDir(ctx, reqPath)
if err != nil {
return result, errors.WithMessagef(err, "failed to makeDir, reqPath: %s", reqPath)
return errors.WithMessagef(err, "failed to makeDir, reqPath: %s", reqPath)
}
} else {
return result, gofakes3.KeyNotFound(objectName)
return gofakes3.KeyNotFound(objectName)
}
}
if isDir {
return result, nil
return nil
}
var ti time.Time
@@ -295,7 +310,7 @@ func (b *s3Backend) PutObject(
}
// Check if system file should be ignored
if setting.GetBool(conf.IgnoreSystemFiles) && utils.IsSystemFile(obj.Name) {
return result, errs.IgnoredSystemFile
return errs.IgnoredSystemFile
}
stream := &stream.FileStream{
Obj: &obj,
@@ -305,18 +320,12 @@ func (b *s3Backend) PutObject(
err = fs.PutDirectly(ctx, reqPath, stream)
if err != nil {
return result, err
return err
}
// if err := stream.Close(); err != nil {
// // remove file when close error occurred (FsPutErr)
// _ = fs.Remove(ctx, fp)
// return result, err
// }
b.meta.Store(fp, meta)
return result, nil
return nil
}
// DeleteMulti deletes multiple objects in a single request.
+272
View File
@@ -0,0 +1,272 @@
// 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"
"crypto/md5"
"encoding/hex"
"fmt"
"io"
"os"
"path/filepath"
"strings"
"sync"
"time"
"github.com/google/uuid"
"github.com/OpenListTeam/OpenList/v4/internal/conf"
"github.com/OpenListTeam/OpenList/v4/pkg/utils"
"github.com/OpenListTeam/gofakes3"
log "github.com/sirupsen/logrus"
)
// Compile-time assertions that s3Backend implements both the base Backend and
// the optional MultipartBackend interface from gofakes3.
var (
_ gofakes3.Backend = (*s3Backend)(nil)
_ gofakes3.MultipartBackend = (*s3Backend)(nil)
)
// multipartPart records a single uploaded part on disk.
type multipartPart struct {
path string
size int64
md5hex string // unquoted lowercase hex
updated time.Time
}
// multipartState tracks one in-progress multipart upload.
//
// gofakes3 does not serialize multipart operations for a given UploadID, so
// the parts map is guarded by mu. Each part is written to its own file inside
// dir, so concurrent UploadPart calls for different part numbers are safe.
type multipartState struct {
bucket string
object string
meta map[string]string
dir string
created time.Time
mu sync.Mutex
parts map[int]*multipartPart
}
// CreateMultipartUpload begins a new multipart upload. Parts are streamed to
// local temp files so that large uploads do not need to be buffered in memory.
//
// It implements gofakes3.MultipartBackend.
func (b *s3Backend) CreateMultipartUpload(ctx context.Context, bucket, object string, meta map[string]string) (gofakes3.UploadID, error) {
if _, err := getBucketByName(bucket); err != nil {
return "", err
}
tempDir := conf.Conf.TempDir
if tempDir == "" {
tempDir = os.TempDir()
}
dir, err := os.MkdirTemp(tempDir, "s3-multipart-*")
if err != nil {
return "", fmt.Errorf("create multipart upload dir: %w", err)
}
uploadID := gofakes3.UploadID(strings.ReplaceAll(uuid.NewString(), "-", ""))
state := &multipartState{
bucket: bucket,
object: object,
meta: meta,
dir: dir,
created: time.Now(),
parts: map[int]*multipartPart{},
}
b.uploads.Store(uploadID, state)
log.Debugf("s3 multipart: created upload %s for %s/%s", uploadID, bucket, object)
return uploadID, nil
}
// UploadPart writes a single part to disk and returns its (quoted) MD5 etag.
//
// It implements gofakes3.MultipartBackend. The body must contain exactly
// contentLength bytes; a short read is reported as ErrIncompleteBody so a
// truncated client request is never silently stored.
func (b *s3Backend) UploadPart(ctx context.Context, bucket, object string, uploadID gofakes3.UploadID, partNumber int, contentLength int64, body io.Reader) (string, error) {
if partNumber <= 0 || partNumber > gofakes3.MaxUploadPartNumber {
return "", gofakes3.ErrInvalidPart
}
val, ok := b.uploads.Load(uploadID)
if !ok {
return "", gofakes3.ErrNoSuchUpload
}
state := val.(*multipartState)
partPath := filepath.Join(state.dir, fmt.Sprintf("part-%05d", partNumber))
f, err := os.Create(partPath)
if err != nil {
return "", fmt.Errorf("create part file: %w", err)
}
// Remove a half-written file on any failure path.
partFailed := true
defer func() {
if partFailed {
_ = f.Close()
_ = os.Remove(partPath)
}
}()
hash := md5.New()
// io.TeeReader feeds the hasher while the part is streamed to disk, so the
// etag costs no extra pass over the data.
n, err := utils.CopyWithBuffer(io.MultiWriter(f, hash), io.LimitReader(body, contentLength))
if err != nil {
return "", fmt.Errorf("write part %d: %w", partNumber, err)
}
if err := f.Close(); err != nil {
return "", fmt.Errorf("close part %d: %w", partNumber, err)
}
if n != contentLength {
// The client under-delivered (truncated/aborted request). gofakes3 does
// not validate this for streaming backends, so we must.
return "", gofakes3.ErrIncompleteBody
}
md5hex := hex.EncodeToString(hash.Sum(nil))
etag := fmt.Sprintf("%q", md5hex)
state.mu.Lock()
if old := state.parts[partNumber]; old != nil && old.path != partPath {
_ = os.Remove(old.path)
}
state.parts[partNumber] = &multipartPart{
path: partPath,
size: n,
md5hex: md5hex,
updated: time.Now(),
}
state.mu.Unlock()
partFailed = false
log.Debugf("s3 multipart: stored part %d for %s (%d bytes)", partNumber, uploadID, n)
return etag, nil
}
// CompleteMultipartUpload assembles the uploaded parts in ascending part-number
// order and streams the result into storage via the shared putStream path.
//
// It implements gofakes3.MultipartBackend. Part ordering and etags are
// validated against the parts actually received.
func (b *s3Backend) CompleteMultipartUpload(ctx context.Context, bucket, object string, uploadID gofakes3.UploadID, input *gofakes3.CompleteMultipartUploadRequest) (gofakes3.VersionID, string, error) {
val, ok := b.uploads.Load(uploadID)
if !ok {
return "", "", gofakes3.ErrNoSuchUpload
}
state := val.(*multipartState)
if input == nil || len(input.Parts) == 0 {
return "", "", gofakes3.ErrorMessagef(gofakes3.ErrMalformedXML, "complete multipart upload has no parts")
}
// S3 requires the parts in a CompleteMultipartUpload request to be listed
// in ascending part-number order.
for i := 1; i < len(input.Parts); i++ {
if input.Parts[i].PartNumber <= input.Parts[i-1].PartNumber {
return "", "", gofakes3.ErrInvalidPartOrder
}
}
// Validate every requested part exists with a matching etag, and collect
// them in the order requested by the client (which is sorted ascending).
state.mu.Lock()
ordered := make([]*multipartPart, 0, len(input.Parts))
var concat []byte
for _, p := range input.Parts {
stored := state.parts[p.PartNumber]
if stored == nil {
state.mu.Unlock()
return "", "", gofakes3.ErrorMessagef(gofakes3.ErrInvalidPart, "unexpected part number %d in complete request", p.PartNumber)
}
if strings.Trim(p.ETag, "\"") != stored.md5hex {
state.mu.Unlock()
return "", "", gofakes3.ErrorMessagef(gofakes3.ErrInvalidPart, "unexpected part etag for number %d in complete request", p.PartNumber)
}
ordered = append(ordered, stored)
// S3 multipart etag = hex(md5(concat(part_md5_digests)))-N
concat = append(concat, stored.md5Bytes()...)
}
// Hold the lock until the part files are opened so an abort racing with
// complete cannot delete them out from under us.
readers := make([]io.Reader, 0, len(ordered))
closers := make([]io.Closer, 0, len(ordered))
var total int64
for _, part := range ordered {
f, err := os.Open(part.path)
if err != nil {
for _, c := range closers {
_ = c.Close()
}
state.mu.Unlock()
return "", "", fmt.Errorf("open part %s: %w", part.path, err)
}
readers = append(readers, f)
closers = append(closers, f)
total += part.size
}
state.mu.Unlock()
combined := utils.NewReadCloser(io.MultiReader(readers...), func() error {
var firstErr error
for _, c := range closers {
if err := c.Close(); err != nil && firstErr == nil {
firstErr = err
}
}
return firstErr
})
err := b.putStream(ctx, bucket, object, state.meta, combined, total)
_ = combined.Close()
if err != nil {
// Leave the upload in place so the client may retry completion, per
// the gofakes3 MultipartBackend contract.
return "", "", err
}
// Success: drop bookkeeping and clean up part files.
b.removeUpload(uploadID)
sum := md5.Sum(concat)
etag := fmt.Sprintf("%q", fmt.Sprintf("%s-%d", hex.EncodeToString(sum[:]), len(ordered)))
log.Debugf("s3 multipart: completed upload %s -> %s/%s (%d bytes)", uploadID, bucket, object, total)
return "", etag, nil
}
// AbortMultipartUpload discards an in-progress upload and its parts.
//
// It implements gofakes3.MultipartBackend and is idempotent: aborting an
// unknown upload succeeds so retries do not fail.
func (b *s3Backend) AbortMultipartUpload(ctx context.Context, bucket, object string, uploadID gofakes3.UploadID) error {
b.removeUpload(uploadID)
return nil
}
// removeUpload deletes the upload's temp directory and drops its bookkeeping.
// Missing uploads are ignored to keep abort/complete idempotent.
func (b *s3Backend) removeUpload(uploadID gofakes3.UploadID) {
val, ok := b.uploads.LoadAndDelete(uploadID)
if !ok {
return
}
state := val.(*multipartState)
if state.dir != "" {
if err := os.RemoveAll(state.dir); err != nil {
log.Warnf("s3 multipart: failed to clean up %s: %v", state.dir, err)
}
}
}
// md5Bytes returns the raw 16-byte MD5 digest of the part.
func (p *multipartPart) md5Bytes() []byte {
b, _ := hex.DecodeString(p.md5hex)
return b
}
+245
View File
@@ -0,0 +1,245 @@
package s3
import (
"bytes"
"context"
"errors"
"os"
"path/filepath"
"strings"
"testing"
_ "github.com/OpenListTeam/OpenList/v4/drivers/local"
"github.com/OpenListTeam/OpenList/v4/internal/conf"
"github.com/OpenListTeam/OpenList/v4/internal/db"
"github.com/OpenListTeam/OpenList/v4/internal/model"
"github.com/OpenListTeam/OpenList/v4/internal/op"
"github.com/OpenListTeam/gofakes3"
"github.com/glebarez/sqlite"
"gorm.io/gorm"
)
func init() {
dataDir, err := os.MkdirTemp("", "openlist-s3-mp-*")
if err != nil {
panic(err)
}
conf.Conf = conf.DefaultConfig(dataDir)
if err := os.MkdirAll(conf.Conf.TempDir, 0o755); err != nil {
panic("mkdir temp dir: " + err.Error())
}
dB, err := gorm.Open(sqlite.Open("file::memory:?cache=shared"), &gorm.Config{})
if err != nil {
panic("failed to connect database: " + err.Error())
}
db.Init(dB)
}
// s3ErrorCode extracts the gofakes3 ErrorCode from an error returned by the
// MultipartBackend methods.
func s3ErrorCode(err error) gofakes3.ErrorCode {
if err == nil {
return gofakes3.ErrNone
}
var s3err interface{ ErrorCode() gofakes3.ErrorCode }
if errors.As(err, &s3err) {
return s3err.ErrorCode()
}
return gofakes3.ErrNone
}
// setupMultipartBackend prepares a Local storage mounted at /mpbucket and an
// s3Backend with an "mp" bucket pointing at it. It returns the backend, the
// local root directory on disk, and a cleanup function.
func setupMultipartBackend(t *testing.T) (*s3Backend, string) {
t.Helper()
ctx := context.Background()
// Unique mount path and bucket per test: the in-memory sqlite is shared
// across tests in this package, so a fixed mount path would clash.
mount := "/" + sanitizeTestName(t.Name())
bucket := "mp"
localRoot, err := os.MkdirTemp("", "openlist-s3-local-*")
if err != nil {
t.Fatalf("mkdir local root: %v", err)
}
t.Cleanup(func() { _ = os.RemoveAll(localRoot) })
_, err = op.CreateStorage(ctx, model.Storage{
Driver: "Local",
MountPath: mount,
Addition: `{"root_folder_path":"` + localRoot + `","thumbnail":false}`,
})
if err != nil {
t.Fatalf("create local storage: %+v", err)
}
if err := op.SaveSettingItem(&model.SettingItem{
Key: conf.S3Buckets,
Value: `[{"name":"` + bucket + `","path":"` + mount + `"}]`,
}); err != nil {
t.Fatalf("save s3 buckets setting: %+v", err)
}
return newBackend().(*s3Backend), localRoot
}
func sanitizeTestName(name string) string {
r := strings.NewReplacer("/", "_", " ", "_")
return r.Replace(name)
}
func TestMultipartUploadEndToEnd(t *testing.T) {
ctx := context.Background()
b, localRoot := setupMultipartBackend(t)
meta := map[string]string{"Content-Type": "text/plain"}
uploadID, err := b.CreateMultipartUpload(ctx, "mp", "dir/hello.txt", meta)
if err != nil {
t.Fatalf("CreateMultipartUpload: %+v", err)
}
if uploadID == "" {
t.Fatal("empty upload id")
}
part := func(n int, body string) string {
t.Helper()
etag, err := b.UploadPart(ctx, "mp", "dir/hello.txt", uploadID, n, int64(len(body)), strings.NewReader(body))
if err != nil {
t.Fatalf("UploadPart %d: %+v", n, err)
}
return etag
}
etag1 := part(1, "Hello, ")
etag2 := part(2, "multipart ")
etag3 := part(3, "world!")
// Re-uploading the same part number overwrites it and returns a fresh etag.
if e := part(2, "multipart "); e != etag2 {
t.Fatalf("re-upload part 2 etag = %q, want %q", e, etag2)
}
// Short read (body smaller than declared Content-Length) must fail.
_, shortErr := b.UploadPart(ctx, "mp", "dir/hello.txt", uploadID, 4, 10, strings.NewReader("abc"))
shortErr = s3ErrorCode(shortErr)
if shortErr != gofakes3.ErrIncompleteBody {
t.Fatalf("short read error = %v, want IncompleteBody", shortErr)
}
// Unknown upload id.
_, uerr := b.UploadPart(ctx, "mp", "dir/hello.txt", "does-not-exist", 1, 1, strings.NewReader("x"))
if code := s3ErrorCode(uerr); code != gofakes3.ErrNoSuchUpload {
t.Fatalf("unknown upload error = %v, want NoSuchUpload", code)
}
// Out-of-range part number.
_, perr := b.UploadPart(ctx, "mp", "dir/hello.txt", uploadID, 0, 1, strings.NewReader("x"))
if code := s3ErrorCode(perr); code != gofakes3.ErrInvalidPart {
t.Fatalf("part 0 error = %v, want InvalidPart", code)
}
// Parts out of order.
_, _, err = b.CompleteMultipartUpload(ctx, "mp", "dir/hello.txt", uploadID, &gofakes3.CompleteMultipartUploadRequest{
Parts: []gofakes3.CompletedPart{
{PartNumber: 2, ETag: etag2},
{PartNumber: 1, ETag: etag1},
},
})
if code := s3ErrorCode(err); code != gofakes3.ErrInvalidPartOrder {
t.Fatalf("out-of-order complete error = %v, want InvalidPartOrder", code)
}
// Wrong etag.
_, _, err = b.CompleteMultipartUpload(ctx, "mp", "dir/hello.txt", uploadID, &gofakes3.CompleteMultipartUploadRequest{
Parts: []gofakes3.CompletedPart{
{PartNumber: 1, ETag: etag1},
{PartNumber: 2, ETag: `"deadbeef"`},
{PartNumber: 3, ETag: etag3},
},
})
if code := s3ErrorCode(err); code != gofakes3.ErrInvalidPart {
t.Fatalf("wrong etag complete error = %v, want InvalidPart", code)
}
// Missing part number.
_, _, err = b.CompleteMultipartUpload(ctx, "mp", "dir/hello.txt", uploadID, &gofakes3.CompleteMultipartUploadRequest{
Parts: []gofakes3.CompletedPart{
{PartNumber: 1, ETag: etag1},
{PartNumber: 99, ETag: etag3},
},
})
if code := s3ErrorCode(err); code != gofakes3.ErrInvalidPart {
t.Fatalf("missing part complete error = %v, want InvalidPart", code)
}
// The failed completes must leave the upload available for retry.
if _, ok := b.uploads.Load(uploadID); !ok {
t.Fatal("upload was removed after a failed complete")
}
// Successful complete.
_, etag, err := b.CompleteMultipartUpload(ctx, "mp", "dir/hello.txt", uploadID, &gofakes3.CompleteMultipartUploadRequest{
Parts: []gofakes3.CompletedPart{
{PartNumber: 1, ETag: etag1},
{PartNumber: 2, ETag: etag2},
{PartNumber: 3, ETag: etag3},
},
})
if err != nil {
t.Fatalf("CompleteMultipartUpload: %+v", err)
}
if !strings.HasSuffix(etag, `-3"`) || !strings.HasPrefix(etag, `"`) {
t.Fatalf("complete etag = %q, want a quoted \"<hex>-3\" multipart etag", etag)
}
// The object must exist on disk with the concatenated content.
got, err := os.ReadFile(filepath.Join(localRoot, "dir", "hello.txt"))
if err != nil {
t.Fatalf("read resulting file: %+v", err)
}
want := []byte("Hello, multipart world!")
if !bytes.Equal(got, want) {
t.Fatalf("resulting file content = %q, want %q", got, want)
}
// Bookkeeping and temp files must be cleaned up on success.
if _, ok := b.uploads.Load(uploadID); ok {
t.Fatal("upload still tracked after successful complete")
}
}
func TestMultipartAbort(t *testing.T) {
ctx := context.Background()
b, _ := setupMultipartBackend(t)
uploadID, err := b.CreateMultipartUpload(ctx, "mp", "abort.txt", nil)
if err != nil {
t.Fatalf("CreateMultipartUpload: %+v", err)
}
if _, err := b.UploadPart(ctx, "mp", "abort.txt", uploadID, 1, 3, strings.NewReader("abc")); err != nil {
t.Fatalf("UploadPart: %+v", err)
}
state, _ := b.uploads.Load(uploadID)
dir := state.(*multipartState).dir
if _, err := os.Stat(dir); err != nil {
t.Fatalf("temp dir missing before abort: %+v", err)
}
if err := b.AbortMultipartUpload(ctx, "mp", "abort.txt", uploadID); err != nil {
t.Fatalf("AbortMultipartUpload: %+v", err)
}
if _, ok := b.uploads.Load(uploadID); ok {
t.Fatal("upload still tracked after abort")
}
if _, err := os.Stat(dir); !os.IsNotExist(err) {
t.Fatalf("temp dir still exists after abort (err=%v)", err)
}
// Aborting an unknown upload must be idempotent.
if err := b.AbortMultipartUpload(ctx, "mp", "abort.txt", "nope"); err != nil {
t.Fatalf("abort unknown upload returned error: %+v", err)
}
}