diff --git a/backend/common.go b/backend/common.go index 1f421c7e..c4d9132a 100644 --- a/backend/common.go +++ b/backend/common.go @@ -102,6 +102,32 @@ func ParseRange(size int64, acceptRange string) (int64, int64, error) { return startOffset, endOffset - startOffset + 1, nil } +// ParseCopySource parses x-amz-copy-source header and returns source bucket, +// source object, versionId, error respectively +func ParseCopySource(copySourceHeader string) (string, string, string, error) { + if copySourceHeader[0] == '/' { + copySourceHeader = copySourceHeader[1:] + } + + cSplitted := strings.Split(copySourceHeader, "?") + copySource := cSplitted[0] + var versionId string + if len(cSplitted) > 1 { + versionIdParts := strings.Split(cSplitted[1], "=") + if len(versionIdParts) != 2 || versionIdParts[0] != "versionId" { + return "", "", "", s3err.GetAPIError(s3err.ErrInvalidRequest) + } + versionId = versionIdParts[1] + } + + srcBucket, srcObject, ok := strings.Cut(copySource, "/") + if !ok { + return "", "", "", s3err.GetAPIError(s3err.ErrInvalidCopySource) + } + + return srcBucket, srcObject, versionId, nil +} + func CreateExceedingRangeErr(objSize int64) s3err.APIError { return s3err.APIError{ Code: "InvalidArgument", diff --git a/backend/posix/posix.go b/backend/posix/posix.go index 0e45d369..5448b436 100644 --- a/backend/posix/posix.go +++ b/backend/posix/posix.go @@ -1173,11 +1173,37 @@ func (p *Posix) CompleteMultipartUpload(ctx context.Context, input *s3.CompleteM return nil, err } } + + vEnabled, err := p.isBucketVersioningEnabled(ctx, bucket) + if err != nil { + return nil, err + } + + d, err := os.Stat(objname) + + // if the versioninng is enabled first create the file object version + if p.versioningEnabled() && vEnabled && err == nil && !d.IsDir() { + _, err := p.createObjVersion(bucket, object, d.Size(), acct) + if err != nil { + return nil, fmt.Errorf("create object version: %w", err) + } + } + err = f.link() if err != nil { return nil, fmt.Errorf("link object in namespace: %w", err) } + // if the versioning is enabled, generate a new versionID for the object + var versionID string + if p.versioningEnabled() && vEnabled { + versionID = ulid.Make().String() + + if err := p.meta.StoreAttribute(bucket, object, versionIdKey, []byte(versionID)); err != nil { + return nil, fmt.Errorf("set versionId attr: %w", err) + } + } + for k, v := range userMetaData { err = p.meta.StoreAttribute(bucket, object, fmt.Sprintf("%v.%v", metaHdr, k), []byte(v)) if err != nil { @@ -1261,9 +1287,10 @@ func (p *Posix) CompleteMultipartUpload(ctx context.Context, input *s3.CompleteM os.Remove(filepath.Join(bucket, objdir)) return &s3.CompleteMultipartUploadOutput{ - Bucket: &bucket, - ETag: &s3MD5, - Key: &object, + Bucket: &bucket, + ETag: &s3MD5, + Key: &object, + VersionId: &versionID, }, nil } @@ -1758,14 +1785,11 @@ func (p *Posix) UploadPartCopy(ctx context.Context, upi *s3.UploadPartCopyInput) partPath := filepath.Join(objdir, *upi.UploadId, fmt.Sprintf("%v", *upi.PartNumber)) - substrs := strings.SplitN(*upi.CopySource, "/", 2) - if len(substrs) != 2 { - return s3response.CopyObjectResult{}, s3err.GetAPIError(s3err.ErrInvalidCopySource) + srcBucket, srcObject, srcVersionId, err := backend.ParseCopySource(*upi.CopySource) + if err != nil { + return s3response.CopyObjectResult{}, err } - srcBucket := substrs[0] - srcObject := substrs[1] - _, err = os.Stat(srcBucket) if errors.Is(err, fs.ErrNotExist) { return s3response.CopyObjectResult{}, s3err.GetAPIError(s3err.ErrNoSuchBucket) @@ -1774,9 +1798,35 @@ func (p *Posix) UploadPartCopy(ctx context.Context, upi *s3.UploadPartCopyInput) return s3response.CopyObjectResult{}, fmt.Errorf("stat bucket: %w", err) } + vEnabled, err := p.isBucketVersioningEnabled(ctx, srcBucket) + if err != nil { + return s3response.CopyObjectResult{}, err + } + + if srcVersionId != "" { + if !p.versioningEnabled() || !vEnabled { + return s3response.CopyObjectResult{}, s3err.GetAPIError(s3err.ErrInvalidVersionId) + } + vId, err := p.meta.RetrieveAttribute(srcBucket, srcObject, versionIdKey) + if errors.Is(err, fs.ErrNotExist) { + return s3response.CopyObjectResult{}, s3err.GetAPIError(s3err.ErrNoSuchKey) + } + if err != nil && !errors.Is(err, meta.ErrNoSuchKey) { + return s3response.CopyObjectResult{}, fmt.Errorf("get src object version id: %w", err) + } + + if string(vId) != srcVersionId { + srcBucket = filepath.Join(p.versioningDir, srcBucket) + srcObject = filepath.Join(genObjVersionKey(srcObject), srcVersionId) + } + } + objPath := filepath.Join(srcBucket, srcObject) fi, err := os.Stat(objPath) if errors.Is(err, fs.ErrNotExist) { + if p.versioningEnabled() && vEnabled { + return s3response.CopyObjectResult{}, s3err.GetAPIError(s3err.ErrNoSuchVersion) + } return s3response.CopyObjectResult{}, s3err.GetAPIError(s3err.ErrNoSuchKey) } if errors.Is(err, syscall.ENAMETOOLONG) { @@ -1848,8 +1898,9 @@ func (p *Posix) UploadPartCopy(ctx context.Context, upi *s3.UploadPartCopyInput) } return s3response.CopyObjectResult{ - ETag: etag, - LastModified: fi.ModTime(), + ETag: etag, + LastModified: fi.ModTime(), + CopySourceVersionId: srcVersionId, }, nil } @@ -2760,30 +2811,14 @@ func (p *Posix) CopyObject(ctx context.Context, input *s3.CopyObjectInput) (*s3. return nil, s3err.GetAPIError(s3err.ErrInvalidRequest) } - copySourceHdr := *input.CopySource - if copySourceHdr[0] == '/' { - copySourceHdr = copySourceHdr[1:] - } - - cSplitted := strings.Split(copySourceHdr, "?") - copySource := cSplitted[0] - var srcVersionId string - if len(cSplitted) > 1 { - versionIdParts := strings.Split(cSplitted[1], "=") - if len(versionIdParts) != 2 || versionIdParts[0] != "versionId" { - return nil, s3err.GetAPIError(s3err.ErrInvalidRequest) - } - srcVersionId = versionIdParts[1] - } - - srcBucket, srcObject, ok := strings.Cut(copySource, "/") - if !ok { - return nil, s3err.GetAPIError(s3err.ErrInvalidCopySource) + srcBucket, srcObject, srcVersionId, err := backend.ParseCopySource(*input.CopySource) + if err != nil { + return nil, err } dstBucket := *input.Bucket dstObject := *input.Key - _, err := os.Stat(srcBucket) + _, err = os.Stat(srcBucket) if errors.Is(err, fs.ErrNotExist) { return nil, s3err.GetAPIError(s3err.ErrNoSuchBucket) } @@ -2812,7 +2847,6 @@ func (p *Posix) CopyObject(ctx context.Context, input *s3.CopyObjectInput) (*s3. srcBucket = filepath.Join(p.versioningDir, srcBucket) srcObject = filepath.Join(genObjVersionKey(srcObject), srcVersionId) } - } _, err = os.Stat(dstBucket) diff --git a/s3api/controllers/base.go b/s3api/controllers/base.go index 969fdd57..03ee2ccc 100644 --- a/s3api/controllers/base.go +++ b/s3api/controllers/base.go @@ -1850,6 +1850,14 @@ func (c S3ApiController) PutActions(ctx *fiber.Ctx) error { ExpectedBucketOwner: &bucketOwner, CopySourceRange: ©SrcRange, }) + if err == nil && resp.CopySourceVersionId != "" { + utils.SetResponseHeaders(ctx, []utils.CustomHeader{ + { + Key: "x-amz-copy-source-version-id", + Value: resp.CopySourceVersionId, + }, + }) + } return SendXMLResponse(ctx, resp, err, &MetaOpts{ Logger: c.logger, @@ -3144,6 +3152,14 @@ func (c S3ApiController) CreateActions(ctx *fiber.Ctx) error { }, }) if err == nil { + if getstring(res.VersionId) != "" { + utils.SetResponseHeaders(ctx, []utils.CustomHeader{ + { + Key: "x-amz-version-id", + Value: getstring(res.VersionId), + }, + }) + } return SendXMLResponse(ctx, res, err, &MetaOpts{ Logger: c.logger, diff --git a/s3response/s3response.go b/s3response/s3response.go index 865e8006..9d8b9b0f 100644 --- a/s3response/s3response.go +++ b/s3response/s3response.go @@ -307,9 +307,10 @@ type CanonicalUser struct { } type CopyObjectResult struct { - XMLName xml.Name `xml:"http://s3.amazonaws.com/doc/2006-03-01/ CopyObjectResult" json:"-"` - LastModified time.Time - ETag string + XMLName xml.Name `xml:"http://s3.amazonaws.com/doc/2006-03-01/ CopyObjectResult" json:"-"` + LastModified time.Time + ETag string + CopySourceVersionId string `xml:"-"` } func (r CopyObjectResult) MarshalXML(e *xml.Encoder, start xml.StartElement) error { diff --git a/tests/integration/group-tests.go b/tests/integration/group-tests.go index ad4fd811..c39b1914 100644 --- a/tests/integration/group-tests.go +++ b/tests/integration/group-tests.go @@ -507,21 +507,27 @@ func TestAccessControl(s *S3Conf) { } func TestVersioning(s *S3Conf) { + // PutBucketVersioning action PutBucketVersioning_non_existing_bucket(s) PutBucketVersioning_invalid_status(s) PutBucketVersioning_success(s) + // GetBucketVersioning action GetBucketVersioning_non_existing_bucket(s) GetBucketVersioning_success(s) Versioning_PutObject_success(s) + // CopyObject action Versioning_CopyObject_success(s) Versioning_CopyObject_non_existing_version_id(s) Versioning_CopyObject_from_an_object_version(s) + // HeadObject action Versioning_HeadObject_invalid_versionId(s) Versioning_HeadObject_success(s) Versioning_HeadObject_delete_marker(s) + // GetObject action Versioning_GetObject_invalid_versionId(s) Versioning_GetObject_success(s) Versioning_GetObject_delete_marker(s) + // DeleteObject(s) actions Versioning_DeleteObject_delete_object_version(s) Versioning_DeleteObject_delete_a_delete_marker(s) Versioning_DeleteObjects_success(s) @@ -532,6 +538,11 @@ func TestVersioning(s *S3Conf) { ListObjectVersions_list_multiple_object_versions(s) ListObjectVersions_multiple_object_versions_truncated(s) ListObjectVersions_with_delete_markers(s) + // Multipart upload + Versioning_Multipart_Upload_success(s) + Versioning_Multipart_Upload_overwrite_an_object(s) + Versioning_UploadPartCopy_non_existing_versionId(s) + Versioning_UploadPartCopy_from_an_object_version(s) } type IntTests map[string]func(s *S3Conf) error @@ -867,5 +878,9 @@ func GetIntTests() IntTests { "ListObjectVersions_list_multiple_object_versions": ListObjectVersions_list_multiple_object_versions, "ListObjectVersions_multiple_object_versions_truncated": ListObjectVersions_multiple_object_versions_truncated, "ListObjectVersions_with_delete_markers": ListObjectVersions_with_delete_markers, + "Versioning_Multipart_Upload_success": Versioning_Multipart_Upload_success, + "Versioning_Multipart_Upload_overwrite_an_object": Versioning_Multipart_Upload_overwrite_an_object, + "Versioning_UploadPartCopy_non_existing_versionId": Versioning_UploadPartCopy_non_existing_versionId, + "Versioning_UploadPartCopy_from_an_object_version": Versioning_UploadPartCopy_from_an_object_version, } } diff --git a/tests/integration/tests.go b/tests/integration/tests.go index 86d30044..f3b11295 100644 --- a/tests/integration/tests.go +++ b/tests/integration/tests.go @@ -11175,3 +11175,278 @@ func ListObjectVersions_with_delete_markers(s *S3Conf) error { return nil }, withVersioning()) } + +func Versioning_Multipart_Upload_success(s *S3Conf) error { + testName := "Versioning_Multipart_Upload_success" + return actionHandler(s, testName, func(s3client *s3.Client, bucket string) error { + obj := "my-obj" + out, err := createMp(s3client, bucket, obj) + if err != nil { + return err + } + + objSize := 5 * 1024 * 1024 + parts, err := uploadParts(s3client, objSize, 5, bucket, obj, *out.UploadId) + if err != nil { + return err + } + + compParts := []types.CompletedPart{} + for _, el := range parts { + compParts = append(compParts, types.CompletedPart{ + ETag: el.ETag, + PartNumber: el.PartNumber, + }) + } + + ctx, cancel := context.WithTimeout(context.Background(), shortTimeout) + res, err := s3client.CompleteMultipartUpload(ctx, &s3.CompleteMultipartUploadInput{ + Bucket: &bucket, + Key: &obj, + UploadId: out.UploadId, + MultipartUpload: &types.CompletedMultipartUpload{ + Parts: compParts, + }, + }) + cancel() + if err != nil { + return err + } + + if *res.Key != obj { + return fmt.Errorf("expected object key to be %v, instead got %v", obj, *res.Key) + } + if *res.Bucket != bucket { + return fmt.Errorf("expected the bucket name to be %v, instead got %v", bucket, *res.Bucket) + } + if res.ETag == nil || *res.ETag == "" { + return fmt.Errorf("expected non-empty ETag") + } + if res.VersionId == nil || *res.VersionId == "" { + return fmt.Errorf("expected non-empty versionId") + } + + ctx, cancel = context.WithTimeout(context.Background(), shortTimeout) + resp, err := s3client.HeadObject(ctx, &s3.HeadObjectInput{ + Bucket: &bucket, + Key: &obj, + VersionId: res.VersionId, + }) + cancel() + if err != nil { + return err + } + + if *resp.ETag != *res.ETag { + return fmt.Errorf("expected the uploaded object etag to be %v, instead got %v", *res.ETag, *resp.ETag) + } + if *resp.ContentLength != int64(objSize) { + return fmt.Errorf("expected the uploaded object size to be %v, instead got %v", objSize, resp.ContentLength) + } + if *resp.VersionId != *res.VersionId { + return fmt.Errorf("expected the versionId to be %v, instead got %v", *res.VersionId, *resp.VersionId) + } + + return nil + }, withVersioning()) +} + +func Versioning_Multipart_Upload_overwrite_an_object(s *S3Conf) error { + testName := "Versioning_Multipart_Upload_overwrite_an_object" + return actionHandler(s, testName, func(s3client *s3.Client, bucket string) error { + obj := "my-obj" + + objVersions, err := createObjVersions(s3client, bucket, obj, 2) + if err != nil { + return err + } + out, err := createMp(s3client, bucket, obj) + if err != nil { + return err + } + + objSize := 5 * 1024 * 1024 + parts, err := uploadParts(s3client, objSize, 5, bucket, obj, *out.UploadId) + if err != nil { + return err + } + + compParts := []types.CompletedPart{} + for _, el := range parts { + compParts = append(compParts, types.CompletedPart{ + ETag: el.ETag, + PartNumber: el.PartNumber, + }) + } + + ctx, cancel := context.WithTimeout(context.Background(), shortTimeout) + res, err := s3client.CompleteMultipartUpload(ctx, &s3.CompleteMultipartUploadInput{ + Bucket: &bucket, + Key: &obj, + UploadId: out.UploadId, + MultipartUpload: &types.CompletedMultipartUpload{ + Parts: compParts, + }, + }) + cancel() + if err != nil { + return err + } + + if *res.Key != obj { + return fmt.Errorf("expected object key to be %v, instead got %v", obj, *res.Key) + } + if *res.Bucket != bucket { + return fmt.Errorf("expected the bucket name to be %v, instead got %v", bucket, *res.Bucket) + } + if res.ETag == nil || *res.ETag == "" { + return fmt.Errorf("expected non-empty ETag") + } + if res.VersionId == nil || *res.VersionId == "" { + return fmt.Errorf("expected non-empty versionId") + } + + ctx, cancel = context.WithTimeout(context.Background(), shortTimeout) + resp, err := s3client.ListObjectVersions(ctx, &s3.ListObjectVersionsInput{ + Bucket: &bucket, + }) + cancel() + if err != nil { + return err + } + + size := int64(objSize) + + objVersions[0].IsLatest = getBoolPtr(false) + versions := append([]types.ObjectVersion{ + { + Key: &obj, + VersionId: res.VersionId, + ETag: res.ETag, + IsLatest: getBoolPtr(true), + Size: &size, + }, + }, objVersions...) + + if !compareVersions(resp.Versions, versions) { + return fmt.Errorf("expected the resulting versions to be %v, instead got %v", versions, resp.Versions) + } + + return nil + }, withVersioning()) +} + +func Versioning_UploadPartCopy_non_existing_versionId(s *S3Conf) error { + testName := "Versioning_UploadPartCopy_non_existing_versionId" + return actionHandler(s, testName, func(s3client *s3.Client, bucket string) error { + dstBucket, dstObj, srcObj := getBucketName(), "dst-obj", "src-obj" + + lgth := int64(100) + _, err := putObjectWithData(lgth, &s3.PutObjectInput{ + Bucket: &bucket, + Key: &srcObj, + }, s3client) + if err != nil { + return err + } + + if err := setup(s, dstBucket); err != nil { + return err + } + + mp, err := createMp(s3client, dstBucket, dstObj) + if err != nil { + return err + } + + pNumber := int32(1) + ctx, cancel := context.WithTimeout(context.Background(), shortTimeout) + _, err = s3client.UploadPartCopy(ctx, &s3.UploadPartCopyInput{ + Bucket: &dstBucket, + Key: &dstObj, + UploadId: mp.UploadId, + PartNumber: &pNumber, + CopySource: getPtr(fmt.Sprintf("%v/%v?versionId=invalid_versionId", bucket, srcObj)), + }) + cancel() + if err := checkApiErr(err, s3err.GetAPIError(s3err.ErrNoSuchVersion)); err != nil { + return err + } + + if err := teardown(s, dstBucket); err != nil { + return err + } + + return nil + }, withVersioning()) +} + +func Versioning_UploadPartCopy_from_an_object_version(s *S3Conf) error { + testName := "Versioning_UploadPartCopy_from_an_object_version" + return actionHandler(s, testName, func(s3client *s3.Client, bucket string) error { + srcObj, dstBucket, obj := "my-obj", getBucketName(), "dst-obj" + err := setup(s, dstBucket) + if err != nil { + return err + } + + srcObjVersions, err := createObjVersions(s3client, bucket, srcObj, 1) + if err != nil { + return err + } + srcObjVersion := srcObjVersions[0] + + out, err := createMp(s3client, dstBucket, obj) + if err != nil { + return err + } + + partNumber := int32(1) + ctx, cancel := context.WithTimeout(context.Background(), shortTimeout) + copyOut, err := s3client.UploadPartCopy(ctx, &s3.UploadPartCopyInput{ + Bucket: &dstBucket, + CopySource: getPtr(fmt.Sprintf("%v/%v?versionId=%v", bucket, srcObj, *srcObjVersion.VersionId)), + UploadId: out.UploadId, + Key: &obj, + PartNumber: &partNumber, + }) + cancel() + if err != nil { + return err + } + + if *copyOut.CopySourceVersionId != *srcObjVersion.VersionId { + return fmt.Errorf("expected the copy-source-version-id to be %v, instead got %v", *srcObjVersion.VersionId, *copyOut.CopySourceVersionId) + } + + ctx, cancel = context.WithTimeout(context.Background(), shortTimeout) + res, err := s3client.ListParts(ctx, &s3.ListPartsInput{ + Bucket: &dstBucket, + Key: &obj, + UploadId: out.UploadId, + }) + cancel() + if err != nil { + return err + } + + if len(res.Parts) != 1 { + return fmt.Errorf("expected parts to be 1, instead got %v", len(res.Parts)) + } + if *res.Parts[0].PartNumber != partNumber { + return fmt.Errorf("expected part-number to be %v, instead got %v", partNumber, res.Parts[0].PartNumber) + } + if *res.Parts[0].Size != *srcObjVersion.Size { + return fmt.Errorf("expected part size to be %v, instead got %v", *srcObjVersion.Size, res.Parts[0].Size) + } + if *res.Parts[0].ETag != *copyOut.CopyPartResult.ETag { + return fmt.Errorf("expected part etag to be %v, instead got %v", *copyOut.CopyPartResult.ETag, *res.Parts[0].ETag) + } + + if err := teardown(s, dstBucket); err != nil { + return err + } + + return nil + }, withVersioning()) +}