o
    Qœ_D  ã                   @   sî   d dl Z d dlZd dlZd dlmZmZ d dlmZmZ d dl	m
Z
mZmZmZmZ d dl	Z	e	jr9d dl	mZmZ g 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G dd„ deƒZdS )é    N)ÚgenÚioloop)ÚFutureÚ"future_set_result_unless_cancelled)ÚUnionÚOptionalÚTypeÚAnyÚ	Awaitable)ÚDequeÚSet)Ú	ConditionÚEventÚ	SemaphoreÚBoundedSemaphoreÚLockc                   @   s$   e Zd ZdZddd„Zddd„ZdS )	Ú_TimeoutGarbageCollectorzáBase class for objects that periodically clean up timed-out waiters.

    Avoids memory leak in a common pattern like:

        while True:
            yield condition.wait(short_timeout)
            print('looping....')
    ÚreturnNc                 C   s   t  ¡ | _d| _d S )Nr   )ÚcollectionsÚdequeÚ_waitersÚ	_timeouts©Úself© r   ú//usr/lib/python3/dist-packages/tornado/locks.pyÚ__init__)   s   

z!_TimeoutGarbageCollector.__init__c                 C   s>   |  j d7  _ | j dkrd| _ t dd„ | jD ƒ¡| _d S d S )Né   éd   r   c                 s   s   � | ]	}|  ¡ s|V  qd S ©N)Údone)Ú.0Úwr   r   r   Ú	<genexpr>2   s   € z<_TimeoutGarbageCollector._garbage_collect.<locals>.<genexpr>)r   r   r   r   r   r   r   r   Ú_garbage_collect-   s
   
þz)_TimeoutGarbageCollector._garbage_collect©r   N)Ú__name__Ú
__module__Ú__qualname__Ú__doc__r   r$   r   r   r   r   r      s    
	r   c                       sv   e Zd ZdZd‡ fdd„Zdefdd„Z	ddeee	e
jf  dee fd	d
„Zddeddfdd„Zddd„Z‡  ZS )r   aÂ  A condition allows one or more coroutines to wait until notified.

    Like a standard `threading.Condition`, but does not need an underlying lock
    that is acquired and released.

    With a `Condition`, coroutines can wait to be notified by other coroutines:

    .. testcode::

        from tornado import gen
        from tornado.ioloop import IOLoop
        from tornado.locks import Condition

        condition = Condition()

        async def waiter():
            print("I'll wait right here")
            await condition.wait()
            print("I'm done waiting")

        async def notifier():
            print("About to notify")
            condition.notify()
            print("Done notifying")

        async def runner():
            # Wait for waiter() and notifier() in parallel
            await gen.multi([waiter(), notifier()])

        IOLoop.current().run_sync(runner)

    .. testoutput::

        I'll wait right here
        About to notify
        Done notifying
        I'm done waiting

    `wait` takes an optional ``timeout`` argument, which is either an absolute
    timestamp::

        io_loop = IOLoop.current()

        # Wait up to 1 second for a notification.
        await condition.wait(timeout=io_loop.time() + 1)

    ...or a `datetime.timedelta` for a timeout relative to the current time::

        # Wait up to 1 second.
        await condition.wait(timeout=datetime.timedelta(seconds=1))

    The method returns False if there's no notification before the deadline.

    .. versionchanged:: 5.0
       Previously, waiters could be notified synchronously from within
       `notify`. Now, the notification will always be received on the
       next iteration of the `.IOLoop`.
    r   Nc                    s   t ƒ  ¡  tj ¡ | _d S r   )Úsuperr   r   ÚIOLoopÚcurrentÚio_loopr   ©Ú	__class__r   r   r   q   s   
