From 4cbb074ccdd3d1fc2453e68f6a18d76abd16e253 Mon Sep 17 00:00:00 2001 From: mr javad seydi Date: Tue, 15 Sep 2026 23:17:55 +0330 Subject: [PATCH] fix: make multipart upload listing S3-compatible Signed-off-by: mr javad seydi --- CHANGELOG.md | 9 + cmd/api-response.go | 2 +- cmd/bucket-handlers.go | 8 - cmd/bucket-handlers_test.go | 6 +- cmd/erasure-multipart.go | 319 ++++++++++++++++++++-- cmd/erasure-server-pool.go | 34 ++- cmd/erasure-sets.go | 30 +- cmd/list-multipart-uploads-compat_test.go | 283 +++++++++++++++++++ cmd/object-api-input-checks.go | 3 +- 9 files changed, 651 insertions(+), 43 deletions(-) create mode 100644 cmd/list-multipart-uploads-compat_test.go diff --git a/CHANGELOG.md b/CHANGELOG.md index 8dad9a0c4..ccd19dfa0 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -22,6 +22,15 @@ and [complete commit range](https://github.com/pgsty/silo/compare/RELEASE.2026-0 ### Object storage and replication +- Make `ListMultipartUploads` discover quorum-valid uploads from durable state + across pools, erasure sets and drives, then apply S3 prefix, delimiter, + marker, ordering and 1,000-entry pagination semantics globally. New uploads + store their canonical bucket and key as reserved fields in the existing + quorum-written `xl.meta`; completion removes those upload-only fields. During + rolling upgrades, detection of any legacy keyless upload retains the prior + listing behavior until those uploads drain. See [issue #79](https://github.com/pgsty/silo/issues/79) + and its [design record](https://silo.pgsty.com/blog/design/list-multipart-uploads/). + - Evaluate conditional multipart completion against the logical current object across all pools while holding the existing object lock. A stale `If-Match` can no longer replace newer data in another pool, and the current ETag is no diff --git a/cmd/api-response.go b/cmd/api-response.go index 0750e4f67..27749b9f5 100644 --- a/cmd/api-response.go +++ b/cmd/api-response.go @@ -41,7 +41,7 @@ import ( const ( maxObjectList = 1000 // Limit number of objects in a listObjectsResponse/listObjectsVersionsResponse. maxDeleteList = 1000 // Limit number of objects deleted in a delete call. - maxUploadsList = 10000 // Limit number of uploads in a listUploadsResponse. + maxUploadsList = 1000 // Limit number of uploads in a listUploadsResponse. maxPartsList = 10000 // Limit number of parts in a listPartsResponse. ) diff --git a/cmd/bucket-handlers.go b/cmd/bucket-handlers.go index 93e14d837..e79a7841f 100644 --- a/cmd/bucket-handlers.go +++ b/cmd/bucket-handlers.go @@ -277,14 +277,6 @@ func (api objectAPIHandlers) ListMultipartUploadsHandler(w http.ResponseWriter, return } - if keyMarker != "" { - // Marker not common with prefix is not implemented. - if !HasPrefix(keyMarker, prefix) { - writeErrorResponse(ctx, w, errorCodes.ToAPIErr(ErrNotImplemented), r.URL) - return - } - } - listMultipartsInfo, err := objectAPI.ListMultipartUploads(ctx, bucket, prefix, keyMarker, uploadIDMarker, delimiter, maxUploads) if err != nil { writeErrorResponse(ctx, w, toAPIError(ctx, err), r.URL) diff --git a/cmd/bucket-handlers_test.go b/cmd/bucket-handlers_test.go index 32c041dff..337435683 100644 --- a/cmd/bucket-handlers_test.go +++ b/cmd/bucket-handlers_test.go @@ -425,7 +425,7 @@ func testListMultipartUploadsHandler(obj ObjectLayer, instanceType, bucketName s shouldPass: true, }, // Test case - 4. - // Setting Invalid prefix and marker combination. + // A key marker outside the prefix is valid and produces an empty page. { bucket: bucketName, prefix: "asia", @@ -435,8 +435,8 @@ func testListMultipartUploadsHandler(obj ObjectLayer, instanceType, bucketName s maxUploads: "0", accessKey: credentials.AccessKey, secretKey: credentials.SecretKey, - expectedRespStatus: http.StatusNotImplemented, - shouldPass: false, + expectedRespStatus: http.StatusOK, + shouldPass: true, }, // Test case - 5. // Invalid upload id and marker combination. diff --git a/cmd/erasure-multipart.go b/cmd/erasure-multipart.go index e78d55899..0a1e70718 100644 --- a/cmd/erasure-multipart.go +++ b/cmd/erasure-multipart.go @@ -44,6 +44,15 @@ import ( "github.com/pgsty/silo-pkg/v3/sync/errgroup" ) +const ( + multipartMetaBucket = ReservedMetadataPrefixLower + "multipart-v1-bucket" + multipartMetaObject = ReservedMetadataPrefixLower + "multipart-v1-object" + + // ponytail: keep scan concurrency fixed until the multipart-list benchmark + // establishes a better adaptive limit. + multipartMetadataScanConcurrency = 4 +) + func (er erasureObjects) getUploadIDDir(bucket, object, uploadID string) string { uploadUUID := uploadID uploadBytes, err := base64.RawURLEncoding.DecodeString(uploadID) @@ -251,14 +260,278 @@ func (er erasureObjects) cleanupStaleUploadsOnDisk(ctx context.Context, disk Sto }) } -// ListMultipartUploads - lists all the pending multipart -// uploads for a particular object in a bucket. -// -// Implements minimal S3 compatible ListMultipartUploads API. We do -// not support prefix based listing, this is a deliberate attempt -// towards simplification of multipart APIs. -// The resulting ListMultipartsInfo structure is unmarshalled directly as XML. -func (er erasureObjects) ListMultipartUploads(ctx context.Context, bucket, object, keyMarker, uploadIDMarker, delimiter string, maxUploads int) (result ListMultipartsInfo, err error) { +func multipartUploadInfo(bucket, object, uploadUUID string, fallback time.Time) MultipartInfo { + initiated := fallback + if i := strings.LastIndexByte(uploadUUID, 'x'); i >= 0 { + if parsed, err := strconv.ParseInt(uploadUUID[i+1:], 10, 64); err == nil { + initiated = time.Unix(0, parsed) + } + } + if initiated.IsZero() { + initiated = UTCNow() + } + return MultipartInfo{ + Bucket: bucket, + Object: object, + UploadID: base64.RawURLEncoding.EncodeToString(fmt.Appendf(nil, "%s.%s", globalDeploymentID(), uploadUUID)), + Initiated: initiated, + } +} + +func listMultipartUploadDirs(ctx context.Context, disk StorageAPI, bucket string) ([]string, error) { + hashDirs, err := disk.ListDir(ctx, bucket, minioMetaMultipartBucket, "", -1) + if errors.Is(err, errFileNotFound) { + return nil, nil + } + if err != nil { + return nil, err + } + + var candidates []string + for _, hashDir := range hashDirs { + if !strings.HasSuffix(hashDir, SlashSeparator) { + continue + } + hashDir = strings.TrimSuffix(hashDir, SlashSeparator) + uploadDirs, err := disk.ListDir(ctx, bucket, minioMetaMultipartBucket, hashDir, -1) + if errors.Is(err, errFileNotFound) { + continue + } + if err != nil { + return nil, err + } + for _, uploadDir := range uploadDirs { + if strings.HasSuffix(uploadDir, SlashSeparator) { + candidates = append(candidates, pathJoin(hashDir, strings.TrimSuffix(uploadDir, SlashSeparator))) + } + } + } + return candidates, nil +} + +func (er erasureObjects) readMultipartUploadCandidate(ctx context.Context, bucket, candidate string) (MultipartInfo, bool, bool, error) { + shaDir, uploadUUID, ok := strings.Cut(candidate, SlashSeparator) + if !ok || shaDir == "" || uploadUUID == "" || strings.Contains(uploadUUID, SlashSeparator) { + return MultipartInfo{}, false, false, fmt.Errorf("invalid multipart upload path %q", candidate) + } + + disks := er.getDisks() + partsMetadata, errs := readAllFileInfo(ctx, disks, bucket, minioMetaMultipartBucket, candidate, "", false, false) + readQuorum, _, err := objectQuorumFromMeta(ctx, partsMetadata, errs, er.defaultParityCount) + if err != nil { + return MultipartInfo{}, false, false, err + } + _, modTime, etag := listOnlineDisks(disks, partsMetadata, errs, readQuorum) + if err = reduceReadQuorumErrs(ctx, errs, objectOpIgnoredErrs, readQuorum); err != nil { + return MultipartInfo{}, false, false, err + } + fi, err := pickValidFileInfo(ctx, partsMetadata, modTime, etag, readQuorum) + if err != nil { + return MultipartInfo{}, false, false, err + } + + storedBucket, hasBucket := fi.Metadata[multipartMetaBucket] + storedObject, hasObject := fi.Metadata[multipartMetaObject] + if !hasBucket || !hasObject || storedBucket == "" || storedObject == "" { + return MultipartInfo{}, false, true, nil + } + if storedBucket != bucket { + return MultipartInfo{}, false, false, nil + } + if !IsValidObjectPrefix(storedObject) || er.getMultipartSHADir(storedBucket, storedObject) != shaDir { + return MultipartInfo{}, false, false, fmt.Errorf("multipart upload %q has invalid stored identity", candidate) + } + return multipartUploadInfo(storedBucket, storedObject, uploadUUID, fi.ModTime), true, false, nil +} + +// scanMultipartUploads discovers durable multipart state from every online +// drive in this erasure set and accepts only quorum-valid upload metadata. +func (er erasureObjects) scanMultipartUploads(ctx context.Context, bucket string) ([]MultipartInfo, bool, error) { + disks := er.getOnlineDisks() + if len(disks) < (er.setDriveCount+1)/2 { + return nil, false, errErasureReadQuorum + } + + candidateSet := make(map[string]struct{}) + var candidateMu sync.Mutex + g := errgroup.WithNErrs(len(disks)) + for i := range disks { + g.Go(func() error { + candidates, err := listMultipartUploadDirs(ctx, disks[i], bucket) + if err != nil { + return err + } + candidateMu.Lock() + for _, candidate := range candidates { + candidateSet[candidate] = struct{}{} + } + candidateMu.Unlock() + return nil + }, i) + } + if err := reduceReadQuorumErrs(ctx, g.Wait(), nil, (er.setDriveCount+1)/2); err != nil { + return nil, false, err + } + + candidates := make([]string, 0, len(candidateSet)) + for candidate := range candidateSet { + candidates = append(candidates, candidate) + } + sort.Strings(candidates) + if len(candidates) == 0 { + return nil, false, nil + } + + jobs := make(chan string) + workerCount := min(multipartMetadataScanConcurrency, len(candidates)) + var wg sync.WaitGroup + var resultMu sync.Mutex + var uploads []MultipartInfo + var keyless bool + var firstErr error + for range workerCount { + wg.Add(1) + go func() { + defer wg.Done() + for candidate := range jobs { + resultMu.Lock() + stopped := firstErr != nil + resultMu.Unlock() + if stopped { + continue + } + + upload, found, legacy, err := er.readMultipartUploadCandidate(ctx, bucket, candidate) + if errors.Is(err, errFileNotFound) || errors.Is(err, errFileVersionNotFound) { + continue + } + resultMu.Lock() + switch { + case err != nil && firstErr == nil: + firstErr = err + case legacy: + keyless = true + case found: + uploads = append(uploads, upload) + } + resultMu.Unlock() + } + }() + } + for _, candidate := range candidates { + jobs <- candidate + } + close(jobs) + wg.Wait() + return uploads, keyless, firstErr +} + +type multipartListEntry struct { + upload *MultipartInfo + commonPrefix string +} + +// paginateMultipartUploads applies the S3 ordering, prefix, delimiter, marker, +// and page rules exactly once after all pools and sets have been merged. +func paginateMultipartUploads(uploads []MultipartInfo, prefix, keyMarker, uploadIDMarker, delimiter string, maxUploads int) ListMultipartsInfo { + if maxUploads > maxUploadsList { + maxUploads = maxUploadsList + } + result := ListMultipartsInfo{ + MaxUploads: maxUploads, + KeyMarker: keyMarker, + UploadIDMarker: uploadIDMarker, + Prefix: prefix, + Delimiter: delimiter, + } + + deduplicated := make([]MultipartInfo, 0, len(uploads)) + seenUploads := make(map[string]struct{}, len(uploads)) + for _, upload := range uploads { + identity := upload.Bucket + "\x00" + upload.Object + "\x00" + upload.UploadID + if _, ok := seenUploads[identity]; ok { + continue + } + seenUploads[identity] = struct{}{} + deduplicated = append(deduplicated, upload) + } + sort.Slice(deduplicated, func(i, j int) bool { + if deduplicated[i].Object != deduplicated[j].Object { + return deduplicated[i].Object < deduplicated[j].Object + } + if !deduplicated[i].Initiated.Equal(deduplicated[j].Initiated) { + return deduplicated[i].Initiated.Before(deduplicated[j].Initiated) + } + return deduplicated[i].UploadID < deduplicated[j].UploadID + }) + + markerPassed := keyMarker == "" + seenPrefixes := make(map[string]struct{}) + entries := make([]multipartListEntry, 0, len(deduplicated)) + for i := range deduplicated { + upload := &deduplicated[i] + if !strings.HasPrefix(upload.Object, prefix) { + continue + } + if !markerPassed { + switch strings.Compare(upload.Object, keyMarker) { + case -1: + continue + case 0: + if uploadIDMarker != "" && upload.UploadID == uploadIDMarker { + markerPassed = true + } + continue + default: + markerPassed = true + } + } + + if delimiter != "" { + remainder := strings.TrimPrefix(upload.Object, prefix) + if i := strings.Index(remainder, delimiter); i >= 0 { + commonPrefix := prefix + remainder[:i+len(delimiter)] + if keyMarker != "" && commonPrefix <= keyMarker { + continue + } + if _, ok := seenPrefixes[commonPrefix]; ok { + continue + } + seenPrefixes[commonPrefix] = struct{}{} + entries = append(entries, multipartListEntry{commonPrefix: commonPrefix}) + continue + } + } + entries = append(entries, multipartListEntry{upload: upload}) + } + + if maxUploads <= 0 { + return result + } + pageSize := min(maxUploads, len(entries)) + for _, entry := range entries[:pageSize] { + if entry.upload != nil { + result.Uploads = append(result.Uploads, *entry.upload) + continue + } + result.CommonPrefixes = append(result.CommonPrefixes, entry.commonPrefix) + } + result.IsTruncated = pageSize < len(entries) + if result.IsTruncated && pageSize > 0 { + last := entries[pageSize-1] + if last.upload != nil { + result.NextKeyMarker = last.upload.Object + result.NextUploadIDMarker = last.upload.UploadID + } else { + result.NextKeyMarker = last.commonPrefix + } + } + return result +} + +// listMultipartUploadsExact preserves the hashed exact-object lookup used by +// multipart write placement and by rolling-upgrade legacy mode. +func (er erasureObjects) listMultipartUploadsExact(ctx context.Context, bucket, object, keyMarker, uploadIDMarker, delimiter string, maxUploads int) (result ListMultipartsInfo, err error) { auditObjectErasureSet(ctx, "ListMultipartUploads", object, &er) result.MaxUploads = maxUploads @@ -311,20 +584,7 @@ func (er erasureObjects) ListMultipartUploads(ctx context.Context, bucket, objec if populatedUploadIDs.Contains(uploadID) { continue } - // If present, use time stored in ID. - startTime := time.Now() - if split := strings.Split(uploadID, "x"); len(split) == 2 { - t, err := strconv.ParseInt(split[1], 10, 64) - if err == nil { - startTime = time.Unix(0, t) - } - } - uploads = append(uploads, MultipartInfo{ - Bucket: bucket, - Object: object, - UploadID: base64.RawURLEncoding.EncodeToString(fmt.Appendf(nil, "%s.%s", globalDeploymentID(), uploadID)), - Initiated: startTime, - }) + uploads = append(uploads, multipartUploadInfo(bucket, object, uploadID, time.Time{})) populatedUploadIDs.Add(uploadID) } @@ -365,6 +625,17 @@ func (er erasureObjects) ListMultipartUploads(ctx context.Context, bucket, objec return result, nil } +func (er erasureObjects) ListMultipartUploads(ctx context.Context, bucket, prefix, keyMarker, uploadIDMarker, delimiter string, maxUploads int) (ListMultipartsInfo, error) { + if err := checkListMultipartArgs(ctx, bucket, prefix, keyMarker, uploadIDMarker, delimiter); err != nil { + return ListMultipartsInfo{}, err + } + uploads, _, err := er.scanMultipartUploads(ctx, bucket) + if err != nil { + return ListMultipartsInfo{}, err + } + return paginateMultipartUploads(uploads, prefix, keyMarker, uploadIDMarker, delimiter, maxUploads), nil +} + // newMultipartUpload - wrapper for initializing a new multipart // request; returns a unique upload id. // @@ -402,6 +673,8 @@ func (er erasureObjects) newMultipartUpload(ctx context.Context, bucket string, } userDefined := cloneMSS(opts.UserDefined) + userDefined[multipartMetaBucket] = bucket + userDefined[multipartMetaObject] = object if opts.PreserveETag != "" { userDefined["etag"] = opts.PreserveETag } @@ -1459,6 +1732,8 @@ func (er erasureObjects) CompleteMultipartUpload(ctx context.Context, bucket str // Remove superfluous internal headers. delete(fi.Metadata, hash.MinIOMultipartChecksum) delete(fi.Metadata, hash.MinIOMultipartChecksumType) + delete(fi.Metadata, multipartMetaBucket) + delete(fi.Metadata, multipartMetaObject) // Save the final object size and modtime. fi.Size = objectSize diff --git a/cmd/erasure-server-pool.go b/cmd/erasure-server-pool.go index 555986bc0..ff55fcb9b 100644 --- a/cmd/erasure-server-pool.go +++ b/cmd/erasure-server-pool.go @@ -1857,7 +1857,34 @@ func (z *erasureServerPools) ListMultipartUploads(ctx context.Context, bucket, p if err := checkListMultipartArgs(ctx, bucket, prefix, keyMarker, uploadIDMarker, delimiter); err != nil { return ListMultipartsInfo{}, err } + if _, err := z.GetBucketInfo(ctx, bucket, BucketOptions{}); err != nil { + return ListMultipartsInfo{}, toObjectErr(err, bucket) + } + var uploads []MultipartInfo + var keyless bool + for idx, pool := range z.serverPools { + if z.IsSuspended(idx) { + continue + } + poolUploads, poolKeyless, err := pool.scanMultipartUploads(ctx, bucket) + if err != nil { + return ListMultipartsInfo{}, err + } + uploads = append(uploads, poolUploads...) + keyless = keyless || poolKeyless + } + + // Old writers did not persist the bucket and object key. Until every such + // upload has drained, retain the old response behavior instead of silently + // claiming that a partial durable scan is complete. + if keyless { + return z.listMultipartUploadsLegacy(ctx, bucket, prefix, keyMarker, uploadIDMarker, delimiter, maxUploads) + } + return paginateMultipartUploads(uploads, prefix, keyMarker, uploadIDMarker, delimiter, maxUploads), nil +} + +func (z *erasureServerPools) listMultipartUploadsLegacy(ctx context.Context, bucket, prefix, keyMarker, uploadIDMarker, delimiter string, maxUploads int) (ListMultipartsInfo, error) { poolResult := ListMultipartsInfo{} poolResult.MaxUploads = maxUploads poolResult.KeyMarker = keyMarker @@ -1883,15 +1910,14 @@ func (z *erasureServerPools) ListMultipartUploads(ctx context.Context, bucket, p } if z.SinglePool() { - return z.serverPools[0].ListMultipartUploads(ctx, bucket, prefix, keyMarker, uploadIDMarker, delimiter, maxUploads) + return z.serverPools[0].getHashedSet(prefix).listMultipartUploadsExact(ctx, bucket, prefix, keyMarker, uploadIDMarker, delimiter, maxUploads) } for idx, pool := range z.serverPools { if z.IsSuspended(idx) { continue } - result, err := pool.ListMultipartUploads(ctx, bucket, prefix, keyMarker, uploadIDMarker, - delimiter, maxUploads) + result, err := pool.getHashedSet(prefix).listMultipartUploadsExact(ctx, bucket, prefix, keyMarker, uploadIDMarker, delimiter, maxUploads) if err != nil { return result, err } @@ -1927,7 +1953,7 @@ func (z *erasureServerPools) NewMultipartUpload(ctx context.Context, bucket, obj continue } - result, err := pool.ListMultipartUploads(ctx, bucket, object, "", "", "", maxUploadsList) + result, err := pool.listMultipartUploadsExact(ctx, bucket, object) if err != nil { return nil, err } diff --git a/cmd/erasure-sets.go b/cmd/erasure-sets.go index 6ad9ece65..20b730451 100644 --- a/cmd/erasure-sets.go +++ b/cmd/erasure-sets.go @@ -880,10 +880,32 @@ func (s *erasureSets) CopyObject(ctx context.Context, srcBucket, srcObject, dstB } func (s *erasureSets) ListMultipartUploads(ctx context.Context, bucket, prefix, keyMarker, uploadIDMarker, delimiter string, maxUploads int) (result ListMultipartsInfo, err error) { - // In list multipart uploads we are going to treat input prefix as the object, - // this means that we are not supporting directory navigation. - set := s.getHashedSet(prefix) - return set.ListMultipartUploads(ctx, bucket, prefix, keyMarker, uploadIDMarker, delimiter, maxUploads) + if err := checkListMultipartArgs(ctx, bucket, prefix, keyMarker, uploadIDMarker, delimiter); err != nil { + return ListMultipartsInfo{}, err + } + uploads, _, err := s.scanMultipartUploads(ctx, bucket) + if err != nil { + return ListMultipartsInfo{}, err + } + return paginateMultipartUploads(uploads, prefix, keyMarker, uploadIDMarker, delimiter, maxUploads), nil +} + +func (s *erasureSets) scanMultipartUploads(ctx context.Context, bucket string) ([]MultipartInfo, bool, error) { + var uploads []MultipartInfo + var keyless bool + for _, set := range s.sets { + setUploads, setKeyless, err := set.scanMultipartUploads(ctx, bucket) + if err != nil { + return nil, false, err + } + uploads = append(uploads, setUploads...) + keyless = keyless || setKeyless + } + return uploads, keyless, nil +} + +func (s *erasureSets) listMultipartUploadsExact(ctx context.Context, bucket, object string) (ListMultipartsInfo, error) { + return s.getHashedSet(object).listMultipartUploadsExact(ctx, bucket, object, "", "", "", maxUploadsList) } // Initiate a new multipart upload on a hashedSet based on object name. diff --git a/cmd/list-multipart-uploads-compat_test.go b/cmd/list-multipart-uploads-compat_test.go new file mode 100644 index 000000000..db0f49d6b --- /dev/null +++ b/cmd/list-multipart-uploads-compat_test.go @@ -0,0 +1,283 @@ +// Copyright (c) 2026 mr javad seydi +// +// This file is part of Silo Object Storage stack +// +// This program is free software: you can redistribute it and/or modify +// it under the terms of the GNU Affero General Public License as published by +// the Free Software Foundation, either version 3 of the License, or +// (at your option) any later version. +// +// This program is distributed in the hope that it will be useful +// but WITHOUT ANY WARRANTY; without even the implied warranty of +// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +// GNU Affero General Public License for more details. +// +// You should have received a copy of the GNU Affero General Public License +// along with this program. If not, see . + +package cmd + +import ( + "bytes" + "fmt" + "slices" + "testing" + "time" +) + +func multipartUploadKeys(uploads []MultipartInfo) []string { + keys := make([]string, len(uploads)) + for i := range uploads { + keys[i] = uploads[i].Object + } + return keys +} + +func requireMultipartUploadKeys(t *testing.T, got ListMultipartsInfo, want ...string) { + t.Helper() + if keys := multipartUploadKeys(got.Uploads); !slices.Equal(keys, want) { + t.Fatalf("uploads = %v, want %v", keys, want) + } +} + +func TestListMultipartUploadsS3Compatibility(t *testing.T) { + obj, dirs, err := prepareErasureSets32(t.Context()) + if err != nil { + t.Fatal(err) + } + z := obj.(*erasureServerPools) + t.Cleanup(func() { + z.Shutdown(t.Context()) + removeRoots(dirs) + }) + + const bucket = "multipart-list-compat" + if err = z.MakeBucket(t.Context(), bucket, MakeBucketOptions{}); err != nil { + t.Fatal(err) + } + + objects := []string{"t/a_b/p1", "t/a_b/p2", "t/c_d/p1", "u/x"} + sets := z.serverPools[0] + firstSet := sets.getHashedSetIndex(objects[0]) + if !slices.ContainsFunc(objects[1:], func(object string) bool { + return sets.getHashedSetIndex(object) != firstSet + }) { + for n := 0; ; n++ { + object := fmt.Sprintf("v/cross-set-%d", n) + if sets.getHashedSetIndex(object) != firstSet { + objects = append(objects, object) + break + } + } + } + + uploadIDs := make(map[string]string, len(objects)) + for _, object := range objects { + mp, err := z.NewMultipartUpload(t.Context(), bucket, object, ObjectOptions{}) + if err != nil { + t.Fatalf("NewMultipartUpload(%q): %v", object, err) + } + uploadIDs[object] = mp.UploadID + } + + // Durable multipart metadata, rather than this node-local cache, must be + // authoritative after a restart or when another node handles the request. + z.mpCache.Range(func(uploadID string, _ MultipartInfo) bool { + z.mpCache.Delete(uploadID) + return true + }) + + all, err := z.ListMultipartUploads(t.Context(), bucket, "", "", "", "", 100) + if err != nil { + t.Fatal(err) + } + requireMultipartUploadKeys(t, all, objects...) + if all.IsTruncated || all.NextKeyMarker != "" || all.NextUploadIDMarker != "" { + t.Fatalf("complete listing has truncation state: %+v", all) + } + + first, err := z.ListMultipartUploads(t.Context(), bucket, "", "", "", "", 1) + if err != nil { + t.Fatal(err) + } + requireMultipartUploadKeys(t, first, objects[0]) + if !first.IsTruncated || first.NextKeyMarker != objects[0] || first.NextUploadIDMarker != uploadIDs[objects[0]] { + t.Fatalf("first page markers = (%q, %q, %t), want (%q, %q, true)", + first.NextKeyMarker, first.NextUploadIDMarker, first.IsTruncated, + objects[0], uploadIDs[objects[0]]) + } + + rest, err := z.ListMultipartUploads(t.Context(), bucket, "", first.NextKeyMarker, first.NextUploadIDMarker, "", 100) + if err != nil { + t.Fatal(err) + } + requireMultipartUploadKeys(t, rest, objects[1:]...) + + afterKey, err := z.ListMultipartUploads(t.Context(), bucket, "", objects[1], "", "", 100) + if err != nil { + t.Fatal(err) + } + requireMultipartUploadKeys(t, afterKey, objects[2:]...) + + prefixed, err := z.ListMultipartUploads(t.Context(), bucket, "t/", "", "", "", 100) + if err != nil { + t.Fatal(err) + } + requireMultipartUploadKeys(t, prefixed, objects[:3]...) + + nested, err := z.ListMultipartUploads(t.Context(), bucket, "t/a_b/", "", "", "", 100) + if err != nil { + t.Fatal(err) + } + requireMultipartUploadKeys(t, nested, objects[:2]...) + + grouped, err := z.ListMultipartUploads(t.Context(), bucket, "t/", "", "", SlashSeparator, 100) + if err != nil { + t.Fatal(err) + } + requireMultipartUploadKeys(t, grouped) + if want := []string{"t/a_b/", "t/c_d/"}; !slices.Equal(grouped.CommonPrefixes, want) { + t.Fatalf("common prefixes = %v, want %v", grouped.CommonPrefixes, want) + } + + groupPage, err := z.ListMultipartUploads(t.Context(), bucket, "t/", "", "", SlashSeparator, 1) + if err != nil { + t.Fatal(err) + } + if want := []string{"t/a_b/"}; !slices.Equal(groupPage.CommonPrefixes, want) { + t.Fatalf("first common-prefix page = %v, want %v", groupPage.CommonPrefixes, want) + } + if !groupPage.IsTruncated || groupPage.NextKeyMarker != "t/a_b/" || groupPage.NextUploadIDMarker != "" { + t.Fatalf("common-prefix page markers = (%q, %q, %t)", + groupPage.NextKeyMarker, groupPage.NextUploadIDMarker, groupPage.IsTruncated) + } + + groupRest, err := z.ListMultipartUploads(t.Context(), bucket, "t/", groupPage.NextKeyMarker, "", SlashSeparator, 1) + if err != nil { + t.Fatal(err) + } + if want := []string{"t/c_d/"}; !slices.Equal(groupRest.CommonPrefixes, want) { + t.Fatalf("second common-prefix page = %v, want %v", groupRest.CommonPrefixes, want) + } + if groupRest.IsTruncated { + t.Fatalf("last common-prefix page is truncated: %+v", groupRest) + } + + part, err := z.PutObjectPart(t.Context(), bucket, objects[0], uploadIDs[objects[0]], 1, + mustGetPutObjReader(t, bytes.NewBufferString("part"), 4, "", ""), ObjectOptions{}) + if err != nil { + t.Fatal(err) + } + completed, err := z.CompleteMultipartUpload(t.Context(), bucket, objects[0], uploadIDs[objects[0]], + []CompletePart{{PartNumber: 1, ETag: part.ETag}}, ObjectOptions{}) + if err != nil { + t.Fatal(err) + } + for _, key := range []string{multipartMetaBucket, multipartMetaObject} { + if _, ok := completed.UserDefined[key]; ok { + t.Errorf("completed object retained upload-only metadata %q", key) + } + } + + // Simulate an upload written by a pre-upgrade server. Its key cannot be + // recovered by scanning the hashed namespace, so detection must retain the + // exact-key legacy path until such uploads have drained. + legacyObject := objects[1] + er := sets.getHashedSet(legacyObject) + fi, metadata, err := er.checkUploadIDExists(t.Context(), bucket, legacyObject, uploadIDs[legacyObject], true) + if err != nil { + t.Fatal(err) + } + for i := range metadata { + delete(metadata[i].Metadata, multipartMetaBucket) + delete(metadata[i].Metadata, multipartMetaObject) + } + if _, err = writeAllMetadata(t.Context(), er.getDisks(), bucket, minioMetaMultipartBucket, + er.getUploadIDDir(bucket, legacyObject, uploadIDs[legacyObject]), metadata, fi.WriteQuorum(er.defaultWQuorum())); err != nil { + t.Fatal(err) + } + legacy, err := z.ListMultipartUploads(t.Context(), bucket, legacyObject, "", "", "", 100) + if err != nil { + t.Fatal(err) + } + requireMultipartUploadKeys(t, legacy, legacyObject) +} + +func TestPaginateMultipartUploads(t *testing.T) { + base := time.Unix(100, 0) + uploads := []MultipartInfo{ + {Bucket: "bucket", Object: "b", UploadID: "b1", Initiated: base}, + {Bucket: "bucket", Object: "a", UploadID: "a2", Initiated: base.Add(time.Second)}, + {Bucket: "bucket", Object: "a", UploadID: "a1", Initiated: base}, + {Bucket: "bucket", Object: "a", UploadID: "a1", Initiated: base}, // duplicate discovery + } + + first := paginateMultipartUploads(uploads, "", "", "", "", 1) + requireMultipartUploadKeys(t, first, "a") + if !first.IsTruncated || first.NextKeyMarker != "a" || first.NextUploadIDMarker != "a1" { + t.Fatalf("first page = %+v", first) + } + + second := paginateMultipartUploads(uploads, "", first.NextKeyMarker, first.NextUploadIDMarker, "", 1) + if len(second.Uploads) != 1 || second.Uploads[0].Object != "a" || second.Uploads[0].UploadID != "a2" { + t.Fatalf("second page uploads = %+v", second.Uploads) + } + if !second.IsTruncated || second.NextKeyMarker != "a" || second.NextUploadIDMarker != "a2" { + t.Fatalf("second page = %+v", second) + } + + last := paginateMultipartUploads(uploads, "", second.NextKeyMarker, second.NextUploadIDMarker, "", 1) + requireMultipartUploadKeys(t, last, "b") + if last.IsTruncated || last.NextKeyMarker != "" || last.NextUploadIDMarker != "" { + t.Fatalf("last page = %+v", last) + } + + missingUploadMarker := paginateMultipartUploads(uploads, "", "a", "missing", "", 10) + requireMultipartUploadKeys(t, missingUploadMarker, "b") + + if err := checkListMultipartArgs(t.Context(), "bucket", "", "", "not-base64=", ""); err != nil { + t.Fatalf("upload-id-marker without key-marker must be ignored: %v", err) + } + + overLimit := make([]MultipartInfo, maxUploadsList+1) + for i := range overLimit { + overLimit[i] = MultipartInfo{Bucket: "bucket", Object: fmt.Sprintf("%04d", i), UploadID: fmt.Sprint(i)} + } + capped := paginateMultipartUploads(overLimit, "", "", "", "", maxUploadsList+1) + if capped.MaxUploads != maxUploadsList || len(capped.Uploads) != maxUploadsList || !capped.IsTruncated { + t.Fatalf("over-limit page = MaxUploads %d, uploads %d, truncated %t", + capped.MaxUploads, len(capped.Uploads), capped.IsTruncated) + } +} + +func TestListMultipartUploadsGlobalPageAcrossPools(t *testing.T) { + z, bucket := consistencyPools(t) + objects := []string{"a/one", "b/two", "c/three", "d/four"} + for i, object := range objects { + if _, err := z.serverPools[i%len(z.serverPools)].NewMultipartUpload(t.Context(), bucket, object, ObjectOptions{}); err != nil { + t.Fatalf("NewMultipartUpload(%q): %v", object, err) + } + } + z.mpCache.Range(func(uploadID string, _ MultipartInfo) bool { + z.mpCache.Delete(uploadID) + return true + }) + + page, err := z.ListMultipartUploads(t.Context(), bucket, "", "", "", "", 2) + if err != nil { + t.Fatal(err) + } + requireMultipartUploadKeys(t, page, objects[:2]...) + if !page.IsTruncated || page.NextKeyMarker != objects[1] { + t.Fatalf("first global page = %+v", page) + } + + rest, err := z.ListMultipartUploads(t.Context(), bucket, "", page.NextKeyMarker, page.NextUploadIDMarker, "", 2) + if err != nil { + t.Fatal(err) + } + requireMultipartUploadKeys(t, rest, objects[2:]...) + if rest.IsTruncated { + t.Fatalf("last global page is truncated: %+v", rest) + } +} diff --git a/cmd/object-api-input-checks.go b/cmd/object-api-input-checks.go index 9c8b213e4..7ffdd617b 100644 --- a/cmd/object-api-input-checks.go +++ b/cmd/object-api-input-checks.go @@ -83,7 +83,8 @@ func checkListMultipartArgs(ctx context.Context, bucket, prefix, keyMarker, uplo if err := checkListObjsArgs(ctx, bucket, prefix, keyMarker); err != nil { return err } - if uploadIDMarker != "" { + // S3 ignores upload-id-marker when key-marker is absent. + if uploadIDMarker != "" && keyMarker != "" { if HasSuffix(keyMarker, SlashSeparator) { return InvalidUploadIDKeyCombination{ UploadIDMarker: uploadIDMarker,