U
    ´.¬b»v  ã                   @   sî   d dl Z d dlmZ d dlmZmZ d dlmZ d dlm	Z	m
Z
mZmZ d dlmZmZmZmZ G dd„ dƒZG d	d
„ d
ƒZG dd„ dƒZG dd„ deƒZG dd„ deƒZG dd„ deƒZG dd„ deƒZG dd„ deƒZG dd„ deƒZdS )é    N)ÚBytesIO©ÚreadableÚseekable)ÚIN_MEMORY_UPLOAD_TAG)ÚCompleteMultipartUploadTaskÚCreateMultipartUploadTaskÚSubmissionTaskÚTask)ÚChunksizeAdjusterÚDeferredOpenFileÚget_callbacksÚget_filtered_dictc                   @   s.   e Zd Zddd„Zdd„ Zdd„ Zdd	„ Zd
S )ÚAggregatedProgressCallbacké   c                 C   s   || _ || _d| _dS )aØ  Aggregates progress updates for every provided progress callback

        :type callbacks: A list of functions that accepts bytes_transferred
            as a single argument
        :param callbacks: The callbacks to invoke when threshold is reached

        :type threshold: int
        :param threshold: The progress threshold in which to take the
            aggregated progress and invoke the progress callback with that
            aggregated progress total
        r   N)Ú