zCondition.__init__c                 C   s.   d| j jf }| jr|dt| jƒ 7 }|d S )Nz<%sz waiters[%s]ú>)r/   r&   r   Úlen)r   Úresultr   r   r   Ú__repr__u   s   zCondition.__repr__Útimeoutc                    sT   t ƒ ‰ˆj ˆ¡ |r(d‡‡fdd„}tj ¡ ‰ ˆ  ||¡‰ˆ ‡ ‡fdd„¡ ˆS )z”Wait for `.notify`.

        Returns a `.Future` that resolves ``True`` if the condition is notified,
        or ``False`` after a timeout.
        r   Nc                      s   ˆ  ¡ s	tˆdƒ ˆ  ¡  d S ©NF)r    r   r$   r   ©r   Úwaiterr   r   Ú
on_timeout‡   s   
z"Condition.wait.<locals>.on_timeoutc                    ó
   ˆ   ˆ¡S r   ©Úremove_timeout©Ú_©r-   Útimeout_handler   r   Ú<lambda>Ž   ó   
 z Condition.wait.<locals>.<lambda>r%   )r   r   Úappendr   r+   r,   Úadd_timeoutÚadd_done_callback©r   r4   r8   r   ©r-   r   r?   r7   r   Úwait{   s   
zCondition.waitr   Únc                 C   sT   g }|r| j r| j  ¡ }| ¡ s|d8 }| |¡ |r| j s|D ]}t|dƒ q dS )zWake ``n`` waiters.r   TN)r   Úpopleftr    rB   r   )r   rH   Úwaitersr7   r   r   r   Únotify‘   s   



üÿzCondition.notifyc                 C   s   |   t| jƒ¡ dS )zWake all waiters.N)rK   r1   r   r   r   r   r   Ú
notify_all�   s   zCondition.notify_allr%   r   ©r   )r&   r'   r(   r)   r   Ústrr3   r   r   ÚfloatÚdatetimeÚ	timedeltar
   ÚboolrG   ÚintrK   rL   Ú__classcell__r   r   r.   r   r   5   s    ;ÿÿ
þr   c                   @   sr   e Zd ZdZddd„Zdefdd„Zdefdd	„Zdd
d„Z	ddd„Z
	ddeeeejf  ded fdd„ZdS )r   a½  An event blocks coroutines until its internal flag is set to True.

    Similar to `threading.Event`.

    A coroutine can wait for an event to be set. Once it is set, calls to
    ``yield event.wait()`` will not block unless the event has been cleared:

    .. testcode::

        from tornado import gen
        from tornado.ioloop import IOLoop
        from tornado.locks import Event

        event = Event()

        async def waiter():
            print("Waiting for event")
            await event.wait()
            print("Not waiting this time")
            await event.wait()
            print("Done")

        async def setter():
            print("About to set the event")
            event.set()

        async def runner():
            await gen.multi([waiter(), setter()])

        IOLoop.current().run_sync(runner)

    .. testoutput::

        Waiting for event
        About to set the event
        Not waiting this time
        Done
    r   Nc                 C   s   d| _ tƒ | _d S r5   )Ú_valueÚsetr   r   r   r   r   r   Ê   s   zEvent.__init__c                 C   s    d| j j|  ¡ rdf S df S )Nz<%s %s>rV   Úclear)r/   r&   Úis_setr   r   r   r   r3   Î   s   
þþzEvent.__repr__c                 C   s   | j S )z-Return ``True`` if the internal flag is true.©rU   r   r   r   r   rX   Ô   s   zEvent.is_setc                 C   s2   | j sd| _ | jD ]}| ¡ s| d¡ q	dS dS )zƒSet the internal flag to ``True``. All waiters are awakened.

        Calling `.wait` once the flag is set will not block.
        TN)rU   r   r    Ú
set_result)r   Úfutr   r   r   rV   Ø   s   

€ûz	Event.setc                 C   s
   d| _ dS )zkReset the internal flag to ``False``.

        Calls to `.wait` will block until `.set` is called.
        FNrY   r   r   r   r   rW   ä   s   
