From 3c17c5146e00fbc56908e41809ffcfacd86585ff Mon Sep 17 00:00:00 2001 From: Chris Lu Date: Mon, 5 Oct 2026 11:58:04 +0800 Subject: [PATCH] S3: fix ListObjectVersions losing keys across page boundaries (#11598) * s3api: thread filer client through the versioned-listing collector findVersionsRecursively now binds one SeaweedFilerClient for the whole recursive walk instead of re-resolving a filer on every list/lookup call, and the collector's list/getEntry/scanLatestVersionEntry/getObjectVersionList helpers go through it. No behavior change; this also lets tests drive collectVersions with a stubbed client. * s3api: keep collecting versions while pending names can sort into the page ListObjectVersions walked the filer in directory-entry name order and stopped as soon as maxKeys+1 items were collected, sorting only that partial set. Filer names do not match key order: "a.copy.versions" sorts before "a.versions" while key "a.copy" sorts after "a", so a page boundary inside the earlier-walked sibling's versions permanently skipped the later key. Track the largest key collected (maxKey) and, once the collector is full, keep walking until entry names pass the ceiling of names that can still resolve to keys at or below it; the ceiling reaches through the prefix versions of maxKey. Entries whose subtree can only hold keys above maxKey are skipped. Versions of an in-bound object are collected in full so its position in the sorted page is exact. Fixes seaweedfs#11594 * s3api: resume versioned listings at the earliest covering name prefix computeStartFrom mapped the key marker straight to an entry name (or cut it at the first '/'), which skips sibling directories that are a prefix of the marker below '0' - for marker "d.x" the listing resumed at name "d.x", skipping directory "d" whose keys "d/*" all sort after it. Resume at the earliest remainder prefix ending at a byte below '0' ('/', '.', '-' and friends), so every directory whose subtree can still hold keys past the marker is revisited; already-returned keys inside are filtered by the existing marker checks as before. * s3api: regression test for versioned-listing pagination order Drive collectVersions against a stubbed filer holding the issue-11594 layout - "a.copy.versions" listing before "a.versions", plus a "d/" subtree next to "d.x" - and assert that every page size from 1 up reproduces the unpaginated ordering with no lost or duplicated entries. Also updates TestComputeStartFrom for the new earliest-prefix resume and gives testFilerClient a LookupDirectoryEntry stub. * s3api: inject list/getEntry functions into the version collector Pinning one SeaweedFilerClient for the whole walk dropped per-call failover: previously each s3a.list resolved a filer through WithFilerClient, so a mid-walk filer failure could fall back to a healthy peer. Inject s3a.list/s3a.getEntry as function fields instead - production keeps the failover behavior, tests can still stub. * s3api: keep scanning marker for a later covering prefix A leading byte below '0' (marker .hidden/file) has no non-empty prefix at index 0, but a deeper separator still does - resuming at .hidden/file skipped the .hidden directory and its keys after file. Continue the scan instead of bailing on the first byte. --- weed/s3api/s3api_object_handlers_list_test.go | 9 ++ ...api_object_handlers_list_versioned_test.go | 8 +- weed/s3api/s3api_object_versioning.go | 133 +++++++++++---- ...s3api_object_versioning_pagination_test.go | 152 ++++++++++++++++++ 4 files changed, 268 insertions(+), 34 deletions(-) create mode 100644 weed/s3api/s3api_object_versioning_pagination_test.go diff --git a/weed/s3api/s3api_object_handlers_list_test.go b/weed/s3api/s3api_object_handlers_list_test.go index 6c6fb95d8..ecbde4671 100644 --- a/weed/s3api/s3api_object_handlers_list_test.go +++ b/weed/s3api/s3api_object_handlers_list_test.go @@ -72,6 +72,15 @@ func (c *testFilerClient) ListEntries(ctx context.Context, in *filer_pb.ListEntr return &testListEntriesStream{entries: entries}, nil } +func (c *testFilerClient) LookupDirectoryEntry(ctx context.Context, in *filer_pb.LookupDirectoryEntryRequest, opts ...grpc.CallOption) (*filer_pb.LookupDirectoryEntryResponse, error) { + for _, e := range c.entriesByDir[in.Directory] { + if e.Name == in.Name { + return &filer_pb.LookupDirectoryEntryResponse{Entry: e}, nil + } + } + return nil, filer_pb.ErrNotFound +} + type markerEchoFilerClient struct { filer_pb.SeaweedFilerClient entriesByDir map[string][]*filer_pb.Entry diff --git a/weed/s3api/s3api_object_handlers_list_versioned_test.go b/weed/s3api/s3api_object_handlers_list_versioned_test.go index 43ec65f27..3bb97a342 100644 --- a/weed/s3api/s3api_object_handlers_list_versioned_test.go +++ b/weed/s3api/s3api_object_handlers_list_versioned_test.go @@ -641,7 +641,7 @@ func TestComputeStartFrom(t *testing.T) { }{ {"empty marker", "", "", "", false}, {"empty marker with path", "", "dir", "", false}, - {"root level file", "file1.txt", "", "file1.txt", true}, + {"root level file", "file1.txt", "", "file1", true}, {"root level with subpath", "Mailboxes/5ac/file1", "", "Mailboxes", true}, {"matching subdir", "Mailboxes/5ac/file1", "Mailboxes", "5ac", true}, {"deeper subdir", "Mailboxes/5ac/ItemsData/file1", "Mailboxes/5ac", "ItemsData", true}, @@ -649,6 +649,12 @@ func TestComputeStartFrom(t *testing.T) { {"unrelated directory", "other/path", "Mailboxes", "", false}, {"marker equals relativePath", "Mailboxes", "Mailboxes", "", false}, {"marker before directory", "aaa/file", "zzz", "", false}, + {"dir prefix below slash", "d.x", "", "d", true}, + {"dir prefix below slash nested", "a-b/file", "", "a", true}, + {"dir prefix at marker end", "d-", "", "d", true}, + {"no qualifying prefix", "zzz", "", "zzz", true}, + {"leading below-slash byte", ".hidden", "", ".hidden", true}, + {"leading below-slash byte nested", ".hidden/file", "", ".hidden", true}, } for _, tt := range tests { diff --git a/weed/s3api/s3api_object_versioning.go b/weed/s3api/s3api_object_versioning.go index da70a7bef..1f6484d46 100644 --- a/weed/s3api/s3api_object_versioning.go +++ b/weed/s3api/s3api_object_versioning.go @@ -346,14 +346,8 @@ func (s3a *S3ApiServer) listObjectVersions(bucket, prefix, keyMarker, versionIdM // Pass keyMarker and versionIdMarker to enable efficient pagination (skip entries before marker) bucketPath := s3a.bucketDir(bucket) - // Memory optimization: limit collection to maxKeys+1 versions. - // This works correctly for objects using the NEW inverted-timestamp format, where - // filesystem order (lexicographic) matches sorted order (newest-first). - // For OLD format objects (raw timestamps), filesystem order is oldest-first, so - // limiting collection may return older versions instead of newest. However: - // - New objects going forward use the new format - // - The alternative (collecting all) causes memory issues for buckets with many versions - // - Pagination continues correctly; users can page through to see all versions + // Memory optimization: limit collection to maxKeys+1 versions, extended past + // entry names that can still resolve to keys sorting into the page. maxCollect := maxKeys + 1 // +1 to detect truncation err := s3a.findVersionsRecursively(bucketPath, "", &allVersions, processedObjects, seenVersionIds, bucket, prefix, keyMarker, versionIdMarker, delimiter, commonPrefixes, maxCollect) if err != nil { @@ -479,6 +473,8 @@ func (s3a *S3ApiServer) splitIntoResult(combinedList []versionListItem, bucket, // versionCollector holds state for collecting object versions during recursive traversal type versionCollector struct { s3a *S3ApiServer + list entryLister + getEntry func(parentDirectoryPath, entryName string) (entry *filer_pb.Entry, err error) bucket string prefix string keyMarker string @@ -489,6 +485,7 @@ type versionCollector struct { seenVersionIds map[string]bool delimiter string commonPrefixes map[string]bool + maxKey string } // isFull returns true if we've collected enough versions @@ -503,6 +500,58 @@ func (vc *versionCollector) isFull() bool { return currentCount >= vc.maxCollect } +// recordKey tracks the largest key collected so far, so collectVersions can +// tell when an entry name still pending can resolve to a key that sorts into +// the collected page. +func (vc *versionCollector) recordKey(key string) { + if key > vc.maxKey { + vc.maxKey = key + } +} + +// nameCeiling returns the largest directory-entry name that can still resolve +// to a key at or below key. Name order is not key order: "a.copy.versions" +// sorts before "a.versions" while key "a.copy" sorts after "a", so the names +// that map to keys <= key reach past key+".versions" - every prefix of key can +// host a .versions directory whose name is the real ceiling. +func nameCeiling(key string) string { + ceiling := key + s3_constants.VersionsFolder + for i := 0; i < len(key); i++ { + if n := key[:i] + s3_constants.VersionsFolder; n > ceiling { + ceiling = n + } + } + return ceiling +} + +// collectBound returns the largest entry name under relativePath that can still +// resolve to a key at or below maxKey, and whether the walk may stop past it. +// The bound is absent when nothing was collected yet, or when every key under +// relativePath already sorts below maxKey and the whole subtree must be read. +func (vc *versionCollector) collectBound(relativePath string) (string, bool) { + if vc.maxKey == "" { + return "", false + } + if relativePath == "" { + return nameCeiling(vc.maxKey), true + } + if rel, ok := strings.CutPrefix(vc.maxKey, relativePath+"/"); ok { + return nameCeiling(rel), true + } + return "", vc.maxKey < relativePath+"/" +} + +// mayCollect reports whether entry name under relativePath can still resolve +// to a key that sorts into the collected page; once the collector is full and +// the name is past the bound, every later name can only produce keys above it. +func (vc *versionCollector) mayCollect(relativePath, name string) bool { + if !vc.isFull() { + return true + } + bound, bounded := vc.collectBound(relativePath) + return !bounded || name <= bound +} + // matchesPrefixFilter checks if an entry path matches the prefix filter func (vc *versionCollector) matchesPrefixFilter(entryPath string, isDirectory bool) bool { if vc.prefix == "" { @@ -544,8 +593,14 @@ func (vc *versionCollector) computeStartFrom(relativePath string) (startFrom str return "", false } - if idx := strings.Index(remainder, "/"); idx >= 0 { - return remainder[:idx], true + // A sibling directory that is a prefix of the marker up to a byte below '0' + // still holds keys that sort after it - "d" for marker "d.x" contains "d/*". + // Resuming at the marker's name would skip those keys, so the walk resumes + // at the earliest such prefix instead. + for i := 1; i < len(remainder); i++ { + if remainder[i] < '0' { + return remainder[:i], true + } } return remainder, true } @@ -613,6 +668,7 @@ func (vc *versionCollector) shouldSkipVersionForMarker(objectKey, versionId stri // addVersion adds a version or delete marker to results func (vc *versionCollector) addVersion(version *ObjectVersion, objectKey string) { + vc.recordKey(objectKey) if version.IsDeleteMarker { deleteMarker := &DeleteMarkerEntry{ Key: objectKey, @@ -655,17 +711,13 @@ func (vc *versionCollector) processVersionsDirectory(entryPath string, versionsE glog.V(2).Infof("processVersionsDirectory: found object %s", normalizedObjectKey) - versions, err := vc.s3a.getObjectVersionList(vc.bucket, normalizedObjectKey, versionsEntry) + versions, err := vc.getObjectVersionList(normalizedObjectKey, versionsEntry) if err != nil { glog.Warningf("processVersionsDirectory: failed to get versions for %s: %v", normalizedObjectKey, err) return nil // Continue with other entries } for _, version := range versions { - if vc.isFull() { - return nil - } - versionKey := normalizedObjectKey + ":" + version.VersionId if vc.seenVersionIds[versionKey] { continue @@ -704,6 +756,7 @@ func (vc *versionCollector) processExplicitDirectory(entryPath string, entry *fi return } + vc.recordKey(directoryKey) versionEntry := &VersionEntry{ Key: directoryKey, VersionId: "null", @@ -742,10 +795,10 @@ func (vc *versionCollector) processRegularFile(currentPath, entryPath string, en // Check if a .versions directory exists for this object versionsEntryName := entry.Name + s3_constants.VersionsFolder - versionsDirEntry, versionsErr := vc.s3a.getEntry(currentPath, versionsEntryName) + versionsDirEntry, versionsErr := vc.getEntry(currentPath, versionsEntryName) if versionsErr == nil && !hasVersionMeta { // .versions exists but file has no version metadata - check for null version in .versions - versions, err := vc.s3a.getObjectVersionList(vc.bucket, normalizedObjectKey, versionsDirEntry) + versions, err := vc.getObjectVersionList(normalizedObjectKey, versionsDirEntry) if err == nil { for _, v := range versions { if v.VersionId == "null" { @@ -775,12 +828,13 @@ func (vc *versionCollector) processRegularFile(currentPath, entryPath string, en if len(versionsDirEntry.Extended[s3_constants.ExtLatestVersionIdKey]) > 0 { isLatest = false } else if !nullVersionIsLatest(versionsDirEntry) { - if latestVersion, _, _, _, scanErr := vc.s3a.scanLatestVersionEntry(currentPath + "/" + versionsEntryName); scanErr == nil && latestVersion != nil && !nullObjectWins(entry, latestVersion) { + if latestVersion, _, _, _, scanErr := scanLatestVersionEntry(vc.list, currentPath+"/"+versionsEntryName); scanErr == nil && latestVersion != nil && !nullObjectWins(entry, latestVersion) { isLatest = false } } } + vc.recordKey(normalizedObjectKey) versionEntry := &VersionEntry{ Key: normalizedObjectKey, VersionId: "null", @@ -801,6 +855,8 @@ func (vc *versionCollector) processRegularFile(currentPath, entryPath string, en func (s3a *S3ApiServer) findVersionsRecursively(currentPath, relativePath string, allVersions *[]interface{}, processedObjects map[string]bool, seenVersionIds map[string]bool, bucket, prefix, keyMarker, versionIdMarker, delimiter string, commonPrefixes map[string]bool, maxCollect int) error { vc := &versionCollector{ s3a: s3a, + list: s3a.list, + getEntry: s3a.getEntry, bucket: bucket, prefix: prefix, keyMarker: keyMarker, @@ -828,11 +884,11 @@ func (vc *versionCollector) collectVersions(currentPath, relativePath string) er } listPrefix := vc.computeListPrefix(relativePath) for { - if vc.isFull() { + if bound, bounded := vc.collectBound(relativePath); vc.isFull() && bounded && startFrom >= bound { return nil } - entries, isLast, err := vc.s3a.list(currentPath, listPrefix, startFrom, inclusive, filer.PaginationSize) + entries, isLast, err := vc.list(currentPath, listPrefix, startFrom, inclusive, filer.PaginationSize) // After the first batch, use exclusive mode for standard pagination inclusive = false if err != nil { @@ -840,7 +896,7 @@ func (vc *versionCollector) collectVersions(currentPath, relativePath string) er } for _, entry := range entries { - if vc.isFull() { + if !vc.mayCollect(relativePath, entry.Name) { return nil } startFrom = entry.Name @@ -885,10 +941,8 @@ func (vc *versionCollector) collectVersions(currentPath, relativePath string) er if vc.keyMarker != "" && commonPrefix <= vc.keyMarker { continue } - if vc.isFull() { - return nil - } vc.commonPrefixes[commonPrefix] = true + vc.recordKey(commonPrefix) } // The prefix rolled up here belongs to the keys nested under this @@ -957,6 +1011,12 @@ func (vc *versionCollector) processDirectory(currentPath, entryPath string, entr return nil } + // Once the page is full, a subtree whose keys all sort above maxKey cannot + // contribute to it. + if bound, bounded := vc.collectBound(entryPath); vc.isFull() && bounded && bound == "" { + return nil + } + // Recursively search subdirectory fullPath := path.Join(currentPath, entry.Name) if err := vc.collectVersions(fullPath, entryPath); err != nil { @@ -971,7 +1031,7 @@ func (vc *versionCollector) processDirectory(currentPath, entryPath string, entr // already holds from listing the parent directory - re-fetching it here would // cost one extra filer round-trip per object listed. // Uses pagination to handle objects with more than 1000 versions -func (s3a *S3ApiServer) getObjectVersionList(bucket, object string, versionsEntry *filer_pb.Entry) ([]*ObjectVersion, error) { +func (vc *versionCollector) getObjectVersionList(object string, versionsEntry *filer_pb.Entry) ([]*ObjectVersion, error) { var versions []*ObjectVersion // A nil entry means the .versions directory is absent: no versions, the @@ -981,9 +1041,9 @@ func (s3a *S3ApiServer) getObjectVersionList(bucket, object string, versionsEntr } // All versions are now stored in the .versions directory only - bucketDir := s3a.bucketDir(bucket) + bucketDir := vc.s3a.bucketDir(vc.bucket) versionsObjectPath := object + s3_constants.VersionsFolder - glog.V(2).Infof("getObjectVersionList: looking for versions of %s/%s in %s", bucket, object, versionsObjectPath) + glog.V(2).Infof("getObjectVersionList: looking for versions of %s/%s in %s", vc.bucket, object, versionsObjectPath) // Get the latest version info from directory metadata var latestVersionId string @@ -1004,7 +1064,7 @@ func (s3a *S3ApiServer) getObjectVersionList(bucket, object string, versionsEntr totalEntries := 0 for { - entries, isLast, err := s3a.list(versionsDir, "", startFrom, false, pageSize) + entries, isLast, err := vc.list(versionsDir, "", startFrom, false, pageSize) if err != nil { glog.Warningf("getObjectVersionList: failed to list version files in %s: %v", versionsDir, err) return nil, err @@ -1031,7 +1091,7 @@ func (s3a *S3ApiServer) getObjectVersionList(bucket, object string, versionsEntr // Check for duplicate version IDs and skip if already seen if seenVersionIds[versionId] { - glog.Warningf("getObjectVersionList: duplicate version ID %s detected for object %s/%s, skipping", versionId, bucket, object) + glog.Warningf("getObjectVersionList: duplicate version ID %s detected for object %s/%s, skipping", versionId, vc.bucket, object) continue } seenVersionIds[versionId] = true @@ -1056,7 +1116,7 @@ func (s3a *S3ApiServer) getObjectVersionList(bucket, object string, versionsEntr IsDeleteMarker: isDeleteMarker, LastModified: time.Unix(entry.Attributes.Mtime, 0), OwnerID: ownerID, - StorageClass: s3a.getStorageClassFromExtended(entry.Extended), + StorageClass: vc.s3a.getStorageClassFromExtended(entry.Extended), } if !isDeleteMarker { @@ -1068,7 +1128,7 @@ func (s3a *S3ApiServer) getObjectVersionList(bucket, object string, versionsEntr } } else { // Fallback: calculate ETag from chunks - version.ETag = s3a.calculateETagFromChunks(entry.Chunks) + version.ETag = vc.s3a.calculateETagFromChunks(entry.Chunks) } version.Size = int64(entry.Attributes.FileSize) } @@ -1087,7 +1147,7 @@ func (s3a *S3ApiServer) getObjectVersionList(bucket, object string, versionsEntr // Don't sort here - let the main listObjectVersions function handle sorting consistently - glog.V(2).Infof("getObjectVersionList: returning %d total versions for %s/%s (after deduplication from %d entries)", len(versions), bucket, object, totalEntries) + glog.V(2).Infof("getObjectVersionList: returning %d total versions for %s/%s (after deduplication from %d entries)", len(versions), vc.bucket, object, totalEntries) for i, version := range versions { glog.V(2).Infof("getObjectVersionList: version %d: %s (isLatest=%v, isDeleteMarker=%v)", i, version.VersionId, version.IsLatest, version.IsDeleteMarker) } @@ -2038,6 +2098,9 @@ func selectLatestVersion(entries []*filer_pb.Entry) (latestEntry *filer_pb.Entry return } +// entryLister is the shared signature of s3a.list and versionCollector.list. +type entryLister func(parentDirectoryPath, prefix, startFrom string, inclusive bool, limit uint32) (entries []*filer_pb.Entry, isLast bool, err error) + // scanLatestVersionEntry paginates a .versions/ directory and returns the // chronologically newest version entry (including delete markers; see // selectLatestVersion). A single-shot list would miss the true latest when @@ -2045,9 +2108,13 @@ func selectLatestVersion(entries []*filer_pb.Entry) (latestEntry *filer_pb.Entry // order is lexicographic-ascending = oldest-first for that format. latestEntry // is nil when the directory holds no version entries. func (s3a *S3ApiServer) scanLatestVersionEntry(versionsDir string) (latestEntry *filer_pb.Entry, latestVersionId, latestVersionFileName string, isDeleteMarker bool, err error) { + return scanLatestVersionEntry(s3a.list, versionsDir) +} + +func scanLatestVersionEntry(list entryLister, versionsDir string) (latestEntry *filer_pb.Entry, latestVersionId, latestVersionFileName string, isDeleteMarker bool, err error) { startFrom := "" for { - entries, isLast, listErr := s3a.list(versionsDir, "", startFrom, false, filer.PaginationSize) + entries, isLast, listErr := list(versionsDir, "", startFrom, false, filer.PaginationSize) if listErr != nil { return nil, "", "", false, fmt.Errorf("list %s: %w", versionsDir, listErr) } diff --git a/weed/s3api/s3api_object_versioning_pagination_test.go b/weed/s3api/s3api_object_versioning_pagination_test.go new file mode 100644 index 000000000..91389ff89 --- /dev/null +++ b/weed/s3api/s3api_object_versioning_pagination_test.go @@ -0,0 +1,152 @@ +package s3api + +import ( + "context" + "testing" + "time" + + "github.com/seaweedfs/seaweedfs/weed/pb/filer_pb" + "github.com/seaweedfs/seaweedfs/weed/s3api/s3_constants" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// versionedDir builds the .versions directory entry of a versioned object whose +// latest version has the given id. +func versionedDir(object, latestVersionId string) *filer_pb.Entry { + return &filer_pb.Entry{ + Name: object + s3_constants.VersionsFolder, + IsDirectory: true, + Attributes: &filer_pb.FuseAttributes{Mtime: time.Now().Unix()}, + Extended: map[string][]byte{ + s3_constants.ExtLatestVersionIdKey: []byte(latestVersionId), + }, + } +} + +// versionFile builds one version entry inside a .versions directory. +func versionFile(versionId string, deleteMarker bool) *filer_pb.Entry { + e := &filer_pb.Entry{ + Name: "v_" + versionId, + Attributes: &filer_pb.FuseAttributes{Mtime: time.Now().Unix()}, + Extended: map[string][]byte{ + s3_constants.ExtVersionIdKey: []byte(versionId), + }, + } + if deleteMarker { + e.Extended[s3_constants.ExtDeleteMarkerKey] = []byte("true") + } + return e +} + +func listVersionsPage(t *testing.T, s3a *S3ApiServer, client filer_pb.SeaweedFilerClient, bucketDir string, maxKeys int, keyMarker, versionIdMarker string) ([]versionListItem, string, string, bool) { + t.Helper() + var allVersions []interface{} + vc := &versionCollector{ + s3a: s3a, + list: clientLister(client), + getEntry: clientLookup(client), + bucket: "b", + keyMarker: keyMarker, + versionIdMarker: versionIdMarker, + maxCollect: maxKeys + 1, + allVersions: &allVersions, + processedObjects: map[string]bool{}, + seenVersionIds: map[string]bool{}, + commonPrefixes: map[string]bool{}, + } + require.NoError(t, vc.collectVersions(bucketDir, "")) + combined := s3a.buildSortedCombinedList(allVersions, vc.commonPrefixes) + return s3a.truncateAndSetMarkers(combined, maxKeys) +} + +// clientLister adapts a stubbed filer client to the collector's list function. +func clientLister(client filer_pb.SeaweedFilerClient) entryLister { + return func(parentDirectoryPath, prefix, startFrom string, inclusive bool, limit uint32) (entries []*filer_pb.Entry, isLast bool, err error) { + err = filer_pb.SeaweedList(context.Background(), client, parentDirectoryPath, prefix, func(entry *filer_pb.Entry, isLastEntry bool) error { + entries = append(entries, entry) + if isLastEntry { + isLast = true + } + return nil + }, startFrom, inclusive, limit) + if len(entries) == 0 { + isLast = true + } + return + } +} + +// clientLookup adapts a stubbed filer client to the collector's getEntry. +func clientLookup(client filer_pb.SeaweedFilerClient) func(string, string) (*filer_pb.Entry, error) { + return func(dir, name string) (*filer_pb.Entry, error) { + resp, err := filer_pb.LookupEntry(context.Background(), client, &filer_pb.LookupDirectoryEntryRequest{Directory: dir, Name: name}) + if err != nil { + return nil, err + } + return resp.Entry, nil + } +} + +func itemId(item versionListItem) string { + return item.key + ":" + item.versionId +} + +// TestListObjectVersionsPagination covers issue 11594: "a.copy.versions" sorts +// before "a.versions" in the filer even though key "a.copy" sorts after "a", so +// cutting the walk at maxKeys+1 can put the later key on the first page and the +// marker then skips the earlier key for good. Directory "d" and file "d.x" +// exercise the mirror-image resume case: a marker on "d.x" must not skip the +// subtree of "d", whose keys all sort after it. +func TestListObjectVersionsPagination(t *testing.T) { + client := &testFilerClient{ + entriesByDir: map[string][]*filer_pb.Entry{ + "/buckets/b": { + newDir(".hidden"), + versionedDir("a.copy", "c2"), + versionedDir("a", "a2"), + newDir("d"), + {Name: "d.x", Attributes: &filer_pb.FuseAttributes{Mtime: time.Now().Unix()}}, + }, + "/buckets/b/.hidden": { + versionedDir("file", "h1"), + versionedDir("z", "z1"), + }, + "/buckets/b/.hidden/file.versions": {versionFile("h1", false)}, + "/buckets/b/.hidden/z.versions": {versionFile("z1", false)}, + "/buckets/b/a.copy.versions": {versionFile("c1", false), versionFile("c2", true)}, + "/buckets/b/a.versions": {versionFile("a1", false), versionFile("a2", false)}, + "/buckets/b/d": { + versionedDir("f", "f1"), + versionedDir("g", "g1"), + }, + "/buckets/b/d/f.versions": {versionFile("f1", false)}, + "/buckets/b/d/g.versions": {versionFile("g1", false)}, + }, + } + s3a := &S3ApiServer{option: &S3ApiServerOption{BucketsPath: "/buckets"}} + + want, _, _, truncated := listVersionsPage(t, s3a, client, "/buckets/b", 100, "", "") + require.False(t, truncated) + wantIds := make([]string, 0, len(want)) + for _, item := range want { + wantIds = append(wantIds, itemId(item)) + } + require.Equal(t, []string{".hidden/file:h1", ".hidden/z:z1", "a:a2", "a:a1", "a.copy:c2", "a.copy:c1", "d.x:null", "d/f:f1", "d/g:g1"}, wantIds) + + for maxKeys := 1; maxKeys <= len(wantIds)+1; maxKeys++ { + var got []string + keyMarker, versionIdMarker := "", "" + for i := 0; i < 20; i++ { + page, nextKey, nextVersion, trunc := listVersionsPage(t, s3a, client, "/buckets/b", maxKeys, keyMarker, versionIdMarker) + for _, item := range page { + got = append(got, itemId(item)) + } + if !trunc { + break + } + keyMarker, versionIdMarker = nextKey, nextVersion + } + assert.Equal(t, wantIds, got, "maxKeys=%d", maxKeys) + } +}