_callbacksÚ
_thresholdÚ_bytes_seen)ÚselfÚ	callbacksÚ	threshold© r   ú5/tmp/pip-unpacked-wheel-hq9fdgne/s3transfer/upload.pyÚ__init__!   s    z#AggregatedProgressCallback.__init__c                 C   s&   |  j |7  _ | j | jkr"|  ¡  d S ©N)r   r   Ú_trigger_callbacks)r   Úbytes_transferredr   r   r   Ú__call__1   s    z#AggregatedProgressCallback.__call__c                 C   s   | j dkr|  ¡  dS )z@Flushes out any progress that has not been sent to its callbacksr   N)r   r   ©r   r   r   r   Úflush6   s    
z AggregatedProgressCallback.flushc                 C   s"   | j D ]}|| jd� qd| _d S )N)r   r   )r   r   )r   Úcallbackr   r   r   r   ;   s    
z-AggregatedProgressCallback._trigger_callbacksN)r   )Ú__name__Ú
__module__Ú__qualname__r   r   r   r   r   r   r   r   r       s   
r   c                   @   sL   e Zd ZdZdd„ Zddd„Zddd	„Zd
d„ Zdd„ Zdd„ Z	dd„ Z
dS )ÚInterruptReaderaÏ  Wrapper that can interrupt reading using an error

    It uses a transfer coordinator to propagate an error if it notices
    that a read is being made while the file is being read from.

    :type fileobj: file-like obj
    :param fileobj: The file-like object to read from

    :type transfer_coordinator: s3transfer.futures.TransferCoordinator
    :param transfer_coordinator: The transfer coordinator to use if the
        reader needs to be interrupted.
    c                 C   s   || _ || _d S r   )Ú_fileobjÚ_transfer_coordinator)r   ÚfileobjÚtransfer_coordinatorr   r   r   r   O   s    zInterruptReader.__init__Nc                 C   s   | j jr| j j‚| j |¡S r   )r&   Ú	exceptionr%   Úread)r   Úamountr   r   r   r*   S   s    zInterruptReader.readr   c                 C   s   | j  ||¡ d S r   )r%   Úseek)r   ÚwhereÚwhencer   r   r   r,   ]   s    zInterruptReader.seekc                 C   s
   | j  ¡ S r   )r%   Útellr   r   r   r   r/   `   s    zInterruptReader.tellc                 C   s   | j  ¡  d S r   )r%   Úcloser   r   r   r   r0   c   s    zInterruptReader.closec                 C   s   | S r   r   r   r   r   r   Ú	__enter__f   s    zInterruptReader.__enter__c                 O   s   |   ¡  d S r   )r0   )r   ÚargsÚkwargsr   r   r   Ú__exit__i   s    zInterruptReader.__exit__)N)r   )r!   r"   r#   Ú__doc__r   r*   r,   r/   r0   r1   r4   r   r   r   r   r$   A   s   


r$   c                   @   sf   e Zd ZdZddd„Zedd„ ƒZdd„ Zd	d
„ Zdd„ Z	dd„ Z
dd„ Zdd„ Zdd„ Zdd„ ZdS )ÚUploadInputManageraJ  Base manager class for handling various types of files for uploads

    This class is typically used for the UploadSubmissionTask class to help
    determine the following:

        * How to determine the size of the file
        * How to determine if a multipart upload is required
        * How to retrieve the body for a PutObject
        * How to retrieve the bodies for a set of UploadParts

    The answers/implementations differ for the various types of file inputs
    that may be accepted. All implementations must subclass and override
    public methods from this class.
    Nc                 C   s   || _ || _|| _d S r   )Ú_osutilr&   Ú_bandwidth_limiter©r   Úosutilr(   Úbandwidth_limiterr   r   r   r   }   s    zUploadInputManager.__init__c                 C   s   t dƒ‚dS )a  Determines if the source for the upload is compatible with manager

        :param upload_source: The source for which the upload will pull data
            from.

        :returns: True if the manager can handle the type of source specified
            otherwise returns False.
        zmust implement _is_compatible()N©ÚNotImplementedError©ÚclsZupload_sourcer   r   r   Úis_compatible‚   s    
z UploadInputManager.is_compatiblec                 C   s   t dƒ‚dS )aÛ  Whether the body it provides are stored in-memory

        :type operation_name: str
        :param operation_name: The name of the client operation that the body
            is being used for. Valid operation_names are ``put_object`` and
            ``upload_part``.

        :rtype: boolean
        :returns: True if the body returned by the manager will be stored in
            memory. False if the manager will not directly store the body in
            memory.
        z%must implement store_body_in_memory()Nr<   ©r   Úoperation_namer   r   r   Ústores_body_in_memoryŽ   s    z(UploadInputManager.stores_body_in_memoryc                 C   s   t dƒ‚dS )z¼Provides the transfer size of an upload

        :type transfer_future: s3transfer.futures.TransferFuture
        :param transfer_future: The future associated with upload request
        z&must implement provide_transfer_size()Nr<   ©r   Útransfer_futurer   r   r   Úprovide_transfer_size�   s    z(UploadInputManager.provide_transfer_sizec                 C   s   t dƒ‚dS )aÔ  Determines where a multipart upload is required

        :type transfer_future: s3transfer.futures.TransferFuture
        :param transfer_future: The future associated with upload request

        :type config: s3transfer.manager.TransferConfig
        :param config: The config associated to the transfer manager

        :rtype: boolean
        :returns: True, if the upload should be multipart based on
            configuration and size. False, otherwise.
        z*must implement requires_multipart_upload()Nr<   ©r   rE   Úconfigr   r   r   Úrequires_multipart_upload¥   s    z,UploadInputManager.requires_multipart_uploadc                 C   s   t dƒ‚dS )aÜ  Returns the body to use for PutObject

        :type transfer_future: s3transfer.futures.TransferFuture
        :param transfer_future: The future associated with upload request

        :type config: s3transfer.manager.TransferConfig
        :param config: The config associated to the transfer manager

        :rtype: s3transfer.utils.ReadFileChunk
        :returns: A ReadFileChunk including all progress callbacks
            associated with the transfer future.
        z$must implement get_put_object_body()Nr<   rD   r   r   r   Úget_put_object_body´   s    z&UploadInputManager.get_put_object_bodyc                 C   s   t dƒ‚dS )a  Yields the part number and body to use for each UploadPart

        :type transfer_future: s3transfer.futures.TransferFuture
        :param transfer_future: The future associated with upload request

        :type chunksize: int
        :param chunksize: The chunksize to use for this upload.

        :rtype: int, s3transfer.utils.ReadFileChunk
        :returns: Yields the part number and the ReadFileChunk including all
            progress callbacks associated with the transfer future for that
            specific yielded part.
        z)must implement yield_upload_part_bodies()Nr<   )r   rE   Ú	chunksizer   r   r   Úyield_upload_part_bodiesÃ   s    z+UploadInputManager.yield_upload_part_bodiesc                 C   s*   t || jƒ}| jr&| jj|| jdd�}|S )NF)Úenabled)r$   r&   r8   Zget_bandwith_limited_stream)r   r'   r   r   r   Ú_wrap_fileobjÓ   s      ÿz UploadInputManager._wrap_fileobjc                 C   s   t |dƒ}|rt|ƒgS g S )NÚprogress)r   r   )r   rE   r   r   r   r   Ú_get_progress_callbacksÛ   s    

z*UploadInputManager._get_progress_callbacksc                 C   s   dd„ |D ƒS )Nc                 S   s   g | ]
}|j ‘qS r   )r   )Ú.0r    r   r   r   Ú
<listcomp>å   s     z;UploadInputManager._get_close_callbacks.<locals>.<listcomp>r   )r   Zaggregated_progress_callbacksr   r   r   Ú_get_close_callbacksä   s    z'UploadInputManager._get_close_callbacks)N)r!   r"   r#   r5   r   Úclassmethodr@   rC   rF   rI   rJ   rL   rN   rP   rS   r   r   r   r   r6   m   s   

	r6   c                   @   sd   e Zd ZdZedd„ ƒZdd„ Zdd„ Zdd	„ Zd
d„ Z	dd„ Z
dd„ Zdd„ Zdd„ Zdd„ ZdS )ÚUploadFilenameInputManagerzUpload utility for filenamesc                 C   s
   t |tƒS r   )Ú
isinstanceÚstrr>   r   r   r   r@   ë   s    z(UploadFilenameInputManager.is_compatiblec                 C   s   dS )NFr   rA   r   r   r   rC   ï   s    z0UploadFilenameInputManager.stores_body_in_memoryc                 C   s   |j  | j |j jj¡¡ d S r   )ÚmetarF   r7   Zget_file_sizeÚ	call_argsr'   rD   r   r   r   rF   ò   s    ÿz0UploadFilenameInputManager.provide_transfer_sizec                 C   s   |j j|jkS r   )rX   ÚsizeÚmultipart_thresholdrG   r   r   r   rI   ÷   s    z4UploadFilenameInputManager.requires_multipart_uploadc                 C   sJ   |   |¡\}}|  |¡}|  |¡}|  |¡}|jj}| jj|||||d�S )N©r'   Ú
chunk_sizeÚfull_file_sizer   Úclose_callbacks)Ú&_get_put_object_fileobj_with_full_sizerN   rP   rS   rX   rZ   r7   Ú#open_file_chunk_reader_from_fileobj)r   rE   r'   Ú	full_sizer   r_   rZ   r   r   r   rJ   ú   s    ÿ


ûz.UploadFilenameInputManager.get_put_object_bodyc                 c   s’   |j j}|  ||¡}td|d ƒD ]j}|  |¡}|  |¡}||d  }| j|j jj|||d�\}	}
|  	|	¡}	| j
j|	||
||d�}||fV  q"d S )Né   )Ú
start_byteÚ	part_sizer^   r\   )rX   rZ   Ú_get_num_partsÚrangerP   rS   Ú'_get_upload_part_fileobj_with_full_sizerY   r'   rN   r7   ra   )r   rE   rK   r^   Z	num_partsÚpart_numberr   r_   rd   r'   rb   Zread_file_chunkr   r   r   rL     s*    

ü


ûz3UploadFilenameInputManager.yield_upload_part_bodiesc                 C   s   t ||| jjd�}|S )N)Zopen_function)r   r7   Úopen)r   r'   rd   r   r   r   Ú_get_deferred_open_file1  s      ÿz2UploadFilenameInputManager._get_deferred_open_filec                 C   s"   |j jj}|j j}|  |d¡|fS )Nr   )rX   rY   r'   rZ   rk   ©r   rE   r'   rZ   r   r   r   r`   7  s    
zAUploadFilenameInputManager._get_put_object_fileobj_with_full_sizec                 K   s    |d }|d }|   ||¡|fS )Nrd   r^   )rk   )r   r'   r3   rd   rb   r   r   r   rh   <  s    zBUploadFilenameInputManager._get_upload_part_fileobj_with_full_sizec                 C   s   t t |jjt|ƒ ¡ƒS r   )ÚintÚmathÚceilrX   rZ   Úfloat)r   rE   re   r   r   r   rf   A  s    z)UploadFilenameInputManager._get_num_partsN)r!   r"   r#   r5   rT   r@   rC   rF   rI   rJ   rL   rk   r`   rh   rf   r   r   r   r   rU   è   s   
rU   c                   @   s<   e Zd ZdZedd„ ƒZdd„ Zdd„ Zdd	„ Zd
d„ Z	dS )ÚUploadSeekableInputManagerz&Upload utility for an open file objectc                 C   s   t |ƒot|ƒS r   r   r>   r   r   r   r@   H  s    z(UploadSeekableInputManager.is_compatiblec                 C   s   |dkrdS dS d S )NÚ
put_objectFTr   rA   r   r   r   rC   L  s    z0UploadSeekableInputManager.stores_body_in_memoryc                 C   sD   |j jj}| ¡ }| dd¡ | ¡ }| |¡ |j  || ¡ d S )Nr   é   )rX   rY   r'   r/   r,   rF   )r   rE   r'   Zstart_positionZend_positionr   r   r   rF   R  s    