zEvent.clearr4   c                    sf   t ƒ ‰ ˆjrˆ  d¡ ˆ S ˆj ˆ ¡ ˆ  ‡fdd„¡ |du r"ˆ S t |ˆ ¡}| ‡ fdd„¡ |S )z�Block until the internal flag is true.

        Returns an awaitable, which raises `tornado.util.TimeoutError` after a
        timeout.
        Nc                    s   ˆ j  | ¡S r   )r   Úremove©r[   r   r   r   r@   ø   s    zEvent.wait.<locals>.<lambda>c                    s   ˆ   ¡ sˆ  ¡ S d S r   )r    Úcancel)Útfr]   r   r   r@     s    )r   rU   rZ   r   ÚaddrD   r   Úwith_timeout)r   r4   Útimeout_futr   )r[   r   r   rG   ë   s   

ÿz
Event.waitr%   r   )r&   r'   r(   r)   r   rN   r3   rR   rX   rV   rW   r   r   rO   rP   rQ   r
   rG   r   r   r   r   r   ¢   s    
'

ÿÿþr   c                   @   sP   e Zd ZdZdeddfdd„Zddd„Zd	d
dee dee	j
 ddfdd„ZdS )Ú_ReleasingContextManagerz³Releases a Lock or Semaphore at the end of a "with" statement.

        with (yield semaphore.acquire()):
            pass

        # Now semaphore.release() has been called.
    Úobjr   Nc                 C   s
   || _ d S r   )Ú_obj)r   rd   r   r   r   r     s   
z!_ReleasingContextManager.__init__c                 C   s   d S r   r   r   r   r   r   Ú	__enter__  s   z"_ReleasingContextManager.__enter__Úexc_typeúOptional[Type[BaseException]]Úexc_valÚexc_tbc                 C   s   | j  ¡  d S r   )re   Úrelease)r   rg   ri   rj   r   r   r   Ú__exit__  s   z!_ReleasingContextManager.__exit__r%   )r&   r'   r(   r)   r	   r   rf   r   ÚBaseExceptionÚtypesÚTracebackTyperl   r   r   r   r   rc     s    
þýüûrc   c                       sÌ   e Zd ZdZddeddf‡ fdd„Zdef‡ fdd	„Zdd
d„Z	dde	e
eejf  dee fdd„Zddd„Zddde	e de	ej ddfdd„Zddd„Zddde	e de	ej ddfdd„Z‡  ZS )r   aS  A lock that can be acquired a fixed number of times before blocking.

    A Semaphore manages a counter representing the number of `.release` calls
    minus the number of `.acquire` calls, plus an initial value. The `.acquire`
    method blocks if necessary until it can return without making the counter
    negative.

    Semaphores limit access to a shared resource. To allow access for two
    workers at a time:

    .. testsetup:: semaphore

       from collections import deque

       from tornado import gen
       from tornado.ioloop import IOLoop
       from tornado.concurrent import Future

       # Ensure reliable doctest output: resolve Futures one at a time.
       futures_q = deque([Future() for _ in range(3)])

       async def simulator(futures):
           for f in futures:
               # simulate the asynchronous passage of time
               await gen.sleep(0)
               await gen.sleep(0)
               f.set_result(None)

       IOLoop.current().add_callback(simulator, list(futures_q))

       def use_some_resource():
           return futures_q.popleft()

    .. testcode:: semaphore

        from tornado import gen
        from tornado.ioloop import IOLoop
        from tornado.locks import Semaphore

        sem = Semaphore(2)

        async def worker(worker_id):
            await sem.acquire()
            try:
                print("Worker %d is working" % worker_id)
                await use_some_resource()
            finally:
                print("Worker %d is done" % worker_id)
                sem.release()

        async def runner():
            # Join all workers.
            await gen.multi([worker(i) for i in range(3)])

        IOLoop.current().run_sync(runner)

    .. testoutput:: semaphore

        Worker 0 is working
        Worker 1 is working
        Worker 0 is done
        Worker 2 is working
        Worker 1 is done
        Worker 2 is done

    Workers 0 and 1 are allowed to run concurrently, but worker 2 waits until
    the semaphore has been released once, by worker 0.

    The semaphore can be used as an async context manager::

        async def worker(worker_id):
            async with sem:
                print("Worker %d is working" % worker_id)
                await use_some_resource()

            # Now the semaphore has been released.
            print("Worker %d is done" % worker_id)

    For compatibility with older versions of Python, `.acquire` is a
    context manager, so ``worker`` could also be written as::

        @gen.coroutine
        def worker(worker_id):
            with (yield sem.acquire()):
                print("Worker %d is working" % worker_id)
                yield use_some_resource()

            # Now the semaphore has been released.
            print("Worker %d is done" % worker_id)

    .. versionchanged:: 4.3
       Added ``async with`` support in Python 3.5.

    r   Úvaluer   Nc                    s$   t ƒ  ¡  |dk rtdƒ‚|| _d S )Nr   z$semaphore initial value must be >= 0)r*   r   Ú
