U
    ´.¬b=  ã                   @   s†   d dl Z d dlZG dd„ deƒZG dd„ dƒZG dd„ dƒZG dd	„ d	ƒZG d
d„ dƒZG dd„ dƒZG dd„ dƒZ	G dd„ dƒZ
dS )é    Nc                       s   e Zd Z‡ fdd„Z‡  ZS )ÚRequestExceededExceptionc                    s(   || _ || _d ||¡}tƒ  |¡ dS )a   Error when requested amount exceeds what is allowed

        The request that raised this error should be retried after waiting
        the time specified by ``retry_time``.

        :type requested_amt: int
        :param requested_amt: The originally requested byte amount

        :type retry_time: float
        :param retry_time: The length in time to wait to retry for the
            requested amount
        z<Request amount {} exceeded the amount available. Retry in {}N)Úrequested_amtÚ
retry_timeÚformatÚsuperÚ__init__)Úselfr   r   Úmsg©Ú	__class__© ú8/tmp/pip-unpacked-wheel-hq9fdgne/s3transfer/bandwidth.pyr      s     ÿz!RequestExceededException.__init__)Ú__name__Ú
__module__Ú__qualname__r   Ú__classcell__r   r   r
   r   r      s   r   c                   @   s   e Zd ZdZdS )ÚRequestTokenzDA token to pass as an identifier when consuming from the LeakyBucketN)r   r   r   Ú__doc__r   r   r   r   r   '   s   r   c                   @   s   e Zd Zdd„ Zdd„ ZdS )Ú	TimeUtilsc                 C   s   t   ¡ S )zgGet the current time back

        :rtype: float
        :returns: The current time in seconds
        )Útime©r   r   r   r   r   .   s    zTimeUtils.timec                 C   s
   t  |¡S )zwSleep for a designated time

        :type value: float
        :param value: The time to sleep for in seconds
        )r   Úsleep)r   Úvaluer   r   r   r   6   s    zTimeUtils.sleepN)r   r   r   r   r   r   r   r   r   r   -   s   r   c                   @   s    e Zd Zddd„Zddd„ZdS )	ÚBandwidthLimiterNc                 C   s    || _ || _|dkrtƒ | _dS )a  Limits bandwidth for shared S3 transfers

        :type leaky_bucket: LeakyBucket
        :param leaky_bucket: The leaky bucket to use limit bandwidth

        :type time_utils: TimeUtils
        :param time_utils: Time utility to use for interacting with time.
        N)Ú_leaky_bucketÚ_time_utilsr   )r   Úleaky_bucketÚ
time_utilsr   r   r   r   @   s    	zBandwidthLimiter.__init__Tc                 C   s"   t || j|| jƒ}|s| ¡  |S )aÎ  Wraps a fileobj in a bandwidth limited stream wrapper

        :type fileobj: file-like obj
        :param fileobj: The file-like obj to wrap

        :type transfer_coordinator: s3transfer.futures.TransferCoordinator
        param transfer_coordinator: The coordinator for the general transfer
            that the wrapped stream is a part of

        :type enabled: boolean
        :param enabled: Whether bandwidth limiting should be enabled to start
        )ÚBandwidthLimitedStreamr   r   Údisable_bandwidth_limiting)r   ÚfileobjÚtransfer_coordinatorÚenabledÚstreamr   r   r   Úget_bandwith_limited_streamN   s       ÿz,BandwidthLimiter.get_bandwith_limited_stream)N)T)r   r   r   r   r$   r   r   r   r   r   ?   s   
 ÿr   c                   @   sp   e Zd Zddd„Z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dd„ Zdd„ Zdd„ ZdS )r   Né   c                 C   sF   || _ || _|| _|| _|dkr(tƒ | _d| _tƒ | _d| _|| _	dS )a[  Limits bandwidth for reads on a wrapped stream

        :type fileobj: file-like object
        :param fileobj: The file like object to wrap

        :type leaky_bucket: LeakyBucket
        :param leaky_bucket: The leaky bucket to use to throttle reads on
            the stream

        :type transfer_coordinator: s3transfer.futures.TransferCoordinator
        param transfer_coordinator: The coordinator for the general transfer
            that the wrapped stream is a part of

        :type time_utils: TimeUtils
        :param time_utils: The time utility to use for interacting with time
        NTr   )
Ú_fileobjr   Ú_transfer_coordinatorr   r   Ú_bandwidth_limiting_enabledr   Ú_request_tokenÚ_bytes_seenÚ_bytes_threshold)r   r    r   r!   r   Zbytes_thresholdr   r   r   r   f   s    zBandwidthLimitedStream.__init__c                 C   s
   d| _ dS )z0Enable bandwidth limiting on reads to the streamTN©r(   r   r   r   r   Úenable_bandwidth_limiting‰   s    z0BandwidthLimitedStream.enable_bandwidth_limitingc                 C   s
   d| _ dS )z1Disable bandwidth limiting on reads to the streamFNr,   r   r   r   r   r   �   s    z1BandwidthLimitedStream.disable_bandwidth_limitingc                 C   sL   | j s| j |¡S |  j|7  _| j| jk r8| j |¡S |  ¡  | j |¡S )zhRead a specified amount

        Reads will only be throttled if bandwidth limiting is enabled.
        )r(   r&   Úreadr*   r+   Ú_consume_through_leaky_bucket)r   Úamountr   r   r   r.   ‘   s    zBandwidthLimitedStream.readc              
   C   sf   | j jsZz| j | j| j¡ d| _W d S  tk