ÿz0UploadSeekableInputManager.provide_transfer_sizec                 K   s   |  |d ¡}t|ƒt|ƒfS )Nre   )r*   r   Úlen)r   r'   r3   Údatar   r   r   rh   _  s    zBUploadSeekableInputManager._get_upload_part_fileobj_with_full_sizec                 C   s"   |j jj}| ¡ |j j }||fS r   )rX   rY   r'   r/   rZ   rl   r   r   r   r`   n  s    
zAUploadSeekableInputManager._get_put_object_fileobj_with_full_sizeN)
r!   r"   r#   r5   rT   r@   rC   rF   rh   r`   r   r   r   r   rq   E  s   
rq   c                       sh   e Zd ZdZd‡ fdd„	Zedd„ ƒZdd„ Zd	d
„ Zdd„ Z	dd„ Z
dd„ Zddd„Zdd„ Z‡  ZS )ÚUploadNonSeekableInputManagerz7Upload utility for a file-like object that cannot seek.Nc                    s   t ƒ  |||¡ d| _d S )Nó    )Úsuperr   Ú_initial_datar9   ©Ú	__class__r   r   r   y  s    z&UploadNonSeekableInputManager.__init__c                 C   s   t |ƒS r   )r   r>   r   r   r   r@   }  s    z+UploadNonSeekableInputManager.is_compatiblec                 C   s   dS )NTr   rA   r   r   r   rC   �  s    z3UploadNonSeekableInputManager.stores_body_in_memoryc                 C   s   d S r   r   rD   r   r   r   rF   „  s    z3UploadNonSeekableInputManager.provide_transfer_sizec                 C   sT   |j jd k	r|j j|jkS |j jj}|j}|  ||d¡| _t| jƒ|k rLdS dS d S )NFT)rX   rZ   r[   rY   r'   Ú_readry   rt   )r   rE   rH   r'   r   r   r   r   rI   ‰  s    
z7UploadNonSeekableInputManager.requires_multipart_uploadc                 C   s@   |   |¡}|  |¡}|jjj}|  | j| ¡  ||¡}d | _|S r   )rP   rS   rX   rY   r'   Ú
_wrap_datary   r*   )r   rE   r   r_   r'   Úbodyr   r   r   rJ   š  s    


  ÿz1UploadNonSeekableInputManager.get_put_object_bodyc           	      c   s^   |j jj}d}|  |¡}|  |¡}|d7 }|  ||¡}|s<qZ|  |||¡}d }||fV  qd S )Nr   rc   )rX   rY   r'   rP   rS   r|   r}   )	r   rE   rK   Zfile_objectri   r   r_   Zpart_contentZpart_objectr   r   r   rL   ¨  s    


  ÿz6UploadNonSeekableInputManager.yield_upload_part_bodiesTc                 C   sx   t | jƒdkr| |¡S |t | jƒkrL| jd|… }|rH| j|d… | _|S |t | jƒ }| j| |¡ }|rtd| _|S )a=  
        Reads a specific amount of data from a stream and returns it. If there
        is any data in initial_data, that will be popped out first.

        :type fileobj: A file-like object that implements read
        :param fileobj: The stream to read from.

        :type amount: int
        :param amount: The number of bytes to read from the stream.

        :type truncate: bool
        :param truncate: Whether or not to truncate initial_data after
            reading from it.

        :return: Generator which generates part bodies from the initial data.
        r   Nrw   )rt   ry   r*   )r   r'   r+   Útruncateru   Zamount_to_readr   r   r   r|   ¼  s    