ValueErrorrU   ©r   rp   r.   r   r   r   ~  s   

zSemaphore.__init__c                    sP   t ƒ  ¡ }| jdkrdnd | j¡}| jrd |t| jƒ¡}d |dd… |¡S )Nr   Úlockedzunlocked,value:{0}z{0},waiters:{1}z<{0} [{1}]>r   éÿÿÿÿ)r*   r3   rU   Úformatr   r1   )r   ÚresÚextrar.   r   r   r3   …  s   
ÿzSemaphore.__repr__c                 C   sT   |  j d7  _ | jr(| j ¡ }| ¡ s#|  j d8  _ | t| ƒ¡ dS | js
dS dS )ú*Increment the counter and wake one waiter.r   N)rU   r   rI   r    rZ   rc   r6   r   r   r   rk   Ž  s   
ôzSemaphore.releaser4   c                    s~   t ƒ ‰ˆjdkrˆ jd8  _ˆ tˆƒ¡ ˆS ˆj ˆ¡ |r=d	‡‡fdd„}tj ¡ ‰ ˆ  	||¡‰ˆ 
‡ ‡fdd„¡ ˆS )
z·Decrement the counter. Returns an awaitable.

        Block if the counter is zero and wait for a `.release`. The awaitable
        raises `.TimeoutError` after the deadline.
        r   r   r   Nc                      s"   ˆ  ¡ sˆ t ¡ ¡ ˆ  ¡  d S r   )r    Úset_exceptionr   ÚTimeoutErrorr$   r   r6   r   r   r8   ¯  s   z%Semaphore.acquire.<locals>.on_timeoutc                    r9   r   r:   r<   r>   r   r   r@   ·  rA   z#Semaphore.acquire.<locals>.<lambda>r%   )r   rU   rZ   rc   r   rB   r   r+   r,   rC   rD   rE   r   rF   r   ÚacquireŸ  s   
ó
ÿzSemaphore.acquirec                 C   ó   t dƒ‚)Nz0Use 'async with' instead of 'with' for Semaphore©ÚRuntimeErrorr   r   r   r   rf   »  ó   zSemaphore.__enter__Útyprh   Ú	tracebackc                 C   ó   |   ¡  d S r   ©rf   )r   r€   rp   r�   r   r   r   rl   ¾  ó   zSemaphore.__exit__c                 Ã   ó   �|   ¡ I d H  d S r   ©r{   r   r   r   r   Ú
__aenter__Æ  ó   €zSemaphore.__aenter__Útbc                 Ã   ó   �|   ¡  d S r   ©rk   ©r   r€   rp   r‰   r   r   r   Ú	__aexit__É  ó   €zSemaphore.__aexit__rM   r%   r   )r&   r'   r(   r)   rS   r   rN   r3   rk   r   r   rO   rP   rQ   r
   rc   r{   rf   rm   rn   ro   rl   r‡   r�   rT   r   r   r.   r   r     s>    _
	ÿÿ
þ
þýü
û
þýüûr   c                       s:   e Zd ZdZd
deddf‡ fdd„Zd‡ fdd	„Z‡  ZS )r   a:  A semaphore that prevents release() being called too many times.

    If `.release` would increment the semaphore's value past the initial
    value, it raises `ValueError`. Semaphores are mostly used to guard
    resources with limited capacity, so a semaphore released too many times
    is a sign of a bug.
    r   rp   r   Nc                    s   t ƒ j|d� || _d S )N©rp   )r*   r   Ú_initial_valuerr   r.   r   r   r   Û  s   
