From 39d89ea3feb82bf47b398eeedcb790f1146e9677 Mon Sep 17 00:00:00 2001 From: Casey Bodley Date: Wed, 1 May 2024 13:59:09 -0400 Subject: [PATCH 01/12] add "checksum" marker, since new checksum tests reference it this removes a Pytest warning during execution Signed-off-by: Matt Benjamin --- pytest.ini | 1 + 1 file changed, 1 insertion(+) diff --git a/pytest.ini b/pytest.ini index 73d156332..1a7d9a836 100644 --- a/pytest.ini +++ b/pytest.ini @@ -7,6 +7,7 @@ markers = auth_common bucket_policy bucket_encryption + checksum cloud_transition encryption fails_on_aws From 51c052fe7a4203b9c805ca768405f50094b6772f Mon Sep 17 00:00:00 2001 From: Matt Benjamin Date: Wed, 1 May 2024 14:05:52 -0400 Subject: [PATCH 02/12] test_multipart_upload_sha256: work around failures re-trying complete-multipart As described in https://tracker.ceph.com/issues/65746, retrying complete-multipart after having attempted to complete the same upload with a bad checksum argument fails with an internal error. The status code is 500, but I'm unsure if it can be retried again, or whether the upload can be aborted later. Signed-off-by: Matt Benjamin --- s3tests_boto3/functional/test_s3.py | 88 +++++++++++++++++++++++++++++ 1 file changed, 88 insertions(+) diff --git a/s3tests_boto3/functional/test_s3.py b/s3tests_boto3/functional/test_s3.py index a0c3fac83..db4e0c21f 100644 --- a/s3tests_boto3/functional/test_s3.py +++ b/s3tests_boto3/functional/test_s3.py @@ -13551,3 +13551,91 @@ def test_upload_part_copy_percent_encoded_key(): final_obj = s3_client.get_object(Bucket=bucket_name, Key=key) content = final_obj['Body'].read() assert content == b"foo" + +@pytest.mark.checksum +def test_object_checksum_sha256(): + bucket = get_new_bucket() + client = get_client() + + key = "myobj" + size = 1024 + body = FakeWriteFile(size, 'A') + sha256sum = 'arcu6553sHVAiX4MjW0j7I7vD4w6R+Gz9Ok0Q9lTa+0=' + response = client.put_object(Bucket=bucket, Key=key, Body=body, ChecksumAlgorithm='SHA256', ChecksumSHA256=sha256sum) + assert sha256sum == response['ChecksumSHA256'] + + response = client.head_object(Bucket=bucket, Key=key) + assert 'ChecksumSHA256' not in response + response = client.head_object(Bucket=bucket, Key=key, ChecksumMode='ENABLED') + assert sha256sum == response['ChecksumSHA256'] + + e = assert_raises(ClientError, client.put_object, Bucket=bucket, Key=key, Body=body, ChecksumAlgorithm='SHA256', ChecksumSHA256='bad') + status, error_code = _get_status_and_error_code(e.response) + assert status == 400 + assert error_code == 'InvalidRequest' + +@pytest.mark.checksum +def test_multipart_checksum_sha256(): + bucket = get_new_bucket() + client = get_client() + + key = "mymultipart" + response = client.create_multipart_upload(Bucket=bucket, Key=key, ChecksumAlgorithm='SHA256') + assert 'SHA256' == response['ChecksumAlgorithm'] + upload_id = response['UploadId'] + + size = 1024 + body = FakeWriteFile(size, 'A') + part_sha256sum = 'arcu6553sHVAiX4MjW0j7I7vD4w6R+Gz9Ok0Q9lTa+0=' + response = client.upload_part(UploadId=upload_id, Bucket=bucket, Key=key, PartNumber=1, Body=body, ChecksumAlgorithm='SHA256', ChecksumSHA256=part_sha256sum) + + # should reject the bad request checksum + e = assert_raises(ClientError, client.complete_multipart_upload, Bucket=bucket, Key=key, UploadId=upload_id, ChecksumSHA256='bad', MultipartUpload={'Parts': [ + {'ETag': response['ETag'].strip('"'), 'ChecksumSHA256': response['ChecksumSHA256'], 'PartNumber': 1}]}) + status, error_code = _get_status_and_error_code(e.response) + assert status == 400 + assert error_code == 'InvalidRequest' + + # XXXX re-trying the complete is failing in RGW due to an internal error that appears not caused + # checksums; + # 2024-04-25T17:47:47.991-0400 7f78e3a006c0 0 req 4931907640780566174 0.011000143s s3:complete_multipart check_previously_completed() ERROR: get_obj_attrs() returned ret=-2 + # 2024-04-25T17:47:47.991-0400 7f78e3a006c0 2 req 4931907640780566174 0.011000143s s3:complete_multipart completing + # 2024-04-25T17:47:47.991-0400 7f78e3a006c0 1 req 4931907640780566174 0.011000143s s3:complete_multipart ERROR: either op_ret is negative (execute failed) or target_obj is null, op_ret: -2200 + # -2200 turns into 500, InternalError + + key = "mymultipart2" + response = client.create_multipart_upload(Bucket=bucket, Key=key, ChecksumAlgorithm='SHA256') + assert 'SHA256' == response['ChecksumAlgorithm'] + upload_id = response['UploadId'] + + size = 1024 + body = FakeWriteFile(size, 'A') + part_sha256sum = 'arcu6553sHVAiX4MjW0j7I7vD4w6R+Gz9Ok0Q9lTa+0=' + response = client.upload_part(UploadId=upload_id, Bucket=bucket, Key=key, PartNumber=1, Body=body, ChecksumAlgorithm='SHA256', ChecksumSHA256=part_sha256sum) + + # should reject the missing part checksum + e = assert_raises(ClientError, client.complete_multipart_upload, Bucket=bucket, Key=key, UploadId=upload_id, ChecksumSHA256='bad', MultipartUpload={'Parts': [ + {'ETag': response['ETag'].strip('"'), 'PartNumber': 1}]}) + status, error_code = _get_status_and_error_code(e.response) + assert status == 400 + assert error_code == 'InvalidRequest' + + key = "mymultipart3" + response = client.create_multipart_upload(Bucket=bucket, Key=key, ChecksumAlgorithm='SHA256') + assert 'SHA256' == response['ChecksumAlgorithm'] + upload_id = response['UploadId'] + + size = 1024 + body = FakeWriteFile(size, 'A') + part_sha256sum = 'arcu6553sHVAiX4MjW0j7I7vD4w6R+Gz9Ok0Q9lTa+0=' + response = client.upload_part(UploadId=upload_id, Bucket=bucket, Key=key, PartNumber=1, Body=body, ChecksumAlgorithm='SHA256', ChecksumSHA256=part_sha256sum) + + composite_sha256sum = 'Ok6Cs5b96ux6+MWQkJO7UBT5sKPBeXBLwvj/hK89smg=-1' + response = client.complete_multipart_upload(Bucket=bucket, Key=key, UploadId=upload_id, ChecksumSHA256=composite_sha256sum, MultipartUpload={'Parts': [ + {'ETag': response['ETag'].strip('"'), 'ChecksumSHA256': response['ChecksumSHA256'], 'PartNumber': 1}]}) + assert composite_sha256sum == response['ChecksumSHA256'] + + response = client.head_object(Bucket=bucket, Key=key) + assert 'ChecksumSHA256' not in response + response = client.head_object(Bucket=bucket, Key=key, ChecksumMode='ENABLED') + assert composite_sha256sum == response['ChecksumSHA256'] From 1c584b6b35b2f33ffaeadab9e043737e17831824 Mon Sep 17 00:00:00 2001 From: Matt Benjamin Date: Wed, 1 May 2024 14:15:36 -0400 Subject: [PATCH 03/12] add test_multipart_checksum_3parts tests a full multipart upload cycle with 3 unique parts, which verifies composite checksum computation and the logic to propagate parts_count to ComleteMultipart Signed-off-by: Matt Benjamin --- s3tests_boto3/functional/test_s3.py | 38 +++++++++++++++++++++++++++++ 1 file changed, 38 insertions(+) diff --git a/s3tests_boto3/functional/test_s3.py b/s3tests_boto3/functional/test_s3.py index db4e0c21f..c73818478 100644 --- a/s3tests_boto3/functional/test_s3.py +++ b/s3tests_boto3/functional/test_s3.py @@ -13639,3 +13639,41 @@ def test_multipart_checksum_sha256(): assert 'ChecksumSHA256' not in response response = client.head_object(Bucket=bucket, Key=key, ChecksumMode='ENABLED') assert composite_sha256sum == response['ChecksumSHA256'] + +@pytest.mark.checksum +def test_multipart_checksum_3parts(): + bucket = get_new_bucket() + client = get_client() + + key = "mymultipart3" + response = client.create_multipart_upload(Bucket=bucket, Key=key, ChecksumAlgorithm='SHA256') + assert 'SHA256' == response['ChecksumAlgorithm'] + upload_id = response['UploadId'] + + size = 5 * 1024 * 1024 # each part but the last must be at least 5M + body = FakeWriteFile(size, 'A') + part1_sha256sum = '275VF5loJr1YYawit0XSHREhkFXYkkPKGuoK0x9VKxI=' + response = client.upload_part(UploadId=upload_id, Bucket=bucket, Key=key, PartNumber=1, Body=body, ChecksumAlgorithm='SHA256', ChecksumSHA256=part1_sha256sum) + etag1 = response['ETag'].strip('"') + + body = FakeWriteFile(size, 'B') + part2_sha256sum = 'mrHwOfjTL5Zwfj74F05HOQGLdUb7E5szdCbxgUSq6NM=' + response = client.upload_part(UploadId=upload_id, Bucket=bucket, Key=key, PartNumber=2, Body=body, ChecksumAlgorithm='SHA256', ChecksumSHA256=part2_sha256sum) + etag2 = response['ETag'].strip('"') + + body = FakeWriteFile(size, 'C') + part3_sha256sum = 'Vw7oB/nKQ5xWb3hNgbyfkvDiivl+U+/Dft48nfJfDow=' + response = client.upload_part(UploadId=upload_id, Bucket=bucket, Key=key, PartNumber=3, Body=body, ChecksumAlgorithm='SHA256', ChecksumSHA256=part3_sha256sum) + etag3 = response['ETag'].strip('"') + + composite_sha256sum = 'uWBwpe1dxI4Vw8Gf0X9ynOdw/SS6VBzfWm9giiv1sf4=-3' + response = client.complete_multipart_upload(Bucket=bucket, Key=key, UploadId=upload_id, ChecksumSHA256=composite_sha256sum, MultipartUpload={'Parts': [ + {'ETag': etag1, 'ChecksumSHA256': response['ChecksumSHA256'], 'PartNumber': 1}, + {'ETag': etag2, 'ChecksumSHA256': response['ChecksumSHA256'], 'PartNumber': 2}, + {'ETag': etag3, 'ChecksumSHA256': response['ChecksumSHA256'], 'PartNumber': 3}]}) + assert composite_sha256sum == response['ChecksumSHA256'] + + response = client.head_object(Bucket=bucket, Key=key) + assert 'ChecksumSHA256' not in response + response = client.head_object(Bucket=bucket, Key=key, ChecksumMode='ENABLED') + assert composite_sha256sum == response['ChecksumSHA256'] From c3aedf84589f430ad073ed72f1a6317f869a61ab Mon Sep 17 00:00:00 2001 From: Matt Benjamin Date: Fri, 3 May 2024 16:25:19 -0400 Subject: [PATCH 04/12] add test_post_object_upload_checksum this tests a two-megabyte binary upload with validated (awscli-computed) SHA256 checksum, and also verifies failure when a bad checksum is provided Signed-off-by: Matt Benjamin --- s3tests_boto3/functional/test_s3.py | 53 +++++++++++++++++++++++++++++ 1 file changed, 53 insertions(+) diff --git a/s3tests_boto3/functional/test_s3.py b/s3tests_boto3/functional/test_s3.py index c73818478..c2705391d 100644 --- a/s3tests_boto3/functional/test_s3.py +++ b/s3tests_boto3/functional/test_s3.py @@ -13677,3 +13677,56 @@ def test_multipart_checksum_3parts(): assert 'ChecksumSHA256' not in response response = client.head_object(Bucket=bucket, Key=key, ChecksumMode='ENABLED') assert composite_sha256sum == response['ChecksumSHA256'] + +def test_post_object_upload_checksum(): + megabytes = 1024 * 1024 + min_size = 0 + max_size = 5 * megabytes + test_payload_size = 2 * megabytes + + bucket_name = get_new_bucket() + client = get_client() + + url = _get_post_url(bucket_name) + utc = pytz.utc + expires = datetime.datetime.now(utc) + datetime.timedelta(seconds=+6000) + + policy_document = {"expiration": expires.strftime("%Y-%m-%dT%H:%M:%SZ"),\ + "conditions": [\ + {"bucket": bucket_name},\ + ["starts-with", "$key", "foo_cksum_test"],\ + {"acl": "private"},\ + ["starts-with", "$Content-Type", "text/plain"],\ + ["content-length-range", min_size, max_size],\ + ]\ + } + + test_payload = b'x' * test_payload_size + + json_policy_document = json.JSONEncoder().encode(policy_document) + bytes_json_policy_document = bytes(json_policy_document, 'utf-8') + policy = base64.b64encode(bytes_json_policy_document) + aws_secret_access_key = get_main_aws_secret_key() + aws_access_key_id = get_main_aws_access_key() + + signature = base64.b64encode(hmac.new(bytes(aws_secret_access_key, 'utf-8'), policy, hashlib.sha1).digest()) + + # good checksum payload (checked via upload from awscli) + payload = OrderedDict([ ("key" , "foo_cksum_test.txt"),("AWSAccessKeyId" , aws_access_key_id),\ + ("acl" , "private"),("signature" , signature),("policy" , policy),\ + ("Content-Type" , "text/plain"),\ + ('x-amz-checksum-sha256', 'aTL9MeXa9HObn6eP93eygxsJlcwdCwCTysgGAZAgE7w='),\ + ('file', (test_payload)),]) + + r = requests.post(url, files=payload, verify=get_config_ssl_verify()) + assert r.status_code == 204 + + # bad checksum payload + payload = OrderedDict([ ("key" , "foo_cksum_test.txt"),("AWSAccessKeyId" , aws_access_key_id),\ + ("acl" , "private"),("signature" , signature),("policy" , policy),\ + ("Content-Type" , "text/plain"),\ + ('x-amz-checksum-sha256', 'sailorjerry'),\ + ('file', (test_payload)),]) + + r = requests.post(url, files=payload, verify=get_config_ssl_verify()) + assert r.status_code == 400 From 0ed7979a02289719e64e3635e951345300ee8821 Mon Sep 17 00:00:00 2001 From: Matt Benjamin Date: Sat, 22 Jun 2024 17:42:21 -0400 Subject: [PATCH 05/12] remove duplicate size assigment [rkhudov review] Signed-off-by: Matt Benjamin --- s3tests_boto3/functional/test_s3.py | 2 -- 1 file changed, 2 deletions(-) diff --git a/s3tests_boto3/functional/test_s3.py b/s3tests_boto3/functional/test_s3.py index c2705391d..5355cdc49 100644 --- a/s3tests_boto3/functional/test_s3.py +++ b/s3tests_boto3/functional/test_s3.py @@ -13608,7 +13608,6 @@ def test_multipart_checksum_sha256(): assert 'SHA256' == response['ChecksumAlgorithm'] upload_id = response['UploadId'] - size = 1024 body = FakeWriteFile(size, 'A') part_sha256sum = 'arcu6553sHVAiX4MjW0j7I7vD4w6R+Gz9Ok0Q9lTa+0=' response = client.upload_part(UploadId=upload_id, Bucket=bucket, Key=key, PartNumber=1, Body=body, ChecksumAlgorithm='SHA256', ChecksumSHA256=part_sha256sum) @@ -13625,7 +13624,6 @@ def test_multipart_checksum_sha256(): assert 'SHA256' == response['ChecksumAlgorithm'] upload_id = response['UploadId'] - size = 1024 body = FakeWriteFile(size, 'A') part_sha256sum = 'arcu6553sHVAiX4MjW0j7I7vD4w6R+Gz9Ok0Q9lTa+0=' response = client.upload_part(UploadId=upload_id, Bucket=bucket, Key=key, PartNumber=1, Body=body, ChecksumAlgorithm='SHA256', ChecksumSHA256=part_sha256sum) From a03400469085303e175cee5194c13a9acc6ebf02 Mon Sep 17 00:00:00 2001 From: Matt Benjamin Date: Wed, 3 Jul 2024 09:42:37 -0400 Subject: [PATCH 06/12] mark two tests that fail on dbstore also add @pytest.mark.checksum for new checksum tests Signed-off-by: Matt Benjamin --- s3tests_boto3/functional/test_s3.py | 3 +++ 1 file changed, 3 insertions(+) diff --git a/s3tests_boto3/functional/test_s3.py b/s3tests_boto3/functional/test_s3.py index 5355cdc49..08c792601 100644 --- a/s3tests_boto3/functional/test_s3.py +++ b/s3tests_boto3/functional/test_s3.py @@ -13575,6 +13575,7 @@ def test_object_checksum_sha256(): assert error_code == 'InvalidRequest' @pytest.mark.checksum +@pytest.mark.fails_on_dbstore def test_multipart_checksum_sha256(): bucket = get_new_bucket() client = get_client() @@ -13639,6 +13640,7 @@ def test_multipart_checksum_sha256(): assert composite_sha256sum == response['ChecksumSHA256'] @pytest.mark.checksum +@pytest.mark.fails_on_dbstore def test_multipart_checksum_3parts(): bucket = get_new_bucket() client = get_client() @@ -13676,6 +13678,7 @@ def test_multipart_checksum_3parts(): response = client.head_object(Bucket=bucket, Key=key, ChecksumMode='ENABLED') assert composite_sha256sum == response['ChecksumSHA256'] +@pytest.mark.checksum def test_post_object_upload_checksum(): megabytes = 1024 * 1024 min_size = 0 From 9d54da1707e067e7c72cb433253cc1893f891f63 Mon Sep 17 00:00:00 2001 From: Matt Benjamin Date: Sat, 6 Jul 2024 13:21:55 -0400 Subject: [PATCH 07/12] test get_object_attributes Signed-off-by: Matt Benjamin --- s3tests_boto3/functional/test_s3.py | 1391 +++++++++++++++++++++++++++ 1 file changed, 1391 insertions(+) diff --git a/s3tests_boto3/functional/test_s3.py b/s3tests_boto3/functional/test_s3.py index 08c792601..263e50b8c 100644 --- a/s3tests_boto3/functional/test_s3.py +++ b/s3tests_boto3/functional/test_s3.py @@ -1211,6 +1211,7 @@ def add_unordered(**kwargs): intersect = set(unordered_keys_out).intersection(unordered_keys_out2) assert 0 == len(intersect) + #pdb.set_trace() # verify that unordered used with delimiter results in error e = assert_raises(ClientError, client.list_objects, Bucket=bucket_name, Delimiter="/") @@ -5757,6 +5758,44 @@ def _multipart_upload(bucket_name, key, size, part_size=5*1024*1024, client=None return (upload_id, s, parts) +def _multipart_upload_checksum(bucket_name, key, size, part_size=5*1024*1024, client=None, content_type=None, metadata=None, resend_parts=[]): + """ + generate a multi-part upload for a random file of specifed size, + if requested, generate a list of the parts + return the upload descriptor + """ + if client == None: + client = get_client() + + + if content_type == None and metadata == None: + response = client.create_multipart_upload(Bucket=bucket_name, Key=key, ChecksumAlgorithm='SHA256') + else: + response = client.create_multipart_upload(Bucket=bucket_name, Key=key, Metadata=metadata, ContentType=content_type, + ChecksumAlgorithm='SHA256') + + upload_id = response['UploadId'] + s = '' + parts = [] + part_checksums = [] + for i, part in enumerate(generate_random(size, part_size)): + # part_num is necessary because PartNumber for upload_part and in parts must start at 1 and i starts at 0 + part_num = i+1 + s += part + response = client.upload_part(UploadId=upload_id, Bucket=bucket_name, Key=key, PartNumber=part_num, Body=part, + ChecksumAlgorithm='SHA256') + + parts.append({'ETag': response['ETag'].strip('"'), 'PartNumber': part_num}) + + armored_part_cksum = base64.b64encode(hashlib.sha256(part.encode('utf-8')).digest()) + part_checksums.append(armored_part_cksum.decode()) + + if i in resend_parts: + client.upload_part(UploadId=upload_id, Bucket=bucket_name, Key=key, PartNumber=part_num, Body=part, + ChecksumAlgorithm='SHA256') + + return (upload_id, s, parts, part_checksums) + @pytest.mark.fails_on_dbstore def test_object_copy_versioning_multipart_upload(): bucket_name = get_new_bucket() @@ -13731,3 +13770,1355 @@ def test_post_object_upload_checksum(): r = requests.post(url, files=payload, verify=get_config_ssl_verify()) assert r.status_code == 400 + +def _has_bucket_logging_extension(): + src_bucket_name = get_new_bucket_name() + src_bucket = get_new_bucket_resource(name=src_bucket_name) + log_bucket_name = get_new_bucket_name() + log_bucket = get_new_bucket_resource(name=log_bucket_name) + client = get_client() + logging_enabled = {'TargetBucket': log_bucket_name, 'TargetPrefix': 'log/', 'LoggingType': 'Journal'} + try: + response = client.put_bucket_logging(Bucket=src_bucket_name, BucketLoggingStatus={ + 'LoggingEnabled': logging_enabled, + }) + except ParamValidationError as e: + return False + return True + +def _has_taget_object_key_format(): + src_bucket_name = get_new_bucket_name() + src_bucket = get_new_bucket_resource(name=src_bucket_name) + log_bucket_name = get_new_bucket_name() + log_bucket = get_new_bucket_resource(name=log_bucket_name) + client = get_client() + logging_enabled = {'TargetBucket': log_bucket_name, 'TargetPrefix': 'log/', 'TargetObjectKeyFormat': {'SimplePrefix': {}}} + try: + response = client.put_bucket_logging(Bucket=src_bucket_name, BucketLoggingStatus={ + 'LoggingEnabled': logging_enabled, + }) + except ParamValidationError as e: + return False + return True + + +import shlex + +def _parse_standard_log_record(record): + record = record.replace('[', '"').replace(']', '"') + chunks = shlex.split(record) + assert len(chunks) == 26 + return { + 'BucketOwner': chunks[0], + 'BucketName': chunks[1], + 'RequestDateTime': chunks[2], + 'RemoteIP': chunks[3], + 'Requester': chunks[4], + 'RequestID': chunks[5], + 'Operation': chunks[6], + 'Key': chunks[7], + 'RequestURI': chunks[8], + 'HTTPStatus': chunks[9], + 'ErrorCode': chunks[10], + 'BytesSent': chunks[11], + 'ObjectSize': chunks[12], + 'TotalTime': chunks[13], + 'TurnAroundTime': chunks[14], + 'Referrer': chunks[15], + 'UserAgent': chunks[16], + 'VersionID': chunks[17], + 'HostID': chunks[18], + 'SigVersion': chunks[19], + 'CipherSuite': chunks[20], + 'AuthType': chunks[21], + 'HostHeader': chunks[22], + 'TLSVersion': chunks[23], + 'AccessPointARN': chunks[24], + 'ACLRequired': chunks[25], + } + + +def _parse_journal_log_record(record): + record = record.replace('[', '"').replace(']', '"') + chunks = shlex.split(record) + assert len(chunks) == 8 + return { + 'BucketOwner': chunks[0], + 'BucketName': chunks[1], + 'RequestDateTime': chunks[2], + 'Operation': chunks[3], + 'Key': chunks[4], + 'ObjectSize': chunks[5], + 'VersionID': chunks[6], + 'ETAG': chunks[7], + } + +def _parse_log_record(record, record_type): + if record_type == 'Standard': + return _parse_standard_log_record(record) + elif record_type == 'Journal': + return _parse_journal_log_record(record) + else: + assert False, 'unknown log record type' + +expected_object_roll_time = 5 + +import logging + +logger = logging.getLogger(__name__) + + +def _verify_records(records, bucket_name, event_type, src_keys, record_type, expected_count, exact_match=False): + keys_found = [] + all_keys = [] + for record in iter(records.splitlines()): + parsed_record = _parse_log_record(record, record_type) + logger.info('bucket log record: %s', json.dumps(parsed_record, indent=4)) + if bucket_name in record and event_type in record: + all_keys.append(parsed_record['Key']) + for key in src_keys: + if key in record: + keys_found.append(key) + break + logger.info('keys found in bucket log: %s', str(all_keys)) + logger.info('keys from the source bucket: %s', str(src_keys)) + if exact_match: + return len(keys_found) == expected_count and len(keys_found) == len(all_keys) + return len(keys_found) == expected_count + + +def randcontent(): + letters = string.ascii_lowercase + length = random.randint(10, 1024) + return ''.join(random.choice(letters) for i in range(length)) + + +@pytest.mark.bucket_logging +def test_put_bucket_logging(): + src_bucket_name = get_new_bucket_name() + src_bucket = get_new_bucket_resource(name=src_bucket_name) + log_bucket_name = get_new_bucket_name() + log_bucket = get_new_bucket_resource(name=log_bucket_name) + client = get_client() + has_extensions = _has_bucket_logging_extension() + has_key_format = _has_taget_object_key_format() + + # minimal configuration + logging_enabled = { + 'TargetBucket': log_bucket_name, + 'TargetPrefix': 'log/' + } + + if has_extensions: + logging_enabled['ObjectRollTime'] = expected_object_roll_time + response = client.put_bucket_logging(Bucket=src_bucket_name, BucketLoggingStatus={ + 'LoggingEnabled': logging_enabled, + }) + assert response['ResponseMetadata']['HTTPStatusCode'] == 200 + response = client.get_bucket_logging(Bucket=src_bucket_name) + assert response['ResponseMetadata']['HTTPStatusCode'] == 200 + if has_extensions: + logging_enabled['LoggingType'] = 'Standard' + logging_enabled['RecordsBatchSize'] = 0 + if has_key_format: + # default value for key prefix is returned + logging_enabled['TargetObjectKeyFormat'] = {'SimplePrefix': {}} + assert response['LoggingEnabled'] == logging_enabled + + if has_key_format: + # with simple target object prefix + logging_enabled = { + 'TargetBucket': log_bucket_name, + 'TargetPrefix': 'log/', + 'TargetObjectKeyFormat': { + 'SimplePrefix': {} + } + } + if has_extensions: + logging_enabled['ObjectRollTime'] = expected_object_roll_time + response = client.put_bucket_logging(Bucket=src_bucket_name, BucketLoggingStatus={ + 'LoggingEnabled': logging_enabled, + }) + assert response['ResponseMetadata']['HTTPStatusCode'] == 200 + response = client.get_bucket_logging(Bucket=src_bucket_name) + assert response['ResponseMetadata']['HTTPStatusCode'] == 200 + if has_extensions: + logging_enabled['LoggingType'] = 'Standard' + logging_enabled['RecordsBatchSize'] = 0 + assert response['LoggingEnabled'] == logging_enabled + + # with partitioned target object prefix + logging_enabled = { + 'TargetBucket': log_bucket_name, + 'TargetPrefix': 'log/', + 'TargetObjectKeyFormat': { + 'PartitionedPrefix': { + 'PartitionDateSource': 'DeliveryTime' + } + } + } + if has_extensions: + logging_enabled['ObjectRollTime'] = expected_object_roll_time + response = client.put_bucket_logging(Bucket=src_bucket_name, BucketLoggingStatus={ + 'LoggingEnabled': logging_enabled, + }) + assert response['ResponseMetadata']['HTTPStatusCode'] == 200 + response = client.get_bucket_logging(Bucket=src_bucket_name) + assert response['ResponseMetadata']['HTTPStatusCode'] == 200 + if has_extensions: + logging_enabled['LoggingType'] = 'Standard' + logging_enabled['RecordsBatchSize'] = 0 + assert response['LoggingEnabled'] == logging_enabled + + # with target grant (not implemented in RGW) + main_display_name = get_main_display_name() + main_user_id = get_main_user_id() + logging_enabled = { + 'TargetBucket': log_bucket_name, + 'TargetPrefix': 'log/', + 'TargetGrants': [{'Grantee': {'DisplayName': main_display_name, 'ID': main_user_id,'Type': 'CanonicalUser'},'Permission': 'FULL_CONTROL'}] + } + if has_extensions: + logging_enabled['ObjectRollTime'] = expected_object_roll_time + response = client.put_bucket_logging(Bucket=src_bucket_name, BucketLoggingStatus={ + 'LoggingEnabled': logging_enabled, + }) + assert response['ResponseMetadata']['HTTPStatusCode'] == 200 + response = client.get_bucket_logging(Bucket=src_bucket_name) + assert response['ResponseMetadata']['HTTPStatusCode'] == 200 + if has_extensions: + logging_enabled['LoggingType'] = 'Standard' + logging_enabled['RecordsBatchSize'] = 0 + # target grants are not implemented + logging_enabled.pop('TargetGrants') + if has_key_format: + # default value for key prefix is returned + logging_enabled['TargetObjectKeyFormat'] = {'SimplePrefix': {}} + assert response['LoggingEnabled'] == logging_enabled + + +def _bucket_logging_key_filter(log_type): + src_bucket_name = get_new_bucket_name() + src_bucket = get_new_bucket_resource(name=src_bucket_name) + log_bucket_name = get_new_bucket_name() + log_bucket = get_new_bucket_resource(name=log_bucket_name) + client = get_client() + + logging_enabled = { + 'TargetBucket': log_bucket_name, + 'LoggingType': log_type, + 'TargetPrefix': 'log/', + 'ObjectRollTime': expected_object_roll_time, + 'TargetObjectKeyFormat': {'SimplePrefix': {}}, + 'RecordsBatchSize': 0, + 'Filter': + { + 'Key': { + 'FilterRules': [ + {'Name': 'prefix', 'Value': 'test/'}, + {'Name': 'suffix', 'Value': '.txt'} + ] + } + } + } + + response = client.put_bucket_logging(Bucket=src_bucket_name, BucketLoggingStatus={ + 'LoggingEnabled': logging_enabled, + }) + assert response['ResponseMetadata']['HTTPStatusCode'] == 200 + + response = client.get_bucket_logging(Bucket=src_bucket_name) + assert response['ResponseMetadata']['HTTPStatusCode'] == 200 + if log_type == 'Journal': + assert response['LoggingEnabled'] == logging_enabled + elif log_type == 'Standard': + print('TODO') + else: + assert False, 'unknown log type: %s' % log_type + + names = [] + num_keys = 5 + for j in range(num_keys): + name = 'myobject'+str(j) + if log_type == 'Standard': + # standard log records are not filtered + names.append(name) + client.put_object(Bucket=src_bucket_name, Key=name, Body=randcontent()) + + for j in range(num_keys): + name = 'test/'+'myobject'+str(j)+'.txt' + names.append(name) + client.put_object(Bucket=src_bucket_name, Key=name, Body=randcontent()) + + expected_count = len(names) + + time.sleep(expected_object_roll_time) + client.put_object(Bucket=src_bucket_name, Key='test/dummy.txt', Body='dummy') + + response = client.list_objects_v2(Bucket=log_bucket_name) + keys = _get_keys(response) + assert len(keys) == 1 + + for key in keys: + assert key.startswith('log/') + response = client.get_object(Bucket=log_bucket_name, Key=key) + body = _get_body(response) + assert _verify_records(body, src_bucket_name, 'REST.PUT.OBJECT', names, log_type, expected_count, exact_match=True) + + +@pytest.mark.bucket_logging +@pytest.mark.fails_on_aws +def test_bucket_logging_key_filter_s(): + if not _has_bucket_logging_extension(): + pytest.skip('ceph extension to bucket logging not supported at client') + _bucket_logging_key_filter('Standard') + + +@pytest.mark.bucket_logging +@pytest.mark.fails_on_aws +def test_bucket_logging_key_filter_j(): + if not _has_bucket_logging_extension(): + pytest.skip('ceph extension to bucket logging not supported at client') + _bucket_logging_key_filter('Journal') + + +def _bucket_logging_flush(log_type): + src_bucket_name = get_new_bucket_name() + src_bucket = get_new_bucket_resource(name=src_bucket_name) + log_bucket_name = get_new_bucket_name() + log_bucket = get_new_bucket_resource(name=log_bucket_name) + client = get_client() + + logging_enabled = { + 'TargetBucket': log_bucket_name, + 'LoggingType': log_type, + 'TargetPrefix': 'log/', + 'ObjectRollTime': 300, # 5 minutes + 'TargetObjectKeyFormat': {'SimplePrefix': {}}, + 'RecordsBatchSize': 0, + } + + response = client.put_bucket_logging(Bucket=src_bucket_name, BucketLoggingStatus={ + 'LoggingEnabled': logging_enabled, + }) + assert response['ResponseMetadata']['HTTPStatusCode'] == 200 + + response = client.get_bucket_logging(Bucket=src_bucket_name) + assert response['ResponseMetadata']['HTTPStatusCode'] == 200 + if log_type == 'Journal': + assert response['LoggingEnabled'] == logging_enabled + elif log_type == 'Standard': + print('TODO') + else: + assert False, 'unknown log type: %s' % log_type + + num_keys = 5 + for j in range(num_keys): + name = 'myobject'+str(j) + client.put_object(Bucket=src_bucket_name, Key=name, Body=randcontent()) + + response = client.list_objects_v2(Bucket=src_bucket_name) + src_keys = _get_keys(response) + + response = client.post_bucket_logging(Bucket=src_bucket_name) + assert response['ResponseMetadata']['HTTPStatusCode'] == 200 + + expected_count = num_keys + + response = client.list_objects_v2(Bucket=log_bucket_name) + keys = _get_keys(response) + assert len(keys) == 1 + + for key in keys: + assert key.startswith('log/') + response = client.get_object(Bucket=log_bucket_name, Key=key) + body = _get_body(response) + assert _verify_records(body, src_bucket_name, 'REST.PUT.OBJECT', src_keys, log_type, expected_count) + + +@pytest.mark.bucket_logging +@pytest.mark.fails_on_aws +def test_bucket_logging_flush_j(): + if not _has_bucket_logging_extension(): + pytest.skip('ceph extension to bucket logging not supported at client') + _bucket_logging_flush('Journal') + + +@pytest.mark.bucket_logging +@pytest.mark.fails_on_aws +def test_bucket_logging_flush_s(): + if not _has_bucket_logging_extension(): + pytest.skip('ceph extension to bucket logging not supported at client') + _bucket_logging_flush('Standard') + + +@pytest.mark.bucket_logging +def test_put_bucket_logging_errors(): + src_bucket_name = get_new_bucket_name() + src_bucket = get_new_bucket_resource(name=src_bucket_name) + log_bucket_name1 = get_new_bucket_name() + log_bucket1 = get_new_bucket_resource(name=log_bucket_name1) + client = get_client() + + # invalid source bucket + try: + response = client.put_bucket_logging(Bucket=src_bucket_name+'kaboom', BucketLoggingStatus={ + 'LoggingEnabled': {'TargetBucket': log_bucket_name1, 'TargetPrefix': 'log/'}, + }) + assert False, 'expected failure' + except ClientError as e: + assert e.response['Error']['Code'] == 'NoSuchBucket' + + # invalid log bucket + try: + response = client.put_bucket_logging(Bucket=src_bucket_name, BucketLoggingStatus={ + 'LoggingEnabled': {'TargetBucket': log_bucket_name1+'kaboom', 'TargetPrefix': 'log/'}, + }) + assert False, 'expected failure' + except ClientError as e: + assert e.response['Error']['Code'] == 'NoSuchKey' + + # log bucket has bucket logging + log_bucket_name2 = get_new_bucket_name() + log_bucket2 = get_new_bucket_resource(name=log_bucket_name2) + response = client.put_bucket_logging(Bucket=log_bucket_name2, BucketLoggingStatus={ + 'LoggingEnabled': {'TargetBucket': log_bucket_name1, 'TargetPrefix': 'log/'}, + }) + assert response['ResponseMetadata']['HTTPStatusCode'] == 200 + try: + response = client.put_bucket_logging(Bucket=src_bucket_name, BucketLoggingStatus={ + 'LoggingEnabled': {'TargetBucket': log_bucket_name2, 'TargetPrefix': 'log/'}, + }) + assert False, 'expected failure' + except ClientError as e: + assert e.response['Error']['Code'] == 'InvalidArgument' + + if _has_taget_object_key_format(): + # invalid partition prefix + logging_enabled = { + 'TargetBucket': log_bucket_name1, + 'TargetPrefix': 'log/', + 'TargetObjectKeyFormat': { + 'PartitionedPrefix': { + 'PartitionDateSource': 'kaboom' + } + } + } + try: + response = client.put_bucket_logging(Bucket=src_bucket_name, BucketLoggingStatus={ + 'LoggingEnabled': logging_enabled, + }) + assert False, 'expected failure' + except ClientError as e: + assert e.response['Error']['Code'] == 'MalformedXML' + + # TODO: log bucket is encrypted + #_put_bucket_encryption_s3(client, log_bucket_name) + #try: + # response = client.put_bucket_logging(Bucket=src_bucket_name, BucketLoggingStatus={ + # 'LoggingEnabled': {'TargetBucket': log_bucket_name, 'TargetPrefix': 'log/'}, + # }) + # assert False, 'expected failure' + #except ClientError as e: + # assert e.response['Error']['Code'] == 'InvalidArgument' + + if _has_bucket_logging_extension(): + try: + response = client.put_bucket_logging(Bucket=src_bucket_name, BucketLoggingStatus={ + 'LoggingEnabled': {'TargetBucket': log_bucket_name1, 'TargetPrefix': 'log/', 'LoggingType': 'kaboom'}, + }) + assert False, 'expected failure' + except ClientError as e: + assert e.response['Error']['Code'] == 'MalformedXML' + + +@pytest.mark.bucket_logging +def test_rm_bucket_logging(): + src_bucket_name = get_new_bucket_name() + src_bucket = get_new_bucket_resource(name=src_bucket_name) + log_bucket_name = get_new_bucket_name() + log_bucket = get_new_bucket_resource(name=log_bucket_name) + client = get_client() + logging_enabled = {'TargetBucket': log_bucket_name, 'TargetPrefix': 'log/'} + response = client.put_bucket_logging(Bucket=src_bucket_name, BucketLoggingStatus={ + 'LoggingEnabled': logging_enabled, + }) + assert response['ResponseMetadata']['HTTPStatusCode'] == 200 + + response = client.put_bucket_logging(Bucket=src_bucket_name, BucketLoggingStatus={}) + assert response['ResponseMetadata']['HTTPStatusCode'] == 200 + + response = client.get_bucket_logging(Bucket=src_bucket_name) + assert response['ResponseMetadata']['HTTPStatusCode'] == 200 + assert not 'LoggingEnabled' in response + + +@pytest.mark.bucket_logging +@pytest.mark.fails_on_aws +def test_put_bucket_logging_extensions(): + if not _has_bucket_logging_extension(): + pytest.skip('ceph extension to bucket logging not supported at client') + src_bucket_name = get_new_bucket_name() + src_bucket = get_new_bucket_resource(name=src_bucket_name) + log_bucket_name = get_new_bucket_name() + log_bucket = get_new_bucket_resource(name=log_bucket_name) + client = get_client() + logging_enabled = {'TargetBucket': log_bucket_name, + 'TargetPrefix': 'log/', + 'LoggingType': 'Standard', + 'ObjectRollTime': expected_object_roll_time, + 'RecordsBatchSize': 0 + } + response = client.put_bucket_logging(Bucket=src_bucket_name, BucketLoggingStatus={ + 'LoggingEnabled': logging_enabled, + }) + assert response['ResponseMetadata']['HTTPStatusCode'] == 200 + + response = client.get_bucket_logging(Bucket=src_bucket_name) + assert response['ResponseMetadata']['HTTPStatusCode'] == 200 + + logging_enabled['TargetObjectKeyFormat'] = {'SimplePrefix': {}} + assert response['LoggingEnabled'] == logging_enabled + + +def _bucket_logging_put_objects(versioned): + src_bucket_name = get_new_bucket() + if versioned: + check_configure_versioning_retry(src_bucket_name, "Enabled", "Enabled") + log_bucket_name = get_new_bucket_name() + log_bucket = get_new_bucket_resource(name=log_bucket_name) + client = get_client() + has_extensions = _has_bucket_logging_extension() + + # minimal configuration + logging_enabled = {'TargetBucket': log_bucket_name, 'TargetPrefix': 'log/'} + if has_extensions: + logging_enabled['ObjectRollTime'] = expected_object_roll_time + logging_enabled['LoggingType'] = 'Journal' + response = client.put_bucket_logging(Bucket=src_bucket_name, BucketLoggingStatus={ + 'LoggingEnabled': logging_enabled, + }) + assert response['ResponseMetadata']['HTTPStatusCode'] == 200 + + num_keys = 5 + for j in range(num_keys): + name = 'myobject'+str(j) + client.put_object(Bucket=src_bucket_name, Key=name, Body=randcontent()) + if versioned: + client.put_object(Bucket=src_bucket_name, Key=name, Body=randcontent()) + + if versioned: + expected_count = 2*num_keys + else: + expected_count = num_keys + + response = client.list_objects_v2(Bucket=src_bucket_name) + src_keys = _get_keys(response) + + time.sleep(expected_object_roll_time) + client.put_object(Bucket=src_bucket_name, Key='dummy', Body='dummy') + + response = client.list_objects_v2(Bucket=log_bucket_name) + keys = _get_keys(response) + assert len(keys) == 1 + + record_type = 'Standard' if not has_extensions else 'Journal' + + for key in keys: + assert key.startswith('log/') + response = client.get_object(Bucket=log_bucket_name, Key=key) + body = _get_body(response) + assert _verify_records(body, src_bucket_name, 'REST.PUT.OBJECT', src_keys, record_type, expected_count) + + +@pytest.mark.bucket_logging +def test_bucket_logging_put_objects(): + _bucket_logging_put_objects(False) + + +@pytest.mark.bucket_logging +def test_bucket_logging_put_objects_versioned(): + _bucket_logging_put_objects(True) + + +@pytest.mark.bucket_logging +def test_bucket_logging_put_concurrency(): + src_bucket_name = get_new_bucket() + log_bucket_name = get_new_bucket_name() + log_bucket = get_new_bucket_resource(name=log_bucket_name) + client = get_client(client_config=botocore.config.Config(max_pool_connections=50)) + has_extensions = _has_bucket_logging_extension() + + # minimal configuration + logging_enabled = {'TargetBucket': log_bucket_name, 'TargetPrefix': 'log/'} + if has_extensions: + logging_enabled['ObjectRollTime'] = expected_object_roll_time + logging_enabled['LoggingType'] = 'Journal' + response = client.put_bucket_logging(Bucket=src_bucket_name, BucketLoggingStatus={ + 'LoggingEnabled': logging_enabled, + }) + assert response['ResponseMetadata']['HTTPStatusCode'] == 200 + + num_keys = 50 + t = [] + for i in range(num_keys): + name = 'myobject'+str(i) + thr = threading.Thread(target = client.put_object, + kwargs={'Bucket': src_bucket_name, 'Key': name, 'Body': randcontent()}) + thr.start() + t.append(thr) + _do_wait_completion(t) + + response = client.list_objects_v2(Bucket=src_bucket_name) + src_keys = _get_keys(response) + + time.sleep(expected_object_roll_time) + t = [] + for i in range(num_keys): + thr = threading.Thread(target = client.put_object, + kwargs={'Bucket': src_bucket_name, 'Key': 'dummy', 'Body': 'dummy'}) + thr.start() + t.append(thr) + _do_wait_completion(t) + + response = client.list_objects_v2(Bucket=log_bucket_name) + keys = _get_keys(response) + assert len(keys) == 1 + + record_type = 'Standard' if not has_extensions else 'Journal' + + for key in keys: + logger.info('logging object: %s', key) + assert key.startswith('log/') + response = client.get_object(Bucket=log_bucket_name, Key=key) + body = _get_body(response) + assert _verify_records(body, src_bucket_name, 'REST.PUT.OBJECT', src_keys, record_type, num_keys) + + +def _bucket_logging_delete_objects(versioned): + src_bucket_name = get_new_bucket() + if versioned: + check_configure_versioning_retry(src_bucket_name, "Enabled", "Enabled") + log_bucket_name = get_new_bucket_name() + log_bucket = get_new_bucket_resource(name=log_bucket_name) + client = get_client() + has_extensions = _has_bucket_logging_extension() + + num_keys = 5 + for j in range(num_keys): + name = 'myobject'+str(j) + client.put_object(Bucket=src_bucket_name, Key=name, Body=randcontent()) + if versioned: + client.put_object(Bucket=src_bucket_name, Key=name, Body=randcontent()) + + # minimal configuration + logging_enabled = {'TargetBucket': log_bucket_name, 'TargetPrefix': 'log/'} + if has_extensions: + logging_enabled['ObjectRollTime'] = expected_object_roll_time + logging_enabled['LoggingType'] = 'Journal' + response = client.put_bucket_logging(Bucket=src_bucket_name, BucketLoggingStatus={ + 'LoggingEnabled': logging_enabled, + }) + + assert response['ResponseMetadata']['HTTPStatusCode'] == 200 + response = client.list_objects_v2(Bucket=src_bucket_name) + src_keys = _get_keys(response) + for key in src_keys: + if versioned: + response = client.head_object(Bucket=src_bucket_name, Key=key) + client.delete_object(Bucket=src_bucket_name, Key=key, VersionId=response['VersionId']) + client.delete_object(Bucket=src_bucket_name, Key=key) + + time.sleep(expected_object_roll_time) + client.put_object(Bucket=src_bucket_name, Key='dummy', Body='dummy') + + response = client.list_objects_v2(Bucket=log_bucket_name) + keys = _get_keys(response) + assert len(keys) == 1 + + if versioned: + expected_count = 2*num_keys + else: + expected_count = num_keys + + key = keys[0] + assert key.startswith('log/') + response = client.get_object(Bucket=log_bucket_name, Key=key) + body = _get_body(response) + record_type = 'Standard' if not has_extensions else 'Journal' + assert _verify_records(body, src_bucket_name, 'REST.DELETE.OBJECT', src_keys, record_type, expected_count) + + +@pytest.mark.bucket_logging +def test_bucket_logging_delete_objects(): + _bucket_logging_delete_objects(False) + + +@pytest.mark.bucket_logging +def test_bucket_logging_delete_objects_versioned(): + _bucket_logging_delete_objects(True) + + +@pytest.mark.bucket_logging +def _bucket_logging_get_objects(versioned): + src_bucket_name = get_new_bucket() + if versioned: + check_configure_versioning_retry(src_bucket_name, "Enabled", "Enabled") + log_bucket_name = get_new_bucket_name() + log_bucket = get_new_bucket_resource(name=log_bucket_name) + client = get_client() + has_extensions = _has_bucket_logging_extension() + + num_keys = 5 + for j in range(num_keys): + name = 'myobject'+str(j) + client.put_object(Bucket=src_bucket_name, Key=name, Body=randcontent()) + if versioned: + client.put_object(Bucket=src_bucket_name, Key=name, Body=randcontent()) + + # minimal configuration + logging_enabled = {'TargetBucket': log_bucket_name, 'TargetPrefix': 'log/'} + if has_extensions: + logging_enabled['ObjectRollTime'] = expected_object_roll_time + logging_enabled['LoggingType'] = 'Standard' + response = client.put_bucket_logging(Bucket=src_bucket_name, BucketLoggingStatus={ + 'LoggingEnabled': logging_enabled, + }) + + assert response['ResponseMetadata']['HTTPStatusCode'] == 200 + response = client.list_objects_v2(Bucket=src_bucket_name) + src_keys = _get_keys(response) + for key in src_keys: + if versioned: + response = client.head_object(Bucket=src_bucket_name, Key=key) + client.get_object(Bucket=src_bucket_name, Key=key, VersionId=response['VersionId']) + client.get_object(Bucket=src_bucket_name, Key=key) + + time.sleep(expected_object_roll_time) + client.put_object(Bucket=src_bucket_name, Key='dummy', Body='dummy') + + response = client.list_objects_v2(Bucket=log_bucket_name) + keys = _get_keys(response) + assert len(keys) == 1 + + if versioned: + expected_count = 2*num_keys + else: + expected_count = num_keys + + key = keys[0] + assert key.startswith('log/') + response = client.get_object(Bucket=log_bucket_name, Key=key) + body = _get_body(response) + assert _verify_records(body, src_bucket_name, 'REST.GET.OBJECT', src_keys, 'Standard', expected_count) + + +@pytest.mark.bucket_logging +def test_bucket_logging_get_objects(): + _bucket_logging_get_objects(False) + + +@pytest.mark.bucket_logging +def test_bucket_logging_get_objects_versioned(): + _bucket_logging_get_objects(True) + + +@pytest.mark.bucket_logging +def _bucket_logging_copy_objects(versioned, another_bucket): + src_bucket_name = get_new_bucket() + if another_bucket: + dst_bucket_name = get_new_bucket() + else: + dst_bucket_name = src_bucket_name + if versioned: + check_configure_versioning_retry(src_bucket_name, "Enabled", "Enabled") + log_bucket_name = get_new_bucket_name() + log_bucket = get_new_bucket_resource(name=log_bucket_name) + client = get_client() + has_extensions = _has_bucket_logging_extension() + + num_keys = 5 + for j in range(num_keys): + name = 'myobject'+str(j) + client.put_object(Bucket=src_bucket_name, Key=name, Body=randcontent()) + if versioned: + client.put_object(Bucket=src_bucket_name, Key=name, Body=randcontent()) + + # minimal configuration + logging_enabled = {'TargetBucket': log_bucket_name, 'TargetPrefix': 'log/'} + if has_extensions: + logging_enabled['ObjectRollTime'] = expected_object_roll_time + logging_enabled['LoggingType'] = 'Journal' + response = client.put_bucket_logging(Bucket=src_bucket_name, BucketLoggingStatus={ + 'LoggingEnabled': logging_enabled, + }) + assert response['ResponseMetadata']['HTTPStatusCode'] == 200 + if another_bucket: + response = client.put_bucket_logging(Bucket=dst_bucket_name, BucketLoggingStatus={ + 'LoggingEnabled': logging_enabled, + }) + assert response['ResponseMetadata']['HTTPStatusCode'] == 200 + + response = client.list_objects_v2(Bucket=src_bucket_name) + src_keys = _get_keys(response) + dst_keys = [] + for key in src_keys: + dst_keys.append('copy_of_'+key) + if another_bucket: + client.copy_object(Bucket=dst_bucket_name, Key='copy_of_'+key, CopySource={'Bucket': src_bucket_name, 'Key': key}) + else: + client.copy_object(Bucket=src_bucket_name, Key='copy_of_'+key, CopySource={'Bucket': src_bucket_name, 'Key': key}) + + time.sleep(expected_object_roll_time) + client.put_object(Bucket=src_bucket_name, Key='dummy', Body='dummy') + + response = client.list_objects_v2(Bucket=log_bucket_name) + keys = _get_keys(response) + assert len(keys) == 1 + + key = keys[0] + assert key.startswith('log/') + response = client.get_object(Bucket=log_bucket_name, Key=key) + body = _get_body(response) + record_type = 'Standard' if not has_extensions else 'Journal' + assert _verify_records(body, dst_bucket_name, 'REST.PUT.OBJECT', dst_keys, record_type, num_keys) + + +@pytest.mark.bucket_logging +def test_bucket_logging_copy_objects(): + _bucket_logging_copy_objects(False, False) + + +@pytest.mark.bucket_logging +def test_bucket_logging_copy_objects_versioned(): + _bucket_logging_copy_objects(True, False) + + +@pytest.mark.bucket_logging +def test_bucket_logging_copy_objects_bucket(): + _bucket_logging_copy_objects(False, True) + + +@pytest.mark.bucket_logging +def test_bucket_logging_copy_objects_bucket_versioned(): + _bucket_logging_copy_objects(True, True) + + +@pytest.mark.bucket_logging +def _bucket_logging_head_objects(versioned): + src_bucket_name = get_new_bucket() + if versioned: + check_configure_versioning_retry(src_bucket_name, "Enabled", "Enabled") + log_bucket_name = get_new_bucket_name() + log_bucket = get_new_bucket_resource(name=log_bucket_name) + client = get_client() + has_extensions = _has_bucket_logging_extension() + + num_keys = 5 + for j in range(num_keys): + name = 'myobject'+str(j) + client.put_object(Bucket=src_bucket_name, Key=name, Body=randcontent()) + + logging_enabled = {'TargetBucket': log_bucket_name, 'TargetPrefix': 'log/'} + if has_extensions: + logging_enabled['ObjectRollTime'] = expected_object_roll_time + logging_enabled['LoggingType'] = 'Standard' + response = client.put_bucket_logging(Bucket=src_bucket_name, BucketLoggingStatus={ + 'LoggingEnabled': logging_enabled, + }) + + assert response['ResponseMetadata']['HTTPStatusCode'] == 200 + response = client.list_objects_v2(Bucket=src_bucket_name) + src_keys = _get_keys(response) + for key in src_keys: + if versioned: + response = client.head_object(Bucket=src_bucket_name, Key=key) + client.head_object(Bucket=src_bucket_name, Key=key, VersionId=response['VersionId']) + else: + client.head_object(Bucket=src_bucket_name, Key=key) + + time.sleep(expected_object_roll_time) + client.put_object(Bucket=src_bucket_name, Key='dummy', Body='dummy') + + response = client.list_objects_v2(Bucket=log_bucket_name) + keys = _get_keys(response) + assert len(keys) == 1 + + if versioned: + expected_count = 2*num_keys + else: + expected_count = num_keys + + key = keys[0] + assert key.startswith('log/') + response = client.get_object(Bucket=log_bucket_name, Key=key) + body = _get_body(response) + assert _verify_records(body, src_bucket_name, 'REST.HEAD.OBJECT', src_keys, 'Standard', expected_count) + + +@pytest.mark.bucket_logging +def test_bucket_logging_head_objects(): + _bucket_logging_head_objects(False) + + +@pytest.mark.bucket_logging +def test_bucket_logging_head_objects_versioned(): + _bucket_logging_head_objects(True) + + +@pytest.mark.bucket_logging +def _bucket_logging_mpu(versioned): + src_bucket_name = get_new_bucket() + if versioned: + check_configure_versioning_retry(src_bucket_name, "Enabled", "Enabled") + log_bucket_name = get_new_bucket_name() + log_bucket = get_new_bucket_resource(name=log_bucket_name) + client = get_client() + has_extensions = _has_bucket_logging_extension() + + # minimal configuration + logging_enabled = {'TargetBucket': log_bucket_name, 'TargetPrefix': 'log/'} + if has_extensions: + logging_enabled['ObjectRollTime'] = expected_object_roll_time + logging_enabled['LoggingType'] = 'Journal' + response = client.put_bucket_logging(Bucket=src_bucket_name, BucketLoggingStatus={ + 'LoggingEnabled': logging_enabled, + }) + assert response['ResponseMetadata']['HTTPStatusCode'] == 200 + + src_key = "myobject" + objlen = 30 * 1024 * 1024 + (upload_id, data, parts) = _multipart_upload(bucket_name=src_bucket_name, key=src_key, size=objlen) + client.complete_multipart_upload(Bucket=src_bucket_name, Key=src_key, UploadId=upload_id, MultipartUpload={'Parts': parts}) + if versioned: + (upload_id, data, parts) = _multipart_upload(bucket_name=src_bucket_name, key=src_key, size=objlen) + client.complete_multipart_upload(Bucket=src_bucket_name, Key=src_key, UploadId=upload_id, MultipartUpload={'Parts': parts}) + + time.sleep(expected_object_roll_time) + client.put_object(Bucket=src_bucket_name, Key='dummy', Body='dummy') + + response = client.list_objects_v2(Bucket=log_bucket_name) + keys = _get_keys(response) + assert len(keys) == 1 + + if versioned: + expected_count = 4 if not has_extensions else 2 + else: + expected_count = 2 if not has_extensions else 1 + + key = keys[0] + assert key.startswith('log/') + response = client.get_object(Bucket=log_bucket_name, Key=key) + body = _get_body(response) + record_type = 'Standard' if not has_extensions else 'Journal' + assert _verify_records(body, src_bucket_name, 'REST.POST.UPLOAD', [src_key, src_key], record_type, expected_count) + + +@pytest.mark.bucket_logging +def test_bucket_logging_mpu(): + _bucket_logging_mpu(False) + + +@pytest.mark.bucket_logging +def test_bucket_logging_mpu_versioned(): + _bucket_logging_mpu(True) + + +@pytest.mark.bucket_logging +def _bucket_logging_mpu_copy(versioned): + src_bucket_name = get_new_bucket() + if versioned: + check_configure_versioning_retry(src_bucket_name, "Enabled", "Enabled") + log_bucket_name = get_new_bucket_name() + log_bucket = get_new_bucket_resource(name=log_bucket_name) + client = get_client() + has_extensions = _has_bucket_logging_extension() + + src_key = "myobject" + objlen = 30 * 1024 * 1024 + (upload_id, data, parts) = _multipart_upload(bucket_name=src_bucket_name, key=src_key, size=objlen) + client.complete_multipart_upload(Bucket=src_bucket_name, Key=src_key, UploadId=upload_id, MultipartUpload={'Parts': parts}) + if versioned: + (upload_id, data, parts) = _multipart_upload(bucket_name=src_bucket_name, key=src_key, size=objlen) + client.complete_multipart_upload(Bucket=src_bucket_name, Key=src_key, UploadId=upload_id, MultipartUpload={'Parts': parts}) + + # minimal configuration + logging_enabled = {'TargetBucket': log_bucket_name, 'TargetPrefix': 'log/'} + if has_extensions: + logging_enabled['ObjectRollTime'] = expected_object_roll_time + logging_enabled['LoggingType'] = 'Journal' + response = client.put_bucket_logging(Bucket=src_bucket_name, BucketLoggingStatus={ + 'LoggingEnabled': logging_enabled, + }) + assert response['ResponseMetadata']['HTTPStatusCode'] == 200 + + client.copy_object(Bucket=src_bucket_name, Key='copy_of_'+src_key, CopySource={'Bucket': src_bucket_name, 'Key': src_key}) + + time.sleep(expected_object_roll_time) + client.put_object(Bucket=src_bucket_name, Key='dummy', Body='dummy') + + response = client.list_objects_v2(Bucket=log_bucket_name) + keys = _get_keys(response) + assert len(keys) == 1 + + key = keys[0] + assert key.startswith('log/') + response = client.get_object(Bucket=log_bucket_name, Key=key) + body = _get_body(response) + record_type = 'Standard' if not has_extensions else 'Journal' + assert _verify_records(body, src_bucket_name, 'REST.PUT.OBJECT', ['copy_of_'+src_key], record_type, 1) + + +@pytest.mark.bucket_logging +def test_bucket_logging_mpu_copy(): + _bucket_logging_mpu_copy(False) + + +@pytest.mark.bucket_logging +def test_bucket_logging_mpu_copy_versioned(): + _bucket_logging_mpu_copy(True) + + +def _bucket_logging_multi_delete(versioned): + src_bucket_name = get_new_bucket() + if versioned: + check_configure_versioning_retry(src_bucket_name, "Enabled", "Enabled") + log_bucket_name = get_new_bucket_name() + log_bucket = get_new_bucket_resource(name=log_bucket_name) + client = get_client() + has_extensions = _has_bucket_logging_extension() + + num_keys = 5 + for j in range(num_keys): + name = 'myobject'+str(j) + client.put_object(Bucket=src_bucket_name, Key=name, Body=randcontent()) + if versioned: + client.put_object(Bucket=src_bucket_name, Key=name, Body=randcontent()) + + # minimal configuration + logging_enabled = {'TargetBucket': log_bucket_name, 'TargetPrefix': 'log/'} + if has_extensions: + logging_enabled['ObjectRollTime'] = expected_object_roll_time + logging_enabled['LoggingType'] = 'Journal' + response = client.put_bucket_logging(Bucket=src_bucket_name, BucketLoggingStatus={ + 'LoggingEnabled': logging_enabled, + }) + + assert response['ResponseMetadata']['HTTPStatusCode'] == 200 + response = client.list_objects_v2(Bucket=src_bucket_name) + src_keys = _get_keys(response) + if versioned: + response = client.list_object_versions(Bucket=src_bucket_name) + objs_list = [] + for version in response['Versions']: + obj_dict = {'Key': version['Key'], 'VersionId': version['VersionId']} + objs_list.append(obj_dict) + objs_dict = {'Objects': objs_list} + client.delete_objects(Bucket=src_bucket_name, Delete=objs_dict) + else: + objs_dict = _make_objs_dict(key_names=src_keys) + client.delete_objects(Bucket=src_bucket_name, Delete=objs_dict) + + time.sleep(expected_object_roll_time) + client.put_object(Bucket=src_bucket_name, Key='dummy', Body='dummy') + + response = client.list_objects_v2(Bucket=log_bucket_name) + keys = _get_keys(response) + assert len(keys) == 1 + + if versioned: + expected_count = 2*num_keys + else: + expected_count = num_keys + + key = keys[0] + assert key.startswith('log/') + response = client.get_object(Bucket=log_bucket_name, Key=key) + body = _get_body(response) + record_type = 'Standard' if not has_extensions else 'Journal' + assert _verify_records(body, src_bucket_name, "REST.POST.DELETE_MULTI_OBJECT", src_keys, record_type, expected_count) + + +@pytest.mark.bucket_logging +def test_bucket_logging_multi_delete(): + _bucket_logging_multi_delete(False) + + +@pytest.mark.bucket_logging +def test_bucket_logging_multi_delete_versioned(): + _bucket_logging_multi_delete(True) + + +def _bucket_logging_type(logging_type): + src_bucket_name = get_new_bucket_name() + src_bucket = get_new_bucket_resource(name=src_bucket_name) + log_bucket_name = get_new_bucket_name() + log_bucket = get_new_bucket_resource(name=log_bucket_name) + client = get_client() + logging_enabled = { + 'TargetBucket': log_bucket_name, + 'TargetPrefix': 'log/', + 'ObjectRollTime': expected_object_roll_time, + 'LoggingType': logging_type + } + response = client.put_bucket_logging(Bucket=src_bucket_name, BucketLoggingStatus={ + 'LoggingEnabled': logging_enabled, + }) + num_keys = 5 + for j in range(num_keys): + name = 'myobject'+str(j) + client.put_object(Bucket=src_bucket_name, Key=name, Body=randcontent()) + client.head_object(Bucket=src_bucket_name, Key=name) + + response = client.list_objects_v2(Bucket=src_bucket_name) + src_keys = _get_keys(response) + + time.sleep(expected_object_roll_time) + client.put_object(Bucket=src_bucket_name, Key='dummy', Body='dummy') + client.head_object(Bucket=src_bucket_name, Key='dummy') + + response = client.list_objects_v2(Bucket=log_bucket_name) + keys = _get_keys(response) + assert len(keys) == 1 + + key = keys[0] + assert key.startswith('log/') + response = client.get_object(Bucket=log_bucket_name, Key=key) + body = _get_body(response) + if logging_type == 'Journal': + assert _verify_records(body, src_bucket_name, 'REST.PUT.OBJECT', src_keys, 'Journal', num_keys) + assert _verify_records(body, src_bucket_name, 'REST.HEAD.OBJECT', src_keys, 'Journal', num_keys) == False + elif logging_type == 'Standard': + assert _verify_records(body, src_bucket_name, 'REST.HEAD.OBJECT', src_keys, 'Standard', num_keys) + assert _verify_records(body, src_bucket_name, 'REST.PUT.OBJECT', src_keys, 'Standard', num_keys) + else: + assert False, 'invalid logging type:'+logging_type + + +@pytest.mark.bucket_logging +@pytest.mark.fails_on_aws +def test_bucket_logging_event_type_j(): + if not _has_bucket_logging_extension(): + pytest.skip('ceph extension to bucket logging not supported at client') + _bucket_logging_type('Journal') + + +@pytest.mark.bucket_logging +@pytest.mark.fails_on_aws +def test_bucket_logging_event_type_s(): + if not _has_bucket_logging_extension(): + pytest.skip('ceph extension to bucket logging not supported at client') + _bucket_logging_type('Standard') + + +@pytest.mark.bucket_logging +@pytest.mark.fails_on_aws +def test_bucket_logging_roll_time(): + if not _has_bucket_logging_extension(): + pytest.skip('ceph extension to bucket logging not supported at client') + src_bucket_name = get_new_bucket_name() + src_bucket = get_new_bucket_resource(name=src_bucket_name) + log_bucket_name = get_new_bucket_name() + log_bucket = get_new_bucket_resource(name=log_bucket_name) + client = get_client() + + roll_time = 10 + logging_enabled = {'TargetBucket': log_bucket_name, 'TargetPrefix': 'log/', 'ObjectRollTime': roll_time} + response = client.put_bucket_logging(Bucket=src_bucket_name, BucketLoggingStatus={ + 'LoggingEnabled': logging_enabled, + }) + assert response['ResponseMetadata']['HTTPStatusCode'] == 200 + + num_keys = 5 + for j in range(num_keys): + name = 'myobject'+str(j) + client.put_object(Bucket=src_bucket_name, Key=name, Body=randcontent()) + + response = client.list_objects_v2(Bucket=src_bucket_name) + src_keys = _get_keys(response) + + time.sleep(roll_time/2) + client.put_object(Bucket=src_bucket_name, Key='myobject', Body=randcontent()) + + response = client.list_objects_v2(Bucket=log_bucket_name) + keys = _get_keys(response) + assert len(keys) == 0 + + time.sleep(roll_time/2) + client.put_object(Bucket=src_bucket_name, Key='myobject', Body=randcontent()) + + response = client.list_objects_v2(Bucket=log_bucket_name) + keys = _get_keys(response) + len(keys) == 1 + + key = keys[0] + assert key.startswith('log/') + response = client.get_object(Bucket=log_bucket_name, Key=key) + body = _get_body(response) + assert _verify_records(body, src_bucket_name, 'REST.PUT.OBJECT', src_keys, 'Standard', num_keys) + client.delete_object(Bucket=log_bucket_name, Key=key) + + num_keys = 25 + for j in range(num_keys): + name = 'myobject'+str(j) + client.put_object(Bucket=src_bucket_name, Key=name, Body=randcontent()) + time.sleep(1) + + response = client.list_objects_v2(Bucket=src_bucket_name) + src_keys = _get_keys(response) + + time.sleep(roll_time) + client.put_object(Bucket=src_bucket_name, Key='myobject', Body=randcontent()) + + response = client.list_objects_v2(Bucket=log_bucket_name) + keys = _get_keys(response) + assert len(keys) > 1 + + body = '' + for key in keys: + assert key.startswith('log/') + response = client.get_object(Bucket=log_bucket_name, Key=key) + body += _get_body(response) + assert _verify_records(body, src_bucket_name, 'REST.PUT.OBJECT', src_keys, 'Standard', num_keys+1) + + +@pytest.mark.bucket_logging +def test_bucket_logging_multiple_prefixes(): + log_bucket_name = get_new_bucket_name() + log_bucket = get_new_bucket_resource(name=log_bucket_name) + client = get_client() + has_extensions = _has_bucket_logging_extension() + + num_buckets = 5 + buckets = [] + bucket_name_prefix = get_new_bucket_name() + for j in range(num_buckets): + src_bucket_name = bucket_name_prefix+str(j) + src_bucket = get_new_bucket_resource(name=src_bucket_name) + logging_enabled = {'TargetBucket': log_bucket_name, 'TargetPrefix': src_bucket_name+'/'} + if has_extensions: + logging_enabled['ObjectRollTime'] = expected_object_roll_time + response = client.put_bucket_logging(Bucket=src_bucket_name, BucketLoggingStatus={ + 'LoggingEnabled': logging_enabled, + }) + assert response['ResponseMetadata']['HTTPStatusCode'] == 200 + buckets.append(src_bucket_name) + + num_keys = 5 + for src_bucket_name in buckets: + for j in range(num_keys): + name = 'myobject'+str(j) + client.put_object(Bucket=src_bucket_name, Key=name, Body=randcontent()) + + time.sleep(expected_object_roll_time) + for src_bucket_name in buckets: + client.head_object(Bucket=src_bucket_name, Key='myobject0') + + response = client.list_objects_v2(Bucket=log_bucket_name) + keys = _get_keys(response) + assert len(keys) >= num_buckets + + for key in keys: + response = client.get_object(Bucket=log_bucket_name, Key=key) + body = _get_body(response) + found = False + for src_bucket_name in buckets: + if key.startswith(src_bucket_name): + found = True + response = client.list_objects_v2(Bucket=src_bucket_name) + src_keys = _get_keys(response) + assert _verify_records(body, src_bucket_name, 'REST.PUT.OBJECT', src_keys, 'Standard', num_keys) + assert found + + +@pytest.mark.bucket_logging +def test_bucket_logging_single_prefix(): + log_bucket_name = get_new_bucket_name() + log_bucket = get_new_bucket_resource(name=log_bucket_name) + client = get_client() + has_extensions = _has_bucket_logging_extension() + + num_buckets = 5 + buckets = [] + bucket_name_prefix = get_new_bucket_name() + for j in range(num_buckets): + src_bucket_name = bucket_name_prefix+str(j) + src_bucket = get_new_bucket_resource(name=src_bucket_name) + # minimal configuration + logging_enabled = {'TargetBucket': log_bucket_name, 'TargetPrefix': 'log/'} + if has_extensions: + logging_enabled['ObjectRollTime'] = expected_object_roll_time + response = client.put_bucket_logging(Bucket=src_bucket_name, BucketLoggingStatus={ + 'LoggingEnabled': logging_enabled, + }) + assert response['ResponseMetadata']['HTTPStatusCode'] == 200 + buckets.append(src_bucket_name) + + num_keys = 5 + bucket_ind = 0 + for src_bucket_name in buckets: + bucket_ind += 1 + for j in range(num_keys): + name = 'myobject'+str(bucket_ind)+str(j) + client.put_object(Bucket=src_bucket_name, Key=name, Body=randcontent()) + + time.sleep(expected_object_roll_time) + client.put_object(Bucket=buckets[0], Key='dummy', Body='dummy') + + response = client.list_objects_v2(Bucket=log_bucket_name) + keys = _get_keys(response) + assert len(keys) == 1 + + key = keys[0] + response = client.get_object(Bucket=log_bucket_name, Key=key) + body = _get_body(response) + found = False + for src_bucket_name in buckets: + response = client.list_objects_v2(Bucket=src_bucket_name) + src_keys = _get_keys(response) + found = _verify_records(body, src_bucket_name, 'REST.PUT.OBJECT', src_keys, 'Standard', num_keys) + assert found + +@pytest.mark.checksum +@pytest.mark.fails_on_dbstore +def test_get_object_attributes(): + bucket_name = get_new_bucket() + client = get_client() + + #pdb.set_trace() + key = "multipart_checksum" + key_metadata = {'foo': 'bar'} + content_type = 'text/plain' + objlen = 64 * 1024 * 1024 + + (upload_id, data, parts, checksums) = \ + _multipart_upload_checksum(bucket_name=bucket_name, key=key, size=objlen, + content_type=content_type, metadata=key_metadata) + response = client.complete_multipart_upload(Bucket=bucket_name, Key=key, + UploadId=upload_id, + MultipartUpload={'Parts': parts}) + upload_checksum = response['ChecksumSHA256'] + + response = client.get_object(Bucket=bucket_name, Key=key) + + request_attributes = ['ETag', 'Checksum', 'ObjectParts', 'StorageClass', 'ObjectSize'] + + response = client.get_object_attributes(Bucket=bucket_name, Key=key, \ + ObjectAttributes=request_attributes) + + # check overall object + nparts = len(parts) + assert response['ObjectSize'] == objlen + assert response['Checksum']['ChecksumSHA256'] == upload_checksum + assert response['ObjectParts']['TotalPartsCount'] == nparts + + # check the parts + partno = 1 + for obj_part in response['ObjectParts']['Parts']: + assert obj_part['PartNumber'] == partno + if partno < len(parts): + assert obj_part['Size'] == 5 * 1024 * 1024 + else: + assert obj_part['Size'] == objlen - ((nparts-1) * (5 * 1024 * 1024)) + assert obj_part['ChecksumSHA256'] == checksums[partno - 1] + partno += 1 From 108c148a4488c9222294d4f69414fdc049b243cc Mon Sep 17 00:00:00 2001 From: Matt Benjamin Date: Sun, 13 Oct 2024 10:47:17 -0400 Subject: [PATCH 08/12] multipart fallback to create-multipart checksum algorithm there seem to be workloads which assume checksum algorithm can be omitted from upload-part Signed-off-by: Matt Benjamin --- s3tests_boto3/functional/test_s3.py | 41 +++++++++++++++++++++++++++++ 1 file changed, 41 insertions(+) diff --git a/s3tests_boto3/functional/test_s3.py b/s3tests_boto3/functional/test_s3.py index 263e50b8c..5abf58180 100644 --- a/s3tests_boto3/functional/test_s3.py +++ b/s3tests_boto3/functional/test_s3.py @@ -13717,6 +13717,47 @@ def test_multipart_checksum_3parts(): response = client.head_object(Bucket=bucket, Key=key, ChecksumMode='ENABLED') assert composite_sha256sum == response['ChecksumSHA256'] +@pytest.mark.checksum +@pytest.mark.fails_on_dbstore +def test_multipart_checksum_upload_fallback(): + bucket = get_new_bucket() + client = get_client() + + key = "mpu_cksum_fallback" + alg = 'SHA256' + + response = client.create_multipart_upload( + Bucket=bucket, Key=key, ChecksumAlgorithm=alg) + assert alg == response['ChecksumAlgorithm'] + upload_id = response['UploadId'] + + nparts = 3 + parts = [] + size = 5 * 1024 * 1024 # each part but the last must be at least 5M + + for ix in range(0,nparts): + body = FakeWriteFile(size, 'A') + part_num = ix + 1 + res = client.upload_part(UploadId=upload_id, Bucket=bucket, + Key=key, PartNumber=part_num, Body=body) + etag = res['ETag'] + part = {'ETag': etag, 'PartNumber': part_num} + parts.append(part) + + res = client.complete_multipart_upload( + Bucket=bucket, Key=key, UploadId=upload_id, + MultipartUpload={'Parts': parts}) + + #pdb.set_trace() + assert res['ResponseMetadata']['HTTPStatusCode'] == 200 + + # not yet merged + #request_attributes = ['ETag', 'Checksum', 'ObjectParts', 'StorageClass', + # 'ObjectSize'] + #res = client.get_object_attributes(Bucket=bucket, Key=key, \ + # ObjectAttributes=request_attributes) + #upload_checksum = res['Checksum']['ChecksumSHA256'] + @pytest.mark.checksum def test_post_object_upload_checksum(): megabytes = 1024 * 1024 From 2254dd31150c4dc4521e13007e590e2f841e54cb Mon Sep 17 00:00:00 2001 From: Casey Bodley Date: Thu, 17 Oct 2024 18:26:07 -0400 Subject: [PATCH 09/12] more tests for GetObjectAttributes * multipart upload without checksums * multipart upload with a single part * pagination of multipart parts * non-multipart upload with/without checksum * versioned object, current and non-current * sse-c encrypted object Signed-off-by: Casey Bodley Signed-off-by: Matt Benjamin --- s3tests_boto3/functional/test_s3.py | 263 +++++++++++++++++++++++++++- 1 file changed, 261 insertions(+), 2 deletions(-) diff --git a/s3tests_boto3/functional/test_s3.py b/s3tests_boto3/functional/test_s3.py index 5abf58180..01dc6ea8c 100644 --- a/s3tests_boto3/functional/test_s3.py +++ b/s3tests_boto3/functional/test_s3.py @@ -15120,9 +15120,16 @@ def test_bucket_logging_single_prefix(): found = _verify_records(body, src_bucket_name, 'REST.PUT.OBJECT', src_keys, 'Standard', num_keys) assert found +def check_parts_count(parts, expected): + # AWS docs disagree on the name of this element + if 'TotalPartsCount' in parts: + assert parts['TotalPartsCount'] == expected + else: + assert parts['PartsCount'] == expected + @pytest.mark.checksum @pytest.mark.fails_on_dbstore -def test_get_object_attributes(): +def test_get_multipart_checksum_object_attributes(): bucket_name = get_new_bucket() client = get_client() @@ -15151,7 +15158,7 @@ def test_get_object_attributes(): nparts = len(parts) assert response['ObjectSize'] == objlen assert response['Checksum']['ChecksumSHA256'] == upload_checksum - assert response['ObjectParts']['TotalPartsCount'] == nparts + check_parts_count(response['ObjectParts'], nparts) # check the parts partno = 1 @@ -15163,3 +15170,255 @@ def test_get_object_attributes(): assert obj_part['Size'] == objlen - ((nparts-1) * (5 * 1024 * 1024)) assert obj_part['ChecksumSHA256'] == checksums[partno - 1] partno += 1 + +@pytest.mark.fails_on_dbstore +def test_get_multipart_object_attributes(): + bucket_name = get_new_bucket() + client = get_client() + + key = "multipart" + part_size = 5*1024*1024 + objlen = 30*1024*1024 + + (upload_id, data, parts) = _multipart_upload(bucket_name, key, objlen, part_size) + response = client.complete_multipart_upload(Bucket=bucket_name, Key=key, + UploadId=upload_id, + MultipartUpload={'Parts': parts}) + etag = response['ETag'].strip('"') + assert len(etag) + + request_attributes = ['ETag', 'Checksum', 'ObjectParts', 'StorageClass', 'ObjectSize'] + response = client.get_object_attributes(Bucket=bucket_name, Key=key, \ + ObjectAttributes=request_attributes) + + # check overall object + nparts = len(parts) + assert response['ObjectSize'] == objlen + check_parts_count(response['ObjectParts'], len(parts)) + assert response['ObjectParts']['IsTruncated'] == False + assert response['ETag'] == etag + assert response['StorageClass'] == 'STANDARD' + + # check the parts + partno = 1 + for obj_part in response['ObjectParts']['Parts']: + assert obj_part['PartNumber'] == partno + assert obj_part['Size'] == part_size + assert 'ChecksumSHA256' not in obj_part + partno += 1 + +@pytest.mark.fails_on_dbstore +def test_get_paginated_multipart_object_attributes(): + bucket_name = get_new_bucket() + client = get_client() + + key = "multipart" + part_size = 5*1024*1024 + objlen = 30*1024*1024 + + (upload_id, data, parts) = _multipart_upload(bucket_name, key, objlen, part_size) + response = client.complete_multipart_upload(Bucket=bucket_name, Key=key, + UploadId=upload_id, + MultipartUpload={'Parts': parts}) + etag = response['ETag'].strip('"') + assert len(etag) + + request_attributes = ['ETag', 'Checksum', 'ObjectParts', 'StorageClass', 'ObjectSize'] + response = client.get_object_attributes(Bucket=bucket_name, Key=key, + ObjectAttributes=request_attributes, + MaxParts=1, PartNumberMarker=3) + + # check overall object + assert response['ObjectSize'] == objlen + check_parts_count(response['ObjectParts'], len(parts)) + assert response['ObjectParts']['MaxParts'] == 1 + assert response['ObjectParts']['PartNumberMarker'] == 3 + assert response['ObjectParts']['IsTruncated'] == True + assert response['ObjectParts']['NextPartNumberMarker'] == 4 + assert response['ETag'] == etag + assert response['StorageClass'] == 'STANDARD' + + # check the part + assert len(response['ObjectParts']['Parts']) == 1 + obj_part = response['ObjectParts']['Parts'][0] + assert obj_part['PartNumber'] == 4 + assert obj_part['Size'] == part_size + assert 'ChecksumSHA256' not in obj_part + + request_attributes = ['ETag', 'Checksum', 'ObjectParts', 'StorageClass', 'ObjectSize'] + response = client.get_object_attributes(Bucket=bucket_name, Key=key, + ObjectAttributes=request_attributes, + MaxParts=10, PartNumberMarker=4) + + # check overall object + assert response['ObjectSize'] == objlen + check_parts_count(response['ObjectParts'], len(parts)) + assert response['ObjectParts']['MaxParts'] == 10 + assert response['ObjectParts']['IsTruncated'] == False + assert response['ObjectParts']['PartNumberMarker'] == 4 + assert response['ETag'] == etag + assert response['StorageClass'] == 'STANDARD' + + # check the parts + assert len(response['ObjectParts']['Parts']) == 2 + partno = 5 + for obj_part in response['ObjectParts']['Parts']: + assert obj_part['PartNumber'] == partno + assert obj_part['Size'] == part_size + assert 'ChecksumSHA256' not in obj_part + partno += 1 + +@pytest.mark.fails_on_dbstore +def test_get_single_multipart_object_attributes(): + bucket_name = get_new_bucket() + client = get_client() + + key = "multipart" + part_size = 5*1024*1024 + part_sizes = [part_size] # just one part + part_count = len(part_sizes) + total_size = sum(part_sizes) + + (upload_id, data, parts) = _multipart_upload(bucket_name, key, total_size, part_size) + response = client.complete_multipart_upload(Bucket=bucket_name, Key=key, + UploadId=upload_id, + MultipartUpload={'Parts': parts}) + etag = response['ETag'].strip('"') + assert len(etag) + + request_attributes = ['ETag', 'Checksum', 'ObjectParts', 'StorageClass', 'ObjectSize'] + response = client.get_object_attributes(Bucket=bucket_name, Key=key, + ObjectAttributes=request_attributes) + + assert response['ObjectSize'] == total_size + check_parts_count(response['ObjectParts'], 1) + assert response['ETag'] == etag + assert response['StorageClass'] == 'STANDARD' + + assert len(response['ObjectParts']['Parts']) == 1 + obj_part = response['ObjectParts']['Parts'][0] + assert obj_part['PartNumber'] == 1 + assert obj_part['Size'] == part_size + assert 'ChecksumSHA256' not in obj_part + +def test_get_checksum_object_attributes(): + bucket_name = get_new_bucket() + client = get_client() + + key = "myobj" + size = 1024 + body = FakeWriteFile(size, 'A') + sha256sum = 'arcu6553sHVAiX4MjW0j7I7vD4w6R+Gz9Ok0Q9lTa+0=' + response = client.put_object(Bucket=bucket_name, Key=key, Body=body, ChecksumAlgorithm='SHA256', ChecksumSHA256=sha256sum) + assert sha256sum == response['ChecksumSHA256'] + etag = response['ETag'].strip('"') + assert len(etag) + + request_attributes = ['ETag', 'Checksum', 'ObjectParts', 'StorageClass', 'ObjectSize'] + response = client.get_object_attributes(Bucket=bucket_name, Key=key, + ObjectAttributes=request_attributes) + + assert response['ObjectSize'] == size + assert response['ETag'] == etag + assert response['StorageClass'] == 'STANDARD' + assert response['Checksum']['ChecksumSHA256'] == sha256sum + assert 'ObjectParts' not in response + +def test_get_versioned_object_attributes(): + bucket_name = get_new_bucket() + check_configure_versioning_retry(bucket_name, "Enabled", "Enabled") + client = get_client() + key = "obj" + objlen = 3 + + response = client.put_object(Bucket=bucket_name, Key=key, Body='foo') + etag = response['ETag'].strip('"') + assert len(etag) + version = response['VersionId'] + assert len(version) + + request_attributes = ['ETag', 'Checksum', 'ObjectParts', 'StorageClass', 'ObjectSize'] + response = client.get_object_attributes(Bucket=bucket_name, Key=key, + ObjectAttributes=request_attributes) + + assert 'DeleteMarker' not in response + assert response['VersionId'] == version + + assert response['ObjectSize'] == 3 + assert response['ETag'] == etag + assert response['StorageClass'] == 'STANDARD' + assert 'ObjectParts' not in response + + # write a new current version + client.put_object(Bucket=bucket_name, Key=key, Body='foo') + + # ask for the original version again + request_attributes = ['ETag', 'Checksum', 'ObjectParts', 'StorageClass', 'ObjectSize'] + response = client.get_object_attributes(Bucket=bucket_name, Key=key, VersionId=version, + ObjectAttributes=request_attributes) + + assert 'DeleteMarker' not in response + assert response['VersionId'] == version + + assert response['ObjectSize'] == 3 + assert response['ETag'] == etag + assert response['StorageClass'] == 'STANDARD' + assert 'ObjectParts' not in response + +@pytest.mark.encryption +def test_get_sse_c_encrypted_object_attributes(): + bucket_name = get_new_bucket() + client = get_client() + key = 'obj' + objlen = 1000 + data = 'A'*objlen + sse_args = { + 'SSECustomerAlgorithm': 'AES256', + 'SSECustomerKey': 'pO3upElrwuEXSoFwCfnZPdSsmt/xWeFa0N9KgDijwVs=', + 'SSECustomerKeyMD5': 'DWygnHRtgiJ77HCm+1rvHw==' + } + attrs = ['ETag', 'Checksum', 'ObjectParts', 'StorageClass', 'ObjectSize'] + + response = client.put_object(Bucket=bucket_name, Key=key, Body=data, **sse_args) + etag = response['ETag'].strip('"') + assert len(etag) + + # GetObjectAttributes fails without sse-c headers + e = assert_raises(ClientError, client.get_object_attributes, + Bucket=bucket_name, Key=key, ObjectAttributes=attrs) + status, error_code = _get_status_and_error_code(e.response) + assert status == 400 + + # and succeeds sse-c headers + response = client.get_object_attributes(Bucket=bucket_name, Key=key, + ObjectAttributes=attrs, **sse_args) + + assert 'DeleteMarker' not in response + assert 'VersionId' not in response + + assert response['ObjectSize'] == objlen + assert response['ETag'] == etag + assert response['StorageClass'] == 'STANDARD' + assert 'ObjectParts' not in response + +def test_get_object_attributes(): + bucket_name = get_new_bucket() + client = get_client() + key = "obj" + objlen = 3 + + response = client.put_object(Bucket=bucket_name, Key=key, Body='foo') + etag = response['ETag'].strip('"') + assert len(etag) + + request_attributes = ['ETag', 'Checksum', 'ObjectParts', 'StorageClass', 'ObjectSize'] + response = client.get_object_attributes(Bucket=bucket_name, Key=key, + ObjectAttributes=request_attributes) + + assert 'DeleteMarker' not in response + assert 'VersionId' not in response + + assert response['ObjectSize'] == 3 + assert response['ETag'] == etag + assert response['StorageClass'] == 'STANDARD' + assert 'ObjectParts' not in response From e0195e9dfc0b945a77ebe3fbf400ae4f07f31cd2 Mon Sep 17 00:00:00 2001 From: Matt Benjamin Date: Fri, 29 Nov 2024 12:31:54 -0500 Subject: [PATCH 10/12] mark attribute tests as failing on dbstore (for now) Signed-off-by: Matt Benjamin --- s3tests_boto3/functional/test_s3.py | 2 ++ 1 file changed, 2 insertions(+) diff --git a/s3tests_boto3/functional/test_s3.py b/s3tests_boto3/functional/test_s3.py index 01dc6ea8c..80796b57b 100644 --- a/s3tests_boto3/functional/test_s3.py +++ b/s3tests_boto3/functional/test_s3.py @@ -15366,6 +15366,7 @@ def test_get_versioned_object_attributes(): assert 'ObjectParts' not in response @pytest.mark.encryption +@pytest.mark.fails_on_dbstore def test_get_sse_c_encrypted_object_attributes(): bucket_name = get_new_bucket() client = get_client() @@ -15401,6 +15402,7 @@ def test_get_sse_c_encrypted_object_attributes(): assert response['StorageClass'] == 'STANDARD' assert 'ObjectParts' not in response +@pytest.mark.fails_on_dbstore def test_get_object_attributes(): bucket_name = get_new_bucket() client = get_client() From 14df06dfec0baa268908b51b9a39418c6ad7ea83 Mon Sep 17 00:00:00 2001 From: Matt Benjamin Date: Tue, 25 Feb 2025 21:04:55 -0500 Subject: [PATCH 11/12] add minimal put-object for CRC64NVME Signed-off-by: Matt Benjamin --- s3tests_boto3/functional/test_s3.py | 22 ++++++++++++++++++++++ 1 file changed, 22 insertions(+) diff --git a/s3tests_boto3/functional/test_s3.py b/s3tests_boto3/functional/test_s3.py index 80796b57b..1b98bb609 100644 --- a/s3tests_boto3/functional/test_s3.py +++ b/s3tests_boto3/functional/test_s3.py @@ -13613,6 +13613,28 @@ def test_object_checksum_sha256(): assert status == 400 assert error_code == 'InvalidRequest' +@pytest.mark.checksum +def test_object_checksum_crc64nvme(): + bucket = get_new_bucket() + client = get_client() + + key = "myobj" + size = 1024 + body = FakeWriteFile(size, 'A') + crc64sum = 'Qeh8oXvGiSo=' + response = client.put_object(Bucket=bucket, Key=key, Body=body, ChecksumAlgorithm='CRC64NVME', ChecksumCRC64NVME=crc64sum) + assert crc64sum == response['ChecksumCRC64NVME'] + + response = client.head_object(Bucket=bucket, Key=key) + assert 'ChecksumCRC64NVME' not in response + response = client.head_object(Bucket=bucket, Key=key, ChecksumMode='ENABLED') + assert crc64sum == response['ChecksumCRC64NVME'] + + e = assert_raises(ClientError, client.put_object, Bucket=bucket, Key=key, Body=body, ChecksumAlgorithm='CRC64NVME', ChecksumCRC64NVME='bad') + status, error_code = _get_status_and_error_code(e.response) + assert status == 400 + assert error_code == 'InvalidRequest' + @pytest.mark.checksum @pytest.mark.fails_on_dbstore def test_multipart_checksum_sha256(): From 5524f4d5616e2b3c26e1c5b3d169b7bec09b4043 Mon Sep 17 00:00:00 2001 From: Matt Benjamin Date: Mon, 3 Mar 2025 13:31:21 -0500 Subject: [PATCH 12/12] enhance additional checksum tests includes tests for CRC64NVME, tests for selecting COMPOSITE and FULL_OBJECT checksums a decomposed matrix of tests for all checksum types also removes the mixed checksum upload case that no longer works in recent boto3 cleanups, add sha1 checksum validation failure (mismatch) returns BadDigest multipart checksum matrix helper now validates checksum and checksum type for all operations which can return them (complete-multipart, head-object, and get-object-attributes) Signed-off-by: Matt Benjamin --- s3tests_boto3/functional/test_s3.py | 212 ++++++++++++++++++++-------- 1 file changed, 151 insertions(+), 61 deletions(-) diff --git a/s3tests_boto3/functional/test_s3.py b/s3tests_boto3/functional/test_s3.py index 1b98bb609..fdf24fb0c 100644 --- a/s3tests_boto3/functional/test_s3.py +++ b/s3tests_boto3/functional/test_s3.py @@ -24,6 +24,7 @@ import socket import dateutil.parser import ssl +import pdb from collections import namedtuple from collections import defaultdict from io import StringIO @@ -13611,7 +13612,7 @@ def test_object_checksum_sha256(): e = assert_raises(ClientError, client.put_object, Bucket=bucket, Key=key, Body=body, ChecksumAlgorithm='SHA256', ChecksumSHA256='bad') status, error_code = _get_status_and_error_code(e.response) assert status == 400 - assert error_code == 'InvalidRequest' + assert error_code == 'BadDigest' @pytest.mark.checksum def test_object_checksum_crc64nvme(): @@ -13633,7 +13634,7 @@ def test_object_checksum_crc64nvme(): e = assert_raises(ClientError, client.put_object, Bucket=bucket, Key=key, Body=body, ChecksumAlgorithm='CRC64NVME', ChecksumCRC64NVME='bad') status, error_code = _get_status_and_error_code(e.response) assert status == 400 - assert error_code == 'InvalidRequest' + assert error_code == 'BadDigest' @pytest.mark.checksum @pytest.mark.fails_on_dbstore @@ -13656,7 +13657,7 @@ def test_multipart_checksum_sha256(): {'ETag': response['ETag'].strip('"'), 'ChecksumSHA256': response['ChecksumSHA256'], 'PartNumber': 1}]}) status, error_code = _get_status_and_error_code(e.response) assert status == 400 - assert error_code == 'InvalidRequest' + assert error_code == 'BadDigest' # XXXX re-trying the complete is failing in RGW due to an internal error that appears not caused # checksums; @@ -13679,7 +13680,7 @@ def test_multipart_checksum_sha256(): {'ETag': response['ETag'].strip('"'), 'PartNumber': 1}]}) status, error_code = _get_status_and_error_code(e.response) assert status == 400 - assert error_code == 'InvalidRequest' + assert error_code == 'BadDigest' key = "mymultipart3" response = client.create_multipart_upload(Bucket=bucket, Key=key, ChecksumAlgorithm='SHA256') @@ -13700,85 +13701,174 @@ def test_multipart_checksum_sha256(): response = client.head_object(Bucket=bucket, Key=key, ChecksumMode='ENABLED') assert composite_sha256sum == response['ChecksumSHA256'] -@pytest.mark.checksum -@pytest.mark.fails_on_dbstore -def test_multipart_checksum_3parts(): +def multipart_checksum_3parts_helper(key=None, checksum_algo=None, checksum_type=None, **kwargs): + bucket = get_new_bucket() client = get_client() - key = "mymultipart3" - response = client.create_multipart_upload(Bucket=bucket, Key=key, ChecksumAlgorithm='SHA256') - assert 'SHA256' == response['ChecksumAlgorithm'] + response = client.create_multipart_upload(Bucket=bucket, Key=key, ChecksumAlgorithm=checksum_algo, ChecksumType=checksum_type) + assert checksum_algo == response['ChecksumAlgorithm'] upload_id = response['UploadId'] - size = 5 * 1024 * 1024 # each part but the last must be at least 5M - body = FakeWriteFile(size, 'A') - part1_sha256sum = '275VF5loJr1YYawit0XSHREhkFXYkkPKGuoK0x9VKxI=' - response = client.upload_part(UploadId=upload_id, Bucket=bucket, Key=key, PartNumber=1, Body=body, ChecksumAlgorithm='SHA256', ChecksumSHA256=part1_sha256sum) + cksum_arg_name = "Checksum" + checksum_algo + + upload_args = {cksum_arg_name : kwargs['part1_cksum']} + response = client.upload_part(UploadId=upload_id, Bucket=bucket, Key=key, PartNumber=1, Body=kwargs['body1'], ChecksumAlgorithm=checksum_algo, **upload_args) etag1 = response['ETag'].strip('"') + cksum1 = response[cksum_arg_name] - body = FakeWriteFile(size, 'B') - part2_sha256sum = 'mrHwOfjTL5Zwfj74F05HOQGLdUb7E5szdCbxgUSq6NM=' - response = client.upload_part(UploadId=upload_id, Bucket=bucket, Key=key, PartNumber=2, Body=body, ChecksumAlgorithm='SHA256', ChecksumSHA256=part2_sha256sum) + upload_args = {cksum_arg_name : kwargs['part2_cksum']} + response = client.upload_part(UploadId=upload_id, Bucket=bucket, Key=key, PartNumber=2, Body=kwargs['body2'], ChecksumAlgorithm=checksum_algo, **upload_args) etag2 = response['ETag'].strip('"') + cksum2 = response[cksum_arg_name] - body = FakeWriteFile(size, 'C') - part3_sha256sum = 'Vw7oB/nKQ5xWb3hNgbyfkvDiivl+U+/Dft48nfJfDow=' - response = client.upload_part(UploadId=upload_id, Bucket=bucket, Key=key, PartNumber=3, Body=body, ChecksumAlgorithm='SHA256', ChecksumSHA256=part3_sha256sum) + upload_args = {cksum_arg_name : kwargs['part3_cksum']} + response = client.upload_part(UploadId=upload_id, Bucket=bucket, Key=key, PartNumber=3, Body=kwargs['body3'], ChecksumAlgorithm=checksum_algo, **upload_args) etag3 = response['ETag'].strip('"') + cksum3 = response[cksum_arg_name] - composite_sha256sum = 'uWBwpe1dxI4Vw8Gf0X9ynOdw/SS6VBzfWm9giiv1sf4=-3' - response = client.complete_multipart_upload(Bucket=bucket, Key=key, UploadId=upload_id, ChecksumSHA256=composite_sha256sum, MultipartUpload={'Parts': [ - {'ETag': etag1, 'ChecksumSHA256': response['ChecksumSHA256'], 'PartNumber': 1}, - {'ETag': etag2, 'ChecksumSHA256': response['ChecksumSHA256'], 'PartNumber': 2}, - {'ETag': etag3, 'ChecksumSHA256': response['ChecksumSHA256'], 'PartNumber': 3}]}) - assert composite_sha256sum == response['ChecksumSHA256'] + upload_args = {cksum_arg_name : kwargs['composite_cksum']} + response = client.complete_multipart_upload(Bucket=bucket, Key=key, UploadId=upload_id, MultipartUpload={'Parts': [ + {'ETag': etag1, cksum_arg_name: cksum1, 'PartNumber': 1}, + {'ETag': etag2, cksum_arg_name: cksum2, 'PartNumber': 2}, + {'ETag': etag3, cksum_arg_name: cksum3, 'PartNumber': 3}]}, + **upload_args) - response = client.head_object(Bucket=bucket, Key=key) - assert 'ChecksumSHA256' not in response - response = client.head_object(Bucket=bucket, Key=key, ChecksumMode='ENABLED') - assert composite_sha256sum == response['ChecksumSHA256'] + assert response['ChecksumType'] == checksum_type + assert response[cksum_arg_name] == kwargs['composite_cksum'] + + response1 = client.head_object(Bucket=bucket, Key=key) + assert cksum_arg_name not in response1 + + response2 = client.head_object(Bucket=bucket, Key=key, ChecksumMode='ENABLED') + assert response2['ChecksumType'] == checksum_type + assert response2[cksum_arg_name] == kwargs['composite_cksum'] + + request_attributes = ['Checksum'] + response3 = client.get_object_attributes(Bucket=bucket, Key=key, \ + ObjectAttributes=request_attributes) + assert response3['Checksum']['ChecksumType'] == checksum_type + assert response3['Checksum'][cksum_arg_name] == kwargs['composite_cksum'] @pytest.mark.checksum @pytest.mark.fails_on_dbstore -def test_multipart_checksum_upload_fallback(): - bucket = get_new_bucket() - client = get_client() +def test_multipart_use_cksum_helper_sha256(): + size = 5 * 1024 * 1024 # each part but the last must be at least 5M - key = "mpu_cksum_fallback" - alg = 'SHA256' + # code to compute checksums for these is in unittest_rgw_cksum + body1 = FakeWriteFile(size, 'A') + body2 = FakeWriteFile(size, 'B') + body3 = FakeWriteFile(size, 'C') + + upload_args = { + "body1" : body1, + "part1_cksum" : '275VF5loJr1YYawit0XSHREhkFXYkkPKGuoK0x9VKxI=', + "body2" : body2, + "part2_cksum" : 'mrHwOfjTL5Zwfj74F05HOQGLdUb7E5szdCbxgUSq6NM=', + "body3" : body3, + "part3_cksum" : 'Vw7oB/nKQ5xWb3hNgbyfkvDiivl+U+/Dft48nfJfDow=', + # the composite OR combined checksum, as appropriate + "composite_cksum" : 'uWBwpe1dxI4Vw8Gf0X9ynOdw/SS6VBzfWm9giiv1sf4=-3', + } - response = client.create_multipart_upload( - Bucket=bucket, Key=key, ChecksumAlgorithm=alg) - assert alg == response['ChecksumAlgorithm'] - upload_id = response['UploadId'] + res = multipart_checksum_3parts_helper(key="mymultipart3", checksum_algo="SHA256", checksum_type="COMPOSITE", **upload_args) + #eof - nparts = 3 - parts = [] +@pytest.mark.checksum +@pytest.mark.fails_on_dbstore +def test_multipart_use_cksum_helper_crc64nvme(): size = 5 * 1024 * 1024 # each part but the last must be at least 5M - for ix in range(0,nparts): - body = FakeWriteFile(size, 'A') - part_num = ix + 1 - res = client.upload_part(UploadId=upload_id, Bucket=bucket, - Key=key, PartNumber=part_num, Body=body) - etag = res['ETag'] - part = {'ETag': etag, 'PartNumber': part_num} - parts.append(part) + # code to compute checksums for these is in unittest_rgw_cksum + body1 = FakeWriteFile(size, 'A') + body2 = FakeWriteFile(size, 'B') + body3 = FakeWriteFile(size, 'C') + + upload_args = { + "body1" : body1, + "part1_cksum" : 'L/E4WYn8v98=', + "body2" : body2, + "part2_cksum" : 'xW1l19VobYM=', + "body3" : body3, + "part3_cksum" : 'cK5MnNaWrW4=', + # the composite OR combined checksum, as appropriate + "composite_cksum" : 'i+6LR0y3eFo=', + } - res = client.complete_multipart_upload( - Bucket=bucket, Key=key, UploadId=upload_id, - MultipartUpload={'Parts': parts}) + res = multipart_checksum_3parts_helper(key="mymultipart3", checksum_algo="CRC64NVME", checksum_type="FULL_OBJECT", **upload_args) + #eof - #pdb.set_trace() - assert res['ResponseMetadata']['HTTPStatusCode'] == 200 - - # not yet merged - #request_attributes = ['ETag', 'Checksum', 'ObjectParts', 'StorageClass', - # 'ObjectSize'] - #res = client.get_object_attributes(Bucket=bucket, Key=key, \ - # ObjectAttributes=request_attributes) - #upload_checksum = res['Checksum']['ChecksumSHA256'] +@pytest.mark.checksum +@pytest.mark.fails_on_dbstore +def test_multipart_use_cksum_helper_crc32(): + size = 5 * 1024 * 1024 # each part but the last must be at least 5M + + # code to compute checksums for these is in unittest_rgw_cksum + body1 = FakeWriteFile(size, 'A') + body2 = FakeWriteFile(size, 'B') + body3 = FakeWriteFile(size, 'C') + + upload_args = { + "body1" : body1, + "part1_cksum" : 'JRTCyQ==', + "body2" : body2, + "part2_cksum" : 'QoZTGg==', + "body3" : body3, + "part3_cksum" : 'YAgjqw==', + # the composite OR combined checksum, as appropriate + "composite_cksum" : 'WgDhBQ==', + } + + res = multipart_checksum_3parts_helper(key="mymultipart3", checksum_algo="CRC32", checksum_type="FULL_OBJECT", **upload_args) + #eof + +@pytest.mark.checksum +@pytest.mark.fails_on_dbstore +def test_multipart_use_cksum_helper_crc32c(): + size = 5 * 1024 * 1024 # each part but the last must be at least 5M + + # code to compute checksums for these is in unittest_rgw_cksum + body1 = FakeWriteFile(size, 'A') + body2 = FakeWriteFile(size, 'B') + body3 = FakeWriteFile(size, 'C') + + upload_args = { + "body1" : body1, + "part1_cksum" : 'MDaLrw==', + "body2" : body2, + "part2_cksum" : 'TH4EZg==', + "body3" : body3, + "part3_cksum" : 'Z7mBIQ==', + # the composite OR combined checksum, as appropriate + "composite_cksum" : 'xU+Krw==', + } + + res = multipart_checksum_3parts_helper(key="mymultipart3", checksum_algo="CRC32C", checksum_type="FULL_OBJECT", **upload_args) + #eof + +@pytest.mark.checksum +@pytest.mark.fails_on_dbstore +def test_multipart_use_cksum_helper_sha1(): + size = 5 * 1024 * 1024 # each part but the last must be at least 5M + + # code to compute checksums for these is in unittest_rgw_cksum + body1 = FakeWriteFile(size, 'A') + body2 = FakeWriteFile(size, 'B') + body3 = FakeWriteFile(size, 'C') + + upload_args = { + "body1" : body1, + "part1_cksum" : 'iIaTCGbm+vdVjNqIMF2S0T7ibMk=', + "body2" : body2, + "part2_cksum" : 'LS/TJ32bAVKEwRu+sE3X7awh/lk=', + "body3" : body3, + "part3_cksum" : '6DDwovUaHwrKNXDMzOGbuvj9kxI=', + # the composite OR combined checksum, as appropriate + "composite_cksum" : 'sizjvY4eud3MrcHdZM3cQ/ol39o=-3', + } + + res = multipart_checksum_3parts_helper(key="mymultipart3", checksum_algo="SHA1", checksum_type="COMPOSITE", **upload_args) + #eof @pytest.mark.checksum def test_post_object_upload_checksum():