Ë
    oj™*  ã                  ó  — d Z ddl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mZmZm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  ed«      Z ed«      Z G d„ de«      Z G d„ de«      Z e«       Z  G d„ d«      Z!y)zº
twisted.python.threadpool: a pool of threads to which we dispatch tasks.

In most cases you can just use C{reactor.callInThread} and friends
instead of creating a thread pool directly.
é    )Úannotations)ÚThreadÚcurrent_thread)ÚAnyÚCallableÚListÚOptionalÚTypeVar)Ú	ParamSpecÚProtocolÚ	TypedDict)Úpool)ÚcontextÚlog)Ú
deprecated)ÚFailure)ÚVersionÚ_PÚ_Rc                  ó   — e Zd Zdd„Zy)Ú_SupportsQsizec                 ó   — y ©N© ©Úselfs    ú;/usr/lib/python3/dist-packages/twisted/python/threadpool.pyÚqsizez_SupportsQsize.qsize   s   € Øó    N©ÚreturnÚint©Ú__name__Ú
__module__Ú__qualname__r   r   r   r   r   r      s   „ ôr   r   c                  ó"   — e Zd ZU ded<   ded<   y)Ú_Stater"   ÚminÚmaxN)r$   r%   r&   Ú__annotations__r   r   r   r(   r(   "   s   … Ø	ƒHØ	„Hr   r(   c                  ób  — e Zd ZdZdZdZdZdZdZe	Z
 e  e edddd	«      d
¬«      e«      «      Z ee«      Z	 d	 	 	 	 	 dd„Zedd„«       Zedd„«       Zedd„«       Zed d„«       Ze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'	 	 	 	 	 d(d„Z!d!d„Z"y))Ú
ThreadPoolaè  
    This class (hopefully) generalizes the functionality of a pool of threads
    to which work can be dispatched.

    L{callInThread} and L{stop} should only be called from a single thread.

    @ivar started: Whether or not the thread pool is currently running.
    @type started: L{bool}

    @ivar threads: List of workers currently running in this thread pool.
    @type threads: L{list}

    @ivar _pool: A hook for testing.
    @type _pool: callable compatible with L{_pool}
    é   é   FNÚTwistedé   é   r   zthreading.current_thread)ÚversionÚreplacementc                ó´   ‡ — |dk\  sJ d«       ‚||k  sJ d«       ‚|‰ _         |‰ _        |‰ _        g ‰ _        dˆ fd„}dˆ fd„}‰ j	                  ||«      ‰ _        y)	ac  
        Create a new threadpool.

        @param minthreads: minimum number of threads in the pool
        @type minthreads: L{int}

        @param maxthreads: maximum number of threads in the pool
        @type maxthreads: L{int}

        @param name: The name to give this threadpool; visible in log messages.
        @type name: native L{str}
        r   úminimum is negativeúminimum is greater than maximumc                 ó‚   •—  ‰j                   | d‰j                  «       i|¤Ž}‰j                  j                  |«       |S )NÚname)ÚthreadFactoryÚ_generateNameÚthreadsÚappend)ÚaÚkwÚthreadr   s      €r   ÚtrackingThreadFactoryz2ThreadPool.__init__.<locals>.trackingThreadFactory`   sJ   ø€ Ø'�T×'Ñ'ØðØ×+Ñ+Ó-ðØ13ñˆFð �L‰L×Ñ Ô'ØˆMr   c                 ó6   •— ‰ j                   sy‰ j                  S )Nr   )Ústartedr*   r   s   €r   ÚcurrentLimitz)ThreadPool.__init__.<locals>.currentLimitg   s   ø€ Ø—<’<ØØ—8‘8ˆOr   N)r>   r   r?   r   r!   r   r    )r)   r*   r9   r<   Ú_poolÚ_team)r   Ú
minthreadsÚ
maxthreadsr9   rA   rD   s   `     r   Ú__init__zThreadPool.__init__J   sf   ø€ ð ˜QŠÐ5Ð 5Ó5ˆØ˜ZÒ'ÐJÐ)JÓJÐ'ØˆŒØˆŒØˆŒ	Ø%'ˆŒõ	õ	ð
 —Z‘Z Ð.CÓDˆ�
r   c                óh   — | j                   j                  «       }|j                  |j                  z   S )a  
        For legacy compatibility purposes, return a total number of workers.

        @return: the current number of workers, both idle and busy (but not
            those that have been quit by L{ThreadPool.adjustPoolsize})
        @rtype: L{int}
        )rF   Ú
statisticsÚidleWorkerCountÚbusyWorkerCount)r   Ústatss     r   ÚworkerszThreadPool.workersn   s-   € ð —
‘
×%Ñ%Ó'ˆØ×$Ñ$ u×'<Ñ'<Ñ<Ð<r   c                óR   — dg| j                   j                  «       j                  z  S )zý
        For legacy compatibility purposes, return the number of busy workers as
        expressed by a list the length of that number.

        @return: the number of workers currently processing a work item.
        @rtype: L{list} of L{None}
        N)rF   rK   rM   r   s    r   ÚworkingzThreadPool.workingz   s$   € ð ˆv˜Ÿ
™
×-Ñ-Ó/×?Ñ?Ñ?Ð?r   c                óR   — dg| j                   j                  «       j                  z  S )a,  
        For legacy compatibility purposes, return the number of idle workers as
        expressed by a list the length of that number.

        @return: the number of workers currently alive (with an allocated
            thread) but waiting for new work.
        @rtype: L{list} of L{None}
        N)rF   rK   rL   r   s    r   ÚwaiterszThreadPool.waiters…   s$   € ð ˆv˜Ÿ
™
×-Ñ-Ó/×?Ñ?Ñ?Ð?r   c                ó*   ‡ —  G ˆ fd„d«      } |«       S )zÙ
        For legacy compatibility purposes, return an object with a C{qsize}
        method that indicates the amount of work not yet allocated to a worker.

        @return: an object with a C{qsize} method.
        c                  ó   •— e Zd Zdˆ fd„Zy)ú$ThreadPool._queue.<locals>.NotAQueuec                óL   •— ‰j                   j                  «       j                  S )a  
                Pretend to be a Python threading Queue and return the
                number of as-yet-unconsumed tasks.

                @return: the amount of backlogged work not yet dispatched to a
                    worker.
                @rtype: L{int}
                )rF   rK   ÚbackloggedWorkCount)Úqr   s    €r   r   z*ThreadPool._queue.<locals>.NotAQueue.qsize›   s   ø€ ð —z‘z×,Ñ,Ó.×BÑBÐBr   Nr    r#   r   s   €r   Ú	NotAQueuerV   š   s	   ø„ ö	Cr   rZ   r   )r   rZ   s   ` r   Ú_queuezThreadPool._queue‘   s   ø€ ÷
	Có 
	Cñ ‹{Ðr   c                óÄ   — d| _         d| _        | j                  «        | j                  j	                  «       j
                  }|r| j                  j                  |«       yy)z'
        Start the threadpool.
        FTN)ÚjoinedrC   ÚadjustPoolsizerF   rK   rX   Úgrow)r   Úbacklogs     r   ÚstartzThreadPool.start«   sN   € ð ˆŒØˆŒà×ÑÔØ—*‘*×'Ñ'Ó)×=Ñ=ˆÙØ�J‰J�O‰O˜GÕ$ð r   c                ó:   — | j                   j                  d«       y)zŒ
        Increase the number of available workers for the thread pool by 1, up
        to the maximum allowed by L{ThreadPool.max}.
        r2   N)rF   r_   r   s    r   ÚstartAWorkerzThreadPool.startAWorker·   s   € ð
 	�
‰
�‰˜Õr   c                óT   — d| j                   xs t        | «      › d| j                  › �S )z‹
        Generate a name for a new pool thread.

        @return: A distinctive name for the thread.
        @rtype: native L{str}
        zPoolThread-ú-)r9   ÚidrO   r   s    r   r;   zThreadPool._generateName¾   s)   € ð ˜TŸY™YÒ2¬"¨T«(Ð3°1°T·\±\°NÐCÐCr   c                ó:   — | j                   j                  d«       y)zn
        Decrease the number of available workers by 1, by quitting one as soon
        as it's idle.
        r2   N)rF   Úshrinkr   s    r   ÚstopAWorkerzThreadPool.stopAWorkerÇ   s   € ð
 	�
