Ë
    oj  ã                  ó¢   — d Z ddlmZ ddlmZ ddlmZmZmZ ddl	m
Z
 ddlmZ ddlmZ dd	lmZ  G d
„ d«      Z e
e«       G d„ d«      «       Zy)zZ
Implementation of a L{Team} of workers; a thread-pool that can allocate work to
workers.
é    )Úannotations)Údeque)ÚCallableÚOptionalÚSet)Úimplementeré   )ÚIWorker)ÚQuit)ÚIExclusiveWorkerc                  ó(   — e Zd ZdZ	 	 	 	 	 	 	 	 dd„Zy)Ú
StatisticsaÓ  
    Statistics about a L{Team}'s current activity.

    @ivar idleWorkerCount: The number of idle workers.
    @type idleWorkerCount: L{int}

    @ivar busyWorkerCount: The number of busy workers.
    @type busyWorkerCount: L{int}

    @ivar backloggedWorkCount: The number of work items passed to L{Team.do}
        which have not yet been sent to a worker to be performed because not
        enough workers are available.
    @type backloggedWorkCount: L{int}
    c                ó.   — || _         || _        || _        y ©N)ÚidleWorkerCountÚbusyWorkerCountÚbackloggedWorkCount)Úselfr   r   r   s       ú8/usr/lib/python3/dist-packages/twisted/_threads/_team.pyÚ__init__zStatistics.__init__%   s   € ð  /ˆÔØ.ˆÔØ#6ˆÕ ó    N)r   Úintr   r   r   r   ÚreturnÚNone)Ú__name__Ú
__module__Ú__qualname__Ú__doc__r   © r   r   r   r      s,   „ ñð7Ø"ð7Ø58ð7ØORð7à	ô7r   r   c                  óh   — e Zd ZdZ	 	 	 	 	 	 dd„Zd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y)ÚTeamax  
    A composite L{IWorker} implementation.

    @ivar _quit: A L{Quit} flag indicating whether this L{Team} has been quit
        yet.  This may be set by an arbitrary thread since L{Team.quit} may be
        called from anywhere.

    @ivar _coordinator: the L{IExclusiveWorker} coordinating access to this
        L{Team}'s internal resources.

    @ivar _createWorker: a callable that will create new workers.

    @ivar _logException: a 0-argument callable called in an exception context
        when there is an unhandled error from a task passed to L{Team.do}

    @ivar _idle: a L{set} of idle workers.

    @ivar _busyCount: the number of workers currently busy.

    @ivar _pending: a C{deque} of tasks - that is, 0-argument callables passed
        to L{Team.do} - that are outstanding.

    @ivar _shouldQuitCoordinator: A flag indicating that the coordinator should
        be quit at the next available opportunity.  Unlike L{Team._quit}, this
        flag is only set by the coordinator.

    @ivar _toShrink: the number of workers to shrink this L{Team} by at the
        next available opportunity; set in the coordinator.
    c                ó²   — t        «       | _        || _        || _        || _        t        «       | _        d| _        t        «       | _	        d| _
        d| _        y)a   
        @param coordinator: an L{IExclusiveWorker} which will coordinate access
            to resources on this L{Team}; that is to say, an
            L{IExclusiveWorker} whose C{do} method ensures that its given work
            will be executed in a mutually exclusive context, not in parallel
            with other work enqueued by C{do} (although possibly in parallel
            with the caller).

        @param createWorker: A 0-argument callable that will create an
            L{IWorker} to perform work.

        @param logException: A 0-argument callable called in an exception
            context when the work passed to C{do} raises an exception.
        r   FN)r   Ú_quitÚ_coordinatorÚ_createWorkerÚ_logExceptionÚsetÚ_idleÚ
_busyCountr   Ú_pendingÚ_shouldQuitCoordinatorÚ	_toShrink)r   ÚcoordinatorÚcreateWorkerÚlogExceptions       r   r   zTeam.__init__M   sO   € ô( “VˆŒ
Ø'ˆÔØ)ˆÔØ)ˆÔô $'£5ˆŒ
ØˆŒÜ8=»ˆŒØ&+ˆÔ#Øˆ�r   c                ó|   — t        t        | j                  «      | j                  t        | j                  «      «      S )z›
        Gather information on the current status of this L{Team}.

        @return: a L{Statistics} describing the current state of this L{Team}.
        )r   Úlenr(   r)   r*   ©r   s    r   Ú
