fix(server/s3): paginate recursive object listings (#2968)

fix(s3): paginate recursive object listings

- stop recursive traversal after filling the requested S3 page
- preserve lexicographic marker ordering and request cancellation
- cover bounded traversal, continuation, prefixes, and cache safety

Co-authored-by: nostalume <nostalucent@gmail.com>
Co-authored-by: Codex <267193182+codex@users.noreply.github.com>
This commit is contained in:
Nostalgia
2026-08-29 01:18:42 +08:00
committed by GitHub
parent 39fcaf488b
commit e0e4de5e82
5 changed files with 413 additions and 86 deletions
+6 -5
View File
@@ -36,13 +36,15 @@ var (
// s3Backend implements the gofacess3.Backend interface to make an S3
// backend for gofakes3
type s3Backend struct {
meta *sync.Map
meta *sync.Map
listDir func(context.Context, string) ([]model.Obj, error)
}
// newBackend creates a new SimpleBucketBackend.
func newBackend() gofakes3.Backend {
return &s3Backend{
meta: new(sync.Map),
meta: new(sync.Map),
listDir: getDirEntries,
}
}
@@ -84,10 +86,9 @@ func (b *s3Backend) ListBucket(ctx context.Context, bucketName string, prefix *g
prefix.HasDelimiter = false
}
response := gofakes3.NewObjectList()
path, remaining := prefixParser(prefix)
err = b.entryListR(bucketPath, path, remaining, prefix.HasDelimiter, response)
response, err := b.listPage(ctx, bucketPath, path, remaining, prefix.HasDelimiter, page)
if err == gofakes3.ErrNoSuchKey {
// AWS just returns an empty list
response = gofakes3.NewObjectList()
@@ -95,7 +96,7 @@ func (b *s3Backend) ListBucket(ctx context.Context, bucketName string, prefix *g
return nil, err
}
return b.pager(response, page)
return response, nil
}
// HeadObject returns the fileinfo for the given object name.
+119 -12
View File
@@ -3,7 +3,10 @@
package s3
import (
"context"
"path"
"slices"
"sort"
"strings"
"time"
@@ -11,16 +14,107 @@ import (
log "github.com/sirupsen/logrus"
)
func (b *s3Backend) entryListR(bucket, fdPath, name string, addPrefix bool, response *gofakes3.ObjectList) error {
fp := path.Join(bucket, fdPath)
// S3 ListObjects responses contain at most 1,000 keys.
const maxListPageKeys int64 = 1000
dirEntries, err := getDirEntries(fp)
type objectPage struct {
result *gofakes3.ObjectList
marker string
maxKeys int64
count int64
lastKey string
}
func newObjectPage(page gofakes3.ListBucketPage) *objectPage {
maxKeys := page.MaxKeys
if maxKeys <= 0 || maxKeys > maxListPageKeys {
maxKeys = maxListPageKeys
}
marker := ""
if page.HasMarker {
marker = page.Marker
}
return &objectPage{
result: gofakes3.NewObjectList(),
marker: marker,
maxKeys: maxKeys,
}
}
func (p *objectPage) addContent(item *gofakes3.Content) bool {
if item.Key <= p.marker {
return true
}
if p.count >= p.maxKeys {
p.result.IsTruncated = true
return false
}
p.result.Add(item)
p.count++
p.lastKey = item.Key
return true
}
func (p *objectPage) addPrefix(prefix string) bool {
if prefix <= p.marker {
return true
}
if p.count >= p.maxKeys {
p.result.IsTruncated = true
return false
}
p.result.AddPrefix(prefix)
p.count++
p.lastKey = prefix
return true
}
func (p *objectPage) finish() *gofakes3.ObjectList {
if p.result.IsTruncated {
p.result.NextMarker = p.lastKey
}
return p.result
}
func (b *s3Backend) listPage(
ctx context.Context,
bucket, fdPath, name string,
addPrefix bool,
page gofakes3.ListBucketPage,
) (*gofakes3.ObjectList, error) {
result := newObjectPage(page)
_, err := b.walkPage(ctx, bucket, fdPath, name, addPrefix, result)
if err != nil {
return err
return nil, err
}
return result.finish(), nil
}
// walkPage returns false after finding an entry beyond the requested page.
func (b *s3Backend) walkPage(
ctx context.Context,
bucket, fdPath, name string,
addPrefix bool,
page *objectPage,
) (bool, error) {
if err := ctx.Err(); err != nil {
return false, err
}
fp := path.Join(bucket, fdPath)
dirEntries, err := b.listDir(ctx, fp)
if err != nil {
return false, err
}
if err := ctx.Err(); err != nil {
return false, err
}
// workaround as s3 can't have empty files in directories, useful in deletions
if len(dirEntries) == 0 {
if !strings.HasPrefix(emptyObjectName, name) {
return true, nil
}
item := &gofakes3.Content{
// Key: gofakes3.URLEncode(path.Join(fdPath, emptyObjectName)),
Key: path.Join(fdPath, emptyObjectName),
@@ -29,11 +123,15 @@ func (b *s3Backend) entryListR(bucket, fdPath, name string, addPrefix bool, resp
Size: 0,
StorageClass: gofakes3.StorageStandard,
}
response.Add(item)
log.Debugf("Adding empty object %s to response", item.Key)
return nil
return page.addContent(item), nil
}
dirEntries = slices.Clone(dirEntries)
sort.Slice(dirEntries, func(i, j int) bool {
return dirEntries[i].GetName() < dirEntries[j].GetName()
})
for _, entry := range dirEntries {
object := entry.GetName()
@@ -47,12 +145,19 @@ func (b *s3Backend) entryListR(bucket, fdPath, name string, addPrefix bool, resp
if entry.IsDir() {
if addPrefix {
// response.AddPrefix(gofakes3.URLEncode(objectPath))
response.AddPrefix(objectPath)
if !page.addPrefix(objectPath) {
return false, nil
}
continue
}
err := b.entryListR(bucket, path.Join(fdPath, object), "", false, response)
if err != nil {
return err
subtreePrefix := objectPath + "/"
// A marker beyond this subtree lets us avoid an upstream directory read.
if subtreePrefix <= page.marker && !strings.HasPrefix(page.marker, subtreePrefix) {
continue
}
keepGoing, err := b.walkPage(ctx, bucket, objectPath, "", false, page)
if err != nil || !keepGoing {
return keepGoing, err
}
} else {
item := &gofakes3.Content{
@@ -63,8 +168,10 @@ func (b *s3Backend) entryListR(bucket, fdPath, name string, addPrefix bool, resp
Size: entry.GetSize(),
StorageClass: gofakes3.StorageStandard,
}
response.Add(item)
if !page.addContent(item) {
return false, nil
}
}
}
return nil
return true, nil
}
+281
View File
@@ -0,0 +1,281 @@
package s3
import (
"context"
"errors"
"fmt"
"sort"
"testing"
"github.com/OpenListTeam/OpenList/v4/internal/model"
"github.com/itsHenry35/gofakes3"
)
func TestListPageBoundsRecursiveTraversal(t *testing.T) {
listCalls := 0
b := &s3Backend{
listDir: func(_ context.Context, dir string) ([]model.Obj, error) {
listCalls++
if dir == "bucket/data" {
entries := make([]model.Obj, 256)
for i := range entries {
entries[i] = &model.Object{Name: fmt.Sprintf("%02x", i), IsFolder: true}
}
return entries, nil
}
return []model.Obj{&model.Object{Name: "pack", Size: 1}}, nil
},
}
got, err := b.listPage(context.Background(), "bucket", "data", "", false, gofakes3.ListBucketPage{MaxKeys: 3})
if err != nil {
t.Fatalf("listPage() error = %v", err)
}
wantKeys := []string{"data/00/pack", "data/01/pack", "data/02/pack"}
if len(got.Contents) != len(wantKeys) {
t.Fatalf("len(Contents) = %d, want %d", len(got.Contents), len(wantKeys))
}
for i, want := range wantKeys {
if got.Contents[i].Key != want {
t.Errorf("Contents[%d].Key = %q, want %q", i, got.Contents[i].Key, want)
}
}
if !got.IsTruncated {
t.Error("IsTruncated = false, want true")
}
if got.NextMarker != wantKeys[len(wantKeys)-1] {
t.Errorf("NextMarker = %q, want %q", got.NextMarker, wantKeys[len(wantKeys)-1])
}
if listCalls > 5 {
t.Errorf("list calls = %d, want at most 5 for a three-key page", listCalls)
}
}
func TestListPageContinuesWithoutGapsOrDuplicates(t *testing.T) {
b := &s3Backend{
listDir: func(_ context.Context, dir string) ([]model.Obj, error) {
if dir == "bucket/data" {
entries := make([]model.Obj, 16)
for i := range entries {
entries[i] = &model.Object{Name: fmt.Sprintf("%02x", 15-i), IsFolder: true}
}
return entries, nil
}
entries := make([]model.Obj, 20)
for i := range entries {
entries[i] = &model.Object{Name: fmt.Sprintf("pack-%02d", 19-i), Size: 1}
}
return entries, nil
},
}
var gotKeys []string
page := gofakes3.ListBucketPage{MaxKeys: 37}
for pageNumber := 0; pageNumber < 10; pageNumber++ {
got, err := b.listPage(context.Background(), "bucket", "data", "", false, page)
if err != nil {
t.Fatalf("listPage() page %d error = %v", pageNumber, err)
}
for _, item := range got.Contents {
gotKeys = append(gotKeys, item.Key)
}
if !got.IsTruncated {
break
}
page.Marker = got.NextMarker
page.HasMarker = true
}
wantKeys := make([]string, 0, 16*20)
for shard := 0; shard < 16; shard++ {
for pack := 0; pack < 20; pack++ {
wantKeys = append(wantKeys, fmt.Sprintf("data/%02x/pack-%02d", shard, pack))
}
}
sort.Strings(wantKeys)
if len(gotKeys) != len(wantKeys) {
t.Fatalf("listed %d keys, want %d", len(gotKeys), len(wantKeys))
}
for i, want := range wantKeys {
if gotKeys[i] != want {
t.Fatalf("key %d = %q, want %q", i, gotKeys[i], want)
}
}
}
func TestListPagePaginatesContentsAndCommonPrefixesTogether(t *testing.T) {
b := &s3Backend{
listDir: func(_ context.Context, _ string) ([]model.Obj, error) {
return []model.Obj{
&model.Object{Name: "d", IsFolder: true},
&model.Object{Name: "c", Size: 1},
&model.Object{Name: "b", IsFolder: true},
&model.Object{Name: "a", Size: 1},
}, nil
},
}
first, err := b.listPage(context.Background(), "bucket", "data", "", true, gofakes3.ListBucketPage{MaxKeys: 2})
if err != nil {
t.Fatalf("first listPage() error = %v", err)
}
if len(first.Contents) != 1 || first.Contents[0].Key != "data/a" {
t.Fatalf("first contents = %#v, want data/a", first.Contents)
}
if len(first.CommonPrefixes) != 1 || first.CommonPrefixes[0].Prefix != "data/b" {
t.Fatalf("first prefixes = %#v, want data/b", first.CommonPrefixes)
}
if !first.IsTruncated || first.NextMarker != "data/b" {
t.Fatalf("first page truncated = %v, next marker = %q", first.IsTruncated, first.NextMarker)
}
second, err := b.listPage(context.Background(), "bucket", "data", "", true, gofakes3.ListBucketPage{
Marker: first.NextMarker,
HasMarker: true,
MaxKeys: 2,
})
if err != nil {
t.Fatalf("second listPage() error = %v", err)
}
if len(second.Contents) != 1 || second.Contents[0].Key != "data/c" {
t.Fatalf("second contents = %#v, want data/c", second.Contents)
}
if len(second.CommonPrefixes) != 1 || second.CommonPrefixes[0].Prefix != "data/d" {
t.Fatalf("second prefixes = %#v, want data/d", second.CommonPrefixes)
}
if second.IsTruncated || second.NextMarker != "" {
t.Fatalf("second page truncated = %v, next marker = %q", second.IsTruncated, second.NextMarker)
}
}
func TestListPageSkipsCompletedSubtrees(t *testing.T) {
var calls []string
b := &s3Backend{
listDir: func(_ context.Context, dir string) ([]model.Obj, error) {
calls = append(calls, dir)
if dir == "bucket/data" {
return []model.Obj{
&model.Object{Name: "c", IsFolder: true},
&model.Object{Name: "b", IsFolder: true},
&model.Object{Name: "a", IsFolder: true},
}, nil
}
return []model.Obj{&model.Object{Name: "pack", Size: 1}}, nil
},
}
got, err := b.listPage(context.Background(), "bucket", "data", "", false, gofakes3.ListBucketPage{
Marker: "data/b/pack",
HasMarker: true,
MaxKeys: 1,
})
if err != nil {
t.Fatalf("listPage() error = %v", err)
}
if len(got.Contents) != 1 || got.Contents[0].Key != "data/c/pack" {
t.Fatalf("contents = %#v, want data/c/pack", got.Contents)
}
for _, call := range calls {
if call == "bucket/data/a" {
t.Fatalf("completed subtree was listed: calls = %v", calls)
}
}
}
func TestListPageUsesS3MaximumAndProbesForTruncation(t *testing.T) {
b := &s3Backend{
listDir: func(_ context.Context, _ string) ([]model.Obj, error) {
entries := make([]model.Obj, maxListPageKeys+1)
for i := range entries {
entries[i] = &model.Object{Name: fmt.Sprintf("pack-%04d", i), Size: 1}
}
return entries, nil
},
}
got, err := b.listPage(context.Background(), "bucket", "data", "", false, gofakes3.ListBucketPage{MaxKeys: 5000})
if err != nil {
t.Fatalf("listPage() error = %v", err)
}
if len(got.Contents) != int(maxListPageKeys) {
t.Fatalf("len(Contents) = %d, want %d", len(got.Contents), maxListPageKeys)
}
if !got.IsTruncated || got.NextMarker != "data/pack-0999" {
t.Fatalf("truncated = %v, next marker = %q", got.IsTruncated, got.NextMarker)
}
}
func TestListPageDoesNotReturnEmptyDirectoryMarkerForDifferentPrefix(t *testing.T) {
b := &s3Backend{
listDir: func(_ context.Context, _ string) ([]model.Obj, error) {
return nil, nil
},
}
got, err := b.listPage(context.Background(), "bucket", "empty", "not-the-marker", false, gofakes3.ListBucketPage{MaxKeys: 10})
if err != nil {
t.Fatalf("listPage() error = %v", err)
}
if len(got.Contents) != 0 {
t.Fatalf("contents = %#v, want empty", got.Contents)
}
}
func TestListPageHonorsCanceledContext(t *testing.T) {
listCalls := 0
b := &s3Backend{
listDir: func(_ context.Context, _ string) ([]model.Obj, error) {
listCalls++
return nil, nil
},
}
ctx, cancel := context.WithCancel(context.Background())
cancel()
_, err := b.listPage(ctx, "bucket", "data", "", false, gofakes3.ListBucketPage{MaxKeys: 10})
if !errors.Is(err, context.Canceled) {
t.Fatalf("listPage() error = %v, want context.Canceled", err)
}
if listCalls != 0 {
t.Fatalf("list calls = %d, want 0", listCalls)
}
}
func TestListPageStopsWhenDirectoryReadCancelsContext(t *testing.T) {
ctx, cancel := context.WithCancel(context.Background())
b := &s3Backend{
listDir: func(_ context.Context, _ string) ([]model.Obj, error) {
cancel()
return []model.Obj{&model.Object{Name: "pack", Size: 1}}, nil
},
}
_, err := b.listPage(ctx, "bucket", "data", "", false, gofakes3.ListBucketPage{MaxKeys: 10})
if !errors.Is(err, context.Canceled) {
t.Fatalf("listPage() error = %v, want context.Canceled", err)
}
}
func TestListPageDoesNotReorderDirectoryCacheEntries(t *testing.T) {
entries := []model.Obj{
&model.Object{Name: "c", Size: 1},
&model.Object{Name: "a", Size: 1},
&model.Object{Name: "b", Size: 1},
}
b := &s3Backend{
listDir: func(_ context.Context, _ string) ([]model.Obj, error) {
return entries, nil
},
}
if _, err := b.listPage(context.Background(), "bucket", "data", "", false, gofakes3.ListBucketPage{MaxKeys: 10}); err != nil {
t.Fatalf("listPage() error = %v", err)
}
wantOrder := []string{"c", "a", "b"}
for i, want := range wantOrder {
if entries[i].GetName() != want {
t.Fatalf("entry %d = %q, want original order %q", i, entries[i].GetName(), want)
}
}
}
-67
View File
@@ -1,67 +0,0 @@
// 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 (
"sort"
"github.com/itsHenry35/gofakes3"
)
// pager splits the object list into multiple pages.
func (db *s3Backend) pager(list *gofakes3.ObjectList, page gofakes3.ListBucketPage) (*gofakes3.ObjectList, error) {
// sort by alphabet
sort.Slice(list.CommonPrefixes, func(i, j int) bool {
return list.CommonPrefixes[i].Prefix < list.CommonPrefixes[j].Prefix
})
// sort by modtime
sort.Slice(list.Contents, func(i, j int) bool {
return list.Contents[i].LastModified.Before(list.Contents[j].LastModified.Time)
})
tokens := page.MaxKeys
if tokens == 0 {
tokens = 1000
}
if page.HasMarker {
for i, obj := range list.Contents {
if obj.Key == page.Marker {
list.Contents = list.Contents[i+1:]
break
}
}
for i, obj := range list.CommonPrefixes {
if obj.Prefix == page.Marker {
list.CommonPrefixes = list.CommonPrefixes[i+1:]
break
}
}
}
response := gofakes3.NewObjectList()
for _, obj := range list.CommonPrefixes {
if tokens <= 0 {
break
}
response.AddPrefix(obj.Prefix)
tokens--
}
for _, obj := range list.Contents {
if tokens <= 0 {
break
}
response.Add(obj)
tokens--
}
if len(list.CommonPrefixes)+len(list.Contents) > int(page.MaxKeys) {
response.IsTruncated = true
if len(response.Contents) > 0 {
response.NextMarker = response.Contents[len(response.Contents)-1].Key
} else {
response.NextMarker = response.CommonPrefixes[len(response.CommonPrefixes)-1].Prefix
}
}
return response, nil
}
+7 -2
View File
@@ -42,10 +42,12 @@ func getBucketByName(name string) (Bucket, error) {
return Bucket{}, gofakes3.BucketNotFound(name)
}
func getDirEntries(path string) ([]model.Obj, error) {
ctx := context.Background()
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 {
@@ -57,6 +59,9 @@ func getDirEntries(path string) ([]model.Obj, error) {
}
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
}