‰
×Ñ˜!Õr   c                ót   — t        | d|«       t        j                  | | j                  | j                  «       y )NÚ__dict__)Úsetattrr-   rI   r)   r*   )r   Ústates     r   Ú__setstate__zThreadPool.__setstate__Î   s(   € Ü��j %Ô(Ü×Ñ˜D $§(¡(¨D¯H©HÕ5r   c                óD   — t        | j                  | j                  ¬«      S )N)r)   r*   )r(   r)   r*   r   s    r   Ú__getstate__zThreadPool.__getstate__Ò   s   € Ü˜$Ÿ(™(¨¯©Ô1Ð1r   c                ó2   —  | j                   d|g|¢­i |¤Ž y)a   
        Call a callable object in a separate thread.

        @param func: callable object to be called in separate thread

        @param args: positional arguments to be passed to C{func}

        @param kw: keyword args to be passed to C{func}
        N)ÚcallInThreadWithCallback)r   ÚfuncÚargsr?   s       r   ÚcallInThreadzThreadPool.callInThreadÕ   s    € ð 	&ˆ×%Ñ% d¨DÐ>°4Ò>¸2Ó>r   c                óè   ‡‡‡‡‡— | j                   ryt        j                  j                  «       j                  d   Šdˆfd„Šˆˆˆˆfd„‰_        |‰_        | j                  j                  ‰«       y)a$  
        Call a callable object in a separate thread and call C{onResult} with
        the return value, or a L{twisted.python.failure.Failure} if the
        callable raises an exception.

        The callable is allowed to block, but the C{onResult} function must not
        block and should perform as little work as possible.

        A typical action for C{onResult} for a threadpool used with a Twisted
        reactor would be to schedule a L{twisted.internet.defer.Deferred} to
        fire in the main reactor thread using C{.callFromThread}.  Note that
        C{onResult} is called inside the separate thread, not inside the
        reactor thread.

        @param onResult: a callable with the signature C{(success, result)}.
            If the callable returns normally, C{onResult} is called with
            C{(True, result)} where C{result} is the return value of the
            callable.  If the callable throws an exception, C{onResult} is
            called with C{(False, failure)}.

            Optionally, C{onResult} may be L{None}, in which case it is not
            called at all.

        @param func: callable object to be called in separate thread

        @param args: positional arguments to be passed to C{func}

        @param kw: keyword arguments to be passed to C{func}
        Néÿÿÿÿc                 óì   •— 	 ‰j                  «       } d}d ‰_         ‰j                  �‰j                  || «       d ‰_        y |st	        j
                  | «       y y # t        $ r t        «       } d}Y Œ]w xY w)NTF)ÚtheWorkÚBaseExceptionr   ÚonResultr   Úerr)ÚresultÚokÚ	inContexts     €r   r   z6ThreadPool.callInThreadWithCallback.<locals>.inContext  sy   ø€ ðØ"×*Ñ*Ó,�Ø�ð
 !%ˆIÔØ×!Ñ!Ð-Ø×"Ñ" 2 vÔ.Ø%)�	Õ"ÙÜ—‘˜•ð øô !ò Ü ›�Ø’ðús   ƒA ÁA3Á2A3c                 ó8   •— t        j                  ‰‰g‰ ¢­i ‰¤ŽS r   )r   Úcall)rt   Úctxrs   r?   s   €€€€r   ú<lambda>z5ThreadPool.callInThreadWithCallback.<locals>.<lambda>  s%   ø€ ¤G§L¡LØ�ð%