statisticszTeam.statisticsm   s(   € ô œ#˜dŸj™j›/¨4¯?©?¼CÀÇÁÓ<NÓOÐOr   c                ó|   ‡ ‡— ‰ j                   j                  «        ‰ j                  j                  dˆˆ fd„«       }y)z—
        Increase the the number of idle workers by C{n}.

        @param n: The number of new idle workers to create.
        @type n: L{int}
        c                 óp   •— t        ‰«      D ]'  } ‰j                  «       }|€ y ‰j                  |«       Œ) y r   )Úranger%   Ú_recycleWorker)ÚxÚworkerÚnr   s     €€r   ÚcreateOneWorkerz"Team.grow.<locals>.createOneWorker~   s:   ø€ ä˜1“Xò ,�Ø×+Ñ+Ó-�Ø�>ÙØ×#Ñ# FÕ+ñ	,r   N©r   r   ©r#   Úcheckr$   Údo)r   r:   r;   s   `` r   Úgrowz	Team.growu   s3   ù€ ð 	�
‰
×ÑÔà	×	Ñ	×	Ñ	õ	,ó 
ñ	,r   Nc                óz   ‡ ‡— ‰ j                   j                  «        ‰ j                  j                  ˆˆ fd„«       y)zß
        Decrease the number of idle workers by C{n}.

        @param n: The number of idle workers to shut down, or L{None} (or
            unspecified) to shut down all workers.
        @type n: L{int} or L{None}
        c                 ó&   •— ‰j                  ‰ «      S r   )Ú_quitIdlers)r:   r   s   €€r   ú<lambda>zTeam.shrink.<locals>.<lambda>�   s   ø€  T×%5Ñ%5°aÓ%8€ r   Nr=   )r   r:   s   ``r   ÚshrinkzTeam.shrink†   s*   ù€ ð 	�
‰
×ÑÔØ×Ñ×ÑÔ8Õ9r   c                ón  — |€"t        | j                  «      | j                  z   }t        |«      D ]L  }| j                  r)| j                  j	                  «       j                  «        Œ8| xj                  dz  c_        ŒN | j                  r+| j                  dk(  r| j                  j                  «        yyy)z|
        The implmentation of C{shrink}, performed by the coordinator worker.

        @param n: see L{Team.shrink}
        Nr	   r   )	r1   r(   r)   r6   ÚpopÚquitr,   r+   r$   )r   r:   r8   s      r   rC   zTeam._quitIdlers‘   s�   € ð ˆ9Ü�D—J‘J“ $§/¡/Ñ1ˆAÜ�q“ò 	$ˆAØ�zŠzØ—
‘
—‘Ó ×%Ñ%Õ'à—’ !Ñ#–ð		$ð
 ×&Ò&¨4¯?©?¸aÒ+?Ø×Ñ×"Ñ"Õ$ð ,@Ð&r   c                óz   ‡ ‡— ‰ j                   j                  «        ‰ j                  j                  ˆ ˆfd„«       y)zu
        Perform some work in a worker created by C{createWorker}.

        @param task: the callable to run
        c                 ó&   •— ‰ j                  ‰«      S r   )Ú_coordinateThisTask©r   Útasks   €€r   rD   zTeam.do.<locals>.<lambda>¨   s   ø€  T×%=Ñ%=¸dÓ%C€ r   Nr=   rL   s   ``r   r?   zTeam.do¡   s*   ù€ ð 	�
‰
×ÑÔØ×Ñ×ÑÔCÕDr   c                ó  ‡ ‡‡— ‰ j                   r‰ j                   j                  «       n‰ j                  «       }|€‰ j                  j	                  ‰«       y|Š‰ xj
                  dz  c_        |j                  dˆˆ ˆfd„«       }y)zø
        Select a worker to dispatch to, either an idle one or a new one, and
        perform it.

        This method should run on the coordinator worker.

        @param task: the task to dispatch
        @type task: 0-argument callable
        Nr	   c                 ó”   •— 	  ‰«        ‰j                  j                  dˆˆfd„«       } y # t         $ r ‰j                  «        Y Œ<w xY w)Nc                 óR   •— ‰xj                   dz  c_         ‰j                  ‰ «       y )Nr	   )r)   r7   )Únot_none_workerr   s   €€r   ÚidleAndPendingz@Team._coordinateThisTask.<locals>.doWork.<locals>.idleAndPendingÄ   s   ø€ à—’ 1Ñ$•Ø×#Ñ# OÕ4r   r<   )ÚBaseExceptionr&   r$   r?   )rR   rQ   r   rM   s    €€€r   ÚdoWorkz(Team._coordinateThisTask.<locals>.doWork½   sJ   ø€ ð%Ù”ð ×Ñ×!Ñ!õ5ó "ñ5øô	 !ò %Ø×"Ñ"Ö$ð%ús   ƒ+ «AÁAr<   )r(   rG   r%   r*   Úappendr)   r?   )r   rM   r9   rT   rQ   s   ``  @r   rK   zTeam._coordinateThisTaskª   sk   ú€ ð &*§Z¢Z�—‘—‘Ô!°T×5GÑ5GÓ5IˆØˆ>ð �M‰M× Ñ  Ô&ØØ ˆØ�Š˜1Ñ�à	�‰ö		5ó 
ñ		5r   c                ó€  — | j                   j                  |«       | j                  r*| j                  | j                  j	                  «       «       y| j
                  r| j                  «        y| j                  dkD  rA| xj                  dz  c_        | j                   j                  |«       |j                  «        yy)zÐ
        Called only from coordinator.

        Recycle the given worker into the idle pool.

        @param worker: a worker created by C{createWorker} and now idle.
        @type worker: L{IWorker}
        r   r	   N)
r(   Úaddr*   rK   Úpopleftr+   rC   r,   ÚremoverH   )r   r9   s     r   r7   zTeam._recycleWorkerÉ   s‡   € ð 	�
‰
�‰�vÔØ�=Š=ð ×$Ñ$ T§]¡]×%:Ñ%:Ó%<Õ=Ø×(Ò(Ø×ÑÕØ�^‰^˜aÒØ�NŠN˜aÑ�NØ�J‰J×Ñ˜fÔ%Ø�K‰K�Mð  r   c                óx   ‡ — ‰ j                   j                  «        ‰ j                  j                  dˆ fd„«       }y)zA
        Stop doing work and shut down all idle workers.
        c                 ó4   •— d‰ _         ‰ j                  «        y )NT)r+   rC   r2   s   €r   ÚstartFinishingz!Team.quit.<locals>.startFinishingå   s   ø€ à*.ˆDÔ'Ø×ÑÕr   Nr<   )r#   r'   r$   r?   )r   r\   s   ` r   rH   z	Team.quitÞ   s3   ø€ ð 	�
‰
�‰Ôð 
×	Ñ	×	Ñ	ô	ó 
ñ	r   )r-   r   r.   zCallable[[], Optional[IWorker]]r/   úCallable[[], None])r   r   )r:   r   r   r   r   )r:   zOptional[int]r   r   )rM   r]   r   r   )rM   zCallable[..., object]r   r   )r9   r
   r   r   r<   )r   r   r   r   r   r3   r@   rE   rC   r?   rK   r7   rH   r   r   r   r!   r!   -   sS   „ ñð<à%ðð 6ðð )ó	ó@Pó,ô"	:ô%ó Eó5ó>ô*
r   r!   N)r   Ú
__future__r   Úcollectionsr   Útypingr   r   r   Úzope.interfacer   Ú r
   Ú_conveniencer   Ú	_ithreadsr   r   r!   r   r   r   ú<module>re      sO   ðñ
õ #å ß *Ñ *å &å Ý Ý '÷7ñ 7ñ0 ˆWÓ÷zð zó ñzr   