
Â�Â­Â�Â­Â�Â­Â�Â­Â�Â­Â�Â­Â�Â­Â�Â­Â�Â­Â�Â­Â�Â­Â�Â­Â�Â­Â�Â­Â�Â­Â�Â­Â�Â­Â�Â­
<!DOCTYPE html>
<html>
U
    ¡ê,aý~  ã                   @   sd  d dgZ ddlZddlZddlZddlZddlZddlZddlZddlZddlm	Z	 ddl
mZ ddl
mZmZ ddlmZ d	Zd
ZdZdZe ¡ Zdd„ Zdd„ ZG dd„ deƒZG dd„ dƒZdd„ ZG dd„ deƒZd+dd„Zdd„ ZG dd „ d eƒZ G d!d „ d e!ƒZ"G d"d#„ d#e!ƒZ#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 ),ÚPoolÚ
ThreadPoolé    N)ÚEmptyé   )Úutil)Úget_contextÚTimeoutError)ÚwaitÚINITÚRUNÚCLOSEÚ	TERMINATEc                 C   s   t t| Ž ƒS ©N)ÚlistÚmap©Úargs© r   ú9/opt/alt/python38/lib64/python3.8/multiprocessing/pool.pyÚmapstar/   s    r   c                 C   s   t t | d | d ¡ƒS )Nr   r   )r   Ú	itertoolsÚstarmapr   r   r   r   Ústarmapstar2   s    r   c                   @   s   e Zd Zdd„ Zdd„ ZdS )ÚRemoteTracebackc                 C   s
   || _ d S r   ©Útb)Úselfr   r   r   r   Ú__init__:   s    zRemoteTraceback.__init__c                 C   s   | j S r   r   ©r   r   r   r   Ú__str__<   s    zRemoteTraceback.__str__N)Ú__name__Ú
__module__Ú__qualname__r   r   r   r   r   r   r   9   s   r   c                   @   s   e Zd Zdd„ Zdd„ ZdS )ÚExceptionWithTracebackc                 C   s0   t  t|ƒ||¡}d |¡}|| _d| | _d S )NÚ z

"""
%s""")Ú	tracebackÚformat_exceptionÚtypeÚjoinÚexcr   )r   r)   r   r   r   r   r   @   s    
zExceptionWithTraceback.__init__c                 C   s   t | j| jffS r   )Úrebuild_excr)   r   r   r   r   r   Ú
__reduce__E   s    z!ExceptionWithTraceback.__reduce__N)r    r!   r"   r   r+   r   r   r   r   r#   ?   s   r#   c                 C   s   t |ƒ| _| S r   )r   Ú	__cause__)r)   r   r   r   r   r*   H   s    
r*   c                       s0   e Zd ZdZ‡ fdd„Zdd„ Zdd„ Z‡  ZS )ÚMaybeEncodingErrorzVWraps possible unpickleable errors, so they can be
    safely sent through the socket.c                    s.   t |ƒ| _t |ƒ| _tt| ƒ | j| j¡ d S r   )Úreprr)   ÚvalueÚsuperr-   r   )r   r)   r/   ©Ú	__class__r   r   r   T   s    

zMaybeEncodingError.__init__c                 C   s   d| j | jf S )Nz(Error sending result: '%s'. Reason: '%s')r/   r)   r   r   r   r   r   Y   s    ÿzMaybeEncodingError.__str__c                 C   s   d| j j| f S )Nz<%s: %s>)r2   r    r   r   r   r   Ú__repr__]   s    zMaybeEncodingError.__repr__)r    r!   r"   Ú__doc__r   r   r3   Ú__classcell__r   r   r1   r   r-   P   s   r-   r   Fc              
   C   sÌ  |d k	r(t |tƒr|dks(td |¡ƒ‚|j}| j}t| dƒrR| j ¡  |j	 ¡  |d k	rb||Ž  d}|d ks~|�rº||k �rºz