Øò%
Ø "ñ%
€ r   ©r!   ÚNone)	r]   r   ÚtheContextTrackerÚcurrentContextÚcontextsry   r{   rF   Údo)r   r{   rs   rt   r?   r‚   r   s     ```@@r   rr   z#ThreadPool.callInThreadWithCallbackã   sX   ü€ ðH �;Š;ØÜ×'Ñ'×6Ñ6Ó8×AÑAÀ"ÑEˆõ	 ö$
ˆ	Ôð &ˆ	Ôà�
‰
�‰�iÕ r   c                ó–   — d| _         d| _        | j                  j                  «        | j                  D ]  }|j                  «        Œ y)z9
        Shutdown the threads in the threadpool.
        TFN)r]   rC   rF   Úquitr<   Újoin)r   r@   s     r   ÚstopzThreadPool.stop$  s<   € ð ˆŒØˆŒØ�
‰
�‰ÔØ—l‘lò 	ˆFØ�K‰K�Mñ	r   c                óÐ  — |€| j                   }|€| j                  }|dk\  sJ d«       ‚||k  sJ d«       ‚|| _         || _        | j                  sy| j                  | j                  kD  r2| j                  j                  | j                  | j                  z
  «       | j                  | j                   k  r3| j                  j                  | j                   | j                  z
  «       yy)zî
        Adjust the number of available threads by setting C{min} and C{max} to
        new values.

        @param minthreads: The new value for L{ThreadPool.min}.

        @param maxthreads: The new value for L{ThreadPool.max}.
        Nr   r6   r7   )r)   r*   rC   rO   rF   rh   r_   )r   rG   rH   s      r   r^   zThreadPool.adjustPoolsize.  sÃ   € ð ÐØŸ™ˆJØÐØŸ™ˆJà˜QŠÐ5Ð 5Ó5ˆØ˜ZÒ'ÐJÐ)JÓJÐ'àˆŒØˆŒØ�|Š|Øð �<‰<˜$Ÿ(™(Ò"Ø�J‰J×Ñ˜dŸl™l¨T¯X©XÑ5Ô6à�<‰<˜$Ÿ(™(Ò"Ø�J‰J�O‰O˜DŸH™H t§|¡|Ñ3Õ4ð #r   c                óÐ   — t        j                  d| j                  › �«       t        j                  d| j                  › �«       t        j                  d| j                  › �«       y)zw
        Dump some plain-text informational messages to the log about the state
        of this L{ThreadPool}.
        z	waiters: z	workers: ztotal: N)r   ÚmsgrS   rQ   r<   r   s    r   Ú	dumpStatszThreadPool.dumpStatsM  sI   € ô
 	�‰�)˜DŸL™L˜>Ð*Ô+Ü�‰�)˜DŸL™L˜>Ð*Ô+Ü�‰�'˜$Ÿ,™,˜Ð(Õ)r   )r.   r/   N)rG   r"   rH   r"   r9   zOptional[str]r    )r!   z
list[None])r!   r   r„   )r!   Ústr)rm   r(   r!   r…   )r!   r(   )rs   zCallable[_P, object]rt   ú_P.argsr?   ú	_P.kwargsr!   r…   )
r{   z&Optional[Callable[[bool, _R], object]]rs   zCallable[_P, _R]rt   r“   r?   r”   r!   r…   )NN)rG   úOptional[int]rH   r•   r!   r…   )#r$   r%   r&   Ú__doc__r)   r*   r]   rC   r9   r   r:   Ústaticmethodr   r   r   ÚcurrentThreadrE   rI   ÚpropertyrO   rQ   rS   r[   rY   ra   rc   r;   ri   rn   rp   ru   rr   r�   r^   r‘   r   r   r   r-   r-   *   s�  „ ñð  €CØ
€CØ€FØ€GØ€Dà€MÙ ð	
‰
Ù˜I r¨1¨aÓ0Ø2ô	
ð ó	ó€Mñ ˜Ó€Eð PTð"EØð"EØ/2ð"EØ?Ló"EðH ò	=ó ð	=ð ò@ó ð@ð ò	@ó ð	@ð òó ðð, 	€Aó
%óóDóó6ó2ð?Ø(ð?Ø18ð?Ø@Ið?à	ó?ð?!à8ð?!ð ð?!ð ð	?!ð
 ð?!ð 
ó?!óBð MQð5Ø'ð5Ø<Ið5à	ó5ô>*r   r-   N)"r–   Ú
__future__r   Ú	threadingr   r   Útypingr   r   r   r	   r
   r   r   r   Útwisted._threadsr   rE   Útwisted.pythonr   r   Útwisted.python.deprecater   Útwisted.python.failurer   Útwisted.python.versionsr   r   r   r   r(   ÚobjectÚ
WorkerStopr-   r   r   r   ú<module>r¤      sl   ðñ
õ #ç ,ß 9Õ 9ç 1Ñ 1å *ß 'Ý /Ý *Ý +áˆtƒ_€ÙˆTƒ]€ô�Xô ô
ˆYô ñ
 ‹X€
÷j*ò j*r   