rV } z| j |j	¡ W 5 d }~X Y q X q | j j‚d S )Nr   )
r'   Ú	exceptionr   Úconsumer*   r)   r   r   r   r   )r   Úer   r   r   r/   ¥   s     ÿ"z4BandwidthLimitedStream._consume_through_leaky_bucketc                 C   s   |   ¡  dS )z6Signal that data being read is being transferred to S3N)r-   r   r   r   r   Úsignal_transferring·   s    z*BandwidthLimitedStream.signal_transferringc                 C   s   |   ¡  dS )z:Signal that data being read is not being transferred to S3N)r   r   r   r   r   Úsignal_not_transferring»   s    z.BandwidthLimitedStream.signal_not_transferringr   c                 C   s   | j  ||¡ d S ©N)r&   Úseek)r   ÚwhereÚwhencer   r   r   r7   ¿   s    zBandwidthLimitedStream.seekc                 C   s
   | j  ¡ S r6   )r&   Útellr   r   r   r   r:   Â   s    zBandwidthLimitedStream.tellc                 C   s"   | j r| jr|  ¡  | j ¡  d S r6   )r(   r*   r/   r&   Úcloser   r   r   r   r;   Å   s    zBandwidthLimitedStream.closec                 C   s   | S r6   r   r   r   r   r   Ú	__enter__Ï   s    z BandwidthLimitedStream.__enter__c                 O   s   |   ¡  d S r6   )r;   )r   ÚargsÚkwargsr   r   r   Ú__exit__Ò   s    zBandwidthLimitedStream.__exit__)Nr%   )r   )r   r   r   r   r-   r   r.   r/   r4   r5   r7   r:   r;   r<   r?   r   r   r   r   r   e   s     ú
#

r   c                   @   s>   e Zd Zddd„Zdd„ Zdd„ Zdd	„ Zd
d„ Zdd„ ZdS )ÚLeakyBucketNc                 C   sZ   t |ƒ| _|| _|dkr tƒ | _t ¡ | _|| _|dkr@tƒ | _|| _	|dkrVt
ƒ | _	dS )a9  A leaky bucket abstraction to limit bandwidth consumption

        :type rate: int
        :type rate: The maximum rate to allow. This rate is in terms of
            bytes per second.

        :type time_utils: TimeUtils
        :param time_utils: The time utility to use for interacting with time

        :type rate_tracker: BandwidthRateTracker
        :param rate_tracker: Tracks bandwidth consumption

        :type consumption_scheduler: ConsumptionScheduler
        :param consumption_scheduler: Schedules consumption retries when
            necessary
        N)ÚfloatÚ	_max_rater   r   Ú	threadingÚLockÚ_lockÚ_rate_trackerÚBandwidthRateTrackerÚ_consumption_schedulerÚConsumptionScheduler)r   Zmax_rater   Zrate_trackerZconsumption_schedulerr   r   r   r   ×   s    

zLeakyBucket.__init__c              
   C   sz   | j �j | j ¡ }| j |¡r8|  |||¡W  5 Q R £ S |  ||¡rT|  |||¡ n|  ||¡W  5 Q R £ S W 5 Q R X dS )ac  Consume an a requested amount

        :type amt: int
        :param amt: The amount of bytes to request to consume

        :type request_token: RequestToken
        :param request_token: The token associated to the consumption
            request that is used to identify the request. So if a
            RequestExceededException is raised the token should be used
            in subsequent retry consume() request.

        :raises RequestExceededException: If the consumption amount would
            exceed the maximum allocated bandwidth

        :rtype: int
        :returns: The amount consumed
        N)	rE   r   r   rH   Úis_scheduledÚ,_release_requested_amt_for_scheduled_requestÚ_projected_to_exceed_max_rateÚ!_raise_request_exceeded_exceptionÚ_release_requested_amt©r   ÚamtÚrequest_tokenÚtime_nowr   r   r   r2   ú   s    
  ÿ  ÿzLeakyBucket.consumec                 C   s   | j  ||¡}|| jkS r6   )rF   Úget_projected_raterB   )r   rP   rR   Zprojected_rater   r   r   rL     s    z)LeakyBucket._projected_to_exceed_max_ratec                 C   s   | j  |¡ |  ||¡S r6   )rH   Úprocess_scheduled_consumptionrN   rO   r   r   r   rK     s    ÿz8LeakyBucket._release_requested_amt_for_scheduled_requestc                 C   s.   |t | jƒ }| j |||¡}t||d�‚d S )N)r   r   )rA   rB   rH   Úschedule_consumptionr   )r   rP   rQ   rR   Zallocated_timer   r   r   r   rM   %  s      ÿ ÿz-LeakyBucket._raise_request_exceeded_exceptionc                 C   s   | j  ||¡ |S r6   )rF   Úrecord_consumption_rate)r   rP   rR   r   r   r   rN   .  s    z"LeakyBucket._release_requested_amt)NNN)	r   r   r   r   r2   rL   rK   rM   rN   r   r   r   r   r@   Ö   s      û