zBoundedSemaphore.__init__c                    s"   | j | jkr
tdƒ‚tƒ  ¡  dS )rx   z!Semaphore released too many timesN)rU   r�   rq   r*   rk   r   r.   r   r   rk   ß  s   zBoundedSemaphore.releaserM   r%   )r&   r'   r(   r)   rS   r   rk   rT   r   r   r.   r   r   Ò  s    r   c                   @   s¶   e Zd ZdZddd„Zdefdd„Z	ddeee	e
jf  dee fd	d
„Zddd„Zddd„Zdddee deej ddfdd„Zddd„Zdddee deej ddfdd„ZdS )r   aÏ  A lock for coroutines.

    A Lock begins unlocked, and `acquire` locks it immediately. While it is
    locked, a coroutine that yields `acquire` waits until another coroutine
    calls `release`.

    Releasing an unlocked lock raises `RuntimeError`.

    A Lock can be used as an async context manager with the ``async
    with`` statement:

    >>> from tornado import locks
    >>> lock = locks.Lock()
    >>>
    >>> async def f():
    ...    async with lock:
    ...        # Do something holding the lock.
    ...        pass
    ...
    ...    # Now the lock is released.

    For compatibility with older versions of Python, the `.acquire`
    method asynchronously returns a regular context manager:

    >>> async def f2():
    ...    with (yield lock.acquire()):
    ...        # Do something holding the lock.
    ...        pass
    ...
    ...    # Now the lock is released.

    .. versionchanged:: 4.3
       Added ``async with`` support in Python 3.5.

    r   Nc                 C   s   t dd�| _d S )Nr   r�   )r   Ú_blockr   r   r   r   r     s   zLock.__init__c                 C   s   d| j j| jf S )Nz<%s _block=%s>)r/   r&   r‘   r   r   r   r   r3     s   zLock.__repr__r4   c                 C   s   | j  |¡S )z�Attempt to lock. Returns an awaitable.

        Returns an awaitable, which raises `tornado.util.TimeoutError` after a
        timeout.
        )r‘   r{   )r   r4   r   r   r   r{     s   zLock.acquirec                 C   s(   z| j  ¡  W dS  ty   tdƒ‚w )zŠUnlock.

        The first coroutine in line waiting for `acquire` gets the lock.

        If not locked, raise a `RuntimeError`.
        zrelease unlocked lockN)r‘   rk   rq   r~   r   r   r   r   rk     s
   ÿzLock.releasec                 C   r|   )Nz+Use `async with` instead of `with` for Lockr}   r   r   r   r   rf   '  r   zLock.__enter__r€   rh   rp   r‰   c                 C   r‚   r   rƒ   rŒ   r   r   r   rl   *  r„   zLock.__exit__c                 Ã   r…   r   r†   r   r   r   r   r‡   2  rˆ   zLock.__aenter__c                 Ã   rŠ   r   r‹   rŒ   r   r   r   r�   5  rŽ   zLock.__aexit__r%   r   )r&   r'   r(   r)   r   rN   r3   r   r   rO   rP   rQ   r
   rc   r{   rk   rf   rm   rn   ro   rl   r‡   r�   r   r   r   r   r   æ  s>    
$ÿÿ
þ


þýü
û
þýüûr   )r   rP   rn   Útornador   r   Útornado.concurrentr   r   Útypingr   r   r   r	   r
   ÚTYPE_CHECKINGr   r   Ú__all__Úobjectr   r   r   rc   r   r   r   r   r   r   r   Ú<module>   s$   md 5