|ƒ }	W n( t
tfk
r°   t d¡ Y �qºY nX |	d krÈt d¡ �qº|	\}
}}}}zd|||Žf}W nH tk
�r0 } z(|�r|tk	�rt||jƒ}d|f}W 5 d }~X Y nX z||
||fƒ W nR tk
�r– } z2t||d ƒ}t d	| ¡ ||
|d|ffƒ W 5 d }~X Y nX d  }	 }
 } } }}|d7 }qft d
| ¡ d S )Nr   zMaxtasks {!r} is not validÚ_writerr   z)worker got EOFError or OSError -- exitingzworker got sentinel -- exitingTFz0Possible encoding error while sending result: %szworker exiting after %d tasks)Ú
isinstanceÚintÚAssertionErrorÚformatÚputÚgetÚhasattrr6   ÚcloseÚ_readerÚEOFErrorÚOSErrorr   ÚdebugÚ	ExceptionÚ_helper_reraises_exceptionr#   Ú__traceback__r-   )ÚinqueueÚoutqueueÚinitializerÚinitargsZmaxtasksÚwrap_exceptionr;   r<   Z	completedÚtaskÚjobÚiÚfuncr   ÚkwdsÚresultÚeÚwrappedr   r   r   Úworkera   sN    ÿ





ÿ$
rS   c                 C   s   | ‚dS )z@Pickle-able helper function for use by _guarded_task_generation.Nr   )Zexr   r   r   rD   Ž   s    rD   c                       s2   e Zd ZdZddœ‡ fdd„
Z‡ fdd„Z‡  ZS )Ú
_PoolCachezò
    Class that implements a cache for the Pool class that will notify
    the pool management threads every time the cache is emptied. The
    notification is done by the use of a queue that is provided when
    instantiating the cache.
    N©Únotifierc                  s   || _ tƒ j||Ž d S r   )rV   r0   r   )r   rV   r   rO   r1   r   r   r   �   s    z_PoolCache.__init__c                    s    t ƒ  |¡ | s| j d ¡ d S r   )r0   Ú__delitem__rV   r;   )r   Úitemr1   r   r   rW   ¡   s    z_PoolCache.__delitem__)r    r!   r"   r4   r   rW   r5   r   r   r1   r   rT   –   s   rT   c                   @   s†  e Zd ZdZdZedd„ ƒZdLdd„Zej	e
fd	d
„Zdd„ Zdd„ Zedd„ ƒZedd„ ƒZdd„ Zedd„ ƒZedd„ ƒZdd„ Zdd„ Zdi fdd„ZdMdd „ZdNd!d"„ZdOd#d$„Zd%d&„ ZdPd(d)„ZdQd*d+„Zdi ddfd,d-„ZdRd.d/„ZdSd0d1„ZedTd2d3„ƒZe d4d5„ ƒZ!ed6d7„ ƒZ"ed8d9„ ƒZ#ed:d;„ ƒZ$d<d=„ Z%d>d?„ Z&d@dA„ Z'dBdC„ Z(edDdE„ ƒZ)e dFdG„ ƒZ*dHdI„ Z+dJdK„ Z,dS )Ur   zS
    Class which supports an async version of applying functions to arguments.
    Tc                 O   s   | j ||ŽS r   ©ÚProcess)Úctxr   rO   r   r   r   rZ   ³   s    zPool.ProcessNr   c                 C   s  g | _ t| _|ptƒ | _|  ¡  t ¡ | _| j ¡ | _	t
| j	d�| _|| _|| _|| _|d krjt ¡ phd}|dk rztdƒ‚|d k	r’t|ƒs’tdƒ‚|| _z|  ¡  W nH tk
rì   | j D ]}|jd krº| ¡  qº| j D ]}| ¡  qØ‚ Y nX |  ¡ }tjtj| j| j| j| j| j| j | j | j!| j| j| j| j"|| j	fd�| _#d| j#_$t%| j#_| j# &¡  tjtj'| j| j(| j!| j | jfd�| _)d| j)_$t%| j)_| j) &¡  tjtj*| j!| j+| jfd�| _,d| j,_$t%| j,_| j, &¡  t-j.| | j/| j| j | j!| j | j	| j#| j)| j,| jf	dd�| _0t%| _d S )	NrU   r   z&Number of processes must be at least 1zinitializer must be a callable©Útargetr   Té   )r   Zexitpriority)1Ú_poolr
   Ú_stater   Ú_ctxÚ_setup_queuesÚqueueÚSimpleQueueÚ