z#UploadNonSeekableInputManager._readc                 C   s,   |   t|ƒ¡}| jj|t|ƒt|ƒ||d�S )a¸  
        Wraps data with the interrupt reader and the file chunk reader.

        :type data: bytes
        :param data: The data to wrap.

        :type callbacks: list
        :param callbacks: The callbacks associated with the transfer future.

        :type close_callbacks: list
        :param close_callbacks: The callbacks to be called when closing the
            wrapper for the data.

        :return: Fully wrapped data.
        r\   )rN   r   r7   ra   rt   )r   ru   r   r_   r'   r   r   r   r}   æ  s    ûz(UploadNonSeekableInputManager._wrap_data)N)T)r!   r"   r#   r5   r   rT   r@   rC   rF   rI   rJ   rL   r|   r}   Ú__classcell__r   r   rz   r   rv   v  s   

*rv   c                   @   sb   e Zd ZdZddddddgZddgZdd	„ Zddd„Zdd„ Zdd„ Z	dd„ Z
dd„ Zdd„ Zd
S )ÚUploadSubmissionTaskz.Task for submitting tasks to execute an uploadÚChecksumAlgorithmZSSECustomerKeyZSSECustomerAlgorithmZSSECustomerKeyMD5ZRequestPayerZExpectedBucketOwnerc                 C   sH   t ttg}|jjj}|D ]}| |¡r|  S qtd |t	|ƒ¡ƒ‚dS )ao  Retrieves a class for managing input for an upload based on file type

        :type transfer_future: s3transfer.futures.TransferFuture
        :param transfer_future: The transfer future for the request

        :rtype: class of UploadInputManager
        :returns: The appropriate class to use for managing a specific type of
            input for uploads.
        z&Input {} of type: {} is not supported.N)
