From 04697cc02ade35bed7e31cbb39c78d3dcdeb61ee Mon Sep 17 00:00:00 2001 From: Nick Craig-Wood Date: Thu, 3 Sep 2026 11:26:38 +0100 Subject: [PATCH] serve s3: list objects lazily so paging a deep hierarchy is fast - fixes #9855 Listing a prefix with no delimiter walked the entire subtree below it into memory and only then sliced out the requested page, so every page of a listing cost a full traversal of the tree. Walk the tree lazily instead, stopping as soon as the page is full and skipping the subtrees an earlier page already returned. A page now costs a number of directory reads proportional to the keys it returns rather than to the size of the subtree, and ETags are computed only for the objects that are actually returned. Entries are emitted in the order their keys have in a flat keyspace, with a directory sorting as if it carried its trailing slash, so that "a.txt" comes before "a/b" as it does in a real S3 bucket. Without this a resumed listing would silently skip keys across a page boundary. --- cmd/serve/s3/backend.go | 4 +- cmd/serve/s3/list.go | 181 +++++++++++++++++++++++-- cmd/serve/s3/list_test.go | 263 +++++++++++++++++++++++++++++++++++++ cmd/serve/s3/pager.go | 66 ---------- cmd/serve/s3/pager_test.go | 32 ----- cmd/serve/s3/s3_test.go | 108 +++++++++++++++ 6 files changed, 540 insertions(+), 114 deletions(-) create mode 100644 cmd/serve/s3/list_test.go delete mode 100644 cmd/serve/s3/pager.go delete mode 100644 cmd/serve/s3/pager_test.go diff --git a/cmd/serve/s3/backend.go b/cmd/serve/s3/backend.go index 337cff5ea..6e5eb92bf 100644 --- a/cmd/serve/s3/backend.go +++ b/cmd/serve/s3/backend.go @@ -111,7 +111,7 @@ func (b *s3Backend) ListBucket(ctx context.Context, bucket string, prefix *gofak response := gofakes3.NewObjectList() path, remaining := prefixParser(prefix) - err = b.entryListR(_vfs, bucket, path, remaining, prefix.HasDelimiter, response) + err = b.listPage(_vfs, bucket, path, remaining, prefix.HasDelimiter, page, response) if err == gofakes3.ErrNoSuchKey { // AWS just returns an empty list response = gofakes3.NewObjectList() @@ -119,7 +119,7 @@ func (b *s3Backend) ListBucket(ctx context.Context, bucket string, prefix *gofak return nil, err } - return b.pager(response, page) + return response, nil } // formatHeaderTime makes an timestamp which is the same as that used by AWS. diff --git a/cmd/serve/s3/list.go b/cmd/serve/s3/list.go index 67e9f21df..3f7a91572 100644 --- a/cmd/serve/s3/list.go +++ b/cmd/serve/s3/list.go @@ -1,10 +1,13 @@ package s3 import ( + "errors" "path" + "slices" "strings" "github.com/rclone/gofakes3" + "github.com/rclone/rclone/fs" "github.com/rclone/rclone/vfs" ) @@ -13,17 +16,139 @@ import ( // (rclone v1.75); leftovers from an older server are still hidden. const legacyMultipartUploadPrefix = ".rclone_multipart_upload_" -func (b *s3Backend) entryListR(_vfs *vfs.VFS, bucketName, fdPath, name string, addPrefix bool, response *gofakes3.ObjectList) error { - fp, err := bucketDirPath(bucketName, fdPath) +// errPageFull is returned by the listing walk once the page has been filled +// and one further key has been seen, so the walk can unwind early. It never +// escapes ListBucket. +var errPageFull = errors.New("listing page full") + +// lister walks a directory tree emitting one page of an S3 listing. +// +// The walk emits keys in flat S3 order and stops as soon as the page is full, +// so the cost of a page is proportional to the keys it returns rather than to +// the size of the subtree below the prefix. +type lister struct { + b *s3Backend + vfs *vfs.VFS + bucket string + marker string // keys at or before this one have already been returned + hasMarker bool + max int // number of keys wanted in this page + response *gofakes3.ObjectList + keys int // keys added to the page so far + lastKey string // the last key added to the page + dirsRead int // directories read, for debug logging +} + +// newLister makes a lister which fills response with at most one page of the +// listing of bucket, as described by page. +func newLister(b *s3Backend, _vfs *vfs.VFS, bucket string, page gofakes3.ListBucketPage, response *gofakes3.ObjectList) *lister { + max := int(page.MaxKeys) + if max <= 0 { + // A missing or zero max-keys means the client didn't ask for a + // particular page size, so use the S3 maximum. + max = 1000 + } + return &lister{ + b: b, + vfs: _vfs, + bucket: bucket, + marker: page.Marker, + hasMarker: page.HasMarker, + max: max, + response: response, + } +} + +// truncate marks the page as truncated after the last key added to it, so +// that the next page resumes from there, and returns errPageFull. +func (l *lister) truncate() error { + l.response.IsTruncated = true + l.response.NextMarker = l.lastKey + return errPageFull +} + +// sortKey returns the S3 key an entry contributes to the listing. +// +// Directories sort as if they carried their trailing slash, which is what +// makes a depth first walk emit keys in the same order as the flat keyspace +// of a real S3 bucket. +func sortKey(entry vfs.Node) string { + if entry.IsDir() { + return entry.Name() + "/" + } + return entry.Name() +} + +// skipKey reports whether key has already been returned by an earlier page. +func (l *lister) skipKey(key string) bool { + return l.hasMarker && key <= l.marker +} + +// skipDir reports whether every key under the directory prefix p was returned +// by an earlier page, so the subtree need not be read at all. +// +// A marker inside the subtree has p as a prefix and is never skipped: that +// subtree is descended and filtered key by key. +func (l *lister) skipDir(p string) bool { + return l.hasMarker && l.marker > p && !strings.HasPrefix(l.marker, p) +} + +// addCommonPrefix adds a common prefix to the page, returning errPageFull if +// the page is now complete. +func (l *lister) addCommonPrefix(prefix string) error { + if l.keys >= l.max { + return l.truncate() + } + l.response.AddPrefix(prefix) + l.keys++ + l.lastKey = prefix + return nil +} + +// addObject adds an object to the page, returning errPageFull if the page is +// now complete. +func (l *lister) addObject(key string, entry vfs.Node) error { + if l.keys >= l.max { + return l.truncate() + } + l.response.Add(&gofakes3.Content{ + Key: key, + LastModified: gofakes3.NewContentTime(entry.ModTime()), + ETag: getFileHash(entry, l.b.s.etagHashType), + Size: entry.Size(), + StorageClass: gofakes3.StorageStandard, + }) + l.keys++ + l.lastKey = key + return nil +} + +// list walks the directory fdPath, adding the entries whose leaf name starts +// with name to the page. +// +// If addPrefix is set, subdirectories are reported as common prefixes, +// otherwise they are descended into. +// +// It returns errPageFull once the page is complete, and gofakes3.ErrNoSuchKey +// if fdPath is not a directory. +func (l *lister) list(fdPath, name string, addPrefix bool) error { + fp, err := bucketDirPath(l.bucket, fdPath) if err != nil { // A listing prefix that can't be represented as a path matches nothing. return gofakes3.ErrNoSuchKey } - dirEntries, err := getDirEntries(fp, _vfs) + dirEntries, err := getDirEntries(fp, l.vfs) if err != nil { return err } + l.dirsRead++ + + // Emit the entries in the order their keys have in a flat keyspace, + // which is not the plain name order getDirEntries returns. + slices.SortFunc(dirEntries, func(a, b vfs.Node) int { + return strings.Compare(sortKey(a), sortKey(b)) + }) for _, entry := range dirEntries { object := entry.Name() @@ -41,25 +166,53 @@ func (b *s3Backend) entryListR(_vfs *vfs.VFS, bucketName, fdPath, name string, a } if entry.IsDir() { + prefixWithTrailingSlash := objectPath + "/" if addPrefix { - prefixWithTrailingSlash := objectPath + "/" - response.AddPrefix(prefixWithTrailingSlash) + if l.skipKey(prefixWithTrailingSlash) { + continue + } + if err := l.addCommonPrefix(prefixWithTrailingSlash); err != nil { + return err + } continue } - err := b.entryListR(_vfs, bucketName, path.Join(fdPath, object), "", false, response) - if err != nil { + if l.skipDir(prefixWithTrailingSlash) { + continue + } + err := l.list(objectPath, "", false) + if errors.Is(err, gofakes3.ErrNoSuchKey) { + // The directory went away while we were listing it, so + // there is nothing below it to report. + continue + } else if err != nil { return err } } else { - item := &gofakes3.Content{ - Key: objectPath, - LastModified: gofakes3.NewContentTime(entry.ModTime()), - ETag: getFileHash(entry, b.s.etagHashType), - Size: entry.Size(), - StorageClass: gofakes3.StorageStandard, + if l.skipKey(objectPath) { + continue + } + if err := l.addObject(objectPath, entry); err != nil { + return err } - response.Add(item) } } return nil } + +// listPage fills response with one page of the listing of the directory +// fdPath in bucket, as described by page. +// +// Only entries whose leaf name starts with name are listed. If addPrefix is +// set, subdirectories are reported as common prefixes, otherwise the whole +// subtree below fdPath is listed. +// +// It returns gofakes3.ErrNoSuchKey if fdPath is not a directory. +func (b *s3Backend) listPage(_vfs *vfs.VFS, bucket, fdPath, name string, addPrefix bool, page gofakes3.ListBucketPage, response *gofakes3.ObjectList) error { + l := newLister(b, _vfs, bucket, page, response) + err := l.list(fdPath, name, addPrefix) + if err != nil && !errors.Is(err, errPageFull) { + return err + } + fs.Debugf("serve s3", "Listed %d keys from %d directories (truncated=%v)", l.keys, l.dirsRead, response.IsTruncated) + return nil +} diff --git a/cmd/serve/s3/list_test.go b/cmd/serve/s3/list_test.go new file mode 100644 index 000000000..4054819b3 --- /dev/null +++ b/cmd/serve/s3/list_test.go @@ -0,0 +1,263 @@ +package s3 + +import ( + "context" + "errors" + "fmt" + "os" + "path/filepath" + "slices" + "testing" + + "github.com/rclone/gofakes3" + _ "github.com/rclone/rclone/backend/local" + "github.com/rclone/rclone/cmd/serve/proxy" + "github.com/rclone/rclone/fs" + "github.com/rclone/rclone/fstest" + "github.com/rclone/rclone/vfs/vfscommon" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// newListBackend serves a temporary directory populated with the given files, +// whose paths are relative to the serve root and so start with a bucket name. +func newListBackend(t *testing.T, files ...string) *s3Backend { + fstest.Initialise() + ctx := context.Background() + root := t.TempDir() + for _, file := range files { + p := filepath.Join(root, filepath.FromSlash(file)) + require.NoError(t, os.MkdirAll(filepath.Dir(p), 0777)) + require.NoError(t, os.WriteFile(p, []byte(file), 0666)) + } + + f, err := fs.NewFs(ctx, root) + require.NoError(t, err) + + opt := Opt + opt.HTTP.ListenAddr = []string{endpoint} + w, err := newServer(ctx, f, &opt, &vfscommon.Opt, &proxy.Opt) + require.NoError(t, err) + + // Use the server's own backend so that a test which drives the server + // and one which calls the backend directly share the same state. + return w.backend +} + +// listPrefix makes a listing prefix, with a delimiter if delimiter is set. +func listPrefix(prefix, delimiter string) *gofakes3.Prefix { + return &gofakes3.Prefix{ + HasPrefix: prefix != "", + Prefix: prefix, + HasDelimiter: delimiter != "", + Delimiter: delimiter, + } +} + +// listOnce returns the keys of one page of a listing, common prefixes and +// object keys merged into the single sorted stream S3 clients see. +func listOnce(t *testing.T, b *s3Backend, prefix *gofakes3.Prefix, page gofakes3.ListBucketPage) (keys []string, list *gofakes3.ObjectList) { + list, err := b.ListBucket(context.Background(), "bucket", prefix, page) + require.NoError(t, err) + for _, commonPrefix := range list.CommonPrefixes { + keys = append(keys, commonPrefix.Prefix) + } + for _, object := range list.Contents { + keys = append(keys, object.Key) + } + slices.Sort(keys) + return keys, list +} + +// listPages pages through a whole listing maxKeys at a time, the way an S3 +// client does, and returns the keys of every page concatenated. +func listPages(t *testing.T, b *s3Backend, prefix *gofakes3.Prefix, maxKeys int) (keys []string) { + page := gofakes3.ListBucketPage{MaxKeys: int64(maxKeys)} + for pages := 0; ; pages++ { + require.Less(t, pages, 1000, "listing did not terminate") + pageKeys, list := listOnce(t, b, prefix, page) + assert.LessOrEqual(t, len(pageKeys), maxKeys, "page has more than max-keys keys") + keys = append(keys, pageKeys...) + if !list.IsTruncated { + return keys + } + require.NotEmpty(t, list.NextMarker, "truncated page has no next marker") + require.Equal(t, list.NextMarker, pageKeys[len(pageKeys)-1], "next marker is not the last key of the page") + page.Marker, page.HasMarker = list.NextMarker, true + } +} + +// TestListFlatKeyOrder checks that keys come back in the order of a flat S3 +// keyspace rather than in the order a directory walk finds them: "a.txt" +// sorts before "a/b" because '.' < '/'. +func TestListFlatKeyOrder(t *testing.T) { + b := newListBackend(t, "bucket/a/b", "bucket/a.txt", "bucket/a0") + + keys, list := listOnce(t, b, listPrefix("", ""), gofakes3.ListBucketPage{MaxKeys: 10}) + assert.Equal(t, []string{"a.txt", "a/b", "a0"}, keys) + assert.False(t, list.IsTruncated) + + // Paging must not skip a key across the "a.txt"/"a/b" boundary + assert.Equal(t, []string{"a.txt", "a/b", "a0"}, listPages(t, b, listPrefix("", ""), 1)) +} + +// TestListPaging checks that paging through a listing returns every key +// exactly once, whatever the page size. +func TestListPaging(t *testing.T) { + var files []string + for dir := range 5 { + for file := range 4 { + files = append(files, fmt.Sprintf("bucket/dir%d/file%d", dir, file)) + } + files = append(files, fmt.Sprintf("bucket/dir%d.txt", dir)) + } + b := newListBackend(t, files...) + + all, list := listOnce(t, b, listPrefix("", ""), gofakes3.ListBucketPage{MaxKeys: 1000}) + require.False(t, list.IsTruncated) + require.Len(t, all, 25) + assert.True(t, slices.IsSorted(all), "unpaged listing is not sorted: %q", all) + + for _, maxKeys := range []int{1, 2, 3, 7, 24, 25, 26} { + t.Run(fmt.Sprintf("MaxKeys%d", maxKeys), func(t *testing.T) { + assert.Equal(t, all, listPages(t, b, listPrefix("", ""), maxKeys)) + }) + } +} + +// TestListPagingWithPrefix checks paging a listing which is filtered by a +// prefix that is not a whole path segment. +func TestListPagingWithPrefix(t *testing.T) { + b := newListBackend(t, + "bucket/apple/1", "bucket/apple/2", "bucket/apricot/1", + "bucket/apt.txt", "bucket/banana/1", + ) + prefix := listPrefix("ap", "") + + all, _ := listOnce(t, b, prefix, gofakes3.ListBucketPage{MaxKeys: 1000}) + assert.Equal(t, []string{"apple/1", "apple/2", "apricot/1", "apt.txt"}, all) + + for _, maxKeys := range []int{1, 2, 3} { + t.Run(fmt.Sprintf("MaxKeys%d", maxKeys), func(t *testing.T) { + assert.Equal(t, all, listPages(t, b, prefix, maxKeys)) + }) + } +} + +// TestListDelimiter checks that a delimited listing pages through common +// prefixes and object keys as one merged stream. +func TestListDelimiter(t *testing.T) { + b := newListBackend(t, + "bucket/a/1", "bucket/b.txt", "bucket/c/1", "bucket/d.txt", "bucket/e/1", + ) + prefix := listPrefix("", "/") + + all, _ := listOnce(t, b, prefix, gofakes3.ListBucketPage{MaxKeys: 1000}) + assert.Equal(t, []string{"a/", "b.txt", "c/", "d.txt", "e/"}, all) + + for _, maxKeys := range []int{1, 2, 3, 4} { + t.Run(fmt.Sprintf("MaxKeys%d", maxKeys), func(t *testing.T) { + assert.Equal(t, all, listPages(t, b, prefix, maxKeys)) + }) + } +} + +// countDirs lists one page directly with a lister so that the number of +// directories the walk had to read can be checked. +func countDirs(t *testing.T, b *s3Backend, page gofakes3.ListBucketPage) (dirsRead int, list *gofakes3.ObjectList) { + _vfs, err := b.s.getVFS(context.Background()) + require.NoError(t, err) + list = gofakes3.NewObjectList() + l := newLister(b, _vfs, "bucket", page, list) + err = l.list("", "", false) + if err != nil && !errors.Is(err, errPageFull) { + require.NoError(t, err) + } + return l.dirsRead, list +} + +// TestListReadsOnlyTheDirectoriesItNeeds checks that a page costs a number of +// directory reads proportional to the keys it returns, not to the size of the +// subtree, and that resuming from a marker does not re-read the subtrees +// earlier pages already returned. +func TestListReadsOnlyTheDirectoriesItNeeds(t *testing.T) { + const dirs = 20 + var files []string + for dir := range dirs { + for file := range 5 { + files = append(files, fmt.Sprintf("bucket/dir%02d/file%d", dir, file)) + } + } + b := newListBackend(t, files...) + + // A full listing has to read the bucket and every directory in it + whole, list := countDirs(t, b, gofakes3.ListBucketPage{MaxKeys: 1000}) + require.False(t, list.IsTruncated) + require.Len(t, list.Contents, dirs*5) + assert.Equal(t, dirs+1, whole) + + // Each page of 5 reads the bucket, the directory holding the marker, + // the directory it returns and the one holding the following key + page := gofakes3.ListBucketPage{MaxKeys: 5} + for pages := range dirs { + dirsRead, list := countDirs(t, b, page) + assert.LessOrEqual(t, dirsRead, 4, "page %d read too many directories", pages) + assert.Less(t, dirsRead, whole, "page %d walked the whole tree", pages) + require.Len(t, list.Contents, 5) + assert.Equal(t, pages == dirs-1, !list.IsTruncated) + page.Marker, page.HasMarker = list.NextMarker, true + } +} + +// TestListHidesTemporaryObjects checks that the temporary objects of +// in-progress uploads stay hidden and do not use up a page. +func TestListHidesTemporaryObjects(t *testing.T) { + b := newListBackend(t, + "bucket/a.txt", + "bucket/"+putObjectPrefix+"hidden", + "bucket/"+legacyMultipartUploadPrefix+"hidden", + "bucket/z.txt", + ) + + keys, list := listOnce(t, b, listPrefix("", ""), gofakes3.ListBucketPage{MaxKeys: 2}) + assert.Equal(t, []string{"a.txt", "z.txt"}, keys) + assert.False(t, list.IsTruncated) +} + +// TestListSortKey checks the key a directory entry sorts under. +func TestListSortKey(t *testing.T) { + b := newListBackend(t, "bucket/dir/file", "bucket/file") + _vfs, err := b.s.getVFS(context.Background()) + require.NoError(t, err) + + entries, err := getDirEntries("bucket", _vfs) + require.NoError(t, err) + got := map[string]string{} + for _, entry := range entries { + got[entry.Name()] = sortKey(entry) + } + assert.Equal(t, map[string]string{"dir": "dir/", "file": "file"}, got) +} + +// TestListSkipDir checks which subtrees a marker allows to be skipped unread. +func TestListSkipDir(t *testing.T) { + for _, test := range []struct { + marker string + dir string + want bool + }{ + {marker: "", dir: "a/", want: false}, // no marker, list everything + {marker: "a", dir: "b/", want: false}, // marker before the subtree + {marker: "b/9", dir: "b/", want: false}, // marker inside the subtree + {marker: "b/", dir: "b/", want: false}, // marker is the subtree + {marker: "b0", dir: "b/", want: true}, // marker past the subtree + {marker: "c", dir: "b/", want: true}, // marker past the subtree + {marker: "b.txt", dir: "b/", want: false}, // '.' sorts before '/' + {marker: "b/a", dir: "b/c/", want: false}, // marker before the subtree + {marker: "b/c/9", dir: "b/c/", want: false}, // marker inside the subtree + {marker: "b/d", dir: "b/c/", want: true}, // marker past the subtree + } { + l := &lister{marker: test.marker, hasMarker: test.marker != ""} + assert.Equal(t, test.want, l.skipDir(test.dir), "marker %q dir %q", test.marker, test.dir) + } +} diff --git a/cmd/serve/s3/pager.go b/cmd/serve/s3/pager.go deleted file mode 100644 index 94c8e3cba..000000000 --- a/cmd/serve/s3/pager.go +++ /dev/null @@ -1,66 +0,0 @@ -// Package s3 implements a fake s3 server for rclone -package s3 - -import ( - "sort" - - "github.com/rclone/gofakes3" -) - -// pager splits the object list into smulitply 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 key name - sort.Slice(list.Contents, func(i, j int) bool { - return list.Contents[i].Key < list.Contents[j].Key - }) - 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 -} diff --git a/cmd/serve/s3/pager_test.go b/cmd/serve/s3/pager_test.go deleted file mode 100644 index e7166b1aa..000000000 --- a/cmd/serve/s3/pager_test.go +++ /dev/null @@ -1,32 +0,0 @@ -package s3 - -import ( - "testing" - "time" - - "github.com/rclone/gofakes3" -) - -func TestPagerSortsContentsByKey(t *testing.T) { - list := gofakes3.NewObjectList() - list.Add(&gofakes3.Content{ - Key: "b.txt", - LastModified: gofakes3.NewContentTime(time.Unix(100, 0)), - }) - list.Add(&gofakes3.Content{ - Key: "a.txt", - LastModified: gofakes3.NewContentTime(time.Unix(200, 0)), - }) - - got, err := (&s3Backend{}).pager(list, gofakes3.ListBucketPage{MaxKeys: 2}) - if err != nil { - t.Fatal(err) - } - if len(got.Contents) != 2 { - t.Fatalf("expected 2 contents, got %d", len(got.Contents)) - } - - if got.Contents[0].Key != "a.txt" || got.Contents[1].Key != "b.txt" { - t.Fatalf("expected lexicographic key order [a.txt b.txt], got [%s %s]", got.Contents[0].Key, got.Contents[1].Key) - } -} diff --git a/cmd/serve/s3/s3_test.go b/cmd/serve/s3/s3_test.go index dbec07a8e..28c2ae7a2 100644 --- a/cmd/serve/s3/s3_test.go +++ b/cmd/serve/s3/s3_test.go @@ -461,6 +461,114 @@ func TestAuthKeyPerServer(t *testing.T) { assert.Equal(t, http.StatusForbidden, signedGet(urlB, keyA, secA), "B with A's key") } +// newMinioClient serves f and returns a minio client connected to it. +func newMinioClient(t *testing.T, f fs.Fs) *minio.Client { + endpoint, keyid, keysec, _ := serveS3(t, f) + testURL, err := url.Parse(endpoint) + require.NoError(t, err) + minioClient, err := minio.New(testURL.Host, &minio.Options{ + Creds: credentials.NewStaticV4(keyid, keysec, ""), + Secure: false, + }) + require.NoError(t, err) + return minioClient +} + +// putTestObject puts an object with the given key into bucket on f. +func putTestObject(t *testing.T, f fs.Fs, bucket, key string) { + contents := "contents" + obji := object.NewStaticObjectInfo(path.Join(bucket, key), time.Now(), int64(len(contents)), true, nil, nil) + _, err := f.Put(context.Background(), bytes.NewBufferString(contents), obji) + require.NoError(t, err) +} + +// listKeysWithMinioClient pages through a listing with the given options and +// returns the keys it yields, common prefixes included. +func listKeysWithMinioClient(t *testing.T, minioClient *minio.Client, bucket string, opts minio.ListObjectsOptions) (keys []string) { + for object := range minioClient.ListObjects(context.Background(), bucket, opts) { + require.NoError(t, object.Err) + keys = append(keys, object.Key) + } + return keys +} + +// TestListObjectsPagingWithMinioClient checks that a client paging through a +// deep hierarchy with a small max-keys sees every key exactly once, in order, +// exercising the continuation tokens end to end. +func TestListObjectsPagingWithMinioClient(t *testing.T) { + fstest.Initialise() + f, _, clean, err := fstest.RandomRemote() + require.NoError(t, err) + defer clean() + + // A chunk store shaped like the one Proxmox Backup Server lists: + // many directories, each holding a few objects, listed with a prefix + // and no delimiter. + var want []string + for dir := range 20 { + for file := range 5 { + key := fmt.Sprintf(".chunks/%04x/chunk%d", dir, file) + want = append(want, key) + putTestObject(t, f, "mybucket", key) + } + } + slices.Sort(want) + + minioClient := newMinioClient(t, f) + for _, useV1 := range []bool{false, true} { + t.Run(fmt.Sprintf("UseV1=%v", useV1), func(t *testing.T) { + got := listKeysWithMinioClient(t, minioClient, "mybucket", minio.ListObjectsOptions{ + Prefix: ".chunks", + Recursive: true, + MaxKeys: 7, + UseV1: useV1, + }) + assert.Equal(t, want, got) + }) + } +} + +// TestListObjectsDelimitedPagingWithMinioClient checks that a client paging +// through a delimited listing sees every common prefix and every object key +// exactly once, whichever side of a page boundary they fall. +func TestListObjectsDelimitedPagingWithMinioClient(t *testing.T) { + fstest.Initialise() + f, _, clean, err := fstest.RandomRemote() + require.NoError(t, err) + defer clean() + + // Directories and objects which interleave in key order, so that a page + // boundary has to fall between a common prefix and an object key. + var want []string + for i := range 20 { + if i%2 == 0 { + putTestObject(t, f, "mybucket", fmt.Sprintf("a%02d/object", i)) + want = append(want, fmt.Sprintf("a%02d/", i)) + } else { + key := fmt.Sprintf("a%02d.txt", i) + putTestObject(t, f, "mybucket", key) + want = append(want, key) + } + } + slices.Sort(want) + + minioClient := newMinioClient(t, f) + for _, maxKeys := range []int{1, 3, 7, 20} { + for _, useV1 := range []bool{false, true} { + t.Run(fmt.Sprintf("MaxKeys%d/UseV1=%v", maxKeys, useV1), func(t *testing.T) { + got := listKeysWithMinioClient(t, minioClient, "mybucket", minio.ListObjectsOptions{ + MaxKeys: maxKeys, + UseV1: useV1, + }) + // The client yields the objects of a page before its common + // prefixes, so only the keys of the whole listing are sorted. + slices.Sort(got) + assert.Equal(t, want, got) + }) + } + } +} + func TestRc(t *testing.T) { servetest.TestRc(t, rc.Params{ "type": "s3",