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) + } +}