diff --git a/test/s3/versioning/s3_copy_versioning_regression_test.go b/test/s3/versioning/s3_copy_versioning_regression_test.go index c92d1eeb9..074aa283a 100644 --- a/test/s3/versioning/s3_copy_versioning_regression_test.go +++ b/test/s3/versioning/s3_copy_versioning_regression_test.go @@ -250,3 +250,59 @@ func TestSelfCopyWithSuspendedVersioningIsRejected(t *testing.T) { assert.Equal(t, "InvalidRequest", apiErr.ErrorCode()) } } + +// A suspended-versioning CopyObject writes the null version at the regular path, so +// like PutObject and multipart completion it has to retire the null delete marker a +// preceding DELETE left in .versions. While the regular-path object owns the null +// slot the leftover marker is shadowed, but it resurfaces as a phantom delete the +// moment that null version goes away. +func TestSuspendedCopyRetiresDeleteMarker(t *testing.T) { + client := getS3Client(t) + bucketName := getNewBucketName() + + createBucket(t, client, bucketName) + defer deleteBucket(t, client, bucketName) + + sourceKey := "suspended-copy-source.txt" + objectKey := "suspended-copy-dest.txt" + + enableVersioning(t, client, bucketName) + putObject(t, client, bucketName, objectKey, "pre-suspension-content") + suspendVersioning(t, client, bucketName) + + putObject(t, client, bucketName, sourceKey, "source-content") + putObject(t, client, bucketName, objectKey, "null-version-content") + deleteKey(t, client, bucketName, objectKey) + + _, err := client.CopyObject(context.TODO(), &s3.CopyObjectInput{ + Bucket: aws.String(bucketName), + Key: aws.String(objectKey), + CopySource: aws.String(versioningCopySource(bucketName, sourceKey)), + }) + require.NoError(t, err) + + getResp, err := client.GetObject(context.TODO(), &s3.GetObjectInput{ + Bucket: aws.String(bucketName), + Key: aws.String(objectKey), + }) + require.NoError(t, err) + defer getResp.Body.Close() + body, err := io.ReadAll(getResp.Body) + require.NoError(t, err) + assert.Equal(t, "source-content", string(body)) + + // Drop the null version the copy just wrote; a retired marker leaves nothing behind. + _, err = client.DeleteObject(context.TODO(), &s3.DeleteObjectInput{ + Bucket: aws.String(bucketName), + Key: aws.String(objectKey), + VersionId: aws.String("null"), + }) + require.NoError(t, err) + + listResp, err := client.ListObjectVersions(context.TODO(), &s3.ListObjectVersionsInput{ + Bucket: aws.String(bucketName), + Prefix: aws.String(objectKey), + }) + require.NoError(t, err) + assert.Empty(t, listResp.DeleteMarkers, "the copy should have retired the null delete marker") +} diff --git a/test/s3/versioning/s3_suspended_delete_marker_regression_test.go b/test/s3/versioning/s3_suspended_delete_marker_regression_test.go index a2d670289..ce7d128a6 100644 --- a/test/s3/versioning/s3_suspended_delete_marker_regression_test.go +++ b/test/s3/versioning/s3_suspended_delete_marker_regression_test.go @@ -8,6 +8,7 @@ import ( "github.com/aws/aws-sdk-go-v2/aws" "github.com/aws/aws-sdk-go-v2/service/s3" + "github.com/aws/aws-sdk-go-v2/service/s3/types" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" ) @@ -80,3 +81,93 @@ func TestSuspendedDeleteCreatesDeleteMarker(t *testing.T) { require.NoError(t, err) assert.Equal(t, "versioned-content", string(body)) } + +// A suspended-versioning completion must drop the null delete marker a preceding +// DELETE left in .versions, or it reports 200 and the object lists while HEAD/GET +// keep resolving the marker and answer NoSuchKey. +func TestSuspendedMultipartOverwritesDeleteMarker(t *testing.T) { + client := getS3Client(t) + bucketName := getNewBucketName() + + createBucket(t, client, bucketName) + defer deleteBucket(t, client, bucketName) + + objectKey := "suspended-multipart-after-delete.bin" + partData := bytes.Repeat([]byte("a"), 5*1024*1024) + + // The cleanup retires only the null version, never the key's real history. + enableVersioning(t, client, bucketName) + realVersion := putObject(t, client, bucketName, objectKey, "pre-suspension-content") + require.NotNil(t, realVersion.VersionId) + suspendVersioning(t, client, bucketName) + + completeSuspendedMultipart := func() { + t.Helper() + createResp, err := client.CreateMultipartUpload(context.TODO(), &s3.CreateMultipartUploadInput{ + Bucket: aws.String(bucketName), + Key: aws.String(objectKey), + }) + require.NoError(t, err) + + uploadResp, err := client.UploadPart(context.TODO(), &s3.UploadPartInput{ + Bucket: aws.String(bucketName), + Key: aws.String(objectKey), + UploadId: createResp.UploadId, + PartNumber: aws.Int32(1), + Body: bytes.NewReader(partData), + }) + require.NoError(t, err) + + _, err = client.CompleteMultipartUpload(context.TODO(), &s3.CompleteMultipartUploadInput{ + Bucket: aws.String(bucketName), + Key: aws.String(objectKey), + UploadId: createResp.UploadId, + MultipartUpload: &types.CompletedMultipartUpload{ + Parts: []types.CompletedPart{{ETag: uploadResp.ETag, PartNumber: aws.Int32(1)}}, + }, + }) + require.NoError(t, err) + } + + completeSuspendedMultipart() + deleteKey(t, client, bucketName, objectKey) + completeSuspendedMultipart() + + headResp := headObject(t, client, bucketName, objectKey) + require.NotNil(t, headResp.ContentLength) + assert.Equal(t, int64(len(partData)), *headResp.ContentLength) + + getResp, err := client.GetObject(context.TODO(), &s3.GetObjectInput{ + Bucket: aws.String(bucketName), + Key: aws.String(objectKey), + }) + require.NoError(t, err) + defer getResp.Body.Close() + written, err := io.Copy(io.Discard, getResp.Body) + require.NoError(t, err) + assert.Equal(t, int64(len(partData)), written) + + listResp, err := client.ListObjectVersions(context.TODO(), &s3.ListObjectVersionsInput{ + Bucket: aws.String(bucketName), + Prefix: aws.String(objectKey), + }) + require.NoError(t, err) + assert.Empty(t, listResp.DeleteMarkers) + + var listedVersionIds []string + for _, version := range listResp.Versions { + listedVersionIds = append(listedVersionIds, aws.ToString(version.VersionId)) + } + assert.ElementsMatch(t, []string{*realVersion.VersionId, "null"}, listedVersionIds) + + realVersionResp, err := client.GetObject(context.TODO(), &s3.GetObjectInput{ + Bucket: aws.String(bucketName), + Key: aws.String(objectKey), + VersionId: realVersion.VersionId, + }) + require.NoError(t, err) + defer realVersionResp.Body.Close() + realVersionBody, err := io.ReadAll(realVersionResp.Body) + require.NoError(t, err) + assert.Equal(t, "pre-suspension-content", string(realVersionBody)) +} diff --git a/weed/s3api/filer_multipart.go b/weed/s3api/filer_multipart.go index 8af6a5dac..bfa2eb3ca 100644 --- a/weed/s3api/filer_multipart.go +++ b/weed/s3api/filer_multipart.go @@ -649,9 +649,10 @@ func (s3a *S3ApiServer) completeMultipartUpload(r *http.Request, input *s3.Compl glog.Errorf("completeMultipartUpload: failed to get versioning state for bucket %s: %v", *input.Bucket, vErr) return s3err.ErrInternalError } + // Full object key, not just entryName, so the right .versions directory is used. + normalizedKey := s3_constants.NormalizeObjectKey(*input.Key) + if versioningState == s3_constants.VersioningEnabled { - // Use full object key (not just entryName) to ensure correct .versions directory is checked - normalizedKey := strings.TrimPrefix(*input.Key, "/") useInvertedFormat := s3a.getVersionIdFormat(*input.Bucket, normalizedKey) versionId := generateVersionId(useInvertedFormat) versionFileName := s3a.getVersionFileName(versionId) @@ -823,6 +824,14 @@ func (s3a *S3ApiServer) completeMultipartUpload(r *http.Request, input *s3.Compl return s3err.ErrInternalError } + // A failed finalize leaves the key reading as deleted, so fail rather than + // return 200 — a non-ErrNone finalize keeps the upload directory, so the + // caller's retry replays. + if err := s3a.finalizeSuspendedNullWrite(owner, *input.Bucket, normalizedKey, s3_constants.SeaweedFSUploadId, *input.UploadId); err != nil { + glog.Errorf("completeMultipartUpload: failed to retire the null delete marker for %s/%s: %v", *input.Bucket, normalizedKey, err) + return s3err.ErrInternalError + } + // Note: Suspended versioning should NOT return VersionId field according to AWS S3 spec output = &CompleteMultipartUploadResult{ Location: aws.String(fmt.Sprintf("%s://%s/%s/%s", getRequestScheme(r), r.Host, url.PathEscape(*input.Bucket), urlPathEscape(*input.Key))), diff --git a/weed/s3api/s3api_object_handlers_copy.go b/weed/s3api/s3api_object_handlers_copy.go index a35eca72b..5ed64ec86 100644 --- a/weed/s3api/s3api_object_handlers_copy.go +++ b/weed/s3api/s3api_object_handlers_copy.go @@ -583,8 +583,9 @@ func (s3a *S3ApiServer) finalizeCopyDestination(dstBucket, dstObject, dstVersion return "", "", err } - if err = s3a.updateIsLatestFlagsForSuspendedVersioning(dstBucket, normalizedObject); err != nil { - glog.Warningf("CopyObjectHandler: failed to update suspended version latest flags for %s/%s: %v", dstBucket, normalizedObject, err) + // mkFile writes through the default filer, so the ownership check reads there too. + if err = s3a.finalizeSuspendedNullWrite("", dstBucket, normalizedObject, s3_constants.ExtETagKey, etag); err != nil { + glog.Warningf("CopyObjectHandler: failed to retire the null delete marker for %s/%s: %v", dstBucket, normalizedObject, err) } return "", etag, nil diff --git a/weed/s3api/s3api_object_routed_read.go b/weed/s3api/s3api_object_routed_read.go index 1dc060726..3bfff77b9 100644 --- a/weed/s3api/s3api_object_routed_read.go +++ b/weed/s3api/s3api_object_routed_read.go @@ -80,6 +80,16 @@ func (s3a *S3ApiServer) ownerRecentlyUnreachable(owner pb.ServerAddress) bool { return false } +// lookupEntryPreferringOwner reads an entry back from a known write owner, so a +// caller that just wrote there sees its own write. Unlike getObjectEntryRoutedByKey +// it never drops the owner for a healthier peer, which would read behind the write. +func (s3a *S3ApiServer) lookupEntryPreferringOwner(owner pb.ServerAddress, dir, name string) (*filer_pb.Entry, error) { + if owner == "" { + return s3a.getEntry(dir, name) + } + return s3a.lookupEntryOnFiler(owner, dir, name) +} + // lookupEntryOnFiler resolves dir/name against a single filer, without failover. func (s3a *S3ApiServer) lookupEntryOnFiler(filer pb.ServerAddress, dir, name string) (*filer_pb.Entry, error) { var entry *filer_pb.Entry diff --git a/weed/s3api/s3api_object_versioned_finalize.go b/weed/s3api/s3api_object_versioned_finalize.go index 354946423..ea0393c8b 100644 --- a/weed/s3api/s3api_object_versioned_finalize.go +++ b/weed/s3api/s3api_object_versioned_finalize.go @@ -1,6 +1,8 @@ package s3api import ( + "errors" + "fmt" "strconv" "time" @@ -194,3 +196,37 @@ func (s3a *S3ApiServer) versionedFinalize(bucket, object, versionId, versionFile }, } } + +// finalizeSuspendedNullWrite retires the null delete marker a suspended DELETE left +// in .versions, so reads resolve the null version the caller just wrote at the +// regular path. Pointer first: clearing the marker while the pointer still names it +// makes reads rescan .versions and promote an older version. Call only once the +// write has committed — retiring the marker for a write that then fails republishes +// the deleted key. +// +// identityKey/identityValue name the extended attribute that marks the entry as the +// caller's write (an upload id, an etag). The cleanup rewrites shared .versions state +// off the object write lock, so it is skipped unless the regular path still holds that +// write: a DELETE that landed in between owns the null slot, and retiring its marker +// would resurrect an older version under a key that was deleted. Narrows that race, +// does not close it. owner, when set, is the filer the write went to, so the check +// reads its own write back rather than a peer that may be behind. +func (s3a *S3ApiServer) finalizeSuspendedNullWrite(owner pb.ServerAddress, bucket, object, identityKey, identityValue string) error { + dir, name := util.FullPath(s3a.toFilerPath(bucket, object)).DirAndName() + current, err := s3a.lookupEntryPreferringOwner(owner, dir, name) + if err != nil && !errors.Is(err, filer_pb.ErrNotFound) { + return fmt.Errorf("re-read %s/%s: %w", bucket, object, err) + } + if current == nil || string(current.Extended[identityKey]) != identityValue { + glog.V(2).Infof("finalizeSuspendedNullWrite: %s/%s superseded by a concurrent write", bucket, object) + return nil + } + + if err := s3a.updateIsLatestFlagsForSuspendedVersioning(bucket, object); err != nil { + return err + } + // Best-effort: with the pointer gone the regular-path object already owns the + // null slot, so a surviving marker is neither read nor listed. + s3a.removeNullVersionFile(bucket, object) + return nil +}