#	r@   c                   @   s,   e Zd Zdd„ Zdd„ Zdd„ Zdd„ Zd	S )
rI   c                 C   s   i | _ d| _dS )z*Schedules when to consume a desired amountr   N)Ú _tokens_to_scheduled_consumptionÚ_total_waitr   r   r   r   r   4  s    zConsumptionScheduler.__init__c                 C   s
   || j kS )zÙIndicates if a consumption request has been scheduled

        :type token: RequestToken
        :param token: The token associated to the consumption
            request that is used to identify the request.
        )rW   )r   Útokenr   r   r   rJ   9  s    z!ConsumptionScheduler.is_scheduledc                 C   s&   |  j |7  _ | j |dœ| j|< | j S )a´  Schedules a wait time to be able to consume an amount

        :type amt: int
        :param amt: The amount of bytes scheduled to be consumed

        :type token: RequestToken
        :param token: The token associated to the consumption
            request that is used to identify the request.

        :type time_to_consume: float
        :param time_to_consume: The desired time it should take for that
            specific request amount to be consumed in regardless of previously
            scheduled consumption requests

        :rtype: float
        :returns: The amount of time to wait for the specific request before
            actually consuming the specified amount.
        )Zwait_durationÚtime_to_consume)rX   rW   )r   rP   rY   rZ   r   r   r   rU   B  s
    þz)ConsumptionScheduler.schedule_consumptionc                 C   s&   | j  |¡}t| j|d  dƒ| _dS )zàProcesses a scheduled consumption request that has completed

        :type token: RequestToken
        :param token: The token associated to the consumption
            request that is used to identify the request.
        rZ   r   N)rW   ÚpopÚmaxrX   )r   rY   Zscheduled_retryr   r   r   rT   \  s
     ÿz2ConsumptionScheduler.process_scheduled_consumptionN)r   r   r   r   rJ   rU   rT   r   r   r   r   rI   3  s   	rI   c                   @   sB   e Zd Zddd„Zedd„ ƒZdd„ Zdd	„ Zd
d„ Zdd„ Z	dS )rG   çš™™™™™é?c                 C   s   || _ d| _d| _dS )a’  Tracks the rate of bandwidth consumption

        :type a: float
        :param a: The constant to use in calculating the exponentional moving
            average of the bandwidth rate. Specifically it is used in the
            following calculation:

            current_rate = alpha * new_rate + (1 - alpha) * current_rate

            This value of this constant should be between 0 and 1.
        N)Ú_alphaÚ
_last_timeÚ_current_rate)r   Úalphar   r   r   r   j  s    zBandwidthRateTracker.__init__c                 C   s   | j dkrdS | jS )zmThe current transfer rate

        :rtype: float
        :returns: The current tracked transfer rate
        Nç        )r_   r`   r   r   r   r   Úcurrent_ratez  s    
z!BandwidthRateTracker.current_ratec                 C   s   | j dkrdS |  ||¡S )aZ  Get the projected rate using a provided amount and time

        :type amt: int
        :param amt: The proposed amount to consume

        :type time_at_consumption: float
        :param time_at_consumption: The proposed time to consume at

        :rtype: float
        :returns: The consumption rate if that amt and time were consumed
        Nrb   )r_   Ú*_calculate_exponential_moving_average_rate©r   rP   Útime_at_consumptionr   r   r   rS   …  s    
 ÿz'BandwidthRateTracker.get_projected_ratec                 C   s2   | j dkr|| _ d| _dS |  ||¡| _|| _ dS )a  Record the consumption rate based off amount and time point

        :type amt: int
        :param amt: The amount that got consumed

        :type time_at_consumption: float
        :param time_at_consumption: The time at which the amount was consumed
        Nrb   )r_   r`   rd   re   r   r   r   rV   —  s    	
 ÿz,BandwidthRateTracker.record_consumption_ratec                 C   s"   || j  }|dkrtdƒS || S )Nr   Úinf)r_   rA   )r   rP   rf   Z
time_deltar   r   r   Ú_calculate_rate©  s    
z$BandwidthRateTracker._calculate_ratec                 C   s&   |   ||¡}| j| d| j | j  S )Né   )rh   r^   r`   )r   rP   rf   Znew_rater   r   r   rd   ³  s    z?BandwidthRateTracker._calculate_exponential_moving_average_rateN)r]   )
r   r   r   r   Úpropertyrc   rS   rV   rh   rd   r   r   r   r   rG   i  s   



rG   )rC   r   Ú	Exceptionr   r   r   r   r   r@   rI   rG   r   r   r   r   Ú<module>   s   &q]6