_taskqueueÚ_change_notifierrT   Ú_cacheÚ_maxtasksperchildÚ_initializerÚ	_initargsÚosÚ	cpu_countÚ
ValueErrorÚcallableÚ	TypeErrorÚ
_processesÚ_repopulate_poolrC   ÚexitcodeÚ	terminater(   Ú_get_sentinelsÚ	threadingZThreadr   Ú_handle_workersrZ   Ú_inqueueÚ	_outqueueÚ_wrap_exceptionÚ_worker_handlerÚdaemonr   ÚstartÚ_handle_tasksÚ
_quick_putÚ_task_handlerÚ_handle_resultsÚ
_quick_getÚ_result_handlerr   ZFinalizeÚ_terminate_poolÚ
_terminate)r   Ú	processesrH   rI   ÚmaxtasksperchildÚcontextÚpÚ	sentinelsr   r   r   r   ·   s–    





       ýþ
 ÿþ
þ
    þûzPool.__init__c                 C   s>   | j |kr:|d| ›�t| d� t| dd ƒd k	r:| j d ¡ d S )Nz&unclosed running multiprocessing pool )Úsourcerf   )r`   ÚResourceWarningÚgetattrrf   r;   )r   Z_warnr   r   r   r   Ú__del__  s    

 ÿzPool.__del__c              	   C   s0   | j }d|j› d|j› d| j› dt| jƒ› d�	S )Nú<Ú.z state=z pool_size=ú>)r2   r!   r"   r`   Úlenr_   )r   Úclsr   r   r   r3     s    zPool.__repr__c                 C   s   | j jg}| jjg}||•S r   )rx   r?   rf   )r   Ztask_queue_sentinelsZself_notifier_sentinelsr   r   r   rt     s    

zPool._get_sentinelsc                 C   s   dd„ | D ƒS )Nc                 S   s   g | ]}t |d ƒr|j‘qS )Úsentinel)r=   r“   )Ú.0rS   r   r   r   Ú
<listcomp>  s    
ÿz.Pool._get_worker_sentinels.<locals>.<listcomp>r   ©Zworkersr   r   r   Ú_get_worker_sentinels  s    ÿzPool._get_worker_sentinelsc                 C   sP   d}t tt| ƒƒƒD ]6}| | }|jdk	rt d| ¡ | ¡  d}| |= q|S )z�Cleanup after any worker processes which have exited due to reaching
        their specified lifetime.  Returns True if any workers were cleaned up.
        FNúcleaning up worker %dT)ÚreversedÚranger‘   rr   r   rB   r(   )ÚpoolZcleanedrM   rS   r   r   r   Ú_join_exited_workers  s    
zPool._join_exited_workersc                 C   s0   |   | j| j| j| j| j| j| j| j| j	| j
¡
S r   )Ú_repopulate_pool_staticra   rZ   rp   r_   rw   rx   ri   rj   rh   ry   r   r   r   r   rq   .  s      úzPool._repopulate_poolc
              
   C   sf   t |t|ƒ ƒD ]P}
|| t||||||	fd�}|j dd¡|_d|_| ¡  | |¡ t 	d¡ qdS )z€Bring the number of pool processes up to the specified number,
        for use after reaping workers which have exited.
        r\   rZ   Z
PoolWorkerTzadded workerN)
rš   r‘   rS   ÚnameÚreplacer{   r|   Úappendr   rB   )r[   rZ   r…   r›   rF   rG   rH   rI   r†   rJ   rM   Úwr   r   r   r�   7  s     ýÿ
zPool._repopulate_pool_staticc
           
      C   s*   t  |¡r&t  | |||||||||	¡
 dS )zEClean up any exited workers and start replacements for them.
        N)r   rœ   r�   )
