mirror of
https://github.com/pgsty/minio.git
synced 2026-09-16 23:44:06 +03:00
fix: make multipart upload listing S3-compatible
Signed-off-by: mr javad seydi <seydi.birjand@gmail.com>
This commit is contained in:
@@ -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
|
||||
|
||||
+1
-1
@@ -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.
|
||||
)
|
||||
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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.
|
||||
|
||||
+297
-22
@@ -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
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
+26
-4
@@ -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.
|
||||
|
||||
@@ -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 <http://www.gnu.org/licenses/>.
|
||||
|
||||
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)
|
||||
}
|
||||
}
|
||||
@@ -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,
|
||||
|
||||
Reference in New Issue
Block a user