o
    Àyð`Ås  ã                   @   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	 d dlm
Z
 d dlmZ d dlmZ d d	lmZ d d
lmZ d dlmZm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G dd„ de
ƒZG dd„ de	ƒZG dd„ de	ƒZdS )é    N)Úsix)ÚseekableÚreadable)ÚIN_MEMORY_UPLOAD_TAG)ÚTask)ÚSubmissionTask)ÚCreateMultipartUploadTask)ÚCompleteMultipartUploadTask)Úget_callbacks)Úget_filtered_dict)ÚDeferredOpenFileÚChunksizeAdjusterc                   @   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   ú3/usr/lib/python3/dist-packages/s3transfer/upload.pyÚ__init__   s   
z#AggregatedProgressCallback.__init__c                 C   s*   |  j |7  _ | j | jkr|  ¡  d S d S ©N)r   r   Ú_trigger_callbacks)r   Úbytes_transferredr   r   r   Ú__call__-   s   ÿz#AggregatedProgressCallback.__call__c                 C   s   | j dkr|  ¡  dS dS )z@Flushes out any progress that has not been sent to its callbacksr   N)r   r   ©r   r   r   r   Úflush2   s   
ÿz AggregatedProgressCallback.flushc                 C   s"   | j D ]}|| jd� qd| _d S )N)r   r   )r   r   )r   Úcallbackr   r   r   r   7   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   J   s   
zInterruptReader.__init__Nc                 C   s   | j jr| j j‚| j |¡S r   )r%   Ú	exceptionr$   Úread)r   Úamountr   r   r   r)   N   s   zInterruptReader.readr   c                 C   s   | j  ||¡ d S r   )r$   Úseek)r   ÚwhereÚwhencer   r   r   r+   X   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   r/   ^   ó   zInterruptReader.closec                 C   s   | S r   r   r   r   r   r   Ú	__enter__a   ó   zInterruptReader.__enter__c                 O   s   |   ¡  d S r   )r/   )r   ÚargsÚkwargsr   r   r   Ú__exit__d   ó   zInterruptReader.__exit__r   )r   )r    r!   r"   Ú__doc__r   r)   r+   r.   r/   r1   r5   r   r   r   r   r#   =   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   w   s   
zUploadInputManager.__init__c                 C   ó   t dƒ‚)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()©ÚNotImplementedError©ÚclsÚupload_sourcer   r   r   Úis_compatible|   s   
z UploadInputManager.is_compatiblec                 C   r>   )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())ÚNotImplemented©r   Úoperation_namer   r   r   Ústores_body_in_memoryˆ   ó   z(UploadInputManager.stores_body_in_memoryc                 C   r>   )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()r?   ©r   Útransfer_futurer   r   r   Úprovide_transfer_size—   s   z(UploadInputManager.provide_transfer_sizec                 C   r>   )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
            configuartion and size. False, otherwise.
        z*must implement requires_multipart_upload()r?   ©r   rK   Úconfigr   r   r   Úrequires_multipart_uploadŸ   rI   z,UploadInputManager.requires_multipart_uploadc                 C   r>   )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()r?   rJ   r   r   r   Úget_put_object_body®   rI   z&UploadInputManager.get_put_object_bodyc                 C   r>   )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()r?   )r   rK   Ú	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%   r:   Ú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   rK   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   Úaggregated_progress_callbacksr   r   r   Ú_get_close_callbacksÝ   r0   z'UploadInputManager._get_close_callbacksr   )r    r!   r"   r7   r   ÚclassmethodrD   rH   rL   rO   rP   rR   rU   rW   r[   r   r   r   r   r8   h   s    

	r8   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jƒS r   )Ú
isinstancer   Ústring_typesrA   r   r   r   rD   ã   s   z(UploadFilenameInputManager.is_compatiblec                 C   ó   dS )NFr   rF   r   r   r   rH   ç   r2   z0UploadFilenameInputManager.stores_body_in_memoryc                 C   s   |j  | j |j jj¡¡ d S r   )ÚmetarL   r9   Úget_file_sizeÚ	call_argsr&   rJ   r   r   r   rL   ê   s
   ÿÿz0UploadFilenameInputManager.provide_transfer_sizec                 C   s   |j j|jkS r   )ra   ÚsizeÚmultipart_thresholdrM   r   r   r   rO   ï   r0   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_sizerU   rW   r[   ra   rd   r9   Ú#open_file_chunk_reader_from_fileobj)r   rK   r&   Ú	full_sizer   ri   rd   r   r   r   rP   ò   s   ÿ


þz.UploadFilenameInputManager.get_put_object_bodyc                 c   s”   � |j j}|  ||¡}td|d ƒD ]5}|  |¡}|  |¡}||d  }| j|j jj|||d�\}	}
|  	|	¡}	| j
j|	||
||d�}||fV  qd S )Né   )Ú
start_byteÚ	part_sizerh   rf   )ra   rd   Ú_get_num_partsÚrangerW   r[   Ú'_get_upload_part_fileobj_with_full_sizerc   r&   rU   r9   rk   )r   rK   rQ   rh   Ú	num_partsÚpart_numberr   ri   rn   r&   rl   Úread_file_chunkr   r   r   rR     s&   €



þ
ýìz3UploadFilenameInputManager.yield_upload_part_bodiesc                 C   s   t ||| jjd�}|S )N)Úopen_function)r   r9   Úopen)r   r&   rn   r   r   r   Ú_get_deferred_open_file  s   
ÿz2UploadFilenameInputManager._get_deferred_open_filec                 C   s"   |j jj}|j j}|  |d¡|fS )Nr   )ra   rc   r&   rd   rx   ©r   rK   r&   rd   r   r   r   rj   #  s   
zAUploadFilenameInputManager._get_put_object_fileobj_with_full_sizec                 K   s    |d }|d }|   ||¡|fS )Nrn   rh   )rx   )r   r&   r4   rn   rl   r   r   r   rr   (  s   zBUploadFilenameInputManager._get_upload_part_fileobj_with_full_sizec                 C   s   t t |jjt|ƒ ¡ƒS r   )ÚintÚmathÚceilra   rd   Úfloat)r   rK   ro   r   r   r   rp   -  s   ÿz)UploadFilenameInputManager._get_num_partsN)r    r!   r"   r7   r\   rD   rH   rL   rO   rP   rR   rx   rj   rr   rp   r   r   r   r   r]   á   s    
r]   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   rA   r   r   r   rD   4  s   z(UploadSeekableInputManager.is_compatiblec                 C   s   |dkrdS dS )NÚ
put_objectFTr   rF   r   r   r   rH   8  s   z0UploadSeekableInputManager.stores_body_in_memoryc                 C   sD   |j jj}| ¡ }| dd¡ | ¡ }| |¡ |j  || ¡ d S )Nr   é   )ra   rc   r&   r.   r+   rL   )r   rK   r&   Ústart_positionÚend_positionr   r   r   rL   >  s   

ÿz0UploadSeekableInputManager.provide_transfer_sizec                 K   s    |  |d ¡}t |¡t|ƒfS )Nro   )r)   r   ÚBytesIOÚlen)r   r&   r4   Údatar   r   r   rr   J  s   zBUploadSeekableInputManager._get_upload_part_fileobj_with_full_sizec                 C   s"   |j jj}| ¡ |j j }||fS r   )ra   rc   r&   r.   rd   ry   r   r   r   rj   Y  s   
zAUploadSeekableInputManager._get_put_object_fileobj_with_full_sizeN)
r    r!   r"   r7   r\   rD   rH   rL   rr   rj   r   r   r   r   r~   2  s    
r~   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 t| ƒ |||¡ d| _d S )Nó    )Úsuperr†   r   Ú_initial_datar;   ©Ú	__class__r   r   r   c  s   
ÿ
z&UploadNonSeekableInputManager.__init__c                 C   s   t |ƒS r   )r   rA   r   r   r   rD   h  s   z+UploadNonSeekableInputManager.is_compatiblec                 C   r`   )NTr   rF   r   r   r   rH   l  r2   z3UploadNonSeekableInputManager.stores_body_in_memoryc                 C   s   d S r   r   rJ   r   r   r   rL   o  s   z3UploadNonSeekableInputManager.provide_transfer_sizec                 C   sP   |j jd ur|j j|jkS |j jj}|j}|  ||d¡| _t| jƒ|k r&dS dS )NFT)ra   rd   re   rc   r&   Ú_readr‰   r„   )r   rK   rN   r&   r   r   r   r   rO   t  s   
z7UploadNonSeekableInputManager.requires_multipart_uploadc                 C   s@   |   |¡}|  |¡}|jjj}|  | j| ¡  ||¡}d | _|S r   )rW   r[   ra   rc   r&   Ú
_wrap_datar‰   r)   )r   rK   r   ri   r&   Úbodyr   r   r   rP   …  s   


