Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes14kdownloads
crt.py860 linesDownload Raw Back to s3transfer
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 
codekingpro/portable-devtools · Team Ai