r[   rZ   r…   r›   rF   rG   rH   rI   r†   rJ   r   r   r   Ú_maintain_poolJ  s    
   ýzPool._maintain_poolc                 C   s4   | j  ¡ | _| j  ¡ | _| jjj| _| jjj| _	d S r   )
ra   rd   rw   rx   r6   Úsendr~   r?   Úrecvr�   r   r   r   r   rb   V  s    zPool._setup_queuesc                 C   s   | j tkrtdƒ‚d S )NzPool not running)r`   r   rm   r   r   r   r   Ú_check_running\  s    
zPool._check_runningc                 C   s   |   |||¡ ¡ S )zT
        Equivalent of `func(*args, **kwds)`.
        Pool must be running.
        )Úapply_asyncr<   )r   rN   r   rO   r   r   r   Úapply`  s    z
Pool.applyc                 C   s   |   ||t|¡ ¡ S )zx
        Apply `func` to each element in `iterable`, collecting the results
        in a list that is returned.
        )Ú
_map_asyncr   r<   ©r   rN   ÚiterableÚ	chunksizer   r   r   r   g  s    zPool.mapc                 C   s   |   ||t|¡ ¡ S )zÌ
        Like `map()` method but the elements of the `iterable` are expected to
        be iterables as well and will be unpacked as arguments. Hence
        `func` and (a, b) becomes func(a, b).
        )r¨   r   r<   r©   r   r   r   r   n  s    zPool.starmapc                 C   s   |   ||t|||¡S )z=
        Asynchronous version of `starmap()` method.
        )r¨   r   ©r   rN   rª   r«   ÚcallbackÚerror_callbackr   r   r   Ústarmap_asyncv  s     ÿzPool.starmap_asyncc              
   c   sj   z,d}t |ƒD ]\}}||||fi fV  qW n8 tk
rd } z||d t|fi fV  W 5 d}~X Y nX dS )zšProvides a generator of tasks for imap and imap_unordered with
        appropriate handling for iterables which throw exceptions during
        iteration.éÿÿÿÿr   N)Ú	enumeraterC   rD   )r   Z
result_jobrN   rª   rM   ÚxrQ   r   r   r   Ú_guarded_task_generation~  s    zPool._guarded_task_generationr   c                 C   s–   |   ¡  |dkr:t| ƒ}| j |  |j||¡|jf¡ |S |dk rPtd |¡ƒ‚t	 
|||¡}t| ƒ}| j |  |jt|¡|jf¡ dd„ |D ƒS dS )zP
        Equivalent of `map()` -- can be MUCH slower than `Pool.map()`.
        r   zChunksize must be 1+, not {0:n}c                 s   s   | ]}|D ]
}|V  q
qd S r   r   ©r”   ÚchunkrX   r   r   r   Ú	<genexpr>¤  s       zPool.imap.<locals>.<genexpr>N)r¥   ÚIMapIteratorre   r;   r³   Ú_jobÚ_set_lengthrm   r:   r   Ú
_get_tasksr   ©r   rN   rª   r«   rP   Útask_batchesr   r   r   Úimap‰  s4    þÿÿÿþüÿz	Pool.imapc                 C   s–   |   ¡  |dkr:t| ƒ}| j |  |j||¡|jf¡ |S |dk rPtd |¡ƒ‚t	 
|||¡}t| ƒ}| j |  |jt|¡|jf¡ dd„ |D ƒS dS )zL
        Like `imap()` method but ordering of results is arbitrary.
        r   zChunksize must be 1+, not {0!r}c                 s   s   | ]}|D ]
}|V  q
qd S r   r   r´   r   r   r   r¶   À  s       z&Pool.imap_unordered.<locals>.<genexpr>N)r¥   ÚIMapUnorderedIteratorre   r;   r³   r¸   r¹   rm   r:   r   rº   r   r»   r   r   r   Úimap_unordered¦  s0    þÿÿþüÿzPool.imap_unorderedc                 C   s6   |   ¡  t| ||ƒ}| j |jd|||fgdf¡ |S )z;
        Asynchronous version of `apply()` method.
        r   N)r¥   ÚApplyResultre   r;   r¸   )r   rN   r   rO   r­   r®   rP   r   r   r   r¦   Â  s    zPool.apply_asyncc                 C   s   |   ||t|||¡S )z9
        Asynchronous version of `map()` method.
        )r¨   r   r¬   r   r   r   Ú	map_asyncÌ  s    ÿzPool.map_asyncc           
      C   sž   |   ¡  t|dƒst|ƒ}|dkrJtt|ƒt| jƒd ƒ\}}|rJ|d7 }t|ƒdkrZd}t |||¡}t| |t|ƒ||d�}	| j	 
|  |	j||¡df¡ |	S )zY
        Helper function to implement map, starmap and their async counterparts.
        Ú__len__Né   r   r   ©r®   )r¥   r=   r   Údivmodr‘   r_   r   rº   Ú	MapResultre   r;   r³   r¸   )
r   rN   rª   Zmapperr«   r­   r®   Zextrar¼   rP   r   r   r   r¨   Ô  s,    
ÿþüÿzPool._map_asyncc                 C   s"   t | |d� | ¡ s| ¡  qd S )N)Útimeout)r	   Úemptyr<   )r‰   Úchange_notifierrÇ   r   r   r   Ú_wait_for_updatesñ  s    zPool._wait_for_updatesc                 C   sp   t  ¡ }|jtks |rX|jtkrX|  |||||||	|
||¡
 |  |¡|•}|  ||¡ q| d ¡ t	 
d¡ d S )Nzworker handler exiting)ru   Úcurrent_threadr`   r   r   r¢   r—   rÊ   r;   r   rB   )r’   ÚcacheÚ	taskqueuer[   rZ   r…   r›   rF   rG   rH   rI   r†   rJ   r‰   rÉ   ÚthreadZcurrent_sentinelsr   r   r   rv   ÷  s       þ
zPool._handle_workersc                 C   sp  t  ¡ }t| jd ƒD ]ê\}}d }zÎ|D ]Š}|jtkrBt d¡  qâz||ƒ W q& tk
r® }
 zB|d d… \}	}z||	  	|d|
f¡ W n t
k
rœ   Y nX W 5 d }
~
X Y q&X q&|rÜt d¡ |rÌ|d nd}||d ƒ W ¢qW ¢
 �q
W 5 d  } }}	X qt d¡ z6t d¡ | d ¡ t d	¡ |D ]}|d ƒ �q.W n  tk
�r`   t d
¡ Y nX t d¡ d S )Nz'task handler found thread._state != RUNé   Fzdoing set_length()r   r°   ztask handler got sentinelz/task handler sending sentinel to result handlerz(task handler sending sentinel to workersz/task handler got OSError when sending sentinelsztask handler exiting)ru   rË   Úiterr<   r`   r   r   rB   rC   Ú_setÚKeyErrorr;   rA   )rÍ   r;   rG   r›   rÌ   rÎ   ZtaskseqZ
set_lengthrK   rL   rQ   Úidxrˆ   r   r   r   r}     sB    






zPool._handle_tasksc              	   C   sÈ  t  ¡ }z
|ƒ }W n$ ttfk
r6   t d¡ Y d S X |jtkr`|jtksTt	dƒ‚t d¡ q¶|d krtt d¡ q¶|\}}}z||  
||¡ W n tk
r¦   Y nX d  } }}q|�rR|jtk�rRz
|ƒ }W n$ ttfk
rö   t d¡ Y d S X |d k�rt d¡ q¶|\}}}z||  
||¡ W n tk
�rB   Y nX d  } }}q¶t| dƒ�r°t d¡ z,tdƒD ]}| j ¡ �sˆ �q’|ƒ  �qrW n ttfk
�r®   Y nX t d	t|ƒ|j¡ d S )
Nz.result handler got EOFError/OSError -- exitingzThread not in TERMINATEz,result handler found thread._state=TERMINATEzresult handler got sentinelz&result handler ignoring extra sentinelr?   z"ensuring that outqueue is not fullé
   z7result handler exiting: len(cache)=%s, thread._state=%s)ru   rË   rA   r@   r   rB   r`   r   r   r9   rÑ   rÒ   r=   rš   r?   Úpollr‘   )rG   r<   rÌ   rÎ   rK   rL   rM   Úobjr   r   r   r€   :  s^    











 ÿzPool._handle_resultsc                 c   s0   t |ƒ}tt ||¡ƒ}|s d S | |fV  qd S r   )rÐ   Útupler   Úislice)rN   ÚitÚsizer²   r   r   r   rº   v  s
    zPool._get_tasksc                 C   s   t dƒ‚d S )Nz:pool objects cannot be passed between processes or pickled)ÚNotImplementedErrorr   r   r   r   r+     s    ÿzPool.__reduce__c                 C   s2   t  d¡ | jtkr.t| _t| j_| j d ¡ d S )Nzclosing pool)r   rB   r`   r   r   rz   rf   r;   r   r   r   r   r>   „  s
    

z
Pool.closec                 C   s   t  d¡ t| _|  ¡  d S )Nzterminating pool)r   rB   r   r`   r„   r   r   r   r   rs   ‹  s    
zPool.terminatec                 C   sj   t  d¡ | jtkrtdƒ‚n| jttfkr4tdƒ‚| j ¡  | j	 ¡  | j
 ¡  | jD ]}| ¡  qXd S )Nzjoining poolzPool is still runningzIn unknown state)r   rB   r`   r   rm   r   r   rz   r(   r   r‚   r_   )r   rˆ   r   r   r   r(   �  s    






z	Pool.joinc                 C   s@   t  d¡ | j ¡  | ¡ r<| j ¡ r<| j ¡  t 	d¡ qd S )Nz7removing tasks from inqueue until task handler finishedr   )
r   rB   Z_rlockÚacquireÚis_aliver?   rÕ   r¤   ÚtimeÚsleep)rF   Útask_handlerrÚ   r   r   r   Ú_help_stuff_finishœ  s
    


zPool._help_stuff_finishc
                 C   sX  t  d¡ t|_| d ¡ t|_t  d¡ |  ||t|ƒ¡ | ¡ sXt|	ƒdkrXtdƒ‚t|_| d ¡ | d ¡ t  d¡ t	 
¡ |k	r�| ¡  |rÈt|d dƒrÈt  d¡ |D ]}
|
jd kr°|
 ¡  q°t  d¡ t	 
¡ |k	ræ| ¡  t  d	¡ t	 
¡ |k	�r| ¡  |�rTt|d dƒ�rTt  d
¡ |D ](}
|
 ¡ �r*t  d|
j ¡ |
 ¡  �q*d S )Nzfinalizing poolz&helping task handler/workers to finishr   z.Cannot have cache with result_hander not alivezjoining worker handlerrs   zterminating workerszjoining task handlerzjoining result handlerzjoining pool workersr˜   )r   rB   r   r`   r;   rá   r‘   rÝ   r9   ru   rË   r(   r=   rr   rs   Úpid)r’   rÍ   rF   rG   r›   rÉ   Zworker_handlerrà   Zresult_handlerrÌ   rˆ   r   r   r   rƒ   ¥  sB    


ÿ









zPool._terminate_poolc                 C   s   |   ¡  | S r   )r¥   r   r   r   r   Ú	__enter__Û  s    zPool.__enter__c                 C   s   |   ¡  d S r   )rs   )r   Úexc_typeZexc_valZexc_tbr   r   r   Ú__exit__ß  s    zPool.__exit__)NNr   NN)N)N)NNN)r   )r   )NNN)NNN)N)-r    r!   r"   r4   ry   ÚstaticmethodrZ   r   ÚwarningsÚwarnr   r�   r3   rt   r—   rœ   rq   r�   r¢   rb   r¥   r§   r   r   r¯   r³   r½   r¿   r¦   rÁ   r¨   rÊ   Úclassmethodrv   r}   r€   rº   r+   r>   rs   r(   rá   rƒ   rã   rå   r   r   r   r   r   ­   sx   
    ÿ
P

	



  ÿ


ÿ

  ÿ
  ÿ


-
;


5c                   @   s@   e Zd Zdd„ Zdd„ Zdd„ Zddd	„Zdd
d„Zdd„ ZdS )rÀ   c                 C   s>   || _ t ¡ | _ttƒ| _|j| _|| _|| _	| | j| j< d S r   )
r_   ru   ZEventÚ_eventÚnextÚjob_counterr¸   rg   Ú	_callbackÚ_error_callback)r   r›   r­   r®   r   r   r   r   è  s    

zApplyResult.__init__c                 C   s
   | j  ¡ S r   )rê   Zis_setr   r   r   r   Úreadyñ  s    zApplyResult.readyc                 C   s   |   ¡ std | ¡ƒ‚| jS )Nz{0!r} not ready)rï   rm   r:   Ú_successr   r   r   r   Ú
successfulô  s    zApplyResult.successfulNc                 C   s   | j  |¡ d S r   )rê   r	   ©r   rÇ   r   r   r   r	   ù  s    zApplyResult.waitc                 C   s,   |   |¡ |  ¡ st‚| jr"| jS | j‚d S r   )r	   rï   r   rð   Ú_valuerò   r   r   r   r<   ü  s    
zApplyResult.getc                 C   sZ   |\| _ | _| jr$| j r$|  | j¡ | jr<| j s<|  | j¡ | j ¡  | j| j= d | _d S r   )	rð   ró   rí   rî   rê   Úsetrg   r¸   r_   ©r   rM   rÖ   r   r   r   rÑ     s    

zApplyResult._set)N)N)	r    r!   r"   r   rï   rñ   r	   r<   rÑ   r   r   r   r   rÀ   æ  s   	

	rÀ   c                   @   s   e Zd Zdd„ Zdd„ ZdS )rÆ   c                 C   sh   t j| |||d� d| _d g| | _|| _|dkrNd| _| j ¡  | j| j	= n|| t
|| ƒ | _d S )NrÄ   Tr   )rÀ   r   rð   ró   Ú
_chunksizeÚ_number_leftrê   rô   rg   r¸   Úbool)r   r›   r«   Úlengthr­   r®   r   r   r   r     s    
ÿ
zMapResult.__init__c                 C   sÆ   |  j d8  _ |\}}|rv| jrv|| j|| j |d | j …< | j dkrÂ| jrZ|  | j¡ | j| j= | j ¡  d | _	nL|sŒ| jrŒd| _|| _| j dkrÂ| j
r¨|  
| j¡ | j| j= | j ¡  d | _	d S )Nr   r   F)r÷   rð   ró   rö   rí   rg   r¸   rê   rô   r_   rî   )r   rM   Zsuccess_resultÚsuccessrP   r   r   r   rÑ   $  s&    







zMapResult._setN)r    r!   r"   r   rÑ   r   r   r   r   rÆ     s   rÆ   c                   @   s:   e Zd Zdd„ Zdd„ Zddd„ZeZdd	„ Zd
d„ ZdS )r·   c                 C   sT   || _ t t ¡ ¡| _ttƒ| _|j| _t	 
¡ | _d| _d | _i | _| | j| j< d S )Nr   )r_   ru   Z	ConditionZLockÚ_condrë   rì   r¸   rg   ÚcollectionsÚdequeÚ_itemsÚ_indexÚ_lengthÚ	_unsorted)r   r›   r   r   r   r   B  s    

zIMapIterator.__init__c                 C   s   | S r   r   r   r   r   r   Ú__iter__M  s    zIMapIterator.__iter__Nc                 C   s´   | j �� z| j ¡ }W nz tk
r�   | j| jkr>d | _td ‚| j  |¡ z| j ¡ }W n2 tk
rŠ   | j| jkr€d | _td ‚t	d ‚Y nX Y nX W 5 Q R X |\}}|r¬|S |‚d S r   )
rû   rþ   ÚpopleftÚ
IndexErrorrÿ   r   r_   ÚStopIterationr	   r   )r   rÇ   rX   rú   r/   r   r   r   rë   P  s&    zIMapIterator.nextc              	   C   s¢   | j �’ | j|krn| j |¡ |  jd7  _| j| jkrb| j | j¡}| j |¡ |  jd7  _q,| j  ¡  n
|| j|< | j| jkr”| j| j	= d | _
W 5 Q R X d S ©Nr   )rû   rÿ   rþ   r    r  ÚpopÚnotifyr   rg   r¸   r_   rõ   r   r   r   rÑ   h  s    


zIMapIterator._setc              	   C   sB   | j �2 || _| j| jkr4| j  ¡  | j| j= d | _W 5 Q R X d S r   )rû   r   rÿ   r  rg   r¸   r_   )r   rù   r   r   r   r¹   y  s    

zIMapIterator._set_length)N)	r    r!   r"   r   r  rë   Ú__next__rÑ   r¹   r   r   r   r   r·   @  s   
r·   c                   @   s   e Zd Zdd„ ZdS )r¾   c              	   C   sV   | j �F | j |¡ |  jd7  _| j  ¡  | j| jkrH| j| j= d | _W 5 Q R X d S r  )	rû   rþ   r    rÿ   r  r   rg   r¸   r_   rõ   r   r   r   rÑ   ‡  s    

zIMapUnorderedIterator._setN)r    r!   r"   rÑ   r   r   r   r   r¾   …  s   r¾   c                   @   sV   e Zd ZdZedd„ ƒZddd„Zdd	„ Zd
d„ Zedd„ ƒZ	edd„ ƒZ
dd„ ZdS )r   Fc                 O   s   ddl m} |||ŽS )Nr   rY   )ZdummyrZ   )r[   r   rO   rZ   r   r   r   rZ   —  s    zThreadPool.ProcessNr   c                 C   s   t  | |||¡ d S r   )r   r   )r   r…   rH   rI   r   r   r   r   œ  s    zThreadPool.__init__c                 C   s,   t  ¡ | _t  ¡ | _| jj| _| jj| _d S r   )rc   rd   rw   rx   r;   r~   r<   r�   r   r   r   r   rb   Ÿ  s    


zThreadPool._setup_queuesc                 C   s
   | j jgS r   )rf   r?   r   r   r   r   rt   ¥  s    zThreadPool._get_sentinelsc                 C   s   g S r   r   r–   r   r   r   r—   ¨  s    z ThreadPool._get_worker_sentinelsc                 C   sF   z| j dd� qW n tjk
r(   Y nX t|ƒD ]}|  d ¡ q2d S )NF)Úblock)r<   rc   r   rš   r;   )rF   rà   rÚ   rM   r   r   r   rá   ¬  s    zThreadPool._help_stuff_finishc                 C   s   t  |¡ d S r   )rÞ   rß   )r   r‰   rÉ   rÇ   r   r   r   rÊ   ·  s    zThreadPool._wait_for_updates)NNr   )r    r!   r"   ry   ræ   rZ   r   rb   rt   r—   rá   rÊ   r   r   r   r   r   ”  s   




)Nr   NF))Ú__all__rü   r   rk   rc   ru   rÞ   r%   rç   r   r$   r   r   r   Z
connectionr	   r
   r   r   r   Úcountrì   r   r   rC   r   r#   r*   r-   rS   rD   ÚdictrT   Úobjectr   rÀ   ZAsyncResultrÆ   r·   r¾   r   r   r   r   r   Ú<module>
   sN   	  ÿ
-    =)+E