ÿz1UploadNonSeekableInputManager.get_put_object_bodyc           	      c   s`   � |j jj}d}	 |  |¡}|  |¡}|d7 }|  ||¡}|s!d S |  |||¡}d }||fV  q	)Nr   Trm   )ra   rc   r&   rW   r[   rŒ   r�   )	r   rK   rQ   Úfile_objectrt   r   ri   Úpart_contentÚpart_objectr   r   r   rR   ’  s    €


ÿ
ôz6UploadNonSeekableInputManager.yield_upload_part_bodiesTc                 C   sx   t | jƒdkr| |¡S |t | jƒkr&| jd|… }|r$| j|d… | _|S |t | jƒ }| j| |¡ }|r:d| _|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   Nr‡   )r„   r‰   r)   )r   r&   r*   Útruncater…   Ú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.
        rf   )rU   r   rƒ   r9   rk   r„   )r   r…   r   ri   r&   r   r   r   r�   Ï  s
   þz(UploadNonSeekableInputManager._wrap_datar   )T)r    r!   r"   r7   r   r\   rD   rH   rL   rO   rP   rR   rŒ   r�   Ú__classcell__r   r   rŠ   r   r†   a  s    

*r†   c                   @   s\   e Zd ZdZg d¢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)ÚSSECustomerKeyÚSSECustomerAlgorithmÚSSECustomerKeyMD5ÚRequestPayerÚExpectedBucketOwnerr™   rš   c                 C   sD   t ttg}|jjj}|D ]}| |¡r|  S qtd|t|ƒf ƒ‚)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 %s of type: %s is not supported.)	r]   r~   r†   ra   rc   r&   rD   ÚRuntimeErrorÚtype)r   rK   Úupload_manager_resolver_chainr&   Úupload_manager_clsr   r   r   Ú_get_upload_input_manager_clsõ  s   ý

ÿÿÿz2UploadSubmissionTask._get_upload_input_manager_clsNc                 C   sf   |   |¡|| j|ƒ}|jjdu r| |¡ | ||¡s'|  ||||||¡ dS |  ||||||¡ 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%   ra   rd   rL   rO   Ú_submit_upload_requestÚ_submit_multipart_request)r   ÚclientrN   r<   Úrequest_executorrK   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 )Nr   )r¢   r&   ÚbucketÚkeyÚ
extra_argsT)r'   Úmain_kwargsÚis_final©Útag)
ra   rc   Ú_get_upload_task_tagr%   ÚsubmitÚPutObjectTaskrP   r¦   r§   r¨   )	r   r¢   rN   r<   r£   rK   r¤   rc   Úput_object_tagr   r   r   r    4  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 ]!\}}|	 | jj|t| j|||j|j||
dœd|id�|d�¡ q<|  |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§   rt   r¨   Ú	upload_id)r'   r©   Úpending_main_kwargsr«   )r²   ÚpartsT)r'   r©   r³   rª   )ra   rc   r%   r®   r   r¦   r§   r¨   Ú_extra_upload_part_argsr­   rd   r   Úadjust_chunksizeÚmultipart_chunksizerR   ÚappendÚUploadPartTaskÚ_extra_complete_multipart_argsr	   )r   r¢   rN   r<   r£   rK   r¤   rc   Úcreate_multipart_futureÚpart_futuresÚextra_part_argsÚupload_part_tagrd   ÚadjusterrQ   Úpart_iteratorrt   r&   Úcomplete_multipart_extra_argsr   r   r   r¡   N  sx   üþþÿÿú	ÿöðÿÿüþôþz.UploadSubmissionTask._submit_multipart_requestc                 C   ó   t || jƒS r   )r   ÚUPLOAD_PART_ARGS©r   r¨   r   r   r   rµ   ›  s   z,UploadSubmissionTask._extra_upload_part_argsc                 C   rÂ   r   )r   ÚCOMPLETE_MULTIPART_ARGSrÄ   r   r   r   rº      r6   z3UploadSubmissionTask._extra_complete_multipart_argsc                 C   s   d }|  |¡r	t}|S r   )rH   r   )r   r¤   rG   r¬   r   r   r   r­   £  s   
z)UploadSubmissionTask._get_upload_task_tagr   )r    r!   r"   r7   rÃ   rÅ   rŸ   r¥   r    r¡   rµ   rº   r­   r   r   r   r   r•   å  s    	þ
ÿ'Mr•   c                   @   ó   e Zd ZdZdd„ ZdS )r¯   z Task to do a nonmultipart uploadc                 C   sB   |�}|j d|||dœ|¤Ž W d  ƒ dS 1 sw   Y  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ÚBodyNr   )r   )r   r¢   r&   r¦   r§   r¨   rŽ   r   r   r   Ú_main¬  s   	"ÿzPutObjectTask._mainN©r    r!   r"   r7   rÊ   r   r   r   r   r¯   ª  ó    r¯   c                   @   rÆ   )r¹   z+Task to upload a part in a multipart uploadc              	   C   sR   |�}|j d|||||dœ|¤Ž}	W d  ƒ n1 sw   Y  |	d }
|
|dœ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È   ÚUploadIdÚ
PartNumberrÉ   NÚETag)rÏ   rÎ   r   )r±   )r   r¢   r&   r¦   r§   r²   rt   r¨   rŽ   ÚresponseÚetagr   r   r   rÊ   »  s   ýýÿ
zUploadPartTask._mainNrË   r   r   r   r   r¹   ¹  rÌ   r¹   )r{   Úbotocore.compatr   Ús3transfer.compatr   r   Ús3transfer.futuresr   Ús3transfer.tasksr   r   r   r	   Ús3transfer.utilsr
   r   r   r   Úobjectr   r#   r8   r]   r~   r†   r•   r¯   r¹   r   r   r   r   Ú<module>   s,   !+yQ/  F