codekingpro/portable-devtools
114k
1# Copyright 2021 Amazon.com, Inc. or its affiliates. All Rights Reserved.2#3# Licensed under the Apache License, Version 2.0 (the "License"). You4# may not use this file except in compliance with the License. A copy of5# the License is located at6#7# http://aws.amazon.com/apache2.0/8#9# or in the "license" file accompanying this file. This file is10# distributed on an "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF11# ANY KIND, either express or implied. See the License for the specific12# language governing permissions and limitations under the License.13import logging14import threading15from io import BytesIO16 17import awscrt.http18import awscrt.s319import botocore.awsrequest20import botocore.session21from awscrt.auth import AwsCredentials, AwsCredentialsProvider22from awscrt.io import (23 ClientBootstrap,24 ClientTlsContext,25 DefaultHostResolver,26 EventLoopGroup,27 TlsContextOptions,28)29from awscrt.s3 import S3Client, S3RequestTlsMode, S3RequestType30from botocore import UNSIGNED31from botocore.compat import urlsplit32from botocore.config import Config33from botocore.exceptions import NoCredentialsError34 35from s3transfer.constants import MB36from s3transfer.exceptions import TransferNotDoneError37from s3transfer.futures import BaseTransferFuture, BaseTransferMeta38from s3transfer.utils import CallArgs, OSUtils, get_callbacks39 40logger = logging.getLogger(__name__)41 42CRT_S3_PROCESS_LOCK = None43 44 45def acquire_crt_s3_process_lock(name):46 # Currently, the CRT S3 client performs best when there is only one47 # instance of it running on a host. This lock allows an application to48 # signal across processes whether there is another process of the same49 # application using the CRT S3 client and prevent spawning more than one50 # CRT S3 clients running on the system for that application.51 #52 # NOTE: When acquiring the CRT process lock, the lock automatically is53 # released when the lock object is garbage collected. So, the CRT process54 # lock is set as a global so that it is not unintentionally garbage55 # collected/released if reference of the lock is lost.56 global CRT_S3_PROCESS_LOCK57 if CRT_S3_PROCESS_LOCK is None:58 crt_lock = awscrt.s3.CrossProcessLock(name)59 try:60 crt_lock.acquire()61 except RuntimeError:62 # If there is another process that is holding the lock, the CRT63 # returns a RuntimeError. We return None here to signal that our64 # current process was not able to acquire the lock.65 return None66 CRT_S3_PROCESS_LOCK = crt_lock67 return CRT_S3_PROCESS_LOCK68 69 70def create_s3_crt_client(71 region,72 crt_credentials_provider=None,73 num_threads=None,74 target_throughput=None,75 part_size=8 * MB,76 use_ssl=True,77 verify=None,78):79 """80 :type region: str81 :param region: The region used for signing82 83 :type crt_credentials_provider:84 Optional[awscrt.auth.AwsCredentialsProvider]85 :param crt_credentials_provider: CRT AWS credentials provider86 to use to sign requests. If not set, requests will not be signed.87 88 :type num_threads: Optional[int]89 :param num_threads: Number of worker threads generated. Default90 is the number of processors in the machine.91 92 :type target_throughput: Optional[int]93 :param target_throughput: Throughput target in bytes per second.94 By default, CRT will automatically attempt to choose a target95 throughput that matches the system's maximum network throughput.96 Currently, if CRT is unable to determine the maximum network97 throughput, a fallback target throughput of ``1_250_000_000`` bytes98 per second (which translates to 10 gigabits per second, or 1.1699 gibibytes per second) is used. To set a specific target100 throughput, set a value for this parameter.101 102 :type part_size: Optional[int]103 :param part_size: Size, in Bytes, of parts that files will be downloaded104 or uploaded in.105 106 :type use_ssl: boolean107 :param use_ssl: Whether or not to use SSL. By default, SSL is used.108 Note that not all services support non-ssl connections.109 110 :type verify: Optional[boolean/string]111 :param verify: Whether or not to verify SSL certificates.112 By default SSL certificates are verified. You can provide the113 following values:114 115 * False - do not validate SSL certificates. SSL will still be116 used (unless use_ssl is False), but SSL certificates117 will not be verified.118 * path/to/cert/bundle.pem - A filename of the CA cert bundle to119 use. Specify this argument if you want to use a custom CA cert120 bundle instead of the default one on your system.121 """122 event_loop_group = EventLoopGroup(num_threads)123 host_resolver = DefaultHostResolver(event_loop_group)124 bootstrap = ClientBootstrap(event_loop_group, host_resolver)125 tls_connection_options = None126 127 tls_mode = (128 S3RequestTlsMode.ENABLED if use_ssl else S3RequestTlsMode.DISABLED129 )130 if verify is not None:131 tls_ctx_options = TlsContextOptions()132 if verify:133 tls_ctx_options.override_default_trust_store_from_path(134 ca_filepath=verify135 )136 else:137 tls_ctx_options.verify_peer = False138 client_tls_option = ClientTlsContext(tls_ctx_options)139 tls_connection_options = client_tls_option.new_connection_options()140 target_gbps = _get_crt_throughput_target_gbps(141 provided_throughput_target_bytes=target_throughput142 )143 return S3Client(144 bootstrap=bootstrap,145 region=region,146 credential_provider=crt_credentials_provider,147 part_size=part_size,148 tls_mode=tls_mode,149 tls_connection_options=tls_connection_options,150 throughput_target_gbps=target_gbps,151 )152 153 154def _get_crt_throughput_target_gbps(provided_throughput_target_bytes=None):155 if provided_throughput_target_bytes is None:156 target_gbps = awscrt.s3.get_recommended_throughput_target_gbps()157 logger.debug(158 'Recommended CRT throughput target in gbps: %s', target_gbps159 )160 if target_gbps is None:161 target_gbps = 10.0162 else:163 # NOTE: The GB constant in s3transfer is technically a gibibyte. The164 # GB constant is not used here because the CRT interprets gigabits165 # for networking as a base power of 10166 # (i.e. 1000 ** 3 instead of 1024 ** 3).167 target_gbps = provided_throughput_target_bytes * 8 / 1_000_000_000168 logger.debug('Using CRT throughput target in gbps: %s', target_gbps)169 return target_gbps170 171 172class CRTTransferManager:173 def __init__(self, crt_s3_client, crt_request_serializer, osutil=None):174 """A transfer manager interface for Amazon S3 on CRT s3 client.175 176 :type crt_s3_client: awscrt.s3.S3Client177 :param crt_s3_client: The CRT s3 client, handling all the178 HTTP requests and functions under then hood179 180 :type crt_request_serializer: s3transfer.crt.BaseCRTRequestSerializer181 :param crt_request_serializer: Serializer, generates unsigned crt HTTP182 request.183 184 :type osutil: s3transfer.utils.OSUtils185 :param osutil: OSUtils object to use for os-related behavior when186 using with transfer manager.187 """188 if osutil is None:189 self._osutil = OSUtils()190 self._crt_s3_client = crt_s3_client191 self._s3_args_creator = S3ClientArgsCreator(192 crt_request_serializer, self._osutil193 )194 self._crt_exception_translator = (195 crt_request_serializer.translate_crt_exception196 )197 self._future_coordinators = []198 self._semaphore = threading.Semaphore(128) # not configurable199 # A counter to create unique id's for each transfer submitted.200 self._id_counter = 0201 202 def __enter__(self):203 return self204 205 def __exit__(self, exc_type, exc_value, *args):206 cancel = False207 if exc_type:208 cancel = True209 self._shutdown(cancel)210 211 def download(212 self, bucket, key, fileobj, extra_args=None, subscribers=None213 ):214 if extra_args is None:215 extra_args = {}216 if subscribers is None:217 subscribers = {}218 callargs = CallArgs(219 bucket=bucket,220 key=key,221 fileobj=fileobj,222 extra_args=extra_args,223 subscribers=subscribers,224 )225 return self._submit_transfer("get_object", callargs)226 227 def upload(self, fileobj, bucket, key, extra_args=None, subscribers=None):228 if extra_args is None:229 extra_args = {}230 if subscribers is None:231 subscribers = {}232 self._validate_checksum_algorithm_supported(extra_args)233 callargs = CallArgs(234 bucket=bucket,235 key=key,236 fileobj=fileobj,237 extra_args=extra_args,238 subscribers=subscribers,239 )240 return self._submit_transfer("put_object", callargs)241 242 def delete(self, bucket, key, extra_args=None, subscribers=None):243 if extra_args is None:244 extra_args = {}245 if subscribers is None:246 subscribers = {}247 callargs = CallArgs(248 bucket=bucket,249 key=key,250 extra_args=extra_args,251 subscribers=subscribers,252 )253 return self._submit_transfer("delete_object", callargs)254 255 def shutdown(self, cancel=False):256 self._shutdown(cancel)257 258 def _validate_checksum_algorithm_supported(self, extra_args):259 checksum_algorithm = extra_args.get('ChecksumAlgorithm')260 if checksum_algorithm is None:261 return262 supported_algorithms = list(awscrt.s3.S3ChecksumAlgorithm.__members__)263 if checksum_algorithm.upper() not in supported_algorithms:264 raise ValueError(265 f'ChecksumAlgorithm: {checksum_algorithm} not supported. '266 f'Supported algorithms are: {supported_algorithms}'267 )268 269 def _cancel_transfers(self):270 for coordinator in self._future_coordinators:271 if not coordinator.done():272 coordinator.cancel()273 274 def _finish_transfers(self):275 for coordinator in self._future_coordinators:276 coordinator.result()277 278 def _wait_transfers_done(self):279 for coordinator in self._future_coordinators:280 coordinator.wait_until_on_done_callbacks_complete()281 282 def _shutdown(self, cancel=False):283 if cancel:284 self._cancel_transfers()285 try:286 self._finish_transfers()287 288 except KeyboardInterrupt:289 self._cancel_transfers()290 except Exception:291 pass292 finally:293 self._wait_transfers_done()294 295 def _release_semaphore(self, **kwargs):296 self._semaphore.release()297 298 def _submit_transfer(self, request_type, call_args):299 on_done_after_calls = [self._release_semaphore]300 coordinator = CRTTransferCoordinator(301 transfer_id=self._id_counter,302 exception_translator=self._crt_exception_translator,303 )304 components = {305 'meta': CRTTransferMeta(self._id_counter, call_args),306 'coordinator': coordinator,307 }308 future = CRTTransferFuture(**components)309 afterdone = AfterDoneHandler(coordinator)310 on_done_after_calls.append(afterdone)311 312 try:313 self._semaphore.acquire()314 on_queued = self._s3_args_creator.get_crt_callback(315 future, 'queued'316 )317 on_queued()318 crt_callargs = self._s3_args_creator.get_make_request_args(319 request_type,320 call_args,321 coordinator,322 future,323 on_done_after_calls,324 )325 crt_s3_request = self._crt_s3_client.make_request(**crt_callargs)326 except Exception as e:327 coordinator.set_exception(e, True)328 on_done = self._s3_args_creator.get_crt_callback(329 future, 'done', after_subscribers=on_done_after_calls330 )331 on_done(error=e)332 else:333 coordinator.set_s3_request(crt_s3_request)334 self._future_coordinators.append(coordinator)335 336 self._id_counter += 1337 return future338 339 340class CRTTransferMeta(BaseTransferMeta):341 """Holds metadata about the CRTTransferFuture"""342 343 def __init__(self, transfer_id=None, call_args=None):344 self._transfer_id = transfer_id345 self._call_args = call_args346 self._user_context = {}347 348 @property349 def call_args(self):350 return self._call_args351 352 @property353 def transfer_id(self):354 return self._transfer_id355 356 @property357 def user_context(self):358 return self._user_context359 360 361class CRTTransferFuture(BaseTransferFuture):362 def __init__(self, meta=None, coordinator=None):363 """The future associated to a submitted transfer request via CRT S3 client364 365 :type meta: s3transfer.crt.CRTTransferMeta366 :param meta: The metadata associated to the transfer future.367 368 :type coordinator: s3transfer.crt.CRTTransferCoordinator369 :param coordinator: The coordinator associated to the transfer future.370 """371 self._meta = meta372 if meta is None:373 self._meta = CRTTransferMeta()374 self._coordinator = coordinator375 376 @property377 def meta(self):378 return self._meta379 380 def done(self):381 return self._coordinator.done()382 383 def result(self, timeout=None):384 self._coordinator.result(timeout)385 386 def cancel(self):387 self._coordinator.cancel()388 389 def set_exception(self, exception):390 """Sets the exception on the future."""391 if not self.done():392 raise TransferNotDoneError(393 'set_exception can only be called once the transfer is '394 'complete.'395 )396 self._coordinator.set_exception(exception, override=True)397 398 399class BaseCRTRequestSerializer:400 def serialize_http_request(self, transfer_type, future):401 """Serialize CRT HTTP requests.402 403 :type transfer_type: string404 :param transfer_type: the type of transfer made,405 e.g 'put_object', 'get_object', 'delete_object'406 407 :type future: s3transfer.crt.CRTTransferFuture408 409 :rtype: awscrt.http.HttpRequest410 :returns: An unsigned HTTP request to be used for the CRT S3 client411 """412 raise NotImplementedError('serialize_http_request()')413 414 def translate_crt_exception(self, exception):415 raise NotImplementedError('translate_crt_exception()')416 417 418class BotocoreCRTRequestSerializer(BaseCRTRequestSerializer):419 def __init__(self, session, client_kwargs=None):420 """Serialize CRT HTTP request using botocore logic421 It also takes into account configuration from both the session422 and any keyword arguments that could be passed to423 `Session.create_client()` when serializing the request.424 425 :type session: botocore.session.Session426 427 :type client_kwargs: Optional[Dict[str, str]])428 :param client_kwargs: The kwargs for the botocore429 s3 client initialization.430 """431 self._session = session432 if client_kwargs is None:433 client_kwargs = {}434 self._resolve_client_config(session, client_kwargs)435 self._client = session.create_client(**client_kwargs)436 self._client.meta.events.register(437 'request-created.s3.*', self._capture_http_request438 )439 self._client.meta.events.register(440 'after-call.s3.*', self._change_response_to_serialized_http_request441 )442 self._client.meta.events.register(443 'before-send.s3.*', self._make_fake_http_response444 )445 446 def _resolve_client_config(self, session, client_kwargs):447 user_provided_config = None448 if session.get_default_client_config():449 user_provided_config = session.get_default_client_config()450 if 'config' in client_kwargs:451 user_provided_config = client_kwargs['config']452 453 client_config = Config(signature_version=UNSIGNED)454 if user_provided_config:455 client_config = user_provided_config.merge(client_config)456 client_kwargs['config'] = client_config457 client_kwargs["service_name"] = "s3"458 459 def _crt_request_from_aws_request(self, aws_request):460 url_parts = urlsplit(aws_request.url)461 crt_path = url_parts.path462 if url_parts.query:463 crt_path = f'{crt_path}?{url_parts.query}'464 headers_list = []465 for name, value in aws_request.headers.items():466 if isinstance(value, str):467 headers_list.append((name, value))468 else:469 headers_list.append((name, str(value, 'utf-8')))470 471 crt_headers = awscrt.http.HttpHeaders(headers_list)472 473 crt_request = awscrt.http.HttpRequest(474 method=aws_request.method,475 path=crt_path,476 headers=crt_headers,477 body_stream=aws_request.body,478 )479 return crt_request480 481 def _convert_to_crt_http_request(self, botocore_http_request):482 # Logic that does CRTUtils.crt_request_from_aws_request483 crt_request = self._crt_request_from_aws_request(botocore_http_request)484 if crt_request.headers.get("host") is None:485 # If host is not set, set it for the request before using CRT s3486 url_parts = urlsplit(botocore_http_request.url)487 crt_request.headers.set("host", url_parts.netloc)488 if crt_request.headers.get('Content-MD5') is not None:489 crt_request.headers.remove("Content-MD5")490 491 # In general, the CRT S3 client expects a content length header. It492 # only expects a missing content length header if the body is not493 # seekable. However, botocore does not set the content length header494 # for GetObject API requests and so we set the content length to zero495 # to meet the CRT S3 client's expectation that the content length496 # header is set even if there is no body.497 if crt_request.headers.get('Content-Length') is None:498 if botocore_http_request.body is None:499 crt_request.headers.add('Content-Length', "0")500 501 # Botocore sets the Transfer-Encoding header when it cannot determine502 # the content length of the request body (e.g. it's not seekable).503 # However, CRT does not support this header, but it supports504 # non-seekable bodies. So we remove this header to not cause issues505 # in the downstream CRT S3 request.506 if crt_request.headers.get('Transfer-Encoding') is not None:507 crt_request.headers.remove('Transfer-Encoding')508 509 return crt_request510 511 def _capture_http_request(self, request, **kwargs):512 request.context['http_request'] = request513 514 def _change_response_to_serialized_http_request(515 self, context, parsed, **kwargs516 ):517 request = context['http_request']518 parsed['HTTPRequest'] = request.prepare()519 520 def _make_fake_http_response(self, request, **kwargs):521 return botocore.awsrequest.AWSResponse(522 None,523 200,524 {},525 FakeRawResponse(b""),526 )527 528 def _get_botocore_http_request(self, client_method, call_args):529 return getattr(self._client, client_method)(530 Bucket=call_args.bucket, Key=call_args.key, **call_args.extra_args531 )['HTTPRequest']532 533 def serialize_http_request(self, transfer_type, future):534 botocore_http_request = self._get_botocore_http_request(535 transfer_type, future.meta.call_args536 )537 crt_request = self._convert_to_crt_http_request(botocore_http_request)538 return crt_request539 540 def translate_crt_exception(self, exception):541 if isinstance(exception, awscrt.s3.S3ResponseError):542 return self._translate_crt_s3_response_error(exception)543 else:544 return None545 546 def _translate_crt_s3_response_error(self, s3_response_error):547 status_code = s3_response_error.status_code548 if status_code < 301:549 # Botocore's exception parsing only550 # runs on status codes >= 301551 return None552 553 headers = {k: v for k, v in s3_response_error.headers}554 operation_name = s3_response_error.operation_name555 if operation_name is not None:556 service_model = self._client.meta.service_model557 shape = service_model.operation_model(operation_name).output_shape558 else:559 shape = None560 561 response_dict = {562 'headers': botocore.awsrequest.HeadersDict(headers),563 'status_code': status_code,564 'body': s3_response_error.body,565 }566 parsed_response = self._client._response_parser.parse(567 response_dict, shape=shape568 )569 570 error_code = parsed_response.get("Error", {}).get("Code")571 error_class = self._client.exceptions.from_code(error_code)572 return error_class(parsed_response, operation_name=operation_name)573 574 575class FakeRawResponse(BytesIO):576 def stream(self, amt=1024, decode_content=None):577 while True:578 chunk = self.read(amt)579 if not chunk:580 break581 yield chunk582 583 584class BotocoreCRTCredentialsWrapper:585 def __init__(self, resolved_botocore_credentials):586 self._resolved_credentials = resolved_botocore_credentials587 588 def __call__(self):589 credentials = self._get_credentials().get_frozen_credentials()590 return AwsCredentials(591 credentials.access_key, credentials.secret_key, credentials.token592 )593 594 def to_crt_credentials_provider(self):595 return AwsCredentialsProvider.new_delegate(self)596 597 def _get_credentials(self):598 if self._resolved_credentials is None:599 raise NoCredentialsError()600 return self._resolved_credentials601 602 603class CRTTransferCoordinator:604 """A helper class for managing CRTTransferFuture"""605 606 def __init__(607 self, transfer_id=None, s3_request=None, exception_translator=None608 ):609 self.transfer_id = transfer_id610 self._exception_translator = exception_translator611 self._s3_request = s3_request612 self._lock = threading.Lock()613 self._exception = None614 self._crt_future = None615 self._done_event = threading.Event()616 617 @property618 def s3_request(self):619 return self._s3_request620 621 def set_done_callbacks_complete(self):622 self._done_event.set()623 624 def wait_until_on_done_callbacks_complete(self, timeout=None):625 self._done_event.wait(timeout)626 627 def set_exception(self, exception, override=False):628 with self._lock:629 if not self.done() or override:630 self._exception = exception631 632 def cancel(self):633 if self._s3_request:634 self._s3_request.cancel()635 636 def result(self, timeout=None):637 if self._exception:638 raise self._exception639 try:640 self._crt_future.result(timeout)641 except KeyboardInterrupt:642 self.cancel()643 self._crt_future.result(timeout)644 raise645 except Exception as e:646 self.handle_exception(e)647 finally:648 if self._s3_request:649 self._s3_request = None650 651 def handle_exception(self, exc):652 translated_exc = None653 if self._exception_translator:654 try:655 translated_exc = self._exception_translator(exc)656 except Exception as e:657 # Bail out if we hit an issue translating658 # and raise the original error.659 logger.debug("Unable to translate exception.", exc_info=e)660 pass661 if translated_exc is not None:662 raise translated_exc from exc663 else:664 raise exc665 666 def done(self):667 if self._crt_future is None:668 return False669 return self._crt_future.done()670 671 def set_s3_request(self, s3_request):672 self._s3_request = s3_request673 self._crt_future = self._s3_request.finished_future674 675 676class S3ClientArgsCreator:677 def __init__(self, crt_request_serializer, os_utils):678 self._request_serializer = crt_request_serializer679 self._os_utils = os_utils680 681 def get_make_request_args(682 self, request_type, call_args, coordinator, future, on_done_after_calls683 ):684 request_args_handler = getattr(685 self,686 f'_get_make_request_args_{request_type}',687 self._default_get_make_request_args,688 )689 return request_args_handler(690 request_type=request_type,691 call_args=call_args,692 coordinator=coordinator,693 future=future,694 on_done_before_calls=[],695 on_done_after_calls=on_done_after_calls,696 )697 698 def get_crt_callback(699 self,700 future,701 callback_type,702 before_subscribers=None,703 after_subscribers=None,704 ):705 def invoke_all_callbacks(*args, **kwargs):706 callbacks_list = []707 if before_subscribers is not None:708 callbacks_list += before_subscribers709 callbacks_list += get_callbacks(future, callback_type)710 if after_subscribers is not None:711 callbacks_list += after_subscribers712 for callback in callbacks_list:713 # The get_callbacks helper will set the first augment714 # by keyword, the other augments need to be set by keyword715 # as well716 if callback_type == "progress":717 callback(bytes_transferred=args[0])718 else:719 callback(*args, **kwargs)720 721 return invoke_all_callbacks722 723 def _get_make_request_args_put_object(724 self,725 request_type,726 call_args,727 coordinator,728 future,729 on_done_before_calls,730 on_done_after_calls,731 ):732 send_filepath = None733 if isinstance(call_args.fileobj, str):734 send_filepath = call_args.fileobj735 data_len = self._os_utils.get_file_size(send_filepath)736 call_args.extra_args["ContentLength"] = data_len737 else:738 call_args.extra_args["Body"] = call_args.fileobj739 740 checksum_algorithm = call_args.extra_args.pop(741 'ChecksumAlgorithm', 'CRC32'742 ).upper()743 checksum_config = awscrt.s3.S3ChecksumConfig(744 algorithm=awscrt.s3.S3ChecksumAlgorithm[checksum_algorithm],745 location=awscrt.s3.S3ChecksumLocation.TRAILER,746 )747 # Suppress botocore's automatic MD5 calculation by setting an override748 # value that will get deleted in the BotocoreCRTRequestSerializer.749 # As part of the CRT S3 request, we request the CRT S3 client to750 # automatically add trailing checksums to its uploads.751 call_args.extra_args["ContentMD5"] = "override-to-be-removed"752 753 make_request_args = self._default_get_make_request_args(754 request_type=request_type,755 call_args=call_args,756 coordinator=coordinator,757 future=future,758 on_done_before_calls=on_done_before_calls,759 on_done_after_calls=on_done_after_calls,760 )761 make_request_args['send_filepath'] = send_filepath762 make_request_args['checksum_config'] = checksum_config763 return make_request_args764 765 def _get_make_request_args_get_object(766 self,767 request_type,768 call_args,769 coordinator,770 future,771 on_done_before_calls,772 on_done_after_calls,773 ):774 recv_filepath = None775 on_body = None776 checksum_config = awscrt.s3.S3ChecksumConfig(validate_response=True)777 if isinstance(call_args.fileobj, str):778 final_filepath = call_args.fileobj779 recv_filepath = self._os_utils.get_temp_filename(final_filepath)780 on_done_before_calls.append(781 RenameTempFileHandler(782 coordinator, final_filepath, recv_filepath, self._os_utils783 )784 )785 else:786 on_body = OnBodyFileObjWriter(call_args.fileobj)787 788 make_request_args = self._default_get_make_request_args(789 request_type=request_type,790 call_args=call_args,791 coordinator=coordinator,792 future=future,793 on_done_before_calls=on_done_before_calls,794 on_done_after_calls=on_done_after_calls,795 )796 make_request_args['recv_filepath'] = recv_filepath797 make_request_args['on_body'] = on_body798 make_request_args['checksum_config'] = checksum_config799 return make_request_args800 801 def _default_get_make_request_args(802 self,803 request_type,804 call_args,805 coordinator,806 future,807 on_done_before_calls,808 on_done_after_calls,809 ):810 return {811 'request': self._request_serializer.serialize_http_request(812 request_type, future813 ),814 'type': getattr(815 S3RequestType, request_type.upper(), S3RequestType.DEFAULT816 ),817 'on_done': self.get_crt_callback(818 future, 'done', on_done_before_calls, on_done_after_calls819 ),820 'on_progress': self.get_crt_callback(future, 'progress'),821 }822 823 824class RenameTempFileHandler:825 def __init__(self, coordinator, final_filename, temp_filename, osutil):826 self._coordinator = coordinator827 self._final_filename = final_filename828 self._temp_filename = temp_filename829 self._osutil = osutil830 831 def __call__(self, **kwargs):832 error = kwargs['error']833 if error:834 self._osutil.remove_file(self._temp_filename)835 else:836 try:837 self._osutil.rename_file(838 self._temp_filename, self._final_filename839 )840 except Exception as e:841 self._osutil.remove_file(self._temp_filename)842 # the CRT future has done already at this point843 self._coordinator.set_exception(e)844 845 846class AfterDoneHandler:847 def __init__(self, coordinator):848 self._coordinator = coordinator849 850 def __call__(self, **kwargs):851 self._coordinator.set_done_callbacks_complete()852 853 854class OnBodyFileObjWriter:855 def __init__(self, fileobj):856 self._fileobj = fileobj857 858 def __call__(self, chunk, **kwargs):859 self._fileobj.write(chunk)860 