rU   rq   rv   rX   rY   r'   r@   ÚRuntimeErrorÚformatÚtype)r   rE   Zupload_manager_resolver_chainr'   Zupload_manager_clsr   r   r   Ú_get_upload_input_manager_cls  s    ý


 ÿÿz2UploadSubmissionTask._get_upload_input_manager_clsNc                 C   sd   |   |¡|| j|ƒ}|jjdkr*| |¡ | ||¡sL|  ||||||¡ n|  ||||||¡ dS )aÒ  
        :param client: The client associated with the transfer manager

        :type config: s3transfer.manager.TransferConfig
        :param config: The transfer config associated with the transfer
            manager

        :type osutil: s3transfer.utils.OSUtil
        :param osutil: The os utility associated to the transfer manager

        :type request_executor: s3transfer.futures.BoundedExecutor
        :param request_executor: The request executor associated with the
            transfer manager

        :type transfer_future: s3transfer.futures.TransferFuture
        :param transfer_future: The transfer future associated with the
            transfer request that tasks are being submitted for
        N)r†   r&   rX   rZ   rF   rI   Ú_submit_upload_requestÚ_submit_multipart_request)r   ÚclientrH   r:   Úrequest_executorrE   r;   Úupload_input_managerr   r   r   Ú_submit(  s:    ÿ  þ
 ÿú	úzUploadSubmissionTask._submitc           	   
   C   sN   |j j}|  |d¡}| jj|t| j|| |¡|j|j|j	dœdd�|d� d S )Nrr   )r‰   r'   ÚbucketÚkeyÚ
extra_argsT)r(   Úmain_kwargsÚis_final©Útag)
rX   rY   Ú_get_upload_task_tagr&   ÚsubmitÚPutObjectTaskrJ   r�   rŽ   r�   )	r   r‰   rH   r:   rŠ   rE   r‹   rY   Zput_object_tagr   r   r   r‡   a  s*    	 ÿÿù	õñz+UploadSubmissionTask._submit_upload_requestc                 C   sü   |j j}| j |t| j||j|j|jdœd�¡}g }	|  |j¡}
|  	|d¡}|j j
}tƒ }| |j|¡}| ||¡}|D ]B\}}|	 | jj|t| j|||j|j||
dœd|id�|d�¡ qx|  |j¡}| j |t| j||j|j|dœ||	dœd	d
�¡ d S )N)r‰   r�   rŽ   r�   )r(   r�   Úupload_part)r‰   r'   r�   rŽ   ri   r�   Ú	upload_id)r(   r�   Úpending_main_kwargsr’   )r˜   ÚpartsT)r(   r�   r™   r‘   )rX   rY   r&   r•   r   r�   rŽ   r�   Ú_extra_upload_part_argsr”   rZ   r   Zadjust_chunksizeZmultipart_chunksizerL   ÚappendÚUploadPartTaskÚ_extra_complete_multipart_argsr   )r   r‰   rH   r:   rŠ   rE   r‹   rY   Zcreate_multipart_futureZpart_futuresZextra_part_argsZupload_part_tagrZ   ZadjusterrK   Zpart_iteratorri   r'   Zcomplete_multipart_extra_argsr   r   r   rˆ   „  s~    	üþþ ÿ ÿú	 ÿöðÿÿüþôþz.UploadSubmissionTask._submit_multipart_requestc                 C   s   t || jƒS r   )r   ÚUPLOAD_PART_ARGS©r   r�   r   r   r   r›   Ú  s    z,UploadSubmissionTask._extra_upload_part_argsc                 C   s   t || jƒS r   )r   ÚCOMPLETE_MULTIPART_ARGSr    r   r   r   rž   ß  s    z3UploadSubmissionTask._extra_complete_multipart_argsc                 C   s   d }|  |¡rt}|S r   )rC   r   )r   r‹   rB   r“   r   r   r   r”   â  s    
z)UploadSubmissionTask._get_upload_task_tag)N)r!   r"   r#   r5   rŸ   r¡   r†   rŒ   r‡   rˆ   r›   rž   r”   r   r   r   r   r�      s"   ú	! ù
9#Vr�   c                   @   s   e Zd ZdZdd„ ZdS )r–   z Task to do a nonmultipart uploadc              	   C   s,   |�}|j f |||dœ|—Ž W 5 Q R X dS )aP  
        :param client: The client to use when calling PutObject
        :param fileobj: The file to upload.
        :param bucket: The name of the bucket to upload to
        :param key: The name of the key to upload to
        :param extra_args: A dictionary of any extra arguments that may be
            used in the upload.
        )ÚBucketÚKeyÚBodyN)rr   )r   r‰   r'   r�   rŽ   r�   r~   r   r   r   Ú_mainì  s    	zPutObjectTask._mainN©r!   r"   r#   r5   r¥   r   r   r   r   r–   é  s   r–   c                   @   s   e Zd ZdZdd„ ZdS )r�   z+Task to upload a part in a multipart uploadc              	   C   st   |�"}|j f |||||dœ|—Ž}	W 5 Q R X |	d }
|
|dœ}d|krp|d  ¡ }d|› �}||	krp|	| ||< |S )aÓ  
        :param client: The client to use when calling PutObject
        :param fileobj: The file to upload.
        :param bucket: The name of the bucket to upload to
        :param key: The name of the key to upload to
        :param upload_id: The id of the upload
        :param part_number: The number representing the part of the multipart
            upload
        :param extra_args: A dictionary of any extra arguments that may be
            used in the upload.

        :rtype: dict
        :returns: A dictionary representing a part::

            {'Etag': etag_value, 'PartNumber': part_number}

            This value can be appended to a list to be used to complete
            the multipart upload.
        )r¢   r£   ZUploadIdÚ
PartNumberr¤   ÚETag)r¨   r§   r‚   ZChecksum)r—   Úupper)r   r‰   r'   r�   rŽ   r˜   ri   r�   r~   ÚresponseÚetagZpart_metadataZalgorithm_nameZchecksum_memberr   r   r   r¥   ü  s$    ûú

zUploadPartTask._mainNr¦   r   r   r   r   r�   ù  s   r�   )rn   Úior   Zs3transfer.compatr   r   Zs3transfer.futuresr   Zs3transfer.tasksr   r   r	   r
   Zs3transfer.utilsr   r   r   r   r   r$   r6   rU   rq   rv   r�   r–   r�   r   r   r   r   Ú<module>   s    !,{]1  j