• Home
  • Features
  • Pricing
  • Docs
  • Announcements
  • Sign In

localstack / localstack / 8d185ffa-7000-47de-8b8d-70e464012974

27 Jan 2025 05:40PM UTC coverage: 86.901% (+0.001%) from 86.9%
8d185ffa-7000-47de-8b8d-70e464012974

push

circleci

web-flow
Fix: validate schedule_expression for EventBridge Scheduler (#12191)

15 of 15 new or added lines in 1 file covered. (100.0%)

14 existing lines in 8 files now uncovered.

61174 of 70395 relevant lines covered (86.9%)

0.87 hits per line

Source File
Press 'n' to go to next uncovered line, 'b' for previous

95.26
/localstack-core/localstack/services/s3/provider.py
1
import base64
1✔
2
import copy
1✔
3
import datetime
1✔
4
import json
1✔
5
import logging
1✔
6
import re
1✔
7
from collections import defaultdict
1✔
8
from io import BytesIO
1✔
9
from operator import itemgetter
1✔
10
from typing import IO, Optional, Union
1✔
11
from urllib import parse as urlparse
1✔
12
from zoneinfo import ZoneInfo
1✔
13

14
from localstack import config
1✔
15
from localstack.aws.api import CommonServiceException, RequestContext, handler
1✔
16
from localstack.aws.api.s3 import (
1✔
17
    MFA,
18
    AbortMultipartUploadOutput,
19
    AccelerateConfiguration,
20
    AccessControlPolicy,
21
    AccessDenied,
22
    AccountId,
23
    AnalyticsConfiguration,
24
    AnalyticsId,
25
    BadDigest,
26
    Body,
27
    Bucket,
28
    BucketAlreadyExists,
29
    BucketAlreadyOwnedByYou,
30
    BucketCannedACL,
31
    BucketLifecycleConfiguration,
32
    BucketLoggingStatus,
33
    BucketName,
34
    BucketNotEmpty,
35
    BucketRegion,
36
    BucketVersioningStatus,
37
    BypassGovernanceRetention,
38
    ChecksumAlgorithm,
39
    ChecksumCRC32,
40
    ChecksumCRC32C,
41
    ChecksumCRC64NVME,
42
    ChecksumSHA1,
43
    ChecksumSHA256,
44
    ChecksumType,
45
    CommonPrefix,
46
    CompletedMultipartUpload,
47
    CompleteMultipartUploadOutput,
48
    ConditionalRequestConflict,
49
    ConfirmRemoveSelfBucketAccess,
50
    ContentMD5,
51
    CopyObjectOutput,
52
    CopyObjectRequest,
53
    CopyObjectResult,
54
    CopyPartResult,
55
    CORSConfiguration,
56
    CreateBucketOutput,
57
    CreateBucketRequest,
58
    CreateMultipartUploadOutput,
59
    CreateMultipartUploadRequest,
60
    CrossLocationLoggingProhibitted,
61
    Delete,
62
    DeletedObject,
63
    DeleteMarkerEntry,
64
    DeleteObjectOutput,
65
    DeleteObjectsOutput,
66
    DeleteObjectTaggingOutput,
67
    Delimiter,
68
    EncodingType,
69
    Error,
70
    Expiration,
71
    FetchOwner,
72
    GetBucketAccelerateConfigurationOutput,
73
    GetBucketAclOutput,
74
    GetBucketAnalyticsConfigurationOutput,
75
    GetBucketCorsOutput,
76
    GetBucketEncryptionOutput,
77
    GetBucketIntelligentTieringConfigurationOutput,
78
    GetBucketInventoryConfigurationOutput,
79
    GetBucketLifecycleConfigurationOutput,
80
    GetBucketLocationOutput,
81
    GetBucketLoggingOutput,
82
    GetBucketOwnershipControlsOutput,
83
    GetBucketPolicyOutput,
84
    GetBucketPolicyStatusOutput,
85
    GetBucketReplicationOutput,
86
    GetBucketRequestPaymentOutput,
87
    GetBucketTaggingOutput,
88
    GetBucketVersioningOutput,
89
    GetBucketWebsiteOutput,
90
    GetObjectAclOutput,
91
    GetObjectAttributesOutput,
92
    GetObjectAttributesParts,
93
    GetObjectAttributesRequest,
94
    GetObjectLegalHoldOutput,
95
    GetObjectLockConfigurationOutput,
96
    GetObjectOutput,
97
    GetObjectRequest,
98
    GetObjectRetentionOutput,
99
    GetObjectTaggingOutput,
100
    GetObjectTorrentOutput,
101
    GetPublicAccessBlockOutput,
102
    HeadBucketOutput,
103
    HeadObjectOutput,
104
    HeadObjectRequest,
105
    IfMatch,
106
    IfMatchInitiatedTime,
107
    IfMatchLastModifiedTime,
108
    IfMatchSize,
109
    IfNoneMatch,
110
    IntelligentTieringConfiguration,
111
    IntelligentTieringId,
112
    InvalidArgument,
113
    InvalidBucketName,
114
    InvalidDigest,
115
    InvalidLocationConstraint,
116
    InvalidObjectState,
117
    InvalidPartNumber,
118
    InvalidPartOrder,
119
    InvalidStorageClass,
120
    InvalidTargetBucketForLogging,
121
    InventoryConfiguration,
122
    InventoryId,
123
    KeyMarker,
124
    LifecycleRules,
125
    ListBucketAnalyticsConfigurationsOutput,
126
    ListBucketIntelligentTieringConfigurationsOutput,
127
    ListBucketInventoryConfigurationsOutput,
128
    ListBucketsOutput,
129
    ListMultipartUploadsOutput,
130
    ListObjectsOutput,
131
    ListObjectsV2Output,
132
    ListObjectVersionsOutput,
133
    ListPartsOutput,
134
    Marker,
135
    MaxBuckets,
136
    MaxKeys,
137
    MaxParts,
138
    MaxUploads,
139
    MethodNotAllowed,
140
    MissingSecurityHeader,
141
    MpuObjectSize,
142
    MultipartUpload,
143
    MultipartUploadId,
144
    NoSuchBucket,
145
    NoSuchBucketPolicy,
146
    NoSuchCORSConfiguration,
147
    NoSuchKey,
148
    NoSuchLifecycleConfiguration,
149
    NoSuchPublicAccessBlockConfiguration,
150
    NoSuchTagSet,
151
    NoSuchUpload,
152
    NoSuchWebsiteConfiguration,
153
    NotificationConfiguration,
154
    Object,
155
    ObjectIdentifier,
156
    ObjectKey,
157
    ObjectLockConfiguration,
158
    ObjectLockConfigurationNotFoundError,
159
    ObjectLockEnabled,
160
    ObjectLockLegalHold,
161
    ObjectLockMode,
162
    ObjectLockRetention,
163
    ObjectLockToken,
164
    ObjectOwnership,
165
    ObjectVersion,
166
    ObjectVersionId,
167
    ObjectVersionStorageClass,
168
    OptionalObjectAttributesList,
169
    Owner,
170
    OwnershipControls,
171
    OwnershipControlsNotFoundError,
172
    Part,
173
    PartNumber,
174
    PartNumberMarker,
175
    Policy,
176
    PostResponse,
177
    PreconditionFailed,
178
    Prefix,
179
    PublicAccessBlockConfiguration,
180
    PutBucketAclRequest,
181
    PutBucketLifecycleConfigurationOutput,
182
    PutObjectAclOutput,
183
    PutObjectAclRequest,
184
    PutObjectLegalHoldOutput,
185
    PutObjectLockConfigurationOutput,
186
    PutObjectOutput,
187
    PutObjectRequest,
188
    PutObjectRetentionOutput,
189
    PutObjectTaggingOutput,
190
    ReplicationConfiguration,
191
    ReplicationConfigurationNotFoundError,
192
    RequestPayer,
193
    RequestPaymentConfiguration,
194
    RestoreObjectOutput,
195
    RestoreRequest,
196
    S3Api,
197
    ServerSideEncryption,
198
    ServerSideEncryptionConfiguration,
199
    SkipValidation,
200
    SSECustomerAlgorithm,
201
    SSECustomerKey,
202
    SSECustomerKeyMD5,
203
    StartAfter,
204
    StorageClass,
205
    Tagging,
206
    Token,
207
    TransitionDefaultMinimumObjectSize,
208
    UploadIdMarker,
209
    UploadPartCopyOutput,
210
    UploadPartCopyRequest,
211
    UploadPartOutput,
212
    UploadPartRequest,
213
    VersionIdMarker,
214
    VersioningConfiguration,
215
    WebsiteConfiguration,
216
)
217
from localstack.aws.api.s3 import NotImplemented as NotImplementedException
1✔
218
from localstack.aws.handlers import (
1✔
219
    modify_service_response,
220
    preprocess_request,
221
    serve_custom_service_request_handlers,
222
)
223
from localstack.constants import AWS_REGION_US_EAST_1
1✔
224
from localstack.services.edge import ROUTER
1✔
225
from localstack.services.plugins import ServiceLifecycleHook
1✔
226
from localstack.services.s3.codec import AwsChunkedDecoder
1✔
227
from localstack.services.s3.constants import (
1✔
228
    ALLOWED_HEADER_OVERRIDES,
229
    ARCHIVES_STORAGE_CLASSES,
230
    CHECKSUM_ALGORITHMS,
231
    DEFAULT_BUCKET_ENCRYPTION,
232
)
233
from localstack.services.s3.cors import S3CorsHandler, s3_cors_request_handler
1✔
234
from localstack.services.s3.exceptions import (
1✔
235
    InvalidBucketOwnerAWSAccountID,
236
    InvalidBucketState,
237
    InvalidRequest,
238
    MalformedPolicy,
239
    MalformedXML,
240
    NoSuchConfiguration,
241
    NoSuchObjectLockConfiguration,
242
    UnexpectedContent,
243
)
244
from localstack.services.s3.models import (
1✔
245
    BucketCorsIndex,
246
    EncryptionParameters,
247
    ObjectLockParameters,
248
    S3Bucket,
249
    S3DeleteMarker,
250
    S3Multipart,
251
    S3Object,
252
    S3Part,
253
    S3Store,
254
    VersionedKeyStore,
255
    s3_stores,
256
)
257
from localstack.services.s3.notifications import NotificationDispatcher, S3EventNotificationContext
1✔
258
from localstack.services.s3.presigned_url import validate_post_policy
1✔
259
from localstack.services.s3.storage.core import LimitedIterableStream, S3ObjectStore
1✔
260
from localstack.services.s3.storage.ephemeral import EphemeralS3ObjectStore
1✔
261
from localstack.services.s3.utils import (
1✔
262
    ObjectRange,
263
    add_expiration_days_to_datetime,
264
    base_64_content_md5_to_etag,
265
    create_redirect_for_post_request,
266
    create_s3_kms_managed_key_for_region,
267
    etag_to_base_64_content_md5,
268
    extract_bucket_key_version_id_from_copy_source,
269
    generate_safe_version_id,
270
    get_canned_acl,
271
    get_class_attrs_from_spec_class,
272
    get_failed_precondition_copy_source,
273
    get_full_default_bucket_location,
274
    get_kms_key_arn,
275
    get_lifecycle_rule_from_object,
276
    get_owner_for_account_id,
277
    get_permission_from_header,
278
    get_retention_from_now,
279
    get_s3_checksum_algorithm_from_request,
280
    get_s3_checksum_algorithm_from_trailing_headers,
281
    get_system_metadata_from_request,
282
    get_unique_key_id,
283
    is_bucket_name_valid,
284
    parse_copy_source_range_header,
285
    parse_post_object_tagging_xml,
286
    parse_range_header,
287
    parse_tagging_header,
288
    s3_response_handler,
289
    serialize_expiration_header,
290
    str_to_rfc_1123_datetime,
291
    validate_dict_fields,
292
    validate_failed_precondition,
293
    validate_kms_key_id,
294
    validate_tag_set,
295
)
296
from localstack.services.s3.validation import (
1✔
297
    parse_grants_in_headers,
298
    validate_acl_acp,
299
    validate_bucket_analytics_configuration,
300
    validate_bucket_intelligent_tiering_configuration,
301
    validate_canned_acl,
302
    validate_checksum_value,
303
    validate_cors_configuration,
304
    validate_inventory_configuration,
305
    validate_lifecycle_configuration,
306
    validate_object_key,
307
    validate_sse_c,
308
    validate_website_configuration,
309
)
310
from localstack.services.s3.website_hosting import register_website_hosting_routes
1✔
311
from localstack.state import AssetDirectory, StateVisitor
1✔
312
from localstack.utils.aws.arns import s3_bucket_name
1✔
313
from localstack.utils.strings import short_uid, to_bytes, to_str
1✔
314

315
LOG = logging.getLogger(__name__)
1✔
316

317
STORAGE_CLASSES = get_class_attrs_from_spec_class(StorageClass)
1✔
318
SSE_ALGORITHMS = get_class_attrs_from_spec_class(ServerSideEncryption)
1✔
319
OBJECT_OWNERSHIPS = get_class_attrs_from_spec_class(ObjectOwnership)
1✔
320

321
DEFAULT_S3_TMP_DIR = "/tmp/localstack-s3-storage"
1✔
322

323

324
class S3Provider(S3Api, ServiceLifecycleHook):
1✔
325
    def __init__(self, storage_backend: S3ObjectStore = None) -> None:
1✔
326
        super().__init__()
1✔
327
        self._storage_backend = storage_backend or EphemeralS3ObjectStore(DEFAULT_S3_TMP_DIR)
1✔
328
        self._notification_dispatcher = NotificationDispatcher()
1✔
329
        self._cors_handler = S3CorsHandler(BucketCorsIndex())
1✔
330

331
        # runtime cache of Lifecycle Expiration headers, as they need to be calculated everytime we fetch an object
332
        # in case the rules have changed
333
        self._expiration_cache: dict[BucketName, dict[ObjectKey, Expiration]] = defaultdict(dict)
1✔
334

335
    def on_after_init(self):
1✔
336
        preprocess_request.append(self._cors_handler)
1✔
337
        serve_custom_service_request_handlers.append(s3_cors_request_handler)
1✔
338
        modify_service_response.append(self.service, s3_response_handler)
1✔
339
        register_website_hosting_routes(router=ROUTER)
1✔
340

341
    def accept_state_visitor(self, visitor: StateVisitor):
1✔
342
        visitor.visit(s3_stores)
×
343
        visitor.visit(AssetDirectory(self.service, self._storage_backend.root_directory))
×
344

345
    def on_before_state_save(self):
1✔
346
        self._storage_backend.flush()
×
347

348
    def on_after_state_reset(self):
1✔
349
        self._cors_handler.invalidate_cache()
×
350

351
    def on_after_state_load(self):
1✔
352
        self._cors_handler.invalidate_cache()
×
353

354
    def on_before_stop(self):
1✔
355
        self._notification_dispatcher.shutdown()
1✔
356
        self._storage_backend.close()
1✔
357

358
    def _notify(
1✔
359
        self,
360
        context: RequestContext,
361
        s3_bucket: S3Bucket,
362
        s3_object: S3Object | S3DeleteMarker = None,
363
        s3_notif_ctx: S3EventNotificationContext = None,
364
    ):
365
        """
366
        :param context: the RequestContext, to retrieve more information about the incoming notification
367
        :param s3_bucket: the S3Bucket object
368
        :param s3_object: the S3Object object if S3EventNotificationContext is not given
369
        :param s3_notif_ctx: S3EventNotificationContext, in case we need specific data only available in the API call
370
        :return:
371
        """
372
        if s3_bucket.notification_configuration:
1✔
373
            if not s3_notif_ctx:
1✔
374
                s3_notif_ctx = S3EventNotificationContext.from_request_context_native(
1✔
375
                    context,
376
                    s3_bucket=s3_bucket,
377
                    s3_object=s3_object,
378
                )
379

380
            self._notification_dispatcher.send_notifications(
1✔
381
                s3_notif_ctx, s3_bucket.notification_configuration
382
            )
383

384
    def _verify_notification_configuration(
1✔
385
        self,
386
        notification_configuration: NotificationConfiguration,
387
        skip_destination_validation: SkipValidation,
388
        context: RequestContext,
389
        bucket_name: str,
390
    ):
391
        self._notification_dispatcher.verify_configuration(
1✔
392
            notification_configuration, skip_destination_validation, context, bucket_name
393
        )
394

395
    def _get_expiration_header(
1✔
396
        self,
397
        lifecycle_rules: LifecycleRules,
398
        bucket: BucketName,
399
        s3_object: S3Object,
400
        object_tags: dict[str, str],
401
    ) -> Expiration:
402
        """
403
        This method will check if the key matches a Lifecycle filter, and return the serializer header if that's
404
        the case. We're caching it because it can change depending on the set rules on the bucket.
405
        We can't use `lru_cache` as the parameters needs to be hashable
406
        :param lifecycle_rules: the bucket LifecycleRules
407
        :param s3_object: S3Object
408
        :param object_tags: the object tags
409
        :return: the Expiration header if there's a rule matching
410
        """
411
        if cached_exp := self._expiration_cache.get(bucket, {}).get(s3_object.key):
1✔
412
            return cached_exp
1✔
413

414
        if lifecycle_rule := get_lifecycle_rule_from_object(
1✔
415
            lifecycle_rules, s3_object.key, s3_object.size, object_tags
416
        ):
417
            expiration_header = serialize_expiration_header(
1✔
418
                lifecycle_rule["ID"],
419
                lifecycle_rule["Expiration"],
420
                s3_object.last_modified,
421
            )
422
            self._expiration_cache[bucket][s3_object.key] = expiration_header
1✔
423
            return expiration_header
1✔
424

425
    def _get_cross_account_bucket(
1✔
426
        self,
427
        context: RequestContext,
428
        bucket_name: BucketName,
429
        *,
430
        expected_bucket_owner: AccountId = None,
431
    ) -> tuple[S3Store, S3Bucket]:
432
        if expected_bucket_owner and not re.fullmatch(r"\w{12}", expected_bucket_owner):
1✔
433
            raise InvalidBucketOwnerAWSAccountID(
1✔
434
                f"The value of the expected bucket owner parameter must be an AWS Account ID... [{expected_bucket_owner}]",
435
            )
436

437
        store = self.get_store(context.account_id, context.region)
1✔
438
        if not (s3_bucket := store.buckets.get(bucket_name)):
1✔
439
            if not (account_id := store.global_bucket_map.get(bucket_name)):
1✔
440
                raise NoSuchBucket("The specified bucket does not exist", BucketName=bucket_name)
1✔
441

442
            store = self.get_store(account_id, context.region)
1✔
443
            if not (s3_bucket := store.buckets.get(bucket_name)):
1✔
444
                raise NoSuchBucket("The specified bucket does not exist", BucketName=bucket_name)
×
445

446
        if expected_bucket_owner and s3_bucket.bucket_account_id != expected_bucket_owner:
1✔
447
            raise AccessDenied("Access Denied")
1✔
448

449
        return store, s3_bucket
1✔
450

451
    @staticmethod
1✔
452
    def get_store(account_id: str, region_name: str) -> S3Store:
1✔
453
        # Use default account id for external access? would need an anonymous one
454
        return s3_stores[account_id][region_name]
1✔
455

456
    @handler("CreateBucket", expand=False)
1✔
457
    def create_bucket(
1✔
458
        self,
459
        context: RequestContext,
460
        request: CreateBucketRequest,
461
    ) -> CreateBucketOutput:
462
        bucket_name = request["Bucket"]
1✔
463

464
        if not is_bucket_name_valid(bucket_name):
1✔
465
            raise InvalidBucketName("The specified bucket is not valid.", BucketName=bucket_name)
1✔
466

467
        # the XML parser returns an empty dict if the body contains the following:
468
        # <CreateBucketConfiguration xmlns="http://s3.amazonaws.com/doc/2006-03-01/" />
469
        # but it also returns an empty dict if the body is fully empty. We need to differentiate the 2 cases by checking
470
        # if the body is empty or not
471
        if context.request.data and (
1✔
472
            (create_bucket_configuration := request.get("CreateBucketConfiguration")) is not None
473
        ):
474
            if not (bucket_region := create_bucket_configuration.get("LocationConstraint")):
1✔
475
                raise MalformedXML()
1✔
476

477
            if context.region == AWS_REGION_US_EAST_1:
1✔
478
                if bucket_region == "us-east-1":
1✔
479
                    raise InvalidLocationConstraint(
1✔
480
                        "The specified location-constraint is not valid",
481
                        LocationConstraint=bucket_region,
482
                    )
483
            elif context.region != bucket_region:
1✔
484
                raise CommonServiceException(
1✔
485
                    code="IllegalLocationConstraintException",
486
                    message=f"The {bucket_region} location constraint is incompatible for the region specific endpoint this request was sent to.",
487
                )
488
        else:
489
            bucket_region = AWS_REGION_US_EAST_1
1✔
490
            if context.region != bucket_region:
1✔
491
                raise CommonServiceException(
1✔
492
                    code="IllegalLocationConstraintException",
493
                    message="The unspecified location constraint is incompatible for the region specific endpoint this request was sent to.",
494
                )
495

496
        store = self.get_store(context.account_id, bucket_region)
1✔
497

498
        if bucket_name in store.global_bucket_map:
1✔
499
            existing_bucket_owner = store.global_bucket_map[bucket_name]
1✔
500
            if existing_bucket_owner != context.account_id:
1✔
501
                raise BucketAlreadyExists()
1✔
502

503
            # if the existing bucket has the same owner, the behaviour will depend on the region
504
            if bucket_region != "us-east-1":
1✔
505
                raise BucketAlreadyOwnedByYou(
1✔
506
                    "Your previous request to create the named bucket succeeded and you already own it.",
507
                    BucketName=bucket_name,
508
                )
509
            else:
510
                # CreateBucket is idempotent in us-east-1
511
                return CreateBucketOutput(Location=f"/{bucket_name}")
1✔
512

513
        if (
1✔
514
            object_ownership := request.get("ObjectOwnership")
515
        ) is not None and object_ownership not in OBJECT_OWNERSHIPS:
516
            raise InvalidArgument(
1✔
517
                f"Invalid x-amz-object-ownership header: {object_ownership}",
518
                ArgumentName="x-amz-object-ownership",
519
            )
520
        # see https://docs.aws.amazon.com/AmazonS3/latest/API/API_Owner.html
521
        owner = get_owner_for_account_id(context.account_id)
1✔
522
        acl = get_access_control_policy_for_new_resource_request(request, owner=owner)
1✔
523
        s3_bucket = S3Bucket(
1✔
524
            name=bucket_name,
525
            account_id=context.account_id,
526
            bucket_region=bucket_region,
527
            owner=owner,
528
            acl=acl,
529
            object_ownership=request.get("ObjectOwnership"),
530
            object_lock_enabled_for_bucket=request.get("ObjectLockEnabledForBucket"),
531
        )
532

533
        store.buckets[bucket_name] = s3_bucket
1✔
534
        store.global_bucket_map[bucket_name] = s3_bucket.bucket_account_id
1✔
535
        self._cors_handler.invalidate_cache()
1✔
536
        self._storage_backend.create_bucket(bucket_name)
1✔
537

538
        # Location is always contained in response -> full url for LocationConstraint outside us-east-1
539
        location = (
1✔
540
            f"/{bucket_name}"
541
            if bucket_region == "us-east-1"
542
            else get_full_default_bucket_location(bucket_name)
543
        )
544
        response = CreateBucketOutput(Location=location)
1✔
545
        return response
1✔
546

547
    def delete_bucket(
1✔
548
        self,
549
        context: RequestContext,
550
        bucket: BucketName,
551
        expected_bucket_owner: AccountId = None,
552
        **kwargs,
553
    ) -> None:
554
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
555

556
        # the bucket still contains objects
557
        if not s3_bucket.objects.is_empty():
1✔
558
            message = "The bucket you tried to delete is not empty"
1✔
559
            if s3_bucket.versioning_status:
1✔
560
                message += ". You must delete all versions in the bucket."
1✔
561
            raise BucketNotEmpty(
1✔
562
                message,
563
                BucketName=bucket,
564
            )
565

566
        store.buckets.pop(bucket)
1✔
567
        store.global_bucket_map.pop(bucket)
1✔
568
        self._cors_handler.invalidate_cache()
1✔
569
        self._expiration_cache.pop(bucket, None)
1✔
570
        # clean up the storage backend
571
        self._storage_backend.delete_bucket(bucket)
1✔
572

573
    def list_buckets(
1✔
574
        self,
575
        context: RequestContext,
576
        max_buckets: MaxBuckets = None,
577
        continuation_token: Token = None,
578
        prefix: Prefix = None,
579
        bucket_region: BucketRegion = None,
580
        **kwargs,
581
    ) -> ListBucketsOutput:
582
        # TODO add support for max_buckets, continuation_token, prefix, and bucket_region
583
        owner = get_owner_for_account_id(context.account_id)
1✔
584
        store = self.get_store(context.account_id, context.region)
1✔
585
        buckets = [
1✔
586
            Bucket(Name=bucket.name, CreationDate=bucket.creation_date)
587
            for bucket in store.buckets.values()
588
        ]
589
        return ListBucketsOutput(Owner=owner, Buckets=buckets)
1✔
590

591
    def head_bucket(
1✔
592
        self,
593
        context: RequestContext,
594
        bucket: BucketName,
595
        expected_bucket_owner: AccountId = None,
596
        **kwargs,
597
    ) -> HeadBucketOutput:
598
        store = self.get_store(context.account_id, context.region)
1✔
599
        if not (s3_bucket := store.buckets.get(bucket)):
1✔
600
            if not (account_id := store.global_bucket_map.get(bucket)):
1✔
601
                # just to return the 404 error message
602
                raise NoSuchBucket()
1✔
603

UNCOV
604
            store = self.get_store(account_id, context.region)
×
UNCOV
605
            if not (s3_bucket := store.buckets.get(bucket)):
×
606
                # just to return the 404 error message
607
                raise NoSuchBucket()
×
608

609
        # TODO: this call is also used to check if the user has access/authorization for the bucket
610
        #  it can return 403
611
        return HeadBucketOutput(BucketRegion=s3_bucket.bucket_region)
1✔
612

613
    def get_bucket_location(
1✔
614
        self,
615
        context: RequestContext,
616
        bucket: BucketName,
617
        expected_bucket_owner: AccountId = None,
618
        **kwargs,
619
    ) -> GetBucketLocationOutput:
620
        """
621
        When implementing the ASF provider, this operation is implemented because:
622
        - The spec defines a root element GetBucketLocationOutput containing a LocationConstraint member, where
623
          S3 actually just returns the LocationConstraint on the root level (only operation so far that we know of).
624
        - We circumvent the root level element here by patching the spec such that this operation returns a
625
          single "payload" (the XML body response), which causes the serializer to directly take the payload element.
626
        - The above "hack" causes the fix in the serializer to not be picked up here as we're passing the XML body as
627
          the payload, which is why we need to manually do this here by manipulating the string.
628
        Botocore implements this hack for parsing the response in `botocore.handlers.py#parse_get_bucket_location`
629
        """
630
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
631

632
        location_constraint = (
1✔
633
            '<?xml version="1.0" encoding="UTF-8"?>\n'
634
            '<LocationConstraint xmlns="http://s3.amazonaws.com/doc/2006-03-01/">{{location}}</LocationConstraint>'
635
        )
636

637
        location = s3_bucket.bucket_region if s3_bucket.bucket_region != "us-east-1" else ""
1✔
638
        location_constraint = location_constraint.replace("{{location}}", location)
1✔
639

640
        response = GetBucketLocationOutput(LocationConstraint=location_constraint)
1✔
641
        return response
1✔
642

643
    @handler("PutObject", expand=False)
1✔
644
    def put_object(
1✔
645
        self,
646
        context: RequestContext,
647
        request: PutObjectRequest,
648
    ) -> PutObjectOutput:
649
        # TODO: validate order of validation
650
        # TODO: still need to handle following parameters
651
        #  request_payer: RequestPayer = None,
652
        bucket_name = request["Bucket"]
1✔
653
        key = request["Key"]
1✔
654
        store, s3_bucket = self._get_cross_account_bucket(context, bucket_name)
1✔
655

656
        if (storage_class := request.get("StorageClass")) is not None and (
1✔
657
            storage_class not in STORAGE_CLASSES or storage_class == StorageClass.OUTPOSTS
658
        ):
659
            raise InvalidStorageClass(
1✔
660
                "The storage class you specified is not valid", StorageClassRequested=storage_class
661
            )
662

663
        if not config.S3_SKIP_KMS_KEY_VALIDATION and (sse_kms_key_id := request.get("SSEKMSKeyId")):
1✔
664
            validate_kms_key_id(sse_kms_key_id, s3_bucket)
1✔
665

666
        validate_object_key(key)
1✔
667

668
        if_match = request.get("IfMatch")
1✔
669
        if (if_none_match := request.get("IfNoneMatch")) and if_match:
1✔
670
            raise NotImplementedException(
671
                "A header you provided implies functionality that is not implemented",
672
                Header="If-Match,If-None-Match",
673
                additionalMessage="Multiple conditional request headers present in the request",
674
            )
675

676
        elif (if_none_match and if_none_match != "*") or (if_match and if_match == "*"):
1✔
677
            header_name = "If-None-Match" if if_none_match else "If-Match"
1✔
678
            raise NotImplementedException(
679
                "A header you provided implies functionality that is not implemented",
680
                Header=header_name,
681
                additionalMessage=f"We don't accept the provided value of {header_name} header for this API",
682
            )
683

684
        system_metadata = get_system_metadata_from_request(request)
1✔
685
        if not system_metadata.get("ContentType"):
1✔
686
            system_metadata["ContentType"] = "binary/octet-stream"
1✔
687

688
        version_id = generate_version_id(s3_bucket.versioning_status)
1✔
689

690
        etag_content_md5 = ""
1✔
691
        if content_md5 := request.get("ContentMD5"):
1✔
692
            # assert that the received ContentMD5 is a properly b64 encoded value that fits a MD5 hash length
693
            etag_content_md5 = base_64_content_md5_to_etag(content_md5)
1✔
694
            if not etag_content_md5:
1✔
695
                raise InvalidDigest(
1✔
696
                    "The Content-MD5 you specified was invalid.",
697
                    Content_MD5=content_md5,
698
                )
699

700
        checksum_algorithm = get_s3_checksum_algorithm_from_request(request)
1✔
701
        checksum_value = (
1✔
702
            request.get(f"Checksum{checksum_algorithm.upper()}") if checksum_algorithm else None
703
        )
704

705
        # TODO: we're not encrypting the object with the provided key for now
706
        sse_c_key_md5 = request.get("SSECustomerKeyMD5")
1✔
707
        validate_sse_c(
1✔
708
            algorithm=request.get("SSECustomerAlgorithm"),
709
            encryption_key=request.get("SSECustomerKey"),
710
            encryption_key_md5=sse_c_key_md5,
711
            server_side_encryption=request.get("ServerSideEncryption"),
712
        )
713

714
        encryption_parameters = get_encryption_parameters_from_request_and_bucket(
1✔
715
            request,
716
            s3_bucket,
717
            store,
718
        )
719

720
        lock_parameters = get_object_lock_parameters_from_bucket_and_request(request, s3_bucket)
1✔
721

722
        acl = get_access_control_policy_for_new_resource_request(request, owner=s3_bucket.owner)
1✔
723

724
        if tagging := request.get("Tagging"):
1✔
725
            tagging = parse_tagging_header(tagging)
1✔
726

727
        s3_object = S3Object(
1✔
728
            key=key,
729
            version_id=version_id,
730
            storage_class=storage_class,
731
            expires=request.get("Expires"),
732
            user_metadata=request.get("Metadata"),
733
            system_metadata=system_metadata,
734
            checksum_algorithm=checksum_algorithm,
735
            checksum_value=checksum_value,
736
            encryption=encryption_parameters.encryption,
737
            kms_key_id=encryption_parameters.kms_key_id,
738
            bucket_key_enabled=encryption_parameters.bucket_key_enabled,
739
            sse_key_hash=sse_c_key_md5,
740
            lock_mode=lock_parameters.lock_mode,
741
            lock_legal_status=lock_parameters.lock_legal_status,
742
            lock_until=lock_parameters.lock_until,
743
            website_redirect_location=request.get("WebsiteRedirectLocation"),
744
            acl=acl,
745
            owner=s3_bucket.owner,  # TODO: for now we only have one owner, but it can depends on Bucket settings
746
        )
747

748
        body = request.get("Body")
1✔
749
        # check if chunked request
750
        headers = context.request.headers
1✔
751
        is_aws_chunked = headers.get("x-amz-content-sha256", "").startswith(
1✔
752
            "STREAMING-"
753
        ) or "aws-chunked" in headers.get("content-encoding", "")
754
        if is_aws_chunked:
1✔
755
            checksum_algorithm = (
1✔
756
                checksum_algorithm
757
                or get_s3_checksum_algorithm_from_trailing_headers(headers.get("x-amz-trailer", ""))
758
            )
759
            if checksum_algorithm:
1✔
760
                s3_object.checksum_algorithm = checksum_algorithm
1✔
761

762
            decoded_content_length = int(headers.get("x-amz-decoded-content-length", 0))
1✔
763
            body = AwsChunkedDecoder(body, decoded_content_length, s3_object=s3_object)
1✔
764

765
            # S3 removes the `aws-chunked` value from ContentEncoding
766
            if content_encoding := s3_object.system_metadata.pop("ContentEncoding", None):
1✔
767
                encodings = [enc for enc in content_encoding.split(",") if enc != "aws-chunked"]
1✔
768
                if encodings:
1✔
769
                    s3_object.system_metadata["ContentEncoding"] = ",".join(encodings)
1✔
770

771
        with self._storage_backend.open(bucket_name, s3_object, mode="w") as s3_stored_object:
1✔
772
            # as we are inside the lock here, if multiple concurrent requests happen for the same object, it's the first
773
            # one to finish to succeed, and subsequent will raise exceptions. Once the first write finishes, we're
774
            # opening the lock and other requests can check this condition
775
            if if_none_match and object_exists_for_precondition_write(s3_bucket, key):
1✔
776
                raise PreconditionFailed(
1✔
777
                    "At least one of the pre-conditions you specified did not hold",
778
                    Condition="If-None-Match",
779
                )
780

781
            elif if_match:
1✔
782
                verify_object_equality_precondition_write(s3_bucket, key, if_match)
1✔
783

784
            s3_stored_object.write(body)
1✔
785

786
            if s3_object.checksum_algorithm:
1✔
787
                if not validate_checksum_value(s3_object.checksum_value, checksum_algorithm):
1✔
788
                    self._storage_backend.remove(bucket_name, s3_object)
1✔
789
                    raise InvalidRequest(
1✔
790
                        f"Value for x-amz-checksum-{s3_object.checksum_algorithm.lower()} header is invalid."
791
                    )
792
                elif s3_object.checksum_value != s3_stored_object.checksum:
1✔
793
                    self._storage_backend.remove(bucket_name, s3_object)
1✔
794
                    raise BadDigest(
1✔
795
                        f"The {checksum_algorithm.upper()} you specified did not match the calculated checksum."
796
                    )
797

798
            # TODO: handle ContentMD5 and ChecksumAlgorithm in a handler for all requests except requests with a
799
            #  streaming body. We can use the specs to verify which operations needs to have the checksum validated
800
            if content_md5:
1✔
801
                calculated_md5 = etag_to_base_64_content_md5(s3_stored_object.etag)
1✔
802
                if calculated_md5 != content_md5:
1✔
803
                    self._storage_backend.remove(bucket_name, s3_object)
1✔
804
                    raise BadDigest(
1✔
805
                        "The Content-MD5 you specified did not match what we received.",
806
                        ExpectedDigest=etag_content_md5,
807
                        CalculatedDigest=calculated_md5,
808
                    )
809

810
            s3_bucket.objects.set(key, s3_object)
1✔
811

812
        # in case we are overriding an object, delete the tags entry
813
        key_id = get_unique_key_id(bucket_name, key, version_id)
1✔
814
        store.TAGS.tags.pop(key_id, None)
1✔
815
        if tagging:
1✔
816
            store.TAGS.tags[key_id] = tagging
1✔
817

818
        # RequestCharged: Optional[RequestCharged]  # TODO
819
        response = PutObjectOutput(
1✔
820
            ETag=s3_object.quoted_etag,
821
        )
822
        if s3_bucket.versioning_status == "Enabled":
1✔
823
            response["VersionId"] = s3_object.version_id
1✔
824

825
        if s3_object.checksum_algorithm:
1✔
826
            response[f"Checksum{s3_object.checksum_algorithm}"] = s3_object.checksum_value
1✔
827
            response["ChecksumType"] = getattr(s3_object, "checksum_type", ChecksumType.FULL_OBJECT)
1✔
828

829
        if s3_bucket.lifecycle_rules:
1✔
830
            if expiration_header := self._get_expiration_header(
1✔
831
                s3_bucket.lifecycle_rules,
832
                bucket_name,
833
                s3_object,
834
                store.TAGS.tags.get(key_id, {}),
835
            ):
836
                # TODO: we either apply the lifecycle to existing objects when we set the new rules, or we need to
837
                #  apply them everytime we get/head an object
838
                response["Expiration"] = expiration_header
1✔
839

840
        add_encryption_to_response(response, s3_object=s3_object)
1✔
841
        if sse_c_key_md5:
1✔
842
            response["SSECustomerAlgorithm"] = "AES256"
1✔
843
            response["SSECustomerKeyMD5"] = sse_c_key_md5
1✔
844

845
        self._notify(context, s3_bucket=s3_bucket, s3_object=s3_object)
1✔
846

847
        return response
1✔
848

849
    @handler("GetObject", expand=False)
1✔
850
    def get_object(
1✔
851
        self,
852
        context: RequestContext,
853
        request: GetObjectRequest,
854
    ) -> GetObjectOutput:
855
        # TODO: missing handling parameters:
856
        #  request_payer: RequestPayer = None,
857
        #  expected_bucket_owner: AccountId = None,
858

859
        bucket_name = request["Bucket"]
1✔
860
        object_key = request["Key"]
1✔
861
        version_id = request.get("VersionId")
1✔
862
        store, s3_bucket = self._get_cross_account_bucket(context, bucket_name)
1✔
863

864
        s3_object = s3_bucket.get_object(
1✔
865
            key=object_key,
866
            version_id=version_id,
867
            http_method="GET",
868
        )
869
        if s3_object.expires and s3_object.expires < datetime.datetime.now(
1✔
870
            tz=s3_object.expires.tzinfo
871
        ):
872
            # TODO: old behaviour was deleting key instantly if expired. AWS cleans up only once a day generally
873
            #  you can still HeadObject on it and you get the expiry time until scheduled deletion
874
            kwargs = {"Key": object_key}
1✔
875
            if version_id:
1✔
876
                kwargs["VersionId"] = version_id
×
877
            raise NoSuchKey("The specified key does not exist.", **kwargs)
1✔
878

879
        if s3_object.storage_class in ARCHIVES_STORAGE_CLASSES and not s3_object.restore:
1✔
880
            raise InvalidObjectState(
1✔
881
                "The operation is not valid for the object's storage class",
882
                StorageClass=s3_object.storage_class,
883
            )
884

885
        if not config.S3_SKIP_KMS_KEY_VALIDATION and s3_object.kms_key_id:
1✔
886
            validate_kms_key_id(kms_key=s3_object.kms_key_id, bucket=s3_bucket)
1✔
887

888
        sse_c_key_md5 = request.get("SSECustomerKeyMD5")
1✔
889
        # we're using getattr access because when restoring, the field might not exist
890
        # TODO: cleanup at next major release
891
        if sse_key_hash := getattr(s3_object, "sse_key_hash", None):
1✔
892
            if sse_key_hash and not sse_c_key_md5:
1✔
893
                raise InvalidRequest(
1✔
894
                    "The object was stored using a form of Server Side Encryption. "
895
                    "The correct parameters must be provided to retrieve the object."
896
                )
897
            elif sse_key_hash != sse_c_key_md5:
1✔
898
                raise AccessDenied(
1✔
899
                    "Requests specifying Server Side Encryption with Customer provided keys must provide the correct secret key."
900
                )
901

902
        validate_sse_c(
1✔
903
            algorithm=request.get("SSECustomerAlgorithm"),
904
            encryption_key=request.get("SSECustomerKey"),
905
            encryption_key_md5=sse_c_key_md5,
906
        )
907

908
        validate_failed_precondition(request, s3_object.last_modified, s3_object.etag)
1✔
909

910
        range_header = request.get("Range")
1✔
911
        part_number = request.get("PartNumber")
1✔
912
        if range_header and part_number:
1✔
913
            raise InvalidRequest("Cannot specify both Range header and partNumber query parameter")
1✔
914
        range_data = None
1✔
915
        if range_header:
1✔
916
            range_data = parse_range_header(range_header, s3_object.size)
1✔
917
        elif part_number:
1✔
918
            range_data = get_part_range(s3_object, part_number)
1✔
919

920
        # we deliberately do not call `.close()` on the s3_stored_object to keep the read lock acquired. When passing
921
        # the object to Werkzeug, the handler will call `.close()` after finishing iterating over `__iter__`.
922
        # this can however lead to deadlocks if an exception happens between the call and returning the object.
923
        # Be careful into adding validation between this call and `return` of `S3Provider.get_object`
924
        s3_stored_object = self._storage_backend.open(bucket_name, s3_object, mode="r")
1✔
925

926
        # this is a hacky way to verify the object hasn't been modified between `s3_object = s3_bucket.get_object`
927
        # and the storage backend call. If it has been modified, now that we're in the read lock, we can safely fetch
928
        # the object again
929
        if s3_stored_object.last_modified != s3_object.internal_last_modified:
1✔
930
            s3_object = s3_bucket.get_object(
1✔
931
                key=object_key,
932
                version_id=version_id,
933
                http_method="GET",
934
            )
935

936
        response = GetObjectOutput(
1✔
937
            AcceptRanges="bytes",
938
            **s3_object.get_system_metadata_fields(),
939
        )
940
        if s3_object.user_metadata:
1✔
941
            response["Metadata"] = s3_object.user_metadata
1✔
942

943
        if s3_object.parts and request.get("PartNumber"):
1✔
944
            response["PartsCount"] = len(s3_object.parts)
1✔
945

946
        if s3_object.version_id:
1✔
947
            response["VersionId"] = s3_object.version_id
1✔
948

949
        if s3_object.website_redirect_location:
1✔
950
            response["WebsiteRedirectLocation"] = s3_object.website_redirect_location
1✔
951

952
        if s3_object.restore:
1✔
953
            response["Restore"] = s3_object.restore
×
954

955
        checksum_value = None
1✔
956
        if checksum_algorithm := s3_object.checksum_algorithm:
1✔
957
            if (request.get("ChecksumMode") or "").upper() == "ENABLED":
1✔
958
                checksum_value = s3_object.checksum_value
1✔
959

960
        if range_data:
1✔
961
            s3_stored_object.seek(range_data.begin)
1✔
962
            response["Body"] = LimitedIterableStream(
1✔
963
                s3_stored_object, max_length=range_data.content_length
964
            )
965
            response["ContentRange"] = range_data.content_range
1✔
966
            response["ContentLength"] = range_data.content_length
1✔
967
            response["StatusCode"] = 206
1✔
968
            if range_data.content_length == s3_object.size and checksum_value:
1✔
969
                response[f"Checksum{checksum_algorithm.upper()}"] = checksum_value
1✔
970
                response["ChecksumType"] = getattr(
1✔
971
                    s3_object, "checksum_type", ChecksumType.FULL_OBJECT
972
                )
973
        else:
974
            response["Body"] = s3_stored_object
1✔
975
            if checksum_value:
1✔
976
                response[f"Checksum{checksum_algorithm.upper()}"] = checksum_value
1✔
977
                response["ChecksumType"] = getattr(
1✔
978
                    s3_object, "checksum_type", ChecksumType.FULL_OBJECT
979
                )
980

981
        add_encryption_to_response(response, s3_object=s3_object)
1✔
982

983
        if object_tags := store.TAGS.tags.get(
1✔
984
            get_unique_key_id(bucket_name, object_key, version_id)
985
        ):
986
            response["TagCount"] = len(object_tags)
1✔
987

988
        if s3_object.is_current and s3_bucket.lifecycle_rules:
1✔
989
            if expiration_header := self._get_expiration_header(
1✔
990
                s3_bucket.lifecycle_rules,
991
                bucket_name,
992
                s3_object,
993
                object_tags,
994
            ):
995
                # TODO: we either apply the lifecycle to existing objects when we set the new rules, or we need to
996
                #  apply them everytime we get/head an object
997
                response["Expiration"] = expiration_header
1✔
998

999
        # TODO: missing returned fields
1000
        #     RequestCharged: Optional[RequestCharged]
1001
        #     ReplicationStatus: Optional[ReplicationStatus]
1002

1003
        if s3_object.lock_mode:
1✔
1004
            response["ObjectLockMode"] = s3_object.lock_mode
×
1005
            if s3_object.lock_until:
×
1006
                response["ObjectLockRetainUntilDate"] = s3_object.lock_until
×
1007
        if s3_object.lock_legal_status:
1✔
1008
            response["ObjectLockLegalHoldStatus"] = s3_object.lock_legal_status
×
1009

1010
        if sse_c_key_md5:
1✔
1011
            response["SSECustomerAlgorithm"] = "AES256"
1✔
1012
            response["SSECustomerKeyMD5"] = sse_c_key_md5
1✔
1013

1014
        for request_param, response_param in ALLOWED_HEADER_OVERRIDES.items():
1✔
1015
            if request_param_value := request.get(request_param):
1✔
1016
                response[response_param] = request_param_value
1✔
1017

1018
        return response
1✔
1019

1020
    @handler("HeadObject", expand=False)
1✔
1021
    def head_object(
1✔
1022
        self,
1023
        context: RequestContext,
1024
        request: HeadObjectRequest,
1025
    ) -> HeadObjectOutput:
1026
        bucket_name = request["Bucket"]
1✔
1027
        object_key = request["Key"]
1✔
1028
        version_id = request.get("VersionId")
1✔
1029
        store, s3_bucket = self._get_cross_account_bucket(context, bucket_name)
1✔
1030

1031
        s3_object = s3_bucket.get_object(
1✔
1032
            key=object_key,
1033
            version_id=version_id,
1034
            http_method="HEAD",
1035
        )
1036

1037
        validate_failed_precondition(request, s3_object.last_modified, s3_object.etag)
1✔
1038

1039
        sse_c_key_md5 = request.get("SSECustomerKeyMD5")
1✔
1040
        if s3_object.sse_key_hash:
1✔
1041
            if not sse_c_key_md5:
1✔
1042
                raise InvalidRequest(
×
1043
                    "The object was stored using a form of Server Side Encryption. "
1044
                    "The correct parameters must be provided to retrieve the object."
1045
                )
1046
            elif s3_object.sse_key_hash != sse_c_key_md5:
1✔
1047
                raise AccessDenied(
1✔
1048
                    "Requests specifying Server Side Encryption with Customer provided keys must provide the correct secret key."
1049
                )
1050

1051
        validate_sse_c(
1✔
1052
            algorithm=request.get("SSECustomerAlgorithm"),
1053
            encryption_key=request.get("SSECustomerKey"),
1054
            encryption_key_md5=sse_c_key_md5,
1055
        )
1056

1057
        response = HeadObjectOutput(
1✔
1058
            AcceptRanges="bytes",
1059
            **s3_object.get_system_metadata_fields(),
1060
        )
1061
        if s3_object.user_metadata:
1✔
1062
            response["Metadata"] = s3_object.user_metadata
1✔
1063

1064
        if checksum_algorithm := s3_object.checksum_algorithm:
1✔
1065
            if (request.get("ChecksumMode") or "").upper() == "ENABLED":
1✔
1066
                response[f"Checksum{checksum_algorithm.upper()}"] = s3_object.checksum_value
1✔
1067

1068
        if s3_object.parts and request.get("PartNumber"):
1✔
1069
            response["PartsCount"] = len(s3_object.parts)
1✔
1070

1071
        if s3_object.version_id:
1✔
1072
            response["VersionId"] = s3_object.version_id
1✔
1073

1074
        if s3_object.website_redirect_location:
1✔
1075
            response["WebsiteRedirectLocation"] = s3_object.website_redirect_location
1✔
1076

1077
        if s3_object.restore:
1✔
1078
            response["Restore"] = s3_object.restore
1✔
1079

1080
        range_header = request.get("Range")
1✔
1081
        part_number = request.get("PartNumber")
1✔
1082
        if range_header and part_number:
1✔
1083
            raise InvalidRequest("Cannot specify both Range header and partNumber query parameter")
×
1084
        range_data = None
1✔
1085
        if range_header:
1✔
1086
            range_data = parse_range_header(range_header, s3_object.size)
×
1087
        elif part_number:
1✔
1088
            range_data = get_part_range(s3_object, part_number)
1✔
1089

1090
        if range_data:
1✔
1091
            response["ContentLength"] = range_data.content_length
1✔
1092
            response["StatusCode"] = 206
1✔
1093

1094
        add_encryption_to_response(response, s3_object=s3_object)
1✔
1095

1096
        # if you specify the VersionId, AWS won't return the Expiration header, even if that's the current version
1097
        if not version_id and s3_bucket.lifecycle_rules:
1✔
1098
            object_tags = store.TAGS.tags.get(
1✔
1099
                get_unique_key_id(bucket_name, object_key, s3_object.version_id)
1100
            )
1101
            if expiration_header := self._get_expiration_header(
1✔
1102
                s3_bucket.lifecycle_rules,
1103
                bucket_name,
1104
                s3_object,
1105
                object_tags,
1106
            ):
1107
                # TODO: we either apply the lifecycle to existing objects when we set the new rules, or we need to
1108
                #  apply them everytime we get/head an object
1109
                response["Expiration"] = expiration_header
1✔
1110

1111
        if s3_object.lock_mode:
1✔
1112
            response["ObjectLockMode"] = s3_object.lock_mode
1✔
1113
            if s3_object.lock_until:
1✔
1114
                response["ObjectLockRetainUntilDate"] = s3_object.lock_until
1✔
1115
        if s3_object.lock_legal_status:
1✔
1116
            response["ObjectLockLegalHoldStatus"] = s3_object.lock_legal_status
1✔
1117

1118
        if sse_c_key_md5:
1✔
1119
            response["SSECustomerAlgorithm"] = "AES256"
1✔
1120
            response["SSECustomerKeyMD5"] = sse_c_key_md5
1✔
1121

1122
        # TODO: missing return fields:
1123
        #  ArchiveStatus: Optional[ArchiveStatus]
1124
        #  RequestCharged: Optional[RequestCharged]
1125
        #  ReplicationStatus: Optional[ReplicationStatus]
1126

1127
        return response
1✔
1128

1129
    def delete_object(
1✔
1130
        self,
1131
        context: RequestContext,
1132
        bucket: BucketName,
1133
        key: ObjectKey,
1134
        mfa: MFA = None,
1135
        version_id: ObjectVersionId = None,
1136
        request_payer: RequestPayer = None,
1137
        bypass_governance_retention: BypassGovernanceRetention = None,
1138
        expected_bucket_owner: AccountId = None,
1139
        if_match: IfMatch = None,
1140
        if_match_last_modified_time: IfMatchLastModifiedTime = None,
1141
        if_match_size: IfMatchSize = None,
1142
        **kwargs,
1143
    ) -> DeleteObjectOutput:
1144
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
1145

1146
        if bypass_governance_retention is not None and not s3_bucket.object_lock_enabled:
1✔
1147
            raise InvalidArgument(
1✔
1148
                "x-amz-bypass-governance-retention is only applicable to Object Lock enabled buckets.",
1149
                ArgumentName="x-amz-bypass-governance-retention",
1150
            )
1151

1152
        if s3_bucket.versioning_status is None:
1✔
1153
            if version_id and version_id != "null":
1✔
1154
                raise InvalidArgument(
1✔
1155
                    "Invalid version id specified",
1156
                    ArgumentName="versionId",
1157
                    ArgumentValue=version_id,
1158
                )
1159

1160
            found_object = s3_bucket.objects.pop(key, None)
1✔
1161
            # TODO: RequestCharged
1162
            if found_object:
1✔
1163
                self._storage_backend.remove(bucket, found_object)
1✔
1164
                self._notify(context, s3_bucket=s3_bucket, s3_object=found_object)
1✔
1165
                store.TAGS.tags.pop(get_unique_key_id(bucket, key, version_id), None)
1✔
1166

1167
            return DeleteObjectOutput()
1✔
1168

1169
        if not version_id:
1✔
1170
            delete_marker_id = generate_version_id(s3_bucket.versioning_status)
1✔
1171
            delete_marker = S3DeleteMarker(key=key, version_id=delete_marker_id)
1✔
1172
            s3_bucket.objects.set(key, delete_marker)
1✔
1173
            s3_notif_ctx = S3EventNotificationContext.from_request_context_native(
1✔
1174
                context,
1175
                s3_bucket=s3_bucket,
1176
                s3_object=delete_marker,
1177
            )
1178
            s3_notif_ctx.event_type = f"{s3_notif_ctx.event_type}MarkerCreated"
1✔
1179
            self._notify(context, s3_bucket=s3_bucket, s3_notif_ctx=s3_notif_ctx)
1✔
1180

1181
            return DeleteObjectOutput(VersionId=delete_marker.version_id, DeleteMarker=True)
1✔
1182

1183
        if key not in s3_bucket.objects:
1✔
1184
            return DeleteObjectOutput()
×
1185

1186
        if not (s3_object := s3_bucket.objects.get(key, version_id)):
1✔
1187
            raise InvalidArgument(
1✔
1188
                "Invalid version id specified",
1189
                ArgumentName="versionId",
1190
                ArgumentValue=version_id,
1191
            )
1192

1193
        if s3_object.is_locked(bypass_governance_retention):
1✔
1194
            raise AccessDenied("Access Denied because object protected by object lock.")
1✔
1195

1196
        s3_bucket.objects.pop(object_key=key, version_id=version_id)
1✔
1197
        response = DeleteObjectOutput(VersionId=s3_object.version_id)
1✔
1198

1199
        if isinstance(s3_object, S3DeleteMarker):
1✔
1200
            response["DeleteMarker"] = True
1✔
1201
        else:
1202
            self._storage_backend.remove(bucket, s3_object)
1✔
1203
            store.TAGS.tags.pop(get_unique_key_id(bucket, key, version_id), None)
1✔
1204
        self._notify(context, s3_bucket=s3_bucket, s3_object=s3_object)
1✔
1205

1206
        return response
1✔
1207

1208
    def delete_objects(
1✔
1209
        self,
1210
        context: RequestContext,
1211
        bucket: BucketName,
1212
        delete: Delete,
1213
        mfa: MFA = None,
1214
        request_payer: RequestPayer = None,
1215
        bypass_governance_retention: BypassGovernanceRetention = None,
1216
        expected_bucket_owner: AccountId = None,
1217
        checksum_algorithm: ChecksumAlgorithm = None,
1218
        **kwargs,
1219
    ) -> DeleteObjectsOutput:
1220
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
1221

1222
        if bypass_governance_retention is not None and not s3_bucket.object_lock_enabled:
1✔
1223
            raise InvalidArgument(
1✔
1224
                "x-amz-bypass-governance-retention is only applicable to Object Lock enabled buckets.",
1225
                ArgumentName="x-amz-bypass-governance-retention",
1226
            )
1227

1228
        objects: list[ObjectIdentifier] = delete.get("Objects")
1✔
1229
        if not objects:
1✔
1230
            raise MalformedXML()
×
1231

1232
        # TODO: max 1000 delete at once? test against AWS?
1233

1234
        quiet = delete.get("Quiet", False)
1✔
1235
        deleted = []
1✔
1236
        errors = []
1✔
1237

1238
        to_remove = []
1✔
1239
        for to_delete_object in objects:
1✔
1240
            object_key = to_delete_object.get("Key")
1✔
1241
            version_id = to_delete_object.get("VersionId")
1✔
1242
            if s3_bucket.versioning_status is None:
1✔
1243
                if version_id and version_id != "null":
1✔
1244
                    errors.append(
1✔
1245
                        Error(
1246
                            Code="NoSuchVersion",
1247
                            Key=object_key,
1248
                            Message="The specified version does not exist.",
1249
                            VersionId=version_id,
1250
                        )
1251
                    )
1252
                    continue
1✔
1253

1254
                found_object = s3_bucket.objects.pop(object_key, None)
1✔
1255
                if found_object:
1✔
1256
                    to_remove.append(found_object)
1✔
1257
                    self._notify(context, s3_bucket=s3_bucket, s3_object=found_object)
1✔
1258
                    store.TAGS.tags.pop(get_unique_key_id(bucket, object_key, version_id), None)
1✔
1259
                # small hack to not create a fake object for nothing
1260
                elif s3_bucket.notification_configuration:
1✔
1261
                    # DeleteObjects is a bit weird, even if the object didn't exist, S3 will trigger a notification
1262
                    # for a non-existing object being deleted
1263
                    self._notify(
1✔
1264
                        context, s3_bucket=s3_bucket, s3_object=S3Object(key=object_key, etag="")
1265
                    )
1266

1267
                if not quiet:
1✔
1268
                    deleted.append(DeletedObject(Key=object_key))
1✔
1269

1270
                continue
1✔
1271

1272
            if not version_id:
1✔
1273
                delete_marker_id = generate_version_id(s3_bucket.versioning_status)
1✔
1274
                delete_marker = S3DeleteMarker(key=object_key, version_id=delete_marker_id)
1✔
1275
                s3_bucket.objects.set(object_key, delete_marker)
1✔
1276
                s3_notif_ctx = S3EventNotificationContext.from_request_context_native(
1✔
1277
                    context,
1278
                    s3_bucket=s3_bucket,
1279
                    s3_object=delete_marker,
1280
                )
1281
                s3_notif_ctx.event_type = f"{s3_notif_ctx.event_type}MarkerCreated"
1✔
1282
                self._notify(context, s3_bucket=s3_bucket, s3_notif_ctx=s3_notif_ctx)
1✔
1283

1284
                if not quiet:
1✔
1285
                    deleted.append(
1✔
1286
                        DeletedObject(
1287
                            DeleteMarker=True,
1288
                            DeleteMarkerVersionId=delete_marker_id,
1289
                            Key=object_key,
1290
                        )
1291
                    )
1292
                continue
1✔
1293

1294
            if not (
1✔
1295
                found_object := s3_bucket.objects.get(object_key=object_key, version_id=version_id)
1296
            ):
1297
                errors.append(
1✔
1298
                    Error(
1299
                        Code="NoSuchVersion",
1300
                        Key=object_key,
1301
                        Message="The specified version does not exist.",
1302
                        VersionId=version_id,
1303
                    )
1304
                )
1305
                continue
1✔
1306

1307
            if found_object.is_locked(bypass_governance_retention):
1✔
1308
                errors.append(
1✔
1309
                    Error(
1310
                        Code="AccessDenied",
1311
                        Key=object_key,
1312
                        Message="Access Denied because object protected by object lock.",
1313
                        VersionId=version_id,
1314
                    )
1315
                )
1316
                continue
1✔
1317

1318
            s3_bucket.objects.pop(object_key=object_key, version_id=version_id)
1✔
1319
            if not quiet:
1✔
1320
                deleted_object = DeletedObject(
1✔
1321
                    Key=object_key,
1322
                    VersionId=version_id,
1323
                )
1324
                if isinstance(found_object, S3DeleteMarker):
1✔
1325
                    deleted_object["DeleteMarker"] = True
1✔
1326
                    deleted_object["DeleteMarkerVersionId"] = found_object.version_id
1✔
1327

1328
                deleted.append(deleted_object)
1✔
1329

1330
            if isinstance(found_object, S3Object):
1✔
1331
                to_remove.append(found_object)
1✔
1332

1333
            self._notify(context, s3_bucket=s3_bucket, s3_object=found_object)
1✔
1334
            store.TAGS.tags.pop(get_unique_key_id(bucket, object_key, version_id), None)
1✔
1335

1336
        # TODO: request charged
1337
        self._storage_backend.remove(bucket, to_remove)
1✔
1338
        response: DeleteObjectsOutput = {}
1✔
1339
        # AWS validated: the list of Deleted objects is unordered, multiple identical calls can return different results
1340
        if errors:
1✔
1341
            response["Errors"] = errors
1✔
1342
        if not quiet:
1✔
1343
            response["Deleted"] = deleted
1✔
1344

1345
        return response
1✔
1346

1347
    @handler("CopyObject", expand=False)
1✔
1348
    def copy_object(
1✔
1349
        self,
1350
        context: RequestContext,
1351
        request: CopyObjectRequest,
1352
    ) -> CopyObjectOutput:
1353
        # request_payer: RequestPayer = None,  # TODO:
1354
        dest_bucket = request["Bucket"]
1✔
1355
        dest_key = request["Key"]
1✔
1356
        validate_object_key(dest_key)
1✔
1357
        store, dest_s3_bucket = self._get_cross_account_bucket(context, dest_bucket)
1✔
1358

1359
        src_bucket, src_key, src_version_id = extract_bucket_key_version_id_from_copy_source(
1✔
1360
            request.get("CopySource")
1361
        )
1362
        _, src_s3_bucket = self._get_cross_account_bucket(context, src_bucket)
1✔
1363

1364
        if not config.S3_SKIP_KMS_KEY_VALIDATION and (sse_kms_key_id := request.get("SSEKMSKeyId")):
1✔
1365
            validate_kms_key_id(sse_kms_key_id, dest_s3_bucket)
1✔
1366

1367
        # if the object is a delete marker, get_object will raise NotFound if no versionId, like AWS
1368
        try:
1✔
1369
            src_s3_object = src_s3_bucket.get_object(key=src_key, version_id=src_version_id)
1✔
1370
        except MethodNotAllowed:
×
1371
            raise InvalidRequest(
×
1372
                "The source of a copy request may not specifically refer to a delete marker by version id."
1373
            )
1374

1375
        if src_s3_object.storage_class in ARCHIVES_STORAGE_CLASSES and not src_s3_object.restore:
1✔
1376
            raise InvalidObjectState(
×
1377
                "Operation is not valid for the source object's storage class",
1378
                StorageClass=src_s3_object.storage_class,
1379
            )
1380

1381
        if failed_condition := get_failed_precondition_copy_source(
1✔
1382
            request, src_s3_object.last_modified, src_s3_object.etag
1383
        ):
1384
            raise PreconditionFailed(
1✔
1385
                "At least one of the pre-conditions you specified did not hold",
1386
                Condition=failed_condition,
1387
            )
1388

1389
        source_sse_c_key_md5 = request.get("CopySourceSSECustomerKeyMD5")
1✔
1390
        if src_s3_object.sse_key_hash:
1✔
1391
            if not source_sse_c_key_md5:
1✔
1392
                raise InvalidRequest(
1✔
1393
                    "The object was stored using a form of Server Side Encryption. "
1394
                    "The correct parameters must be provided to retrieve the object."
1395
                )
1396
            elif src_s3_object.sse_key_hash != source_sse_c_key_md5:
1✔
1397
                raise AccessDenied("Access Denied")
×
1398

1399
        validate_sse_c(
1✔
1400
            algorithm=request.get("CopySourceSSECustomerAlgorithm"),
1401
            encryption_key=request.get("CopySourceSSECustomerKey"),
1402
            encryption_key_md5=source_sse_c_key_md5,
1403
        )
1404

1405
        target_sse_c_key_md5 = request.get("SSECustomerKeyMD5")
1✔
1406
        server_side_encryption = request.get("ServerSideEncryption")
1✔
1407
        # validate target SSE-C parameters
1408
        validate_sse_c(
1✔
1409
            algorithm=request.get("SSECustomerAlgorithm"),
1410
            encryption_key=request.get("SSECustomerKey"),
1411
            encryption_key_md5=target_sse_c_key_md5,
1412
            server_side_encryption=server_side_encryption,
1413
        )
1414

1415
        # TODO validate order of validation
1416
        storage_class = request.get("StorageClass")
1✔
1417
        metadata_directive = request.get("MetadataDirective")
1✔
1418
        website_redirect_location = request.get("WebsiteRedirectLocation")
1✔
1419
        # we need to check for identity of the object, to see if the default one has been changed
1420
        is_default_encryption = (
1✔
1421
            dest_s3_bucket.encryption_rule is DEFAULT_BUCKET_ENCRYPTION
1422
            and src_s3_object.encryption == "AES256"
1423
        )
1424
        if (
1✔
1425
            src_bucket == dest_bucket
1426
            and src_key == dest_key
1427
            and not any(
1428
                (
1429
                    storage_class,
1430
                    server_side_encryption,
1431
                    target_sse_c_key_md5,
1432
                    metadata_directive == "REPLACE",
1433
                    website_redirect_location,
1434
                    dest_s3_bucket.encryption_rule
1435
                    and not is_default_encryption,  # S3 will allow copy in place if the bucket has encryption configured
1436
                    src_s3_object.restore,
1437
                )
1438
            )
1439
        ):
1440
            raise InvalidRequest(
1✔
1441
                "This copy request is illegal because it is trying to copy an object to itself without changing the "
1442
                "object's metadata, storage class, website redirect location or encryption attributes."
1443
            )
1444

1445
        if tagging := request.get("Tagging"):
1✔
1446
            tagging = parse_tagging_header(tagging)
1✔
1447

1448
        if metadata_directive == "REPLACE":
1✔
1449
            user_metadata = request.get("Metadata")
1✔
1450
            system_metadata = get_system_metadata_from_request(request)
1✔
1451
            if not system_metadata.get("ContentType"):
1✔
1452
                system_metadata["ContentType"] = "binary/octet-stream"
1✔
1453
        else:
1454
            user_metadata = src_s3_object.user_metadata
1✔
1455
            system_metadata = src_s3_object.system_metadata
1✔
1456

1457
        dest_version_id = generate_version_id(dest_s3_bucket.versioning_status)
1✔
1458

1459
        encryption_parameters = get_encryption_parameters_from_request_and_bucket(
1✔
1460
            request,
1461
            dest_s3_bucket,
1462
            store,
1463
        )
1464
        lock_parameters = get_object_lock_parameters_from_bucket_and_request(
1✔
1465
            request, dest_s3_bucket
1466
        )
1467

1468
        acl = get_access_control_policy_for_new_resource_request(
1✔
1469
            request, owner=dest_s3_bucket.owner
1470
        )
1471

1472
        s3_object = S3Object(
1✔
1473
            key=dest_key,
1474
            size=src_s3_object.size,
1475
            version_id=dest_version_id,
1476
            storage_class=storage_class,
1477
            expires=request.get("Expires"),
1478
            user_metadata=user_metadata,
1479
            system_metadata=system_metadata,
1480
            checksum_algorithm=request.get("ChecksumAlgorithm") or src_s3_object.checksum_algorithm,
1481
            encryption=encryption_parameters.encryption,
1482
            kms_key_id=encryption_parameters.kms_key_id,
1483
            bucket_key_enabled=request.get(
1484
                "BucketKeyEnabled"
1485
            ),  # CopyObject does not inherit from the bucket here
1486
            sse_key_hash=target_sse_c_key_md5,
1487
            lock_mode=lock_parameters.lock_mode,
1488
            lock_legal_status=lock_parameters.lock_legal_status,
1489
            lock_until=lock_parameters.lock_until,
1490
            website_redirect_location=website_redirect_location,
1491
            expiration=None,  # TODO, from lifecycle
1492
            acl=acl,
1493
            owner=dest_s3_bucket.owner,
1494
        )
1495

1496
        with self._storage_backend.copy(
1✔
1497
            src_bucket=src_bucket,
1498
            src_object=src_s3_object,
1499
            dest_bucket=dest_bucket,
1500
            dest_object=s3_object,
1501
        ) as s3_stored_object:
1502
            s3_object.checksum_value = s3_stored_object.checksum or src_s3_object.checksum_value
1✔
1503
            s3_object.etag = s3_stored_object.etag or src_s3_object.etag
1✔
1504

1505
            dest_s3_bucket.objects.set(dest_key, s3_object)
1✔
1506

1507
        dest_key_id = get_unique_key_id(dest_bucket, dest_key, dest_version_id)
1✔
1508

1509
        if (request.get("TaggingDirective")) == "REPLACE":
1✔
1510
            store.TAGS.tags[dest_key_id] = tagging or {}
1✔
1511
        else:
1512
            src_key_id = get_unique_key_id(src_bucket, src_key, src_s3_object.version_id)
1✔
1513
            src_tags = store.TAGS.tags.get(src_key_id, {})
1✔
1514
            store.TAGS.tags[dest_key_id] = copy.copy(src_tags)
1✔
1515

1516
        copy_object_result = CopyObjectResult(
1✔
1517
            ETag=s3_object.quoted_etag,
1518
            LastModified=s3_object.last_modified,
1519
        )
1520
        if s3_object.checksum_algorithm:
1✔
1521
            copy_object_result[f"Checksum{s3_object.checksum_algorithm.upper()}"] = (
1✔
1522
                s3_object.checksum_value
1523
            )
1524

1525
        response = CopyObjectOutput(
1✔
1526
            CopyObjectResult=copy_object_result,
1527
        )
1528

1529
        if s3_object.version_id:
1✔
1530
            response["VersionId"] = s3_object.version_id
1✔
1531

1532
        if s3_object.expiration:
1✔
1533
            response["Expiration"] = s3_object.expiration  # TODO: properly parse the datetime
×
1534

1535
        add_encryption_to_response(response, s3_object=s3_object)
1✔
1536
        if target_sse_c_key_md5:
1✔
1537
            response["SSECustomerAlgorithm"] = "AES256"
1✔
1538
            response["SSECustomerKeyMD5"] = target_sse_c_key_md5
1✔
1539

1540
        if (
1✔
1541
            src_s3_bucket.versioning_status
1542
            and src_s3_object.version_id
1543
            and src_s3_object.version_id != "null"
1544
        ):
1545
            response["CopySourceVersionId"] = src_s3_object.version_id
1✔
1546

1547
        # RequestCharged: Optional[RequestCharged] # TODO
1548
        self._notify(context, s3_bucket=dest_s3_bucket, s3_object=s3_object)
1✔
1549

1550
        return response
1✔
1551

1552
    def list_objects(
1✔
1553
        self,
1554
        context: RequestContext,
1555
        bucket: BucketName,
1556
        delimiter: Delimiter = None,
1557
        encoding_type: EncodingType = None,
1558
        marker: Marker = None,
1559
        max_keys: MaxKeys = None,
1560
        prefix: Prefix = None,
1561
        request_payer: RequestPayer = None,
1562
        expected_bucket_owner: AccountId = None,
1563
        optional_object_attributes: OptionalObjectAttributesList = None,
1564
        **kwargs,
1565
    ) -> ListObjectsOutput:
1566
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
1567

1568
        common_prefixes = set()
1✔
1569
        count = 0
1✔
1570
        is_truncated = False
1✔
1571
        next_key_marker = None
1✔
1572
        max_keys = max_keys or 1000
1✔
1573
        prefix = prefix or ""
1✔
1574
        delimiter = delimiter or ""
1✔
1575
        if encoding_type:
1✔
1576
            prefix = urlparse.quote(prefix)
1✔
1577
            delimiter = urlparse.quote(delimiter)
1✔
1578

1579
        s3_objects: list[Object] = []
1✔
1580

1581
        all_keys = sorted(s3_bucket.objects.values(), key=lambda r: r.key)
1✔
1582
        last_key = all_keys[-1] if all_keys else None
1✔
1583

1584
        # sort by key
1585
        for s3_object in all_keys:
1✔
1586
            key = urlparse.quote(s3_object.key) if encoding_type else s3_object.key
1✔
1587
            # skip all keys that alphabetically come before key_marker
1588
            if marker:
1✔
1589
                if key <= marker:
1✔
1590
                    continue
1✔
1591

1592
            # Filter for keys that start with prefix
1593
            if prefix and not key.startswith(prefix):
1✔
1594
                continue
×
1595

1596
            # see ListObjectsV2 for the logic comments (shared logic here)
1597
            prefix_including_delimiter = None
1✔
1598
            if delimiter and delimiter in (key_no_prefix := key.removeprefix(prefix)):
1✔
1599
                pre_delimiter, _, _ = key_no_prefix.partition(delimiter)
1✔
1600
                prefix_including_delimiter = f"{prefix}{pre_delimiter}{delimiter}"
1✔
1601

1602
                if prefix_including_delimiter in common_prefixes or (
1✔
1603
                    marker and marker.startswith(prefix_including_delimiter)
1604
                ):
1605
                    continue
1✔
1606

1607
            if prefix_including_delimiter:
1✔
1608
                common_prefixes.add(prefix_including_delimiter)
1✔
1609
            else:
1610
                # TODO: add RestoreStatus if present
1611
                object_data = Object(
1✔
1612
                    Key=key,
1613
                    ETag=s3_object.quoted_etag,
1614
                    Owner=s3_bucket.owner,  # TODO: verify reality
1615
                    Size=s3_object.size,
1616
                    LastModified=s3_object.last_modified,
1617
                    StorageClass=s3_object.storage_class,
1618
                )
1619

1620
                if s3_object.checksum_algorithm:
1✔
1621
                    object_data["ChecksumAlgorithm"] = [s3_object.checksum_algorithm]
1✔
1622
                    object_data["ChecksumType"] = getattr(
1✔
1623
                        s3_object, "checksum_type", ChecksumType.FULL_OBJECT
1624
                    )
1625

1626
                s3_objects.append(object_data)
1✔
1627

1628
            # we just added a CommonPrefix or an Object, increase the counter
1629
            count += 1
1✔
1630
            if count >= max_keys and last_key.key != s3_object.key:
1✔
1631
                is_truncated = True
1✔
1632
                if prefix_including_delimiter:
1✔
1633
                    next_key_marker = prefix_including_delimiter
1✔
1634
                elif s3_objects:
1✔
1635
                    next_key_marker = s3_objects[-1]["Key"]
1✔
1636
                break
1✔
1637

1638
        common_prefixes = [CommonPrefix(Prefix=prefix) for prefix in sorted(common_prefixes)]
1✔
1639

1640
        response = ListObjectsOutput(
1✔
1641
            IsTruncated=is_truncated,
1642
            Name=bucket,
1643
            MaxKeys=max_keys,
1644
            Prefix=prefix or "",
1645
            Marker=marker or "",
1646
        )
1647
        if s3_objects:
1✔
1648
            response["Contents"] = s3_objects
1✔
1649
        if encoding_type:
1✔
1650
            response["EncodingType"] = EncodingType.url
1✔
1651
        if delimiter:
1✔
1652
            response["Delimiter"] = delimiter
1✔
1653
        if common_prefixes:
1✔
1654
            response["CommonPrefixes"] = common_prefixes
1✔
1655
        if delimiter and next_key_marker:
1✔
1656
            response["NextMarker"] = next_key_marker
1✔
1657
        if s3_bucket.bucket_region != "us-east-1":
1✔
UNCOV
1658
            response["BucketRegion"] = s3_bucket.bucket_region
×
1659

1660
        # RequestCharged: Optional[RequestCharged]  # TODO
1661
        return response
1✔
1662

1663
    def list_objects_v2(
1✔
1664
        self,
1665
        context: RequestContext,
1666
        bucket: BucketName,
1667
        delimiter: Delimiter = None,
1668
        encoding_type: EncodingType = None,
1669
        max_keys: MaxKeys = None,
1670
        prefix: Prefix = None,
1671
        continuation_token: Token = None,
1672
        fetch_owner: FetchOwner = None,
1673
        start_after: StartAfter = None,
1674
        request_payer: RequestPayer = None,
1675
        expected_bucket_owner: AccountId = None,
1676
        optional_object_attributes: OptionalObjectAttributesList = None,
1677
        **kwargs,
1678
    ) -> ListObjectsV2Output:
1679
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
1680

1681
        if continuation_token == "":
1✔
1682
            raise InvalidArgument(
1✔
1683
                "The continuation token provided is incorrect",
1684
                ArgumentName="continuation-token",
1685
            )
1686

1687
        common_prefixes = set()
1✔
1688
        count = 0
1✔
1689
        is_truncated = False
1✔
1690
        next_continuation_token = None
1✔
1691
        max_keys = max_keys or 1000
1✔
1692
        prefix = prefix or ""
1✔
1693
        delimiter = delimiter or ""
1✔
1694
        if encoding_type:
1✔
1695
            prefix = urlparse.quote(prefix)
1✔
1696
            delimiter = urlparse.quote(delimiter)
1✔
1697
        decoded_continuation_token = (
1✔
1698
            to_str(base64.urlsafe_b64decode(continuation_token.encode()))
1699
            if continuation_token
1700
            else None
1701
        )
1702

1703
        s3_objects: list[Object] = []
1✔
1704

1705
        # sort by key
1706
        for s3_object in sorted(s3_bucket.objects.values(), key=lambda r: r.key):
1✔
1707
            key = urlparse.quote(s3_object.key) if encoding_type else s3_object.key
1✔
1708

1709
            # skip all keys that alphabetically come before continuation_token
1710
            if continuation_token:
1✔
1711
                if key < decoded_continuation_token:
1✔
1712
                    continue
1✔
1713

1714
            elif start_after:
1✔
1715
                if key <= start_after:
1✔
1716
                    continue
1✔
1717

1718
            # Filter for keys that start with prefix
1719
            if prefix and not key.startswith(prefix):
1✔
1720
                continue
1✔
1721

1722
            # separate keys that contain the same string between the prefix and the first occurrence of the delimiter
1723
            prefix_including_delimiter = None
1✔
1724
            if delimiter and delimiter in (key_no_prefix := key.removeprefix(prefix)):
1✔
1725
                pre_delimiter, _, _ = key_no_prefix.partition(delimiter)
1✔
1726
                prefix_including_delimiter = f"{prefix}{pre_delimiter}{delimiter}"
1✔
1727

1728
                # if the CommonPrefix is already in the CommonPrefixes, it doesn't count towards MaxKey, we can skip
1729
                # the entry without increasing the counter. We need to iterate over all of these entries before
1730
                # returning the next continuation marker, to properly start at the next key after this CommonPrefix
1731
                if prefix_including_delimiter in common_prefixes:
1✔
1732
                    continue
1✔
1733

1734
            # After skipping all entries, verify we're not over the MaxKeys before adding a new entry
1735
            if count >= max_keys:
1✔
1736
                is_truncated = True
1✔
1737
                next_continuation_token = to_str(base64.urlsafe_b64encode(s3_object.key.encode()))
1✔
1738
                break
1✔
1739

1740
            # if we found a new CommonPrefix, add it to the CommonPrefixes
1741
            # else, it means it's a new Object, add it to the Contents
1742
            if prefix_including_delimiter:
1✔
1743
                common_prefixes.add(prefix_including_delimiter)
1✔
1744
            else:
1745
                # TODO: add RestoreStatus if present
1746
                object_data = Object(
1✔
1747
                    Key=key,
1748
                    ETag=s3_object.quoted_etag,
1749
                    Size=s3_object.size,
1750
                    LastModified=s3_object.last_modified,
1751
                    StorageClass=s3_object.storage_class,
1752
                )
1753

1754
                if fetch_owner:
1✔
1755
                    object_data["Owner"] = s3_bucket.owner
×
1756

1757
                if s3_object.checksum_algorithm:
1✔
1758
                    object_data["ChecksumAlgorithm"] = [s3_object.checksum_algorithm]
1✔
1759
                    object_data["ChecksumType"] = getattr(
1✔
1760
                        s3_object, "checksum_type", ChecksumType.FULL_OBJECT
1761
                    )
1762

1763
                s3_objects.append(object_data)
1✔
1764

1765
            # we just added either a CommonPrefix or an Object to the List, increase the counter by one
1766
            count += 1
1✔
1767

1768
        common_prefixes = [CommonPrefix(Prefix=prefix) for prefix in sorted(common_prefixes)]
1✔
1769

1770
        response = ListObjectsV2Output(
1✔
1771
            IsTruncated=is_truncated,
1772
            Name=bucket,
1773
            MaxKeys=max_keys,
1774
            Prefix=prefix or "",
1775
            KeyCount=count,
1776
        )
1777
        if s3_objects:
1✔
1778
            response["Contents"] = s3_objects
1✔
1779
        if encoding_type:
1✔
1780
            response["EncodingType"] = EncodingType.url
1✔
1781
        if delimiter:
1✔
1782
            response["Delimiter"] = delimiter
1✔
1783
        if common_prefixes:
1✔
1784
            response["CommonPrefixes"] = common_prefixes
1✔
1785
        if next_continuation_token:
1✔
1786
            response["NextContinuationToken"] = next_continuation_token
1✔
1787

1788
        if continuation_token:
1✔
1789
            response["ContinuationToken"] = continuation_token
1✔
1790
        elif start_after:
1✔
1791
            response["StartAfter"] = start_after
1✔
1792

1793
        if s3_bucket.bucket_region != "us-east-1":
1✔
1794
            response["BucketRegion"] = s3_bucket.bucket_region
1✔
1795

1796
        # RequestCharged: Optional[RequestCharged]  # TODO
1797
        return response
1✔
1798

1799
    def list_object_versions(
1✔
1800
        self,
1801
        context: RequestContext,
1802
        bucket: BucketName,
1803
        delimiter: Delimiter = None,
1804
        encoding_type: EncodingType = None,
1805
        key_marker: KeyMarker = None,
1806
        max_keys: MaxKeys = None,
1807
        prefix: Prefix = None,
1808
        version_id_marker: VersionIdMarker = None,
1809
        expected_bucket_owner: AccountId = None,
1810
        request_payer: RequestPayer = None,
1811
        optional_object_attributes: OptionalObjectAttributesList = None,
1812
        **kwargs,
1813
    ) -> ListObjectVersionsOutput:
1814
        if version_id_marker and not key_marker:
1✔
1815
            raise InvalidArgument(
1✔
1816
                "A version-id marker cannot be specified without a key marker.",
1817
                ArgumentName="version-id-marker",
1818
                ArgumentValue=version_id_marker,
1819
            )
1820

1821
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
1822
        common_prefixes = set()
1✔
1823
        count = 0
1✔
1824
        is_truncated = False
1✔
1825
        next_key_marker = None
1✔
1826
        next_version_id_marker = None
1✔
1827
        max_keys = max_keys or 1000
1✔
1828
        prefix = prefix or ""
1✔
1829
        delimiter = delimiter or ""
1✔
1830
        if encoding_type:
1✔
1831
            prefix = urlparse.quote(prefix)
1✔
1832
            delimiter = urlparse.quote(delimiter)
1✔
1833
        version_key_marker_found = False
1✔
1834

1835
        object_versions: list[ObjectVersion] = []
1✔
1836
        delete_markers: list[DeleteMarkerEntry] = []
1✔
1837

1838
        all_versions = s3_bucket.objects.values(with_versions=True)
1✔
1839
        # sort by key, and last-modified-date, to get the last version first
1840
        all_versions.sort(key=lambda r: (r.key, -r.last_modified.timestamp()))
1✔
1841
        last_version = all_versions[-1] if all_versions else None
1✔
1842

1843
        for version in all_versions:
1✔
1844
            key = urlparse.quote(version.key) if encoding_type else version.key
1✔
1845
            # skip all keys that alphabetically come before key_marker
1846
            if key_marker:
1✔
1847
                if key < key_marker:
1✔
1848
                    continue
1✔
1849
                elif key == key_marker:
1✔
1850
                    if not version_id_marker:
1✔
1851
                        continue
1✔
1852
                    # as the keys are ordered by time, once we found the key marker, we can return the next one
1853
                    if version.version_id == version_id_marker:
1✔
1854
                        version_key_marker_found = True
1✔
1855
                        continue
1✔
1856
                    elif not version_key_marker_found:
1✔
1857
                        # as long as we have not passed the version_key_marker, skip the versions
1858
                        continue
1✔
1859

1860
            # Filter for keys that start with prefix
1861
            if prefix and not key.startswith(prefix):
1✔
1862
                continue
1✔
1863

1864
            # see ListObjectsV2 for the logic comments (shared logic here)
1865
            prefix_including_delimiter = None
1✔
1866
            if delimiter and delimiter in (key_no_prefix := key.removeprefix(prefix)):
1✔
1867
                pre_delimiter, _, _ = key_no_prefix.partition(delimiter)
1✔
1868
                prefix_including_delimiter = f"{prefix}{pre_delimiter}{delimiter}"
1✔
1869

1870
                if prefix_including_delimiter in common_prefixes or (
1✔
1871
                    key_marker and key_marker.startswith(prefix_including_delimiter)
1872
                ):
1873
                    continue
1✔
1874

1875
            if prefix_including_delimiter:
1✔
1876
                common_prefixes.add(prefix_including_delimiter)
1✔
1877

1878
            elif isinstance(version, S3DeleteMarker):
1✔
1879
                delete_marker = DeleteMarkerEntry(
1✔
1880
                    Key=key,
1881
                    Owner=s3_bucket.owner,
1882
                    VersionId=version.version_id,
1883
                    IsLatest=version.is_current,
1884
                    LastModified=version.last_modified,
1885
                )
1886
                delete_markers.append(delete_marker)
1✔
1887
            else:
1888
                # TODO: add RestoreStatus if present
1889
                object_version = ObjectVersion(
1✔
1890
                    Key=key,
1891
                    ETag=version.quoted_etag,
1892
                    Owner=s3_bucket.owner,  # TODO: verify reality
1893
                    Size=version.size,
1894
                    VersionId=version.version_id or "null",
1895
                    LastModified=version.last_modified,
1896
                    IsLatest=version.is_current,
1897
                    # TODO: verify this, are other class possible?
1898
                    # StorageClass=version.storage_class,
1899
                    StorageClass=ObjectVersionStorageClass.STANDARD,
1900
                )
1901

1902
                if version.checksum_algorithm:
1✔
1903
                    object_version["ChecksumAlgorithm"] = [version.checksum_algorithm]
1✔
1904
                    object_version["ChecksumType"] = getattr(
1✔
1905
                        version, "checksum_type", ChecksumType.FULL_OBJECT
1906
                    )
1907

1908
                object_versions.append(object_version)
1✔
1909

1910
            # we just added a CommonPrefix, an Object or a DeleteMarker, increase the counter
1911
            count += 1
1✔
1912
            if count >= max_keys and last_version.version_id != version.version_id:
1✔
1913
                is_truncated = True
1✔
1914
                if prefix_including_delimiter:
1✔
1915
                    next_key_marker = prefix_including_delimiter
1✔
1916
                else:
1917
                    next_key_marker = version.key
1✔
1918
                    next_version_id_marker = version.version_id
1✔
1919
                break
1✔
1920

1921
        common_prefixes = [CommonPrefix(Prefix=prefix) for prefix in sorted(common_prefixes)]
1✔
1922

1923
        response = ListObjectVersionsOutput(
1✔
1924
            IsTruncated=is_truncated,
1925
            Name=bucket,
1926
            MaxKeys=max_keys,
1927
            Prefix=prefix,
1928
            KeyMarker=key_marker or "",
1929
            VersionIdMarker=version_id_marker or "",
1930
        )
1931
        if object_versions:
1✔
1932
            response["Versions"] = object_versions
1✔
1933
        if encoding_type:
1✔
1934
            response["EncodingType"] = EncodingType.url
1✔
1935
        if delete_markers:
1✔
1936
            response["DeleteMarkers"] = delete_markers
1✔
1937
        if delimiter:
1✔
1938
            response["Delimiter"] = delimiter
1✔
1939
        if common_prefixes:
1✔
1940
            response["CommonPrefixes"] = common_prefixes
1✔
1941
        if next_key_marker:
1✔
1942
            response["NextKeyMarker"] = next_key_marker
1✔
1943
        if next_version_id_marker:
1✔
1944
            response["NextVersionIdMarker"] = next_version_id_marker
1✔
1945

1946
        # RequestCharged: Optional[RequestCharged]  # TODO
1947
        return response
1✔
1948

1949
    @handler("GetObjectAttributes", expand=False)
1✔
1950
    def get_object_attributes(
1✔
1951
        self,
1952
        context: RequestContext,
1953
        request: GetObjectAttributesRequest,
1954
    ) -> GetObjectAttributesOutput:
1955
        bucket_name = request["Bucket"]
1✔
1956
        object_key = request["Key"]
1✔
1957
        store, s3_bucket = self._get_cross_account_bucket(context, bucket_name)
1✔
1958

1959
        s3_object = s3_bucket.get_object(
1✔
1960
            key=object_key,
1961
            version_id=request.get("VersionId"),
1962
            http_method="GET",
1963
        )
1964

1965
        sse_c_key_md5 = request.get("SSECustomerKeyMD5")
1✔
1966
        if s3_object.sse_key_hash:
1✔
1967
            if not sse_c_key_md5:
1✔
1968
                raise InvalidRequest(
×
1969
                    "The object was stored using a form of Server Side Encryption. "
1970
                    "The correct parameters must be provided to retrieve the object."
1971
                )
1972
            elif s3_object.sse_key_hash != sse_c_key_md5:
1✔
1973
                raise AccessDenied("Access Denied")
×
1974

1975
        validate_sse_c(
1✔
1976
            algorithm=request.get("SSECustomerAlgorithm"),
1977
            encryption_key=request.get("SSECustomerKey"),
1978
            encryption_key_md5=sse_c_key_md5,
1979
        )
1980

1981
        object_attrs = request.get("ObjectAttributes", [])
1✔
1982
        response = GetObjectAttributesOutput()
1✔
1983
        if "ETag" in object_attrs:
1✔
1984
            response["ETag"] = s3_object.etag
1✔
1985
        if "StorageClass" in object_attrs:
1✔
1986
            response["StorageClass"] = s3_object.storage_class
1✔
1987
        if "ObjectSize" in object_attrs:
1✔
1988
            response["ObjectSize"] = s3_object.size
1✔
1989
        if "Checksum" in object_attrs and (checksum_algorithm := s3_object.checksum_algorithm):
1✔
1990
            if s3_object.parts:
1✔
1991
                checksum_value = s3_object.checksum_value.split("-")[0]
1✔
1992
            else:
1993
                checksum_value = s3_object.checksum_value
1✔
1994
            response["Checksum"] = {
1✔
1995
                f"Checksum{checksum_algorithm.upper()}": checksum_value,
1996
                "ChecksumType": getattr(s3_object, "checksum_type", ChecksumType.FULL_OBJECT),
1997
            }
1998

1999
        response["LastModified"] = s3_object.last_modified
1✔
2000

2001
        if s3_bucket.versioning_status:
1✔
2002
            response["VersionId"] = s3_object.version_id
1✔
2003

2004
        if "ObjectParts" in object_attrs and s3_object.parts:
1✔
2005
            # TODO: implements ObjectParts, this is basically a simplified `ListParts` call on the object, we might
2006
            #  need to store more data about the Parts once we implement checksums for them
2007
            response["ObjectParts"] = GetObjectAttributesParts(TotalPartsCount=len(s3_object.parts))
1✔
2008

2009
        return response
1✔
2010

2011
    def restore_object(
1✔
2012
        self,
2013
        context: RequestContext,
2014
        bucket: BucketName,
2015
        key: ObjectKey,
2016
        version_id: ObjectVersionId = None,
2017
        restore_request: RestoreRequest = None,
2018
        request_payer: RequestPayer = None,
2019
        checksum_algorithm: ChecksumAlgorithm = None,
2020
        expected_bucket_owner: AccountId = None,
2021
        **kwargs,
2022
    ) -> RestoreObjectOutput:
2023
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
2024

2025
        s3_object = s3_bucket.get_object(
1✔
2026
            key=key,
2027
            version_id=version_id,
2028
            http_method="GET",  # TODO: verify http method
2029
        )
2030
        if s3_object.storage_class not in ARCHIVES_STORAGE_CLASSES:
1✔
2031
            raise InvalidObjectState(StorageClass=s3_object.storage_class)
×
2032

2033
        # TODO: moto was only supported "Days" parameters from RestoreRequest, and was ignoring the others
2034
        # will only implement only the same functionality for now
2035

2036
        # if a request was already done and the object was available, and we're updating it, set the status code to 200
2037
        status_code = 200 if s3_object.restore else 202
1✔
2038
        restore_days = restore_request.get("Days")
1✔
2039
        if not restore_days:
1✔
2040
            LOG.debug("LocalStack does not support restore SELECT requests yet.")
×
2041
            return RestoreObjectOutput()
×
2042

2043
        restore_expiration_date = add_expiration_days_to_datetime(
1✔
2044
            datetime.datetime.now(datetime.UTC), restore_days
2045
        )
2046
        # TODO: add a way to transition from ongoing-request=true to false? for now it is instant
2047
        s3_object.restore = f'ongoing-request="false", expiry-date="{restore_expiration_date}"'
1✔
2048

2049
        s3_notif_ctx_initiated = S3EventNotificationContext.from_request_context_native(
1✔
2050
            context,
2051
            s3_bucket=s3_bucket,
2052
            s3_object=s3_object,
2053
        )
2054
        self._notify(context, s3_bucket=s3_bucket, s3_notif_ctx=s3_notif_ctx_initiated)
1✔
2055
        # But because it's instant in LocalStack, we can directly send the Completed notification as well
2056
        # We just need to copy the context so that we don't mutate the first context while it could be sent
2057
        # And modify its event type from `ObjectRestore:Post` to `ObjectRestore:Completed`
2058
        s3_notif_ctx_completed = copy.copy(s3_notif_ctx_initiated)
1✔
2059
        s3_notif_ctx_completed.event_type = s3_notif_ctx_completed.event_type.replace(
1✔
2060
            "Post", "Completed"
2061
        )
2062
        self._notify(context, s3_bucket=s3_bucket, s3_notif_ctx=s3_notif_ctx_completed)
1✔
2063

2064
        # TODO: request charged
2065
        return RestoreObjectOutput(StatusCode=status_code)
1✔
2066

2067
    @handler("CreateMultipartUpload", expand=False)
1✔
2068
    def create_multipart_upload(
1✔
2069
        self,
2070
        context: RequestContext,
2071
        request: CreateMultipartUploadRequest,
2072
    ) -> CreateMultipartUploadOutput:
2073
        # TODO: handle missing parameters:
2074
        #  request_payer: RequestPayer = None,
2075
        bucket_name = request["Bucket"]
1✔
2076
        store, s3_bucket = self._get_cross_account_bucket(context, bucket_name)
1✔
2077

2078
        if (storage_class := request.get("StorageClass")) is not None and (
1✔
2079
            storage_class not in STORAGE_CLASSES or storage_class == StorageClass.OUTPOSTS
2080
        ):
2081
            raise InvalidStorageClass(
1✔
2082
                "The storage class you specified is not valid", StorageClassRequested=storage_class
2083
            )
2084

2085
        if not config.S3_SKIP_KMS_KEY_VALIDATION and (sse_kms_key_id := request.get("SSEKMSKeyId")):
1✔
2086
            validate_kms_key_id(sse_kms_key_id, s3_bucket)
1✔
2087

2088
        if tagging := request.get("Tagging"):
1✔
2089
            tagging = parse_tagging_header(tagging_header=tagging)
×
2090

2091
        key = request["Key"]
1✔
2092

2093
        system_metadata = get_system_metadata_from_request(request)
1✔
2094
        if not system_metadata.get("ContentType"):
1✔
2095
            system_metadata["ContentType"] = "binary/octet-stream"
1✔
2096

2097
        checksum_algorithm = request.get("ChecksumAlgorithm")
1✔
2098
        if checksum_algorithm and checksum_algorithm not in CHECKSUM_ALGORITHMS:
1✔
2099
            raise InvalidRequest(
1✔
2100
                "Checksum algorithm provided is unsupported. Please try again with any of the valid types: [CRC32, CRC32C, SHA1, SHA256]"
2101
            )
2102

2103
        # TODO: we're not encrypting the object with the provided key for now
2104
        sse_c_key_md5 = request.get("SSECustomerKeyMD5")
1✔
2105
        validate_sse_c(
1✔
2106
            algorithm=request.get("SSECustomerAlgorithm"),
2107
            encryption_key=request.get("SSECustomerKey"),
2108
            encryption_key_md5=sse_c_key_md5,
2109
            server_side_encryption=request.get("ServerSideEncryption"),
2110
        )
2111

2112
        encryption_parameters = get_encryption_parameters_from_request_and_bucket(
1✔
2113
            request,
2114
            s3_bucket,
2115
            store,
2116
        )
2117
        lock_parameters = get_object_lock_parameters_from_bucket_and_request(request, s3_bucket)
1✔
2118

2119
        acl = get_access_control_policy_for_new_resource_request(request, owner=s3_bucket.owner)
1✔
2120

2121
        # validate encryption values
2122
        s3_multipart = S3Multipart(
1✔
2123
            key=key,
2124
            storage_class=storage_class,
2125
            expires=request.get("Expires"),
2126
            user_metadata=request.get("Metadata"),
2127
            system_metadata=system_metadata,
2128
            checksum_algorithm=checksum_algorithm,
2129
            encryption=encryption_parameters.encryption,
2130
            kms_key_id=encryption_parameters.kms_key_id,
2131
            bucket_key_enabled=encryption_parameters.bucket_key_enabled,
2132
            sse_key_hash=sse_c_key_md5,
2133
            lock_mode=lock_parameters.lock_mode,
2134
            lock_legal_status=lock_parameters.lock_legal_status,
2135
            lock_until=lock_parameters.lock_until,
2136
            website_redirect_location=request.get("WebsiteRedirectLocation"),
2137
            expiration=None,  # TODO, from lifecycle, or should it be updated with config?
2138
            acl=acl,
2139
            initiator=get_owner_for_account_id(context.account_id),
2140
            tagging=tagging,
2141
            owner=s3_bucket.owner,
2142
            precondition=object_exists_for_precondition_write(s3_bucket, key),
2143
        )
2144

2145
        s3_bucket.multiparts[s3_multipart.id] = s3_multipart
1✔
2146

2147
        response = CreateMultipartUploadOutput(
1✔
2148
            Bucket=bucket_name, Key=key, UploadId=s3_multipart.id
2149
        )
2150

2151
        if checksum_algorithm:
1✔
2152
            response["ChecksumAlgorithm"] = checksum_algorithm
1✔
2153

2154
        add_encryption_to_response(response, s3_object=s3_multipart.object)
1✔
2155
        if sse_c_key_md5:
1✔
2156
            response["SSECustomerAlgorithm"] = "AES256"
1✔
2157
            response["SSECustomerKeyMD5"] = sse_c_key_md5
1✔
2158

2159
        # TODO: missing response fields we're not currently supporting
2160
        # - AbortDate: lifecycle related,not currently supported, todo
2161
        # - AbortRuleId: lifecycle related, not currently supported, todo
2162
        # - RequestCharged: todo
2163

2164
        return response
1✔
2165

2166
    @handler("UploadPart", expand=False)
1✔
2167
    def upload_part(
1✔
2168
        self,
2169
        context: RequestContext,
2170
        request: UploadPartRequest,
2171
    ) -> UploadPartOutput:
2172
        # TODO: missing following parameters:
2173
        #  content_length: ContentLength = None, ->validate?
2174
        #  content_md5: ContentMD5 = None, -> validate?
2175
        #  request_payer: RequestPayer = None,
2176
        bucket_name = request["Bucket"]
1✔
2177
        store, s3_bucket = self._get_cross_account_bucket(context, bucket_name)
1✔
2178

2179
        upload_id = request.get("UploadId")
1✔
2180
        if not (
1✔
2181
            s3_multipart := s3_bucket.multiparts.get(upload_id)
2182
        ) or s3_multipart.object.key != request.get("Key"):
2183
            raise NoSuchUpload(
1✔
2184
                "The specified upload does not exist. "
2185
                "The upload ID may be invalid, or the upload may have been aborted or completed.",
2186
                UploadId=upload_id,
2187
            )
2188
        elif (part_number := request.get("PartNumber", 0)) < 1 or part_number > 10000:
1✔
2189
            raise InvalidArgument(
1✔
2190
                "Part number must be an integer between 1 and 10000, inclusive",
2191
                ArgumentName="partNumber",
2192
                ArgumentValue=part_number,
2193
            )
2194

2195
        if content_md5 := request.get("ContentMD5"):
1✔
2196
            # assert that the received ContentMD5 is a properly b64 encoded value that fits a MD5 hash length
2197
            if not base_64_content_md5_to_etag(content_md5):
1✔
2198
                raise InvalidDigest(
1✔
2199
                    "The Content-MD5 you specified was invalid.",
2200
                    Content_MD5=content_md5,
2201
                )
2202

2203
        checksum_algorithm = get_s3_checksum_algorithm_from_request(request)
1✔
2204
        checksum_value = (
1✔
2205
            request.get(f"Checksum{checksum_algorithm.upper()}") if checksum_algorithm else None
2206
        )
2207

2208
        # TODO: we're not encrypting the object with the provided key for now
2209
        sse_c_key_md5 = request.get("SSECustomerKeyMD5")
1✔
2210
        validate_sse_c(
1✔
2211
            algorithm=request.get("SSECustomerAlgorithm"),
2212
            encryption_key=request.get("SSECustomerKey"),
2213
            encryption_key_md5=sse_c_key_md5,
2214
        )
2215

2216
        if (s3_multipart.object.sse_key_hash and not sse_c_key_md5) or (
1✔
2217
            sse_c_key_md5 and not s3_multipart.object.sse_key_hash
2218
        ):
2219
            raise InvalidRequest(
1✔
2220
                "The multipart upload initiate requested encryption. "
2221
                "Subsequent part requests must include the appropriate encryption parameters."
2222
            )
2223
        elif (
1✔
2224
            s3_multipart.object.sse_key_hash
2225
            and sse_c_key_md5
2226
            and s3_multipart.object.sse_key_hash != sse_c_key_md5
2227
        ):
2228
            raise InvalidRequest(
1✔
2229
                "The provided encryption parameters did not match the ones used originally."
2230
            )
2231

2232
        s3_part = S3Part(
1✔
2233
            part_number=part_number,
2234
            checksum_algorithm=checksum_algorithm,
2235
            checksum_value=checksum_value,
2236
        )
2237
        body = request.get("Body")
1✔
2238
        headers = context.request.headers
1✔
2239
        is_aws_chunked = headers.get("x-amz-content-sha256", "").startswith(
1✔
2240
            "STREAMING-"
2241
        ) or "aws-chunked" in headers.get("content-encoding", "")
2242
        # check if chunked request
2243
        if is_aws_chunked:
1✔
2244
            checksum_algorithm = (
1✔
2245
                checksum_algorithm
2246
                or get_s3_checksum_algorithm_from_trailing_headers(headers.get("x-amz-trailer", ""))
2247
            )
2248
            if checksum_algorithm:
1✔
2249
                s3_part.checksum_algorithm = checksum_algorithm
×
2250

2251
            decoded_content_length = int(headers.get("x-amz-decoded-content-length", 0))
1✔
2252
            body = AwsChunkedDecoder(body, decoded_content_length, s3_part)
1✔
2253

2254
        if s3_part.checksum_algorithm != s3_multipart.object.checksum_algorithm:
1✔
2255
            error_req_checksum = checksum_algorithm.lower() if checksum_algorithm else "null"
1✔
2256
            error_mp_checksum = (
1✔
2257
                s3_multipart.object.checksum_algorithm.lower()
2258
                if s3_multipart.object.checksum_algorithm
2259
                else "null"
2260
            )
2261
            # TODO: properly fix this, this is to unblock default behavior of boto adding checksums and it being
2262
            #  accepted by AWS
2263
            if not error_mp_checksum == "null":
1✔
2264
                raise InvalidRequest(
1✔
2265
                    f"Checksum Type mismatch occurred, expected checksum Type: {error_mp_checksum}, actual checksum Type: {error_req_checksum}"
2266
                )
2267

2268
        stored_multipart = self._storage_backend.get_multipart(bucket_name, s3_multipart)
1✔
2269
        with stored_multipart.open(s3_part, mode="w") as stored_s3_part:
1✔
2270
            try:
1✔
2271
                stored_s3_part.write(body)
1✔
2272
            except Exception:
1✔
2273
                stored_multipart.remove_part(s3_part)
1✔
2274
                raise
1✔
2275

2276
            if checksum_algorithm and s3_part.checksum_value != stored_s3_part.checksum:
1✔
2277
                stored_multipart.remove_part(s3_part)
×
2278
                # TODO: validate this to be BadDigest as well
2279
                raise InvalidRequest(
×
2280
                    f"Value for x-amz-checksum-{checksum_algorithm.lower()} header is invalid."
2281
                )
2282

2283
            if content_md5:
1✔
2284
                calculated_md5 = etag_to_base_64_content_md5(s3_part.etag)
1✔
2285
                if calculated_md5 != content_md5:
1✔
2286
                    stored_multipart.remove_part(s3_part)
1✔
2287
                    raise BadDigest(
1✔
2288
                        "The Content-MD5 you specified did not match what we received.",
2289
                        ExpectedDigest=content_md5,
2290
                        CalculatedDigest=calculated_md5,
2291
                    )
2292

2293
            s3_multipart.parts[part_number] = s3_part
1✔
2294

2295
        response = UploadPartOutput(
1✔
2296
            ETag=s3_part.quoted_etag,
2297
        )
2298

2299
        add_encryption_to_response(response, s3_object=s3_multipart.object)
1✔
2300
        if sse_c_key_md5:
1✔
2301
            response["SSECustomerAlgorithm"] = "AES256"
1✔
2302
            response["SSECustomerKeyMD5"] = sse_c_key_md5
1✔
2303

2304
        if s3_part.checksum_algorithm:
1✔
2305
            response[f"Checksum{s3_part.checksum_algorithm.upper()}"] = s3_part.checksum_value
1✔
2306

2307
        # TODO: RequestCharged: Optional[RequestCharged]
2308
        return response
1✔
2309

2310
    @handler("UploadPartCopy", expand=False)
1✔
2311
    def upload_part_copy(
1✔
2312
        self,
2313
        context: RequestContext,
2314
        request: UploadPartCopyRequest,
2315
    ) -> UploadPartCopyOutput:
2316
        # TODO: handle following parameters:
2317
        #  copy_source_if_match: CopySourceIfMatch = None,
2318
        #  copy_source_if_modified_since: CopySourceIfModifiedSince = None,
2319
        #  copy_source_if_none_match: CopySourceIfNoneMatch = None,
2320
        #  copy_source_if_unmodified_since: CopySourceIfUnmodifiedSince = None,
2321
        #  request_payer: RequestPayer = None,
2322
        dest_bucket = request["Bucket"]
1✔
2323
        dest_key = request["Key"]
1✔
2324
        store = self.get_store(context.account_id, context.region)
1✔
2325
        # TODO: validate cross-account UploadPartCopy
2326
        if not (dest_s3_bucket := store.buckets.get(dest_bucket)):
1✔
2327
            raise NoSuchBucket("The specified bucket does not exist", BucketName=dest_bucket)
×
2328

2329
        src_bucket, src_key, src_version_id = extract_bucket_key_version_id_from_copy_source(
1✔
2330
            request.get("CopySource")
2331
        )
2332

2333
        if not (src_s3_bucket := store.buckets.get(src_bucket)):
1✔
2334
            raise NoSuchBucket("The specified bucket does not exist", BucketName=src_bucket)
×
2335

2336
        # if the object is a delete marker, get_object will raise NotFound if no versionId, like AWS
2337
        try:
1✔
2338
            src_s3_object = src_s3_bucket.get_object(key=src_key, version_id=src_version_id)
1✔
2339
        except MethodNotAllowed:
×
2340
            raise InvalidRequest(
×
2341
                "The source of a copy request may not specifically refer to a delete marker by version id."
2342
            )
2343

2344
        if src_s3_object.storage_class in ARCHIVES_STORAGE_CLASSES and not src_s3_object.restore:
1✔
2345
            raise InvalidObjectState(
×
2346
                "Operation is not valid for the source object's storage class",
2347
                StorageClass=src_s3_object.storage_class,
2348
            )
2349

2350
        upload_id = request.get("UploadId")
1✔
2351
        if (
1✔
2352
            not (s3_multipart := dest_s3_bucket.multiparts.get(upload_id))
2353
            or s3_multipart.object.key != dest_key
2354
        ):
2355
            raise NoSuchUpload(
×
2356
                "The specified upload does not exist. "
2357
                "The upload ID may be invalid, or the upload may have been aborted or completed.",
2358
                UploadId=upload_id,
2359
            )
2360

2361
        elif (part_number := request.get("PartNumber", 0)) < 1 or part_number > 10000:
1✔
2362
            raise InvalidArgument(
×
2363
                "Part number must be an integer between 1 and 10000, inclusive",
2364
                ArgumentName="partNumber",
2365
                ArgumentValue=part_number,
2366
            )
2367

2368
        source_range = request.get("CopySourceRange")
1✔
2369
        # TODO implement copy source IF (done in ASF provider)
2370

2371
        range_data: Optional[ObjectRange] = None
1✔
2372
        if source_range:
1✔
2373
            range_data = parse_copy_source_range_header(source_range, src_s3_object.size)
1✔
2374

2375
        s3_part = S3Part(part_number=part_number)
1✔
2376

2377
        stored_multipart = self._storage_backend.get_multipart(dest_bucket, s3_multipart)
1✔
2378
        stored_multipart.copy_from_object(s3_part, src_bucket, src_s3_object, range_data)
1✔
2379

2380
        s3_multipart.parts[part_number] = s3_part
1✔
2381

2382
        # TODO: return those fields (checksum not handled currently in moto for parts)
2383
        # ChecksumCRC32: Optional[ChecksumCRC32]
2384
        # ChecksumCRC32C: Optional[ChecksumCRC32C]
2385
        # ChecksumSHA1: Optional[ChecksumSHA1]
2386
        # ChecksumSHA256: Optional[ChecksumSHA256]
2387
        #     RequestCharged: Optional[RequestCharged]
2388

2389
        result = CopyPartResult(
1✔
2390
            ETag=s3_part.quoted_etag,
2391
            LastModified=s3_part.last_modified,
2392
        )
2393

2394
        response = UploadPartCopyOutput(
1✔
2395
            CopyPartResult=result,
2396
        )
2397

2398
        if src_s3_bucket.versioning_status and src_s3_object.version_id:
1✔
2399
            response["CopySourceVersionId"] = src_s3_object.version_id
×
2400

2401
        add_encryption_to_response(response, s3_object=s3_multipart.object)
1✔
2402

2403
        return response
1✔
2404

2405
    def complete_multipart_upload(
1✔
2406
        self,
2407
        context: RequestContext,
2408
        bucket: BucketName,
2409
        key: ObjectKey,
2410
        upload_id: MultipartUploadId,
2411
        multipart_upload: CompletedMultipartUpload = None,
2412
        checksum_crc32: ChecksumCRC32 = None,
2413
        checksum_crc32_c: ChecksumCRC32C = None,
2414
        checksum_crc64_nvme: ChecksumCRC64NVME = None,
2415
        checksum_sha1: ChecksumSHA1 = None,
2416
        checksum_sha256: ChecksumSHA256 = None,
2417
        checksum_type: ChecksumType = None,
2418
        mpu_object_size: MpuObjectSize = None,
2419
        request_payer: RequestPayer = None,
2420
        expected_bucket_owner: AccountId = None,
2421
        if_match: IfMatch = None,
2422
        if_none_match: IfNoneMatch = None,
2423
        sse_customer_algorithm: SSECustomerAlgorithm = None,
2424
        sse_customer_key: SSECustomerKey = None,
2425
        sse_customer_key_md5: SSECustomerKeyMD5 = None,
2426
        **kwargs,
2427
    ) -> CompleteMultipartUploadOutput:
2428
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
2429

2430
        if (
1✔
2431
            not (s3_multipart := s3_bucket.multiparts.get(upload_id))
2432
            or s3_multipart.object.key != key
2433
        ):
2434
            raise NoSuchUpload(
1✔
2435
                "The specified upload does not exist. The upload ID may be invalid, or the upload may have been aborted or completed.",
2436
                UploadId=upload_id,
2437
            )
2438

2439
        if if_none_match and if_match:
1✔
2440
            raise NotImplementedException(
2441
                "A header you provided implies functionality that is not implemented",
2442
                Header="If-Match,If-None-Match",
2443
                additionalMessage="Multiple conditional request headers present in the request",
2444
            )
2445

2446
        elif if_none_match:
1✔
2447
            if if_none_match != "*":
1✔
2448
                raise NotImplementedException(
2449
                    "A header you provided implies functionality that is not implemented",
2450
                    Header="If-None-Match",
2451
                    additionalMessage="We don't accept the provided value of If-None-Match header for this API",
2452
                )
2453
            if object_exists_for_precondition_write(s3_bucket, key):
1✔
2454
                raise PreconditionFailed(
1✔
2455
                    "At least one of the pre-conditions you specified did not hold",
2456
                    Condition="If-None-Match",
2457
                )
2458
            elif s3_multipart.precondition:
1✔
2459
                raise ConditionalRequestConflict(
1✔
2460
                    "The conditional request cannot succeed due to a conflicting operation against this resource.",
2461
                    Condition="If-None-Match",
2462
                    Key=key,
2463
                )
2464

2465
        elif if_match:
1✔
2466
            if if_match == "*":
1✔
2467
                raise NotImplementedException(
2468
                    "A header you provided implies functionality that is not implemented",
2469
                    Header="If-None-Match",
2470
                    additionalMessage="We don't accept the provided value of If-None-Match header for this API",
2471
                )
2472
            verify_object_equality_precondition_write(
1✔
2473
                s3_bucket, key, if_match, initiated=s3_multipart.initiated
2474
            )
2475

2476
        parts = multipart_upload.get("Parts", [])
1✔
2477
        if not parts:
1✔
2478
            raise InvalidRequest("You must specify at least one part")
1✔
2479

2480
        parts_numbers = [part.get("PartNumber") for part in parts]
1✔
2481
        # TODO: it seems that with new S3 data integrity, sorting might not be mandatory depending on checksum type
2482
        # see https://docs.aws.amazon.com/AmazonS3/latest/userguide/checking-object-integrity.html
2483
        # sorted is very fast (fastest) if the list is already sorted, which should be the case
2484
        if sorted(parts_numbers) != parts_numbers:
1✔
2485
            raise InvalidPartOrder(
1✔
2486
                "The list of parts was not in ascending order. Parts must be ordered by part number.",
2487
                UploadId=upload_id,
2488
            )
2489

2490
        # generate the versionId before completing, in case the bucket versioning status has changed between
2491
        # creation and completion? AWS validate this
2492
        version_id = generate_version_id(s3_bucket.versioning_status)
1✔
2493
        s3_multipart.object.version_id = version_id
1✔
2494
        s3_multipart.complete_multipart(parts)
1✔
2495

2496
        stored_multipart = self._storage_backend.get_multipart(bucket, s3_multipart)
1✔
2497
        stored_multipart.complete_multipart(
1✔
2498
            [s3_multipart.parts.get(part_number) for part_number in parts_numbers]
2499
        )
2500

2501
        s3_object = s3_multipart.object
1✔
2502

2503
        s3_bucket.objects.set(key, s3_object)
1✔
2504

2505
        # remove the multipart now that it's complete
2506
        self._storage_backend.remove_multipart(bucket, s3_multipart)
1✔
2507
        s3_bucket.multiparts.pop(s3_multipart.id, None)
1✔
2508

2509
        key_id = get_unique_key_id(bucket, key, version_id)
1✔
2510
        store.TAGS.tags.pop(key_id, None)
1✔
2511
        if s3_multipart.tagging:
1✔
2512
            store.TAGS.tags[key_id] = s3_multipart.tagging
×
2513

2514
        # TODO: validate if you provide wrong checksum compared to the given algorithm? should you calculate it anyway
2515
        #  when you complete? sounds weird, not sure how that works?
2516

2517
        #     ChecksumCRC32: Optional[ChecksumCRC32] ??
2518
        #     ChecksumCRC32C: Optional[ChecksumCRC32C] ??
2519
        #     ChecksumSHA1: Optional[ChecksumSHA1] ??
2520
        #     ChecksumSHA256: Optional[ChecksumSHA256] ??
2521
        #     RequestCharged: Optional[RequestCharged] TODO
2522

2523
        response = CompleteMultipartUploadOutput(
1✔
2524
            Bucket=bucket,
2525
            Key=key,
2526
            ETag=s3_object.quoted_etag,
2527
            Location=f"{get_full_default_bucket_location(bucket)}{key}",
2528
        )
2529

2530
        if s3_object.version_id:
1✔
2531
            response["VersionId"] = s3_object.version_id
×
2532

2533
        # TODO: check this?
2534
        if s3_object.checksum_algorithm:
1✔
2535
            response[f"Checksum{s3_object.checksum_algorithm.upper()}"] = s3_object.checksum_value
1✔
2536

2537
        if s3_object.expiration:
1✔
2538
            response["Expiration"] = s3_object.expiration  # TODO: properly parse the datetime
×
2539

2540
        add_encryption_to_response(response, s3_object=s3_object)
1✔
2541

2542
        self._notify(context, s3_bucket=s3_bucket, s3_object=s3_object)
1✔
2543

2544
        return response
1✔
2545

2546
    def abort_multipart_upload(
1✔
2547
        self,
2548
        context: RequestContext,
2549
        bucket: BucketName,
2550
        key: ObjectKey,
2551
        upload_id: MultipartUploadId,
2552
        request_payer: RequestPayer = None,
2553
        expected_bucket_owner: AccountId = None,
2554
        if_match_initiated_time: IfMatchInitiatedTime = None,
2555
        **kwargs,
2556
    ) -> AbortMultipartUploadOutput:
2557
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
2558

2559
        if (
1✔
2560
            not (s3_multipart := s3_bucket.multiparts.get(upload_id))
2561
            or s3_multipart.object.key != key
2562
        ):
2563
            raise NoSuchUpload(
1✔
2564
                "The specified upload does not exist. "
2565
                "The upload ID may be invalid, or the upload may have been aborted or completed.",
2566
                UploadId=upload_id,
2567
            )
2568
        s3_bucket.multiparts.pop(upload_id, None)
1✔
2569

2570
        self._storage_backend.remove_multipart(bucket, s3_multipart)
1✔
2571
        response = AbortMultipartUploadOutput()
1✔
2572
        # TODO: requestCharged
2573
        return response
1✔
2574

2575
    def list_parts(
1✔
2576
        self,
2577
        context: RequestContext,
2578
        bucket: BucketName,
2579
        key: ObjectKey,
2580
        upload_id: MultipartUploadId,
2581
        max_parts: MaxParts = None,
2582
        part_number_marker: PartNumberMarker = None,
2583
        request_payer: RequestPayer = None,
2584
        expected_bucket_owner: AccountId = None,
2585
        sse_customer_algorithm: SSECustomerAlgorithm = None,
2586
        sse_customer_key: SSECustomerKey = None,
2587
        sse_customer_key_md5: SSECustomerKeyMD5 = None,
2588
        **kwargs,
2589
    ) -> ListPartsOutput:
2590
        # TODO: implement MaxParts
2591
        # TODO: implements PartNumberMarker
2592
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
2593

2594
        if (
1✔
2595
            not (s3_multipart := s3_bucket.multiparts.get(upload_id))
2596
            or s3_multipart.object.key != key
2597
        ):
2598
            raise NoSuchUpload(
1✔
2599
                "The specified upload does not exist. "
2600
                "The upload ID may be invalid, or the upload may have been aborted or completed.",
2601
                UploadId=upload_id,
2602
            )
2603

2604
        #     AbortDate: Optional[AbortDate] TODO: lifecycle
2605
        #     AbortRuleId: Optional[AbortRuleId] TODO: lifecycle
2606
        #     RequestCharged: Optional[RequestCharged]
2607

2608
        count = 0
1✔
2609
        is_truncated = False
1✔
2610
        part_number_marker = part_number_marker or 0
1✔
2611
        max_parts = max_parts or 1000
1✔
2612

2613
        parts = []
1✔
2614
        all_parts = sorted(s3_multipart.parts.items())
1✔
2615
        last_part_number = all_parts[-1][0] if all_parts else None
1✔
2616
        for part_number, part in all_parts:
1✔
2617
            if part_number <= part_number_marker:
1✔
2618
                continue
1✔
2619
            part_item = Part(
1✔
2620
                ETag=part.quoted_etag,
2621
                LastModified=part.last_modified,
2622
                PartNumber=part_number,
2623
                Size=part.size,
2624
            )
2625
            if part.checksum_algorithm:
1✔
2626
                part_item[f"Checksum{part.checksum_algorithm.upper()}"] = part.checksum_value
1✔
2627

2628
            parts.append(part_item)
1✔
2629
            count += 1
1✔
2630

2631
            if count >= max_parts and part.part_number != last_part_number:
1✔
2632
                is_truncated = True
1✔
2633
                break
1✔
2634

2635
        response = ListPartsOutput(
1✔
2636
            Bucket=bucket,
2637
            Key=key,
2638
            UploadId=upload_id,
2639
            Initiator=s3_multipart.initiator,
2640
            Owner=s3_multipart.initiator,
2641
            StorageClass=s3_multipart.object.storage_class,
2642
            IsTruncated=is_truncated,
2643
            MaxParts=max_parts,
2644
            PartNumberMarker=0,
2645
            NextPartNumberMarker=0,
2646
        )
2647
        if parts:
1✔
2648
            response["Parts"] = parts
1✔
2649
            last_part = parts[-1]["PartNumber"]
1✔
2650
            response["NextPartNumberMarker"] = last_part
1✔
2651

2652
        if part_number_marker:
1✔
2653
            response["PartNumberMarker"] = part_number_marker
1✔
2654
        if s3_multipart.object.checksum_algorithm:
1✔
2655
            response["ChecksumAlgorithm"] = s3_multipart.object.checksum_algorithm
1✔
2656

2657
        return response
1✔
2658

2659
    def list_multipart_uploads(
1✔
2660
        self,
2661
        context: RequestContext,
2662
        bucket: BucketName,
2663
        delimiter: Delimiter = None,
2664
        encoding_type: EncodingType = None,
2665
        key_marker: KeyMarker = None,
2666
        max_uploads: MaxUploads = None,
2667
        prefix: Prefix = None,
2668
        upload_id_marker: UploadIdMarker = None,
2669
        expected_bucket_owner: AccountId = None,
2670
        request_payer: RequestPayer = None,
2671
        **kwargs,
2672
    ) -> ListMultipartUploadsOutput:
2673
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
2674

2675
        common_prefixes = set()
1✔
2676
        count = 0
1✔
2677
        is_truncated = False
1✔
2678
        max_uploads = max_uploads or 1000
1✔
2679
        prefix = prefix or ""
1✔
2680
        delimiter = delimiter or ""
1✔
2681
        if encoding_type:
1✔
2682
            prefix = urlparse.quote(prefix)
1✔
2683
            delimiter = urlparse.quote(delimiter)
1✔
2684
        upload_id_marker_found = False
1✔
2685

2686
        if key_marker and upload_id_marker:
1✔
2687
            multipart = s3_bucket.multiparts.get(upload_id_marker)
1✔
2688
            if multipart:
1✔
2689
                key = (
1✔
2690
                    urlparse.quote(multipart.object.key) if encoding_type else multipart.object.key
2691
                )
2692
            else:
2693
                # set key to None so it fails if the multipart is not Found
2694
                key = None
×
2695

2696
            if key_marker != key:
1✔
2697
                raise InvalidArgument(
1✔
2698
                    "Invalid uploadId marker",
2699
                    ArgumentName="upload-id-marker",
2700
                    ArgumentValue=upload_id_marker,
2701
                )
2702

2703
        uploads = []
1✔
2704
        # sort by key and initiated
2705
        all_multiparts = sorted(
1✔
2706
            s3_bucket.multiparts.values(), key=lambda r: (r.object.key, r.initiated.timestamp())
2707
        )
2708
        last_multipart = all_multiparts[-1] if all_multiparts else None
1✔
2709

2710
        for multipart in all_multiparts:
1✔
2711
            key = urlparse.quote(multipart.object.key) if encoding_type else multipart.object.key
1✔
2712
            # skip all keys that are different than key_marker
2713
            if key_marker:
1✔
2714
                if key < key_marker:
1✔
2715
                    continue
1✔
2716
                elif key == key_marker:
1✔
2717
                    if not upload_id_marker:
1✔
2718
                        continue
1✔
2719
                    # as the keys are ordered by time, once we found the key marker, we can return the next one
2720
                    if multipart.id == upload_id_marker:
1✔
2721
                        upload_id_marker_found = True
1✔
2722
                        continue
1✔
2723
                    elif not upload_id_marker_found:
1✔
2724
                        # as long as we have not passed the version_key_marker, skip the versions
2725
                        continue
1✔
2726

2727
            # Filter for keys that start with prefix
2728
            if prefix and not key.startswith(prefix):
1✔
2729
                continue
1✔
2730

2731
            # see ListObjectsV2 for the logic comments (shared logic here)
2732
            prefix_including_delimiter = None
1✔
2733
            if delimiter and delimiter in (key_no_prefix := key.removeprefix(prefix)):
1✔
2734
                pre_delimiter, _, _ = key_no_prefix.partition(delimiter)
1✔
2735
                prefix_including_delimiter = f"{prefix}{pre_delimiter}{delimiter}"
1✔
2736

2737
                if prefix_including_delimiter in common_prefixes or (
1✔
2738
                    key_marker and key_marker.startswith(prefix_including_delimiter)
2739
                ):
2740
                    continue
1✔
2741

2742
            if prefix_including_delimiter:
1✔
2743
                common_prefixes.add(prefix_including_delimiter)
1✔
2744
            else:
2745
                multipart_upload = MultipartUpload(
1✔
2746
                    UploadId=multipart.id,
2747
                    Key=multipart.object.key,
2748
                    Initiated=multipart.initiated,
2749
                    StorageClass=multipart.object.storage_class,
2750
                    Owner=multipart.initiator,  # TODO: check the difference
2751
                    Initiator=multipart.initiator,
2752
                )
2753
                uploads.append(multipart_upload)
1✔
2754

2755
            count += 1
1✔
2756
            if count >= max_uploads and last_multipart.id != multipart.id:
1✔
2757
                is_truncated = True
1✔
2758
                break
1✔
2759

2760
        common_prefixes = [CommonPrefix(Prefix=prefix) for prefix in sorted(common_prefixes)]
1✔
2761

2762
        response = ListMultipartUploadsOutput(
1✔
2763
            Bucket=bucket,
2764
            IsTruncated=is_truncated,
2765
            MaxUploads=max_uploads or 1000,
2766
            KeyMarker=key_marker or "",
2767
            UploadIdMarker=upload_id_marker or "" if key_marker else "",
2768
            NextKeyMarker="",
2769
            NextUploadIdMarker="",
2770
        )
2771
        if uploads:
1✔
2772
            response["Uploads"] = uploads
1✔
2773
            last_upload = uploads[-1]
1✔
2774
            response["NextKeyMarker"] = last_upload["Key"]
1✔
2775
            response["NextUploadIdMarker"] = last_upload["UploadId"]
1✔
2776
        if delimiter:
1✔
2777
            response["Delimiter"] = delimiter
1✔
2778
        if prefix:
1✔
2779
            response["Prefix"] = prefix
1✔
2780
        if encoding_type:
1✔
2781
            response["EncodingType"] = EncodingType.url
1✔
2782
        if common_prefixes:
1✔
2783
            response["CommonPrefixes"] = common_prefixes
1✔
2784

2785
        return response
1✔
2786

2787
    def put_bucket_versioning(
1✔
2788
        self,
2789
        context: RequestContext,
2790
        bucket: BucketName,
2791
        versioning_configuration: VersioningConfiguration,
2792
        content_md5: ContentMD5 = None,
2793
        checksum_algorithm: ChecksumAlgorithm = None,
2794
        mfa: MFA = None,
2795
        expected_bucket_owner: AccountId = None,
2796
        **kwargs,
2797
    ) -> None:
2798
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
2799
        if not (versioning_status := versioning_configuration.get("Status")):
1✔
2800
            raise CommonServiceException(
1✔
2801
                code="IllegalVersioningConfigurationException",
2802
                message="The Versioning element must be specified",
2803
            )
2804

2805
        if versioning_status not in ("Enabled", "Suspended"):
1✔
2806
            raise MalformedXML()
1✔
2807

2808
        if s3_bucket.object_lock_enabled and versioning_status == "Suspended":
1✔
2809
            raise InvalidBucketState(
1✔
2810
                "An Object Lock configuration is present on this bucket, so the versioning state cannot be changed."
2811
            )
2812

2813
        if not s3_bucket.versioning_status:
1✔
2814
            s3_bucket.objects = VersionedKeyStore.from_key_store(s3_bucket.objects)
1✔
2815

2816
        s3_bucket.versioning_status = versioning_status
1✔
2817

2818
    def get_bucket_versioning(
1✔
2819
        self,
2820
        context: RequestContext,
2821
        bucket: BucketName,
2822
        expected_bucket_owner: AccountId = None,
2823
        **kwargs,
2824
    ) -> GetBucketVersioningOutput:
2825
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
2826

2827
        if not s3_bucket.versioning_status:
1✔
2828
            return GetBucketVersioningOutput()
1✔
2829

2830
        return GetBucketVersioningOutput(Status=s3_bucket.versioning_status)
1✔
2831

2832
    def get_bucket_encryption(
1✔
2833
        self,
2834
        context: RequestContext,
2835
        bucket: BucketName,
2836
        expected_bucket_owner: AccountId = None,
2837
        **kwargs,
2838
    ) -> GetBucketEncryptionOutput:
2839
        # AWS now encrypts bucket by default with AES256, see:
2840
        # https://docs.aws.amazon.com/AmazonS3/latest/userguide/default-bucket-encryption.html
2841
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
2842

2843
        if not s3_bucket.encryption_rule:
1✔
2844
            return GetBucketEncryptionOutput()
×
2845

2846
        return GetBucketEncryptionOutput(
1✔
2847
            ServerSideEncryptionConfiguration={"Rules": [s3_bucket.encryption_rule]}
2848
        )
2849

2850
    def put_bucket_encryption(
1✔
2851
        self,
2852
        context: RequestContext,
2853
        bucket: BucketName,
2854
        server_side_encryption_configuration: ServerSideEncryptionConfiguration,
2855
        content_md5: ContentMD5 = None,
2856
        checksum_algorithm: ChecksumAlgorithm = None,
2857
        expected_bucket_owner: AccountId = None,
2858
        **kwargs,
2859
    ) -> None:
2860
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
2861

2862
        if not (rules := server_side_encryption_configuration.get("Rules")):
1✔
2863
            raise MalformedXML()
1✔
2864

2865
        if len(rules) != 1 or not (
1✔
2866
            encryption := rules[0].get("ApplyServerSideEncryptionByDefault")
2867
        ):
2868
            raise MalformedXML()
1✔
2869

2870
        if not (sse_algorithm := encryption.get("SSEAlgorithm")):
1✔
2871
            raise MalformedXML()
×
2872

2873
        if sse_algorithm not in SSE_ALGORITHMS:
1✔
2874
            raise MalformedXML()
×
2875

2876
        if sse_algorithm != ServerSideEncryption.aws_kms and "KMSMasterKeyID" in encryption:
1✔
2877
            raise InvalidArgument(
1✔
2878
                "a KMSMasterKeyID is not applicable if the default sse algorithm is not aws:kms or aws:kms:dsse",
2879
                ArgumentName="ApplyServerSideEncryptionByDefault",
2880
            )
2881
        # elif master_kms_key := encryption.get("KMSMasterKeyID"):
2882
        # TODO: validate KMS key? not currently done in moto
2883
        # You can pass either the KeyId or the KeyArn. If cross-account, it has to be the ARN.
2884
        # It's always saved as the ARN in the bucket configuration.
2885
        # kms_key_arn = get_kms_key_arn(master_kms_key, s3_bucket.bucket_account_id)
2886
        # encryption["KMSMasterKeyID"] = master_kms_key
2887

2888
        s3_bucket.encryption_rule = rules[0]
1✔
2889

2890
    def delete_bucket_encryption(
1✔
2891
        self,
2892
        context: RequestContext,
2893
        bucket: BucketName,
2894
        expected_bucket_owner: AccountId = None,
2895
        **kwargs,
2896
    ) -> None:
2897
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
2898

2899
        s3_bucket.encryption_rule = None
1✔
2900

2901
    def put_bucket_notification_configuration(
1✔
2902
        self,
2903
        context: RequestContext,
2904
        bucket: BucketName,
2905
        notification_configuration: NotificationConfiguration,
2906
        expected_bucket_owner: AccountId = None,
2907
        skip_destination_validation: SkipValidation = None,
2908
        **kwargs,
2909
    ) -> None:
2910
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
2911

2912
        self._verify_notification_configuration(
1✔
2913
            notification_configuration, skip_destination_validation, context, bucket
2914
        )
2915
        s3_bucket.notification_configuration = notification_configuration
1✔
2916

2917
    def get_bucket_notification_configuration(
1✔
2918
        self,
2919
        context: RequestContext,
2920
        bucket: BucketName,
2921
        expected_bucket_owner: AccountId = None,
2922
        **kwargs,
2923
    ) -> NotificationConfiguration:
2924
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
2925

2926
        return s3_bucket.notification_configuration or NotificationConfiguration()
1✔
2927

2928
    def put_bucket_tagging(
1✔
2929
        self,
2930
        context: RequestContext,
2931
        bucket: BucketName,
2932
        tagging: Tagging,
2933
        content_md5: ContentMD5 = None,
2934
        checksum_algorithm: ChecksumAlgorithm = None,
2935
        expected_bucket_owner: AccountId = None,
2936
        **kwargs,
2937
    ) -> None:
2938
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
2939

2940
        if "TagSet" not in tagging:
1✔
2941
            raise MalformedXML()
×
2942

2943
        validate_tag_set(tagging["TagSet"], type_set="bucket")
1✔
2944

2945
        # remove the previous tags before setting the new ones, it overwrites the whole TagSet
2946
        store.TAGS.tags.pop(s3_bucket.bucket_arn, None)
1✔
2947
        store.TAGS.tag_resource(s3_bucket.bucket_arn, tags=tagging["TagSet"])
1✔
2948

2949
    def get_bucket_tagging(
1✔
2950
        self,
2951
        context: RequestContext,
2952
        bucket: BucketName,
2953
        expected_bucket_owner: AccountId = None,
2954
        **kwargs,
2955
    ) -> GetBucketTaggingOutput:
2956
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
2957
        tag_set = store.TAGS.list_tags_for_resource(s3_bucket.bucket_arn, root_name="Tags")["Tags"]
1✔
2958
        if not tag_set:
1✔
2959
            raise NoSuchTagSet(
1✔
2960
                "The TagSet does not exist",
2961
                BucketName=bucket,
2962
            )
2963

2964
        return GetBucketTaggingOutput(TagSet=tag_set)
1✔
2965

2966
    def delete_bucket_tagging(
1✔
2967
        self,
2968
        context: RequestContext,
2969
        bucket: BucketName,
2970
        expected_bucket_owner: AccountId = None,
2971
        **kwargs,
2972
    ) -> None:
2973
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
2974

2975
        store.TAGS.tags.pop(s3_bucket.bucket_arn, None)
1✔
2976

2977
    def put_object_tagging(
1✔
2978
        self,
2979
        context: RequestContext,
2980
        bucket: BucketName,
2981
        key: ObjectKey,
2982
        tagging: Tagging,
2983
        version_id: ObjectVersionId = None,
2984
        content_md5: ContentMD5 = None,
2985
        checksum_algorithm: ChecksumAlgorithm = None,
2986
        expected_bucket_owner: AccountId = None,
2987
        request_payer: RequestPayer = None,
2988
        **kwargs,
2989
    ) -> PutObjectTaggingOutput:
2990
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
2991

2992
        s3_object = s3_bucket.get_object(key=key, version_id=version_id, http_method="PUT")
1✔
2993

2994
        if "TagSet" not in tagging:
1✔
2995
            raise MalformedXML()
×
2996

2997
        validate_tag_set(tagging["TagSet"], type_set="object")
1✔
2998

2999
        key_id = get_unique_key_id(bucket, key, s3_object.version_id)
1✔
3000
        # remove the previous tags before setting the new ones, it overwrites the whole TagSet
3001
        store.TAGS.tags.pop(key_id, None)
1✔
3002
        store.TAGS.tag_resource(key_id, tags=tagging["TagSet"])
1✔
3003
        response = PutObjectTaggingOutput()
1✔
3004
        if s3_object.version_id:
1✔
3005
            response["VersionId"] = s3_object.version_id
1✔
3006

3007
        self._notify(context, s3_bucket=s3_bucket, s3_object=s3_object)
1✔
3008

3009
        return response
1✔
3010

3011
    def get_object_tagging(
1✔
3012
        self,
3013
        context: RequestContext,
3014
        bucket: BucketName,
3015
        key: ObjectKey,
3016
        version_id: ObjectVersionId = None,
3017
        expected_bucket_owner: AccountId = None,
3018
        request_payer: RequestPayer = None,
3019
        **kwargs,
3020
    ) -> GetObjectTaggingOutput:
3021
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
3022

3023
        try:
1✔
3024
            s3_object = s3_bucket.get_object(key=key, version_id=version_id)
1✔
3025
        except NoSuchKey as e:
1✔
3026
            # it seems GetObjectTagging does not work like all other operations, so we need to raise a different
3027
            # exception. As we already need to catch it because of the format of the Key, it is not worth to modify the
3028
            # `S3Bucket.get_object` signature for one operation.
3029
            if s3_bucket.versioning_status and (
1✔
3030
                s3_object_version := s3_bucket.objects.get(key, version_id)
3031
            ):
3032
                raise MethodNotAllowed(
1✔
3033
                    "The specified method is not allowed against this resource.",
3034
                    Method="GET",
3035
                    ResourceType="DeleteMarker",
3036
                    DeleteMarker=True,
3037
                    Allow="DELETE",
3038
                    VersionId=s3_object_version.version_id,
3039
                )
3040

3041
            # There a weird AWS validated bug in S3: the returned key contains the bucket name as well
3042
            # follow AWS on this one
3043
            e.Key = f"{bucket}/{key}"
1✔
3044
            raise e
1✔
3045

3046
        tag_set = store.TAGS.list_tags_for_resource(
1✔
3047
            get_unique_key_id(bucket, key, s3_object.version_id)
3048
        )["Tags"]
3049
        response = GetObjectTaggingOutput(TagSet=tag_set)
1✔
3050
        if s3_object.version_id:
1✔
3051
            response["VersionId"] = s3_object.version_id
1✔
3052

3053
        return response
1✔
3054

3055
    def delete_object_tagging(
1✔
3056
        self,
3057
        context: RequestContext,
3058
        bucket: BucketName,
3059
        key: ObjectKey,
3060
        version_id: ObjectVersionId = None,
3061
        expected_bucket_owner: AccountId = None,
3062
        **kwargs,
3063
    ) -> DeleteObjectTaggingOutput:
3064
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
3065

3066
        s3_object = s3_bucket.get_object(key=key, version_id=version_id, http_method="DELETE")
1✔
3067

3068
        store.TAGS.tags.pop(get_unique_key_id(bucket, key, version_id), None)
1✔
3069
        response = DeleteObjectTaggingOutput()
1✔
3070
        if s3_object.version_id:
1✔
3071
            response["VersionId"] = s3_object.version_id
×
3072

3073
        self._notify(context, s3_bucket=s3_bucket, s3_object=s3_object)
1✔
3074

3075
        return response
1✔
3076

3077
    def put_bucket_cors(
1✔
3078
        self,
3079
        context: RequestContext,
3080
        bucket: BucketName,
3081
        cors_configuration: CORSConfiguration,
3082
        content_md5: ContentMD5 = None,
3083
        checksum_algorithm: ChecksumAlgorithm = None,
3084
        expected_bucket_owner: AccountId = None,
3085
        **kwargs,
3086
    ) -> None:
3087
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
3088
        validate_cors_configuration(cors_configuration)
1✔
3089
        s3_bucket.cors_rules = cors_configuration
1✔
3090
        self._cors_handler.invalidate_cache()
1✔
3091

3092
    def get_bucket_cors(
1✔
3093
        self,
3094
        context: RequestContext,
3095
        bucket: BucketName,
3096
        expected_bucket_owner: AccountId = None,
3097
        **kwargs,
3098
    ) -> GetBucketCorsOutput:
3099
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
3100

3101
        if not s3_bucket.cors_rules:
1✔
3102
            raise NoSuchCORSConfiguration(
1✔
3103
                "The CORS configuration does not exist",
3104
                BucketName=bucket,
3105
            )
3106
        return GetBucketCorsOutput(CORSRules=s3_bucket.cors_rules["CORSRules"])
1✔
3107

3108
    def delete_bucket_cors(
1✔
3109
        self,
3110
        context: RequestContext,
3111
        bucket: BucketName,
3112
        expected_bucket_owner: AccountId = None,
3113
        **kwargs,
3114
    ) -> None:
3115
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
3116

3117
        if s3_bucket.cors_rules:
1✔
3118
            self._cors_handler.invalidate_cache()
1✔
3119
            s3_bucket.cors_rules = None
1✔
3120

3121
    def get_bucket_lifecycle_configuration(
1✔
3122
        self,
3123
        context: RequestContext,
3124
        bucket: BucketName,
3125
        expected_bucket_owner: AccountId = None,
3126
        **kwargs,
3127
    ) -> GetBucketLifecycleConfigurationOutput:
3128
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
3129

3130
        if not s3_bucket.lifecycle_rules:
1✔
3131
            raise NoSuchLifecycleConfiguration(
1✔
3132
                "The lifecycle configuration does not exist",
3133
                BucketName=bucket,
3134
            )
3135

3136
        return GetBucketLifecycleConfigurationOutput(Rules=s3_bucket.lifecycle_rules)
1✔
3137

3138
    def put_bucket_lifecycle_configuration(
1✔
3139
        self,
3140
        context: RequestContext,
3141
        bucket: BucketName,
3142
        checksum_algorithm: ChecksumAlgorithm = None,
3143
        lifecycle_configuration: BucketLifecycleConfiguration = None,
3144
        expected_bucket_owner: AccountId = None,
3145
        transition_default_minimum_object_size: TransitionDefaultMinimumObjectSize = None,
3146
        **kwargs,
3147
    ) -> PutBucketLifecycleConfigurationOutput:
3148
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
3149

3150
        validate_lifecycle_configuration(lifecycle_configuration)
1✔
3151
        # TODO: we either apply the lifecycle to existing objects when we set the new rules, or we need to apply them
3152
        #  everytime we get/head an object
3153
        # for now, we keep a cache and get it everytime we fetch an object
3154
        s3_bucket.lifecycle_rules = lifecycle_configuration["Rules"]
1✔
3155
        self._expiration_cache[bucket].clear()
1✔
3156
        return PutBucketLifecycleConfigurationOutput(
1✔
3157
            TransitionDefaultMinimumObjectSize=transition_default_minimum_object_size
3158
        )
3159

3160
    def delete_bucket_lifecycle(
1✔
3161
        self,
3162
        context: RequestContext,
3163
        bucket: BucketName,
3164
        expected_bucket_owner: AccountId = None,
3165
        **kwargs,
3166
    ) -> None:
3167
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
3168

3169
        s3_bucket.lifecycle_rules = None
1✔
3170
        self._expiration_cache[bucket].clear()
1✔
3171

3172
    def put_bucket_analytics_configuration(
1✔
3173
        self,
3174
        context: RequestContext,
3175
        bucket: BucketName,
3176
        id: AnalyticsId,
3177
        analytics_configuration: AnalyticsConfiguration,
3178
        expected_bucket_owner: AccountId = None,
3179
        **kwargs,
3180
    ) -> None:
3181
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
3182

3183
        validate_bucket_analytics_configuration(
1✔
3184
            id=id, analytics_configuration=analytics_configuration
3185
        )
3186

3187
        s3_bucket.analytics_configurations[id] = analytics_configuration
1✔
3188

3189
    def get_bucket_analytics_configuration(
1✔
3190
        self,
3191
        context: RequestContext,
3192
        bucket: BucketName,
3193
        id: AnalyticsId,
3194
        expected_bucket_owner: AccountId = None,
3195
        **kwargs,
3196
    ) -> GetBucketAnalyticsConfigurationOutput:
3197
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
3198

3199
        if not (analytic_config := s3_bucket.analytics_configurations.get(id)):
1✔
3200
            raise NoSuchConfiguration("The specified configuration does not exist.")
1✔
3201

3202
        return GetBucketAnalyticsConfigurationOutput(AnalyticsConfiguration=analytic_config)
1✔
3203

3204
    def list_bucket_analytics_configurations(
1✔
3205
        self,
3206
        context: RequestContext,
3207
        bucket: BucketName,
3208
        continuation_token: Token = None,
3209
        expected_bucket_owner: AccountId = None,
3210
        **kwargs,
3211
    ) -> ListBucketAnalyticsConfigurationsOutput:
3212
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
3213

3214
        return ListBucketAnalyticsConfigurationsOutput(
1✔
3215
            IsTruncated=False,
3216
            AnalyticsConfigurationList=sorted(
3217
                s3_bucket.analytics_configurations.values(),
3218
                key=itemgetter("Id"),
3219
            ),
3220
        )
3221

3222
    def delete_bucket_analytics_configuration(
1✔
3223
        self,
3224
        context: RequestContext,
3225
        bucket: BucketName,
3226
        id: AnalyticsId,
3227
        expected_bucket_owner: AccountId = None,
3228
        **kwargs,
3229
    ) -> None:
3230
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
3231

3232
        if not s3_bucket.analytics_configurations.pop(id, None):
1✔
3233
            raise NoSuchConfiguration("The specified configuration does not exist.")
1✔
3234

3235
    def put_bucket_intelligent_tiering_configuration(
1✔
3236
        self,
3237
        context: RequestContext,
3238
        bucket: BucketName,
3239
        id: IntelligentTieringId,
3240
        intelligent_tiering_configuration: IntelligentTieringConfiguration,
3241
        **kwargs,
3242
    ) -> None:
3243
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
3244

3245
        validate_bucket_intelligent_tiering_configuration(id, intelligent_tiering_configuration)
1✔
3246

3247
        s3_bucket.intelligent_tiering_configurations[id] = intelligent_tiering_configuration
1✔
3248

3249
    def get_bucket_intelligent_tiering_configuration(
1✔
3250
        self, context: RequestContext, bucket: BucketName, id: IntelligentTieringId, **kwargs
3251
    ) -> GetBucketIntelligentTieringConfigurationOutput:
3252
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
3253

3254
        if not (itier_config := s3_bucket.intelligent_tiering_configurations.get(id)):
1✔
3255
            raise NoSuchConfiguration("The specified configuration does not exist.")
×
3256

3257
        return GetBucketIntelligentTieringConfigurationOutput(
1✔
3258
            IntelligentTieringConfiguration=itier_config
3259
        )
3260

3261
    def delete_bucket_intelligent_tiering_configuration(
1✔
3262
        self, context: RequestContext, bucket: BucketName, id: IntelligentTieringId, **kwargs
3263
    ) -> None:
3264
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
3265

3266
        if not s3_bucket.intelligent_tiering_configurations.pop(id, None):
1✔
3267
            raise NoSuchConfiguration("The specified configuration does not exist.")
1✔
3268

3269
    def list_bucket_intelligent_tiering_configurations(
1✔
3270
        self,
3271
        context: RequestContext,
3272
        bucket: BucketName,
3273
        continuation_token: Token = None,
3274
        **kwargs,
3275
    ) -> ListBucketIntelligentTieringConfigurationsOutput:
3276
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
3277

3278
        return ListBucketIntelligentTieringConfigurationsOutput(
1✔
3279
            IsTruncated=False,
3280
            IntelligentTieringConfigurationList=sorted(
3281
                s3_bucket.intelligent_tiering_configurations.values(),
3282
                key=itemgetter("Id"),
3283
            ),
3284
        )
3285

3286
    def put_bucket_inventory_configuration(
1✔
3287
        self,
3288
        context: RequestContext,
3289
        bucket: BucketName,
3290
        id: InventoryId,
3291
        inventory_configuration: InventoryConfiguration,
3292
        expected_bucket_owner: AccountId = None,
3293
        **kwargs,
3294
    ) -> None:
3295
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
3296

3297
        validate_inventory_configuration(
1✔
3298
            config_id=id, inventory_configuration=inventory_configuration
3299
        )
3300
        s3_bucket.inventory_configurations[id] = inventory_configuration
1✔
3301

3302
    def get_bucket_inventory_configuration(
1✔
3303
        self,
3304
        context: RequestContext,
3305
        bucket: BucketName,
3306
        id: InventoryId,
3307
        expected_bucket_owner: AccountId = None,
3308
        **kwargs,
3309
    ) -> GetBucketInventoryConfigurationOutput:
3310
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
3311

3312
        if not (inv_config := s3_bucket.inventory_configurations.get(id)):
1✔
3313
            raise NoSuchConfiguration("The specified configuration does not exist.")
1✔
3314
        return GetBucketInventoryConfigurationOutput(InventoryConfiguration=inv_config)
1✔
3315

3316
    def list_bucket_inventory_configurations(
1✔
3317
        self,
3318
        context: RequestContext,
3319
        bucket: BucketName,
3320
        continuation_token: Token = None,
3321
        expected_bucket_owner: AccountId = None,
3322
        **kwargs,
3323
    ) -> ListBucketInventoryConfigurationsOutput:
3324
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
3325

3326
        return ListBucketInventoryConfigurationsOutput(
1✔
3327
            IsTruncated=False,
3328
            InventoryConfigurationList=sorted(
3329
                s3_bucket.inventory_configurations.values(), key=itemgetter("Id")
3330
            ),
3331
        )
3332

3333
    def delete_bucket_inventory_configuration(
1✔
3334
        self,
3335
        context: RequestContext,
3336
        bucket: BucketName,
3337
        id: InventoryId,
3338
        expected_bucket_owner: AccountId = None,
3339
        **kwargs,
3340
    ) -> None:
3341
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
3342

3343
        if not s3_bucket.inventory_configurations.pop(id, None):
1✔
3344
            raise NoSuchConfiguration("The specified configuration does not exist.")
×
3345

3346
    def get_bucket_website(
1✔
3347
        self,
3348
        context: RequestContext,
3349
        bucket: BucketName,
3350
        expected_bucket_owner: AccountId = None,
3351
        **kwargs,
3352
    ) -> GetBucketWebsiteOutput:
3353
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
3354

3355
        if not s3_bucket.website_configuration:
1✔
3356
            raise NoSuchWebsiteConfiguration(
1✔
3357
                "The specified bucket does not have a website configuration",
3358
                BucketName=bucket,
3359
            )
3360
        return s3_bucket.website_configuration
1✔
3361

3362
    def put_bucket_website(
1✔
3363
        self,
3364
        context: RequestContext,
3365
        bucket: BucketName,
3366
        website_configuration: WebsiteConfiguration,
3367
        content_md5: ContentMD5 = None,
3368
        checksum_algorithm: ChecksumAlgorithm = None,
3369
        expected_bucket_owner: AccountId = None,
3370
        **kwargs,
3371
    ) -> None:
3372
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
3373

3374
        validate_website_configuration(website_configuration)
1✔
3375
        s3_bucket.website_configuration = website_configuration
1✔
3376

3377
    def delete_bucket_website(
1✔
3378
        self,
3379
        context: RequestContext,
3380
        bucket: BucketName,
3381
        expected_bucket_owner: AccountId = None,
3382
        **kwargs,
3383
    ) -> None:
3384
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
3385
        # does not raise error if the bucket did not have a config, will simply return
3386
        s3_bucket.website_configuration = None
1✔
3387

3388
    def get_object_lock_configuration(
1✔
3389
        self,
3390
        context: RequestContext,
3391
        bucket: BucketName,
3392
        expected_bucket_owner: AccountId = None,
3393
        **kwargs,
3394
    ) -> GetObjectLockConfigurationOutput:
3395
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
3396
        if not s3_bucket.object_lock_enabled:
1✔
3397
            raise ObjectLockConfigurationNotFoundError(
1✔
3398
                "Object Lock configuration does not exist for this bucket",
3399
                BucketName=bucket,
3400
            )
3401

3402
        response = GetObjectLockConfigurationOutput(
1✔
3403
            ObjectLockConfiguration=ObjectLockConfiguration(
3404
                ObjectLockEnabled=ObjectLockEnabled.Enabled
3405
            )
3406
        )
3407
        if s3_bucket.object_lock_default_retention:
1✔
3408
            response["ObjectLockConfiguration"]["Rule"] = {
1✔
3409
                "DefaultRetention": s3_bucket.object_lock_default_retention
3410
            }
3411

3412
        return response
1✔
3413

3414
    def put_object_lock_configuration(
1✔
3415
        self,
3416
        context: RequestContext,
3417
        bucket: BucketName,
3418
        object_lock_configuration: ObjectLockConfiguration = None,
3419
        request_payer: RequestPayer = None,
3420
        token: ObjectLockToken = None,
3421
        content_md5: ContentMD5 = None,
3422
        checksum_algorithm: ChecksumAlgorithm = None,
3423
        expected_bucket_owner: AccountId = None,
3424
        **kwargs,
3425
    ) -> PutObjectLockConfigurationOutput:
3426
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
3427
        if s3_bucket.versioning_status != "Enabled":
1✔
3428
            raise InvalidBucketState(
1✔
3429
                "Versioning must be 'Enabled' on the bucket to apply a Object Lock configuration"
3430
            )
3431

3432
        if (
1✔
3433
            not object_lock_configuration
3434
            or object_lock_configuration.get("ObjectLockEnabled") != "Enabled"
3435
        ):
3436
            raise MalformedXML()
1✔
3437

3438
        if "Rule" not in object_lock_configuration:
1✔
3439
            s3_bucket.object_lock_default_retention = None
1✔
3440
            if not s3_bucket.object_lock_enabled:
1✔
3441
                s3_bucket.object_lock_enabled = True
1✔
3442

3443
            return PutObjectLockConfigurationOutput()
1✔
3444
        elif not (rule := object_lock_configuration["Rule"]) or not (
1✔
3445
            default_retention := rule.get("DefaultRetention")
3446
        ):
3447
            raise MalformedXML()
1✔
3448

3449
        if "Mode" not in default_retention or (
1✔
3450
            ("Days" in default_retention and "Years" in default_retention)
3451
            or ("Days" not in default_retention and "Years" not in default_retention)
3452
        ):
3453
            raise MalformedXML()
1✔
3454

3455
        s3_bucket.object_lock_default_retention = default_retention
1✔
3456
        if not s3_bucket.object_lock_enabled:
1✔
3457
            s3_bucket.object_lock_enabled = True
×
3458

3459
        return PutObjectLockConfigurationOutput()
1✔
3460

3461
    def get_object_legal_hold(
1✔
3462
        self,
3463
        context: RequestContext,
3464
        bucket: BucketName,
3465
        key: ObjectKey,
3466
        version_id: ObjectVersionId = None,
3467
        request_payer: RequestPayer = None,
3468
        expected_bucket_owner: AccountId = None,
3469
        **kwargs,
3470
    ) -> GetObjectLegalHoldOutput:
3471
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
3472
        if not s3_bucket.object_lock_enabled:
1✔
3473
            raise InvalidRequest("Bucket is missing Object Lock Configuration")
1✔
3474

3475
        s3_object = s3_bucket.get_object(
1✔
3476
            key=key,
3477
            version_id=version_id,
3478
            http_method="GET",
3479
        )
3480
        if not s3_object.lock_legal_status:
1✔
3481
            raise NoSuchObjectLockConfiguration(
1✔
3482
                "The specified object does not have a ObjectLock configuration"
3483
            )
3484

3485
        return GetObjectLegalHoldOutput(
1✔
3486
            LegalHold=ObjectLockLegalHold(Status=s3_object.lock_legal_status)
3487
        )
3488

3489
    def put_object_legal_hold(
1✔
3490
        self,
3491
        context: RequestContext,
3492
        bucket: BucketName,
3493
        key: ObjectKey,
3494
        legal_hold: ObjectLockLegalHold = None,
3495
        request_payer: RequestPayer = None,
3496
        version_id: ObjectVersionId = None,
3497
        content_md5: ContentMD5 = None,
3498
        checksum_algorithm: ChecksumAlgorithm = None,
3499
        expected_bucket_owner: AccountId = None,
3500
        **kwargs,
3501
    ) -> PutObjectLegalHoldOutput:
3502
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
3503

3504
        if not legal_hold:
1✔
3505
            raise MalformedXML()
1✔
3506

3507
        if not s3_bucket.object_lock_enabled:
1✔
3508
            raise InvalidRequest("Bucket is missing Object Lock Configuration")
1✔
3509

3510
        s3_object = s3_bucket.get_object(
1✔
3511
            key=key,
3512
            version_id=version_id,
3513
            http_method="PUT",
3514
        )
3515
        # TODO: check casing
3516
        if not (status := legal_hold.get("Status")) or status not in ("ON", "OFF"):
1✔
3517
            raise MalformedXML()
×
3518

3519
        s3_object.lock_legal_status = status
1✔
3520

3521
        # TODO: return RequestCharged
3522
        return PutObjectRetentionOutput()
1✔
3523

3524
    def get_object_retention(
1✔
3525
        self,
3526
        context: RequestContext,
3527
        bucket: BucketName,
3528
        key: ObjectKey,
3529
        version_id: ObjectVersionId = None,
3530
        request_payer: RequestPayer = None,
3531
        expected_bucket_owner: AccountId = None,
3532
        **kwargs,
3533
    ) -> GetObjectRetentionOutput:
3534
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
3535
        if not s3_bucket.object_lock_enabled:
1✔
3536
            raise InvalidRequest("Bucket is missing Object Lock Configuration")
1✔
3537

3538
        s3_object = s3_bucket.get_object(
1✔
3539
            key=key,
3540
            version_id=version_id,
3541
            http_method="GET",
3542
        )
3543
        if not s3_object.lock_mode:
1✔
3544
            raise NoSuchObjectLockConfiguration(
1✔
3545
                "The specified object does not have a ObjectLock configuration"
3546
            )
3547

3548
        return GetObjectRetentionOutput(
1✔
3549
            Retention=ObjectLockRetention(
3550
                Mode=s3_object.lock_mode,
3551
                RetainUntilDate=s3_object.lock_until,
3552
            )
3553
        )
3554

3555
    def put_object_retention(
1✔
3556
        self,
3557
        context: RequestContext,
3558
        bucket: BucketName,
3559
        key: ObjectKey,
3560
        retention: ObjectLockRetention = None,
3561
        request_payer: RequestPayer = None,
3562
        version_id: ObjectVersionId = None,
3563
        bypass_governance_retention: BypassGovernanceRetention = None,
3564
        content_md5: ContentMD5 = None,
3565
        checksum_algorithm: ChecksumAlgorithm = None,
3566
        expected_bucket_owner: AccountId = None,
3567
        **kwargs,
3568
    ) -> PutObjectRetentionOutput:
3569
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
3570
        if not s3_bucket.object_lock_enabled:
1✔
3571
            raise InvalidRequest("Bucket is missing Object Lock Configuration")
1✔
3572

3573
        s3_object = s3_bucket.get_object(
1✔
3574
            key=key,
3575
            version_id=version_id,
3576
            http_method="PUT",
3577
        )
3578

3579
        if retention and not validate_dict_fields(
1✔
3580
            retention, required_fields={"Mode", "RetainUntilDate"}
3581
        ):
3582
            raise MalformedXML()
1✔
3583

3584
        if retention and retention["RetainUntilDate"] < datetime.datetime.now(datetime.UTC):
1✔
3585
            # weirdly, this date is format as following: Tue Dec 31 16:00:00 PST 2019
3586
            # it contains the timezone as PST, even if you target a bucket in Europe or Asia
3587
            pst_datetime = retention["RetainUntilDate"].astimezone(tz=ZoneInfo("US/Pacific"))
1✔
3588
            raise InvalidArgument(
1✔
3589
                "The retain until date must be in the future!",
3590
                ArgumentName="RetainUntilDate",
3591
                ArgumentValue=pst_datetime.strftime("%a %b %d %H:%M:%S %Z %Y"),
3592
            )
3593

3594
        if (
1✔
3595
            not retention
3596
            or (s3_object.lock_until and s3_object.lock_until > retention["RetainUntilDate"])
3597
        ) and not (
3598
            bypass_governance_retention and s3_object.lock_mode == ObjectLockMode.GOVERNANCE
3599
        ):
3600
            raise AccessDenied("Access Denied because object protected by object lock.")
1✔
3601

3602
        s3_object.lock_mode = retention["Mode"] if retention else None
1✔
3603
        s3_object.lock_until = retention["RetainUntilDate"] if retention else None
1✔
3604

3605
        # TODO: return RequestCharged
3606
        return PutObjectRetentionOutput()
1✔
3607

3608
    def put_bucket_request_payment(
1✔
3609
        self,
3610
        context: RequestContext,
3611
        bucket: BucketName,
3612
        request_payment_configuration: RequestPaymentConfiguration,
3613
        content_md5: ContentMD5 = None,
3614
        checksum_algorithm: ChecksumAlgorithm = None,
3615
        expected_bucket_owner: AccountId = None,
3616
        **kwargs,
3617
    ) -> None:
3618
        # TODO: this currently only mock the operation, but its actual effect is not emulated
3619
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
3620

3621
        payer = request_payment_configuration.get("Payer")
1✔
3622
        if payer not in ["Requester", "BucketOwner"]:
1✔
3623
            raise MalformedXML()
1✔
3624

3625
        s3_bucket.payer = payer
1✔
3626

3627
    def get_bucket_request_payment(
1✔
3628
        self,
3629
        context: RequestContext,
3630
        bucket: BucketName,
3631
        expected_bucket_owner: AccountId = None,
3632
        **kwargs,
3633
    ) -> GetBucketRequestPaymentOutput:
3634
        # TODO: this currently only mock the operation, but its actual effect is not emulated
3635
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
3636

3637
        return GetBucketRequestPaymentOutput(Payer=s3_bucket.payer)
1✔
3638

3639
    def get_bucket_ownership_controls(
1✔
3640
        self,
3641
        context: RequestContext,
3642
        bucket: BucketName,
3643
        expected_bucket_owner: AccountId = None,
3644
        **kwargs,
3645
    ) -> GetBucketOwnershipControlsOutput:
3646
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
3647

3648
        if not s3_bucket.object_ownership:
1✔
3649
            raise OwnershipControlsNotFoundError(
1✔
3650
                "The bucket ownership controls were not found",
3651
                BucketName=bucket,
3652
            )
3653

3654
        return GetBucketOwnershipControlsOutput(
1✔
3655
            OwnershipControls={"Rules": [{"ObjectOwnership": s3_bucket.object_ownership}]}
3656
        )
3657

3658
    def put_bucket_ownership_controls(
1✔
3659
        self,
3660
        context: RequestContext,
3661
        bucket: BucketName,
3662
        ownership_controls: OwnershipControls,
3663
        content_md5: ContentMD5 = None,
3664
        expected_bucket_owner: AccountId = None,
3665
        **kwargs,
3666
    ) -> None:
3667
        # TODO: this currently only mock the operation, but its actual effect is not emulated
3668
        #  it for example almost forbid ACL usage when set to BucketOwnerEnforced
3669
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
3670

3671
        if not (rules := ownership_controls.get("Rules")) or len(rules) > 1:
1✔
3672
            raise MalformedXML()
1✔
3673

3674
        rule = rules[0]
1✔
3675
        if (object_ownership := rule.get("ObjectOwnership")) not in OBJECT_OWNERSHIPS:
1✔
3676
            raise MalformedXML()
1✔
3677

3678
        s3_bucket.object_ownership = object_ownership
1✔
3679

3680
    def delete_bucket_ownership_controls(
1✔
3681
        self,
3682
        context: RequestContext,
3683
        bucket: BucketName,
3684
        expected_bucket_owner: AccountId = None,
3685
        **kwargs,
3686
    ) -> None:
3687
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
3688

3689
        s3_bucket.object_ownership = None
1✔
3690

3691
    def get_public_access_block(
1✔
3692
        self,
3693
        context: RequestContext,
3694
        bucket: BucketName,
3695
        expected_bucket_owner: AccountId = None,
3696
        **kwargs,
3697
    ) -> GetPublicAccessBlockOutput:
3698
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
3699

3700
        if not s3_bucket.public_access_block:
1✔
3701
            raise NoSuchPublicAccessBlockConfiguration(
1✔
3702
                "The public access block configuration was not found", BucketName=bucket
3703
            )
3704

3705
        return GetPublicAccessBlockOutput(
1✔
3706
            PublicAccessBlockConfiguration=s3_bucket.public_access_block
3707
        )
3708

3709
    def put_public_access_block(
1✔
3710
        self,
3711
        context: RequestContext,
3712
        bucket: BucketName,
3713
        public_access_block_configuration: PublicAccessBlockConfiguration,
3714
        content_md5: ContentMD5 = None,
3715
        checksum_algorithm: ChecksumAlgorithm = None,
3716
        expected_bucket_owner: AccountId = None,
3717
        **kwargs,
3718
    ) -> None:
3719
        # TODO: this currently only mock the operation, but its actual effect is not emulated
3720
        #  as we do not enforce ACL directly. Also, this should take the most restrictive between S3Control and the
3721
        #  bucket configuration. See s3control
3722
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
3723

3724
        public_access_block_fields = {
1✔
3725
            "BlockPublicAcls",
3726
            "BlockPublicPolicy",
3727
            "IgnorePublicAcls",
3728
            "RestrictPublicBuckets",
3729
        }
3730
        if not validate_dict_fields(
1✔
3731
            public_access_block_configuration,
3732
            required_fields=set(),
3733
            optional_fields=public_access_block_fields,
3734
        ):
3735
            raise MalformedXML()
×
3736

3737
        for field in public_access_block_fields:
1✔
3738
            if public_access_block_configuration.get(field) is None:
1✔
3739
                public_access_block_configuration[field] = False
1✔
3740

3741
        s3_bucket.public_access_block = public_access_block_configuration
1✔
3742

3743
    def delete_public_access_block(
1✔
3744
        self,
3745
        context: RequestContext,
3746
        bucket: BucketName,
3747
        expected_bucket_owner: AccountId = None,
3748
        **kwargs,
3749
    ) -> None:
3750
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
3751

3752
        s3_bucket.public_access_block = None
1✔
3753

3754
    def get_bucket_policy(
1✔
3755
        self,
3756
        context: RequestContext,
3757
        bucket: BucketName,
3758
        expected_bucket_owner: AccountId = None,
3759
        **kwargs,
3760
    ) -> GetBucketPolicyOutput:
3761
        store, s3_bucket = self._get_cross_account_bucket(
1✔
3762
            context, bucket, expected_bucket_owner=expected_bucket_owner
3763
        )
3764
        if not s3_bucket.policy:
1✔
3765
            raise NoSuchBucketPolicy(
1✔
3766
                "The bucket policy does not exist",
3767
                BucketName=bucket,
3768
            )
3769
        return GetBucketPolicyOutput(Policy=s3_bucket.policy)
1✔
3770

3771
    def put_bucket_policy(
1✔
3772
        self,
3773
        context: RequestContext,
3774
        bucket: BucketName,
3775
        policy: Policy,
3776
        content_md5: ContentMD5 = None,
3777
        checksum_algorithm: ChecksumAlgorithm = None,
3778
        confirm_remove_self_bucket_access: ConfirmRemoveSelfBucketAccess = None,
3779
        expected_bucket_owner: AccountId = None,
3780
        **kwargs,
3781
    ) -> None:
3782
        store, s3_bucket = self._get_cross_account_bucket(
1✔
3783
            context, bucket, expected_bucket_owner=expected_bucket_owner
3784
        )
3785

3786
        if not policy or policy[0] != "{":
1✔
3787
            raise MalformedPolicy("Policies must be valid JSON and the first byte must be '{'")
1✔
3788
        try:
1✔
3789
            json_policy = json.loads(policy)
1✔
3790
            if not json_policy:
1✔
3791
                # TODO: add more validation around the policy?
3792
                raise MalformedPolicy("Missing required field Statement")
1✔
3793
        except ValueError:
1✔
3794
            raise MalformedPolicy("Policies must be valid JSON and the first byte must be '{'")
×
3795

3796
        s3_bucket.policy = policy
1✔
3797

3798
    def delete_bucket_policy(
1✔
3799
        self,
3800
        context: RequestContext,
3801
        bucket: BucketName,
3802
        expected_bucket_owner: AccountId = None,
3803
        **kwargs,
3804
    ) -> None:
3805
        store, s3_bucket = self._get_cross_account_bucket(
1✔
3806
            context, bucket, expected_bucket_owner=expected_bucket_owner
3807
        )
3808

3809
        s3_bucket.policy = None
1✔
3810

3811
    def get_bucket_accelerate_configuration(
1✔
3812
        self,
3813
        context: RequestContext,
3814
        bucket: BucketName,
3815
        expected_bucket_owner: AccountId = None,
3816
        request_payer: RequestPayer = None,
3817
        **kwargs,
3818
    ) -> GetBucketAccelerateConfigurationOutput:
3819
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
3820

3821
        response = GetBucketAccelerateConfigurationOutput()
1✔
3822
        if s3_bucket.accelerate_status:
1✔
3823
            response["Status"] = s3_bucket.accelerate_status
1✔
3824

3825
        return response
1✔
3826

3827
    def put_bucket_accelerate_configuration(
1✔
3828
        self,
3829
        context: RequestContext,
3830
        bucket: BucketName,
3831
        accelerate_configuration: AccelerateConfiguration,
3832
        expected_bucket_owner: AccountId = None,
3833
        checksum_algorithm: ChecksumAlgorithm = None,
3834
        **kwargs,
3835
    ) -> None:
3836
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
3837

3838
        if "." in bucket:
1✔
3839
            raise InvalidRequest(
1✔
3840
                "S3 Transfer Acceleration is not supported for buckets with periods (.) in their names"
3841
            )
3842

3843
        if not (status := accelerate_configuration.get("Status")) or status not in (
1✔
3844
            "Enabled",
3845
            "Suspended",
3846
        ):
3847
            raise MalformedXML()
1✔
3848

3849
        s3_bucket.accelerate_status = status
1✔
3850

3851
    def put_bucket_logging(
1✔
3852
        self,
3853
        context: RequestContext,
3854
        bucket: BucketName,
3855
        bucket_logging_status: BucketLoggingStatus,
3856
        content_md5: ContentMD5 = None,
3857
        checksum_algorithm: ChecksumAlgorithm = None,
3858
        expected_bucket_owner: AccountId = None,
3859
        **kwargs,
3860
    ) -> None:
3861
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
3862

3863
        if not (logging_config := bucket_logging_status.get("LoggingEnabled")):
1✔
3864
            s3_bucket.logging = {}
1✔
3865
            return
1✔
3866

3867
        # the target bucket must be in the same account
3868
        if not (target_bucket_name := logging_config.get("TargetBucket")):
1✔
3869
            raise MalformedXML()
×
3870

3871
        if not logging_config.get("TargetPrefix"):
1✔
3872
            logging_config["TargetPrefix"] = ""
×
3873

3874
        # TODO: validate Grants
3875

3876
        if not (target_s3_bucket := store.buckets.get(target_bucket_name)):
1✔
3877
            raise InvalidTargetBucketForLogging(
1✔
3878
                "The target bucket for logging does not exist",
3879
                TargetBucket=target_bucket_name,
3880
            )
3881

3882
        source_bucket_region = s3_bucket.bucket_region
1✔
3883
        if target_s3_bucket.bucket_region != source_bucket_region:
1✔
3884
            raise (
1✔
3885
                CrossLocationLoggingProhibitted(
3886
                    "Cross S3 location logging not allowed. ",
3887
                    TargetBucketLocation=target_s3_bucket.bucket_region,
3888
                )
3889
                if source_bucket_region == AWS_REGION_US_EAST_1
3890
                else CrossLocationLoggingProhibitted(
3891
                    "Cross S3 location logging not allowed. ",
3892
                    SourceBucketLocation=source_bucket_region,
3893
                    TargetBucketLocation=target_s3_bucket.bucket_region,
3894
                )
3895
            )
3896

3897
        s3_bucket.logging = logging_config
1✔
3898

3899
    def get_bucket_logging(
1✔
3900
        self,
3901
        context: RequestContext,
3902
        bucket: BucketName,
3903
        expected_bucket_owner: AccountId = None,
3904
        **kwargs,
3905
    ) -> GetBucketLoggingOutput:
3906
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
3907

3908
        if not s3_bucket.logging:
1✔
3909
            return GetBucketLoggingOutput()
1✔
3910

3911
        return GetBucketLoggingOutput(LoggingEnabled=s3_bucket.logging)
1✔
3912

3913
    def put_bucket_replication(
1✔
3914
        self,
3915
        context: RequestContext,
3916
        bucket: BucketName,
3917
        replication_configuration: ReplicationConfiguration,
3918
        content_md5: ContentMD5 = None,
3919
        checksum_algorithm: ChecksumAlgorithm = None,
3920
        token: ObjectLockToken = None,
3921
        expected_bucket_owner: AccountId = None,
3922
        **kwargs,
3923
    ) -> None:
3924
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
3925
        if not s3_bucket.versioning_status == BucketVersioningStatus.Enabled:
1✔
3926
            raise InvalidRequest(
1✔
3927
                "Versioning must be 'Enabled' on the bucket to apply a replication configuration"
3928
            )
3929

3930
        if not (rules := replication_configuration.get("Rules")):
1✔
3931
            raise MalformedXML()
1✔
3932

3933
        for rule in rules:
1✔
3934
            if "ID" not in rule:
1✔
3935
                rule["ID"] = short_uid()
1✔
3936

3937
            dest_bucket_arn = rule.get("Destination", {}).get("Bucket")
1✔
3938
            dest_bucket_name = s3_bucket_name(dest_bucket_arn)
1✔
3939
            if (
1✔
3940
                not (dest_s3_bucket := store.buckets.get(dest_bucket_name))
3941
                or not dest_s3_bucket.versioning_status == BucketVersioningStatus.Enabled
3942
            ):
3943
                # according to AWS testing the same exception is raised if the bucket does not exist
3944
                # or if versioning was disabled
3945
                raise InvalidRequest("Destination bucket must have versioning enabled.")
1✔
3946

3947
        # TODO more validation on input
3948
        s3_bucket.replication = replication_configuration
1✔
3949

3950
    def get_bucket_replication(
1✔
3951
        self,
3952
        context: RequestContext,
3953
        bucket: BucketName,
3954
        expected_bucket_owner: AccountId = None,
3955
        **kwargs,
3956
    ) -> GetBucketReplicationOutput:
3957
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
3958

3959
        if not s3_bucket.replication:
1✔
3960
            raise ReplicationConfigurationNotFoundError(
1✔
3961
                "The replication configuration was not found",
3962
                BucketName=bucket,
3963
            )
3964

3965
        return GetBucketReplicationOutput(ReplicationConfiguration=s3_bucket.replication)
1✔
3966

3967
    def delete_bucket_replication(
1✔
3968
        self,
3969
        context: RequestContext,
3970
        bucket: BucketName,
3971
        expected_bucket_owner: AccountId = None,
3972
        **kwargs,
3973
    ) -> None:
3974
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
3975

3976
        s3_bucket.replication = None
1✔
3977

3978
    @handler("PutBucketAcl", expand=False)
1✔
3979
    def put_bucket_acl(
1✔
3980
        self,
3981
        context: RequestContext,
3982
        request: PutBucketAclRequest,
3983
    ) -> None:
3984
        bucket = request["Bucket"]
1✔
3985
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
3986
        acp = get_access_control_policy_from_acl_request(
1✔
3987
            request=request, owner=s3_bucket.owner, request_body=context.request.data
3988
        )
3989
        s3_bucket.acl = acp
1✔
3990

3991
    def get_bucket_acl(
1✔
3992
        self,
3993
        context: RequestContext,
3994
        bucket: BucketName,
3995
        expected_bucket_owner: AccountId = None,
3996
        **kwargs,
3997
    ) -> GetBucketAclOutput:
3998
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
3999

4000
        return GetBucketAclOutput(Owner=s3_bucket.acl["Owner"], Grants=s3_bucket.acl["Grants"])
1✔
4001

4002
    @handler("PutObjectAcl", expand=False)
1✔
4003
    def put_object_acl(
1✔
4004
        self,
4005
        context: RequestContext,
4006
        request: PutObjectAclRequest,
4007
    ) -> PutObjectAclOutput:
4008
        bucket = request["Bucket"]
1✔
4009
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
4010

4011
        s3_object = s3_bucket.get_object(
1✔
4012
            key=request["Key"],
4013
            version_id=request.get("VersionId"),
4014
            http_method="PUT",
4015
        )
4016
        acp = get_access_control_policy_from_acl_request(
1✔
4017
            request=request, owner=s3_object.owner, request_body=context.request.data
4018
        )
4019
        previous_acl = s3_object.acl
1✔
4020
        s3_object.acl = acp
1✔
4021

4022
        if previous_acl != acp:
1✔
4023
            self._notify(context, s3_bucket=s3_bucket, s3_object=s3_object)
1✔
4024

4025
        # TODO: RequestCharged
4026
        return PutObjectAclOutput()
1✔
4027

4028
    def get_object_acl(
1✔
4029
        self,
4030
        context: RequestContext,
4031
        bucket: BucketName,
4032
        key: ObjectKey,
4033
        version_id: ObjectVersionId = None,
4034
        request_payer: RequestPayer = None,
4035
        expected_bucket_owner: AccountId = None,
4036
        **kwargs,
4037
    ) -> GetObjectAclOutput:
4038
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
4039

4040
        s3_object = s3_bucket.get_object(
1✔
4041
            key=key,
4042
            version_id=version_id,
4043
        )
4044
        # TODO: RequestCharged
4045
        return GetObjectAclOutput(Owner=s3_object.acl["Owner"], Grants=s3_object.acl["Grants"])
1✔
4046

4047
    def get_bucket_policy_status(
1✔
4048
        self,
4049
        context: RequestContext,
4050
        bucket: BucketName,
4051
        expected_bucket_owner: AccountId = None,
4052
        **kwargs,
4053
    ) -> GetBucketPolicyStatusOutput:
4054
        raise NotImplementedError
4055

4056
    def get_object_torrent(
1✔
4057
        self,
4058
        context: RequestContext,
4059
        bucket: BucketName,
4060
        key: ObjectKey,
4061
        request_payer: RequestPayer = None,
4062
        expected_bucket_owner: AccountId = None,
4063
        **kwargs,
4064
    ) -> GetObjectTorrentOutput:
4065
        raise NotImplementedError
4066

4067
    def post_object(
1✔
4068
        self, context: RequestContext, bucket: BucketName, body: IO[Body] = None, **kwargs
4069
    ) -> PostResponse:
4070
        if "multipart/form-data" not in context.request.headers.get("Content-Type", ""):
1✔
4071
            raise PreconditionFailed(
1✔
4072
                "At least one of the pre-conditions you specified did not hold",
4073
                Condition="Bucket POST must be of the enclosure-type multipart/form-data",
4074
            )
4075
        # see https://docs.aws.amazon.com/AmazonS3/latest/API/RESTObjectPOST.html
4076
        # TODO: signature validation is not implemented for pre-signed POST
4077
        # policy validation is not implemented either, except expiration and mandatory fields
4078
        # This operation is the only one using form for storing the request data. We will have to do some manual
4079
        # parsing here, as no specs are present for this, as no client directly implements this operation.
4080
        store, s3_bucket = self._get_cross_account_bucket(context, bucket)
1✔
4081

4082
        form = context.request.form
1✔
4083
        object_key = context.request.form.get("key")
1✔
4084

4085
        if "file" in form:
1✔
4086
            # in AWS, you can pass the file content as a string in the form field and not as a file object
4087
            file_data = to_bytes(form["file"])
1✔
4088
            object_content_length = len(file_data)
1✔
4089
            stream = BytesIO(file_data)
1✔
4090
        else:
4091
            # this is the default behaviour
4092
            fileobj = context.request.files["file"]
1✔
4093
            stream = fileobj.stream
1✔
4094
            # stream is a SpooledTemporaryFile, so we can seek the stream to know its length, necessary for policy
4095
            # validation
4096
            original_pos = stream.tell()
1✔
4097
            object_content_length = stream.seek(0, 2)
1✔
4098
            # reset the stream and put it back at its original position
4099
            stream.seek(original_pos, 0)
1✔
4100

4101
            if "${filename}" in object_key:
1✔
4102
                # TODO: ${filename} is actually usable in all form fields
4103
                # See https://docs.aws.amazon.com/sdk-for-ruby/v3/api/Aws/S3/PresignedPost.html
4104
                # > The string ${filename} is automatically replaced with the name of the file provided by the user and
4105
                # is recognized by all form fields.
4106
                object_key = object_key.replace("${filename}", fileobj.filename)
×
4107

4108
        # TODO: see if we need to pass additional metadata not contained in the policy from the table under
4109
        # https://docs.aws.amazon.com/AmazonS3/latest/API/sigv4-HTTPPOSTConstructPolicy.html#sigv4-PolicyConditions
4110
        additional_policy_metadata = {
1✔
4111
            "bucket": bucket,
4112
            "content_length": object_content_length,
4113
        }
4114
        validate_post_policy(form, additional_policy_metadata)
1✔
4115

4116
        if canned_acl := form.get("acl"):
1✔
4117
            validate_canned_acl(canned_acl)
×
4118
            acp = get_canned_acl(canned_acl, owner=s3_bucket.owner)
×
4119
        else:
4120
            acp = get_canned_acl(BucketCannedACL.private, owner=s3_bucket.owner)
1✔
4121

4122
        post_system_settable_headers = [
1✔
4123
            "Cache-Control",
4124
            "Content-Type",
4125
            "Content-Disposition",
4126
            "Content-Encoding",
4127
        ]
4128
        system_metadata = {}
1✔
4129
        for system_metadata_field in post_system_settable_headers:
1✔
4130
            if field_value := form.get(system_metadata_field):
1✔
4131
                system_metadata[system_metadata_field.replace("-", "")] = field_value
1✔
4132

4133
        if not system_metadata.get("ContentType"):
1✔
4134
            system_metadata["ContentType"] = "binary/octet-stream"
1✔
4135

4136
        user_metadata = {
1✔
4137
            field.removeprefix("x-amz-meta-").lower(): form.get(field)
4138
            for field in form
4139
            if field.startswith("x-amz-meta-")
4140
        }
4141

4142
        if tagging := form.get("tagging"):
1✔
4143
            # this is weird, as it's direct XML in the form, we need to parse it directly
4144
            tagging = parse_post_object_tagging_xml(tagging)
1✔
4145

4146
        if (storage_class := form.get("x-amz-storage-class")) is not None and (
1✔
4147
            storage_class not in STORAGE_CLASSES or storage_class == StorageClass.OUTPOSTS
4148
        ):
4149
            raise InvalidStorageClass(
1✔
4150
                "The storage class you specified is not valid", StorageClassRequested=storage_class
4151
            )
4152

4153
        encryption_request = {
1✔
4154
            "ServerSideEncryption": form.get("x-amz-server-side-encryption"),
4155
            "SSEKMSKeyId": form.get("x-amz-server-side-encryption-aws-kms-key-id"),
4156
            "BucketKeyEnabled": form.get("x-amz-server-side-encryption-bucket-key-enabled"),
4157
        }
4158

4159
        encryption_parameters = get_encryption_parameters_from_request_and_bucket(
1✔
4160
            encryption_request,
4161
            s3_bucket,
4162
            store,
4163
        )
4164

4165
        checksum_algorithm = form.get("x-amz-checksum-algorithm")
1✔
4166
        checksum_value = (
1✔
4167
            form.get(f"x-amz-checksum-{checksum_algorithm.lower()}") if checksum_algorithm else None
4168
        )
4169
        expires = (
1✔
4170
            str_to_rfc_1123_datetime(expires_str) if (expires_str := form.get("Expires")) else None
4171
        )
4172

4173
        version_id = generate_version_id(s3_bucket.versioning_status)
1✔
4174

4175
        s3_object = S3Object(
1✔
4176
            key=object_key,
4177
            version_id=version_id,
4178
            storage_class=storage_class,
4179
            expires=expires,
4180
            user_metadata=user_metadata,
4181
            system_metadata=system_metadata,
4182
            checksum_algorithm=checksum_algorithm,
4183
            checksum_value=checksum_value,
4184
            encryption=encryption_parameters.encryption,
4185
            kms_key_id=encryption_parameters.kms_key_id,
4186
            bucket_key_enabled=encryption_parameters.bucket_key_enabled,
4187
            website_redirect_location=form.get("x-amz-website-redirect-location"),
4188
            acl=acp,
4189
            owner=s3_bucket.owner,  # TODO: for now we only have one owner, but it can depends on Bucket settings
4190
        )
4191

4192
        with self._storage_backend.open(bucket, s3_object, mode="w") as s3_stored_object:
1✔
4193
            s3_stored_object.write(stream)
1✔
4194

4195
            if checksum_algorithm and s3_object.checksum_value != s3_stored_object.checksum:
1✔
4196
                self._storage_backend.remove(bucket, s3_object)
×
4197
                raise InvalidRequest(
×
4198
                    f"Value for x-amz-checksum-{checksum_algorithm.lower()} header is invalid."
4199
                )
4200

4201
            s3_bucket.objects.set(object_key, s3_object)
1✔
4202

4203
        # in case we are overriding an object, delete the tags entry
4204
        key_id = get_unique_key_id(bucket, object_key, version_id)
1✔
4205
        store.TAGS.tags.pop(key_id, None)
1✔
4206
        if tagging:
1✔
4207
            store.TAGS.tags[key_id] = tagging
1✔
4208

4209
        response = PostResponse()
1✔
4210
        # hacky way to set the etag in the headers as well: two locations for one value
4211
        response["ETagHeader"] = s3_object.quoted_etag
1✔
4212

4213
        if redirect := form.get("success_action_redirect"):
1✔
4214
            # we need to create the redirect, as the parser could not return the moto-calculated one
4215
            try:
1✔
4216
                redirect = create_redirect_for_post_request(
1✔
4217
                    base_redirect=redirect,
4218
                    bucket=bucket,
4219
                    object_key=object_key,
4220
                    etag=s3_object.quoted_etag,
4221
                )
4222
                response["LocationHeader"] = redirect
1✔
4223
                response["StatusCode"] = 303
1✔
4224
            except ValueError:
1✔
4225
                # If S3 cannot interpret the URL, it acts as if the field is not present.
4226
                response["StatusCode"] = form.get("success_action_status", 204)
1✔
4227

4228
        elif status_code := form.get("success_action_status"):
1✔
4229
            response["StatusCode"] = status_code
1✔
4230
        else:
4231
            response["StatusCode"] = 204
1✔
4232

4233
        response["LocationHeader"] = response.get(
1✔
4234
            "LocationHeader", f"{get_full_default_bucket_location(bucket)}{object_key}"
4235
        )
4236

4237
        if s3_bucket.versioning_status == "Enabled":
1✔
4238
            response["VersionId"] = s3_object.version_id
×
4239

4240
        if s3_object.checksum_algorithm:
1✔
4241
            response[f"Checksum{checksum_algorithm.upper()}"] = s3_object.checksum_value
×
4242

4243
        if s3_bucket.lifecycle_rules:
1✔
4244
            if expiration_header := self._get_expiration_header(
×
4245
                s3_bucket.lifecycle_rules,
4246
                bucket,
4247
                s3_object,
4248
                store.TAGS.tags.get(key_id, {}),
4249
            ):
4250
                # TODO: we either apply the lifecycle to existing objects when we set the new rules, or we need to
4251
                #  apply them everytime we get/head an object
4252
                response["Expiration"] = expiration_header
×
4253

4254
        add_encryption_to_response(response, s3_object=s3_object)
1✔
4255

4256
        self._notify(context, s3_bucket=s3_bucket, s3_object=s3_object)
1✔
4257

4258
        if response["StatusCode"] == "201":
1✔
4259
            # if the StatusCode is 201, S3 returns an XML body with additional information
4260
            response["ETag"] = s3_object.quoted_etag
1✔
4261
            response["Bucket"] = bucket
1✔
4262
            response["Key"] = object_key
1✔
4263
            response["Location"] = response["LocationHeader"]
1✔
4264

4265
        return response
1✔
4266

4267

4268
def generate_version_id(bucket_versioning_status: str) -> str | None:
1✔
4269
    if not bucket_versioning_status:
1✔
4270
        return None
1✔
4271
    elif bucket_versioning_status.lower() == "enabled":
1✔
4272
        return generate_safe_version_id()
1✔
4273
    else:
4274
        return "null"
1✔
4275

4276

4277
def add_encryption_to_response(response: dict, s3_object: S3Object):
1✔
4278
    if encryption := s3_object.encryption:
1✔
4279
        response["ServerSideEncryption"] = encryption
1✔
4280
        if encryption == ServerSideEncryption.aws_kms:
1✔
4281
            response["SSEKMSKeyId"] = s3_object.kms_key_id
1✔
4282
            if s3_object.bucket_key_enabled:
1✔
4283
                response["BucketKeyEnabled"] = s3_object.bucket_key_enabled
1✔
4284

4285

4286
def get_encryption_parameters_from_request_and_bucket(
1✔
4287
    request: PutObjectRequest | CopyObjectRequest | CreateMultipartUploadRequest,
4288
    s3_bucket: S3Bucket,
4289
    store: S3Store,
4290
) -> EncryptionParameters:
4291
    if request.get("SSECustomerKey"):
1✔
4292
        # we return early, because ServerSideEncryption does not apply if the request has SSE-C
4293
        return EncryptionParameters(None, None, False)
1✔
4294

4295
    encryption = request.get("ServerSideEncryption")
1✔
4296
    kms_key_id = request.get("SSEKMSKeyId")
1✔
4297
    bucket_key_enabled = request.get("BucketKeyEnabled")
1✔
4298
    if s3_bucket.encryption_rule:
1✔
4299
        bucket_key_enabled = bucket_key_enabled or s3_bucket.encryption_rule.get("BucketKeyEnabled")
1✔
4300
        encryption = (
1✔
4301
            encryption
4302
            or s3_bucket.encryption_rule["ApplyServerSideEncryptionByDefault"]["SSEAlgorithm"]
4303
        )
4304
        if encryption == ServerSideEncryption.aws_kms:
1✔
4305
            key_id = kms_key_id or s3_bucket.encryption_rule[
1✔
4306
                "ApplyServerSideEncryptionByDefault"
4307
            ].get("KMSMasterKeyID")
4308
            kms_key_id = get_kms_key_arn(
1✔
4309
                key_id, s3_bucket.bucket_account_id, s3_bucket.bucket_region
4310
            )
4311
            if not kms_key_id:
1✔
4312
                # if not key is provided, AWS will use an AWS managed KMS key
4313
                # create it if it doesn't already exist, and save it in the store per region
4314
                if not store.aws_managed_kms_key_id:
1✔
4315
                    managed_kms_key_id = create_s3_kms_managed_key_for_region(
1✔
4316
                        s3_bucket.bucket_account_id, s3_bucket.bucket_region
4317
                    )
4318
                    store.aws_managed_kms_key_id = managed_kms_key_id
1✔
4319

4320
                kms_key_id = store.aws_managed_kms_key_id
1✔
4321

4322
    return EncryptionParameters(encryption, kms_key_id, bucket_key_enabled)
1✔
4323

4324

4325
def get_object_lock_parameters_from_bucket_and_request(
1✔
4326
    request: PutObjectRequest | CopyObjectRequest | CreateMultipartUploadRequest,
4327
    s3_bucket: S3Bucket,
4328
):
4329
    # TODO: also validate here?
4330
    lock_mode = request.get("ObjectLockMode")
1✔
4331
    lock_legal_status = request.get("ObjectLockLegalHoldStatus")
1✔
4332
    lock_until = request.get("ObjectLockRetainUntilDate")
1✔
4333

4334
    if default_retention := s3_bucket.object_lock_default_retention:
1✔
4335
        lock_mode = lock_mode or default_retention.get("Mode")
1✔
4336
        if lock_mode and not lock_until:
1✔
4337
            lock_until = get_retention_from_now(
1✔
4338
                days=default_retention.get("Days"),
4339
                years=default_retention.get("Years"),
4340
            )
4341

4342
    return ObjectLockParameters(lock_until, lock_legal_status, lock_mode)
1✔
4343

4344

4345
def get_part_range(s3_object: S3Object, part_number: PartNumber) -> ObjectRange:
1✔
4346
    """
4347
    Calculate the range value from a part Number for an S3 Object
4348
    :param s3_object: S3Object
4349
    :param part_number: the wanted part from the S3Object
4350
    :return: an ObjectRange used to return only a slice of an Object
4351
    """
4352
    if not s3_object.parts:
1✔
4353
        if part_number > 1:
1✔
4354
            raise InvalidPartNumber(
1✔
4355
                "The requested partnumber is not satisfiable",
4356
                PartNumberRequested=part_number,
4357
                ActualPartCount=1,
4358
            )
4359
        return ObjectRange(
1✔
4360
            begin=0,
4361
            end=s3_object.size - 1,
4362
            content_length=s3_object.size,
4363
            content_range=f"bytes 0-{s3_object.size - 1}/{s3_object.size}",
4364
        )
4365
    elif not (part_data := s3_object.parts.get(part_number)):
1✔
4366
        raise InvalidPartNumber(
1✔
4367
            "The requested partnumber is not satisfiable",
4368
            PartNumberRequested=part_number,
4369
            ActualPartCount=len(s3_object.parts),
4370
        )
4371

4372
    begin, part_length = part_data
1✔
4373
    end = begin + part_length - 1
1✔
4374
    return ObjectRange(
1✔
4375
        begin=begin,
4376
        end=end,
4377
        content_length=part_length,
4378
        content_range=f"bytes {begin}-{end}/{s3_object.size}",
4379
    )
4380

4381

4382
def get_acl_headers_from_request(
1✔
4383
    request: Union[
4384
        PutObjectRequest,
4385
        CreateMultipartUploadRequest,
4386
        CopyObjectRequest,
4387
        CreateBucketRequest,
4388
        PutBucketAclRequest,
4389
        PutObjectAclRequest,
4390
    ],
4391
) -> list[tuple[str, str]]:
4392
    permission_keys = [
1✔
4393
        "GrantFullControl",
4394
        "GrantRead",
4395
        "GrantReadACP",
4396
        "GrantWrite",
4397
        "GrantWriteACP",
4398
    ]
4399
    acl_headers = [
1✔
4400
        (permission, grant_header)
4401
        for permission in permission_keys
4402
        if (grant_header := request.get(permission))
4403
    ]
4404
    return acl_headers
1✔
4405

4406

4407
def get_access_control_policy_from_acl_request(
1✔
4408
    request: Union[PutBucketAclRequest, PutObjectAclRequest],
4409
    owner: Owner,
4410
    request_body: bytes,
4411
) -> AccessControlPolicy:
4412
    canned_acl = request.get("ACL")
1✔
4413
    acl_headers = get_acl_headers_from_request(request)
1✔
4414

4415
    # FIXME: this is very dirty, but the parser does not differentiate between an empty body and an empty XML node
4416
    # errors are different depending on that data, so we need to access the context. Modifying the parser for this
4417
    # use case seems dangerous
4418
    is_acp_in_body = request_body
1✔
4419

4420
    if not (canned_acl or acl_headers or is_acp_in_body):
1✔
4421
        raise MissingSecurityHeader(
1✔
4422
            "Your request was missing a required header", MissingHeaderName="x-amz-acl"
4423
        )
4424

4425
    elif canned_acl and acl_headers:
1✔
4426
        raise InvalidRequest("Specifying both Canned ACLs and Header Grants is not allowed")
1✔
4427

4428
    elif (canned_acl or acl_headers) and is_acp_in_body:
1✔
4429
        raise UnexpectedContent("This request does not support content")
1✔
4430

4431
    if canned_acl:
1✔
4432
        validate_canned_acl(canned_acl)
1✔
4433
        acp = get_canned_acl(canned_acl, owner=owner)
1✔
4434

4435
    elif acl_headers:
1✔
4436
        grants = []
1✔
4437
        for permission, grantees_values in acl_headers:
1✔
4438
            permission = get_permission_from_header(permission)
1✔
4439
            partial_grants = parse_grants_in_headers(permission, grantees_values)
1✔
4440
            grants.extend(partial_grants)
1✔
4441

4442
        acp = AccessControlPolicy(Owner=owner, Grants=grants)
1✔
4443
    else:
4444
        acp = request.get("AccessControlPolicy")
1✔
4445
        validate_acl_acp(acp)
1✔
4446
        if (
1✔
4447
            owner.get("DisplayName")
4448
            and acp["Grants"]
4449
            and "DisplayName" not in acp["Grants"][0]["Grantee"]
4450
        ):
4451
            acp["Grants"][0]["Grantee"]["DisplayName"] = owner["DisplayName"]
1✔
4452

4453
    return acp
1✔
4454

4455

4456
def get_access_control_policy_for_new_resource_request(
1✔
4457
    request: Union[
4458
        PutObjectRequest, CreateMultipartUploadRequest, CopyObjectRequest, CreateBucketRequest
4459
    ],
4460
    owner: Owner,
4461
) -> AccessControlPolicy:
4462
    # TODO: this is basic ACL, not taking into account Bucket settings. Revisit once we really implement ACLs.
4463
    canned_acl = request.get("ACL")
1✔
4464
    acl_headers = get_acl_headers_from_request(request)
1✔
4465

4466
    if not (canned_acl or acl_headers):
1✔
4467
        return get_canned_acl(BucketCannedACL.private, owner=owner)
1✔
4468

4469
    elif canned_acl and acl_headers:
1✔
4470
        raise InvalidRequest("Specifying both Canned ACLs and Header Grants is not allowed")
×
4471

4472
    if canned_acl:
1✔
4473
        validate_canned_acl(canned_acl)
1✔
4474
        return get_canned_acl(canned_acl, owner=owner)
1✔
4475

4476
    grants = []
×
4477
    for permission, grantees_values in acl_headers:
×
4478
        permission = get_permission_from_header(permission)
×
4479
        partial_grants = parse_grants_in_headers(permission, grantees_values)
×
4480
        grants.extend(partial_grants)
×
4481

4482
    return AccessControlPolicy(Owner=owner, Grants=grants)
×
4483

4484

4485
def object_exists_for_precondition_write(s3_bucket: S3Bucket, key: ObjectKey) -> bool:
1✔
4486
    return (existing := s3_bucket.objects.get(key)) and not isinstance(existing, S3DeleteMarker)
1✔
4487

4488

4489
def verify_object_equality_precondition_write(
1✔
4490
    s3_bucket: S3Bucket,
4491
    key: ObjectKey,
4492
    etag: str,
4493
    initiated: datetime.datetime | None = None,
4494
) -> None:
4495
    existing = s3_bucket.objects.get(key)
1✔
4496
    if not existing or isinstance(existing, S3DeleteMarker):
1✔
4497
        raise NoSuchKey("The specified key does not exist.", Key=key)
1✔
4498

4499
    if not existing.etag == etag.strip('"'):
1✔
4500
        raise PreconditionFailed(
1✔
4501
            "At least one of the pre-conditions you specified did not hold",
4502
            Condition="If-Match",
4503
        )
4504

4505
    if initiated and initiated < existing.last_modified:
1✔
4506
        raise ConditionalRequestConflict(
1✔
4507
            "The conditional request cannot succeed due to a conflicting operation against this resource.",
4508
            Condition="If-Match",
4509
            Key=key,
4510
        )
STATUS · Troubleshooting · Open an Issue · Sales · Support · CAREERS · ENTERPRISE · START FREE TRIAL · SCHEDULE DEMO
ANNOUNCEMENTS · TWITTER · TOS & SLA · Supported CI Services · What's a CI service? · Automated Testing

© 2026 Coveralls, Inc