o
    éT•j÷  ã                   @   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	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é   )Úutil)Úget_contextÚTimeoutError)ÚwaitÚINITÚRUNÚCLOSEÚ	TERMINATEc                 C   s   t t| Ž ƒS ©N)ÚlistÚmap©Úargs© r   ú+/usr/lib/python3.10/multiprocessing/pool.pyÚmapstar/   ó   r   c                 C   s   t t | d | d ¡ƒS )Nr   r   )r   Ú	itertoolsÚstarmapr   r   r   r   Ústarmapstar2   s   r   c                   @   ó   e Zd Zdd„ Zdd„ ZdS )ÚRemoteTracebackc                 C   s
   || _ d S r   ©Útb)Úselfr   r   r   r   Ú__init__:   ó   
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                   @   r   )Ú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   ó   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Úsuperr0   r   )r   r+   r2   ©Ú	__class__r   r   r   T   s   

zMaybeEncodingError.__init__c                 C   s   d| j | jf S )Nz(Error sending result: '%s'. Reason: '%s')r2   r+   r    r   r   r   r!   Y   s   ÿzMaybeEncodingError.__str__c                 C   s   d| j j| f S )Nz<%s: %s>)r5   r"   r    r   r   r   Ú__repr__]   r.   zMaybeEncodingError.__repr__)r"   r#   r$   Ú__doc__r   r!   r6   Ú__classcell__r   r   r4   r   r0   P   s
    r0   r   Fc              
   C   sÐ  |d urt |tƒr|dkstd |¡ƒ‚|j}| j}t| dƒr)| j ¡  |j	 ¡  |d ur1||Ž  d}|d u s=|rß||k rßz|ƒ }	W n t
tfyR   t d¡ Y n�w |	d u r]t d¡ n‚|	\}
}}}}zd||i |¤Žf}W n" ty‘ } z|rƒ|turƒt||jƒ}d|f}W Y d }~nd }~ww z	||
||fƒ W n) tyÄ } zt||d ƒ}t d	| ¡ ||
|d|ffƒ W Y d }~nd }~ww d  }	 }
 } } }}|d7 }|d u s=|rß||k s=t 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Úhasattrr9   ÚcloseÚ_readerÚEOFErrorÚOSErrorr   ÚdebugÚ	ExceptionÚ_helper_reraises_exceptionr%   Ú__traceback__r0   )ÚinqueueÚoutqueueÚinitializerÚinitargsÚmaxtasksÚwrap_exceptionr>   r?   Ú	completedÚtaskÚjobÚiÚfuncr   ÚkwdsÚresultÚeÚwrappedr   r   r   Úworkera   sX   




þ
€ýÿ€üårX   c                 C   s   | ‚)z@Pickle-able helper function for use by _guarded_task_generation.r   )Úexr   r   r   rG   Ž   ó   rG   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|i |¤Ž d S r   )r]   r3   r   )r   r]   r   rT   r4   r   r   r   �   s   z_PoolCache.__init__c                    s$   t ƒ  |¡ | s| j d ¡ d S d S r   )r3   Ú__delitem__r]   r>   )r   Úitemr4   r   r   r^   ¡   s   ÿz_PoolCache.__delitem__)r"   r#   r$   r7   r   r^   r8   r   r   r4   r   r[   –   s    r[   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dMd!d"„Z		dNd#d$„Zd%d&„ ZdOd(d)„ZdOd*d+„Zdi ddfd,d-„Z		dNd.d/„Z		dNd0d1„ZedMd2d3„ƒ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 )Pr   zS
    Class which supports an async version of applying functions to arguments.
    Tc                 O   s   | j |i |¤ŽS r   ©ÚProcess)Úctxr   rT   r   r   r   ra   ³   s   zPool.ProcessNr   c                 C   s0  g | _ t| _|p
tƒ | _|  ¡  t ¡ | _| j ¡ | _	t
| j	d�| _|| _|| _|| _|d u r5t ¡ p4d}|dk r=tdƒ‚|d urNt|tƒrJ|dkrNtdƒ‚|d urZt|ƒsZtdƒ‚|| _z|  ¡  W n! ty„   | j D ]}|jd u rx| ¡  qm| j D ]}| ¡  q|‚ w |  ¡ }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/j0| | j1| j| j"| j#| j | j	| j%| j+| j.| jf	d	d
�| _2t'| _d S )Nr\   r   z&Number of processes must be at least 1r   z/maxtasksperchild must be a positive int or Nonezinitializer must be a callable©Útargetr   Té   )r   Úexitpriority)3Ú_poolr	   Ú_stater   Ú_ctxÚ_setup_queuesÚqueueÚSimpleQueueÚ
_taskqueueÚ_change_notifierr[   Ú_cacheÚ_maxtasksperchildÚ_initializerÚ	_initargsÚosÚ	cpu_countÚ
ValueErrorr:   r;   ÚcallableÚ	TypeErrorÚ
_processesÚ_repopulate_poolrF   ÚexitcodeÚ	terminater*   Ú_get_sentinelsÚ	threadingÚThreadr   Ú_handle_workersra   Ú_inqueueÚ	_outqueueÚ_wrap_exceptionÚ_worker_handlerÚdaemonr
   ÚstartÚ_handle_tasksÚ
_quick_putÚ_task_handlerÚ_handle_resultsÚ
_quick_getÚ_result_handlerr   ÚFinalizeÚ_terminate_poolÚ
_terminate)r   Ú	processesrK   rL   ÚmaxtasksperchildÚcontextÚpÚ	sentinelsr   r   r   r   ·   sˆ   


€

ú
ýþ
ÿþ
þ
þû
zPool.__init__c                 C   sF   | j |kr|d| ›�t| d� t| dd ƒd ur!| j d ¡ d S d S d S )Nz&unclosed running multiprocessing pool )Úsourcern   )rh   ÚResourceWarningÚgetattrrn   r>   )r   Ú_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=ú>)r5   r#   r$   rh   Úlenrg   )r   Úclsr   r   r   r6     s   ÿþzPool.__repr__c                 C   s    | j jg}| jjg}g |¢|¢S r   )r�   rB   rn   )r   Útask_queue_sentinelsÚself_notifier_sentinelsr   r   r   r|     s   

zPool._get_sentinelsc                 C   s   dd„ | D ƒS )Nc                 S   s   g | ]
}t |d ƒr|j‘qS )Úsentinel)r@   r    )Ú.0rX   r   r   r   Ú
<listcomp>  s    ÿz.Pool._get_worker_sentinels.<locals>.<listcomp>r   ©Úworkersr   r   r   Ú_get_worker_sentinels  s   ÿzPool._get_worker_sentinelsc                 C   sP   d}t tt| ƒƒƒD ]}| | }|jdu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œ   rz   r   rE   r*   )ÚpoolÚcleanedrR   rX   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_staticri   ra   rx   rg   r€   r�   rq   rr   rp   r‚   r    r   r   r   ry   1  s   úzPool._repopulate_poolc
              
   C   sf   t |t|ƒ ƒD ](}
|| 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.
        rc   ra   Ú
PoolWorkerTzadded workerN)
r¨   rœ   rX   ÚnameÚreplacer„   r…   Úappendr   rE   )rb   ra   r�   r©   rI   rJ   rK   rL   r�   rN   rR   Úwr   r   r   r¬   :  s   ýÿ
özPool._repopulate_pool_staticc
           
      C   s.   t  |¡rt  | |||||||||	¡
 dS dS )zEClean up any exited workers and start replacements for them.
        N)r   r«   r¬   )
rb   ra   r�   r©   rI   rJ   rK   rL   r�   rN   r   r   r   Ú_maintain_poolM  s   
ýÿzPool._maintain_poolc                 C   s4   | j  ¡ | _| j  ¡ | _| jjj| _| jjj| _	d S r   )
ri   rl   r€   r�   r9   Úsendr‡   rB   ÚrecvrŠ   r    r   r   r   rj   Y  s   zPool._setup_queuesc                 C   s   | j tkr	tdƒ‚d S )NzPool not running)rh   r
   ru   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   rS   r   rT   r   r   r   Úapplyc  s   z
Pool.applyc                 C   ó   |   ||t|¡ ¡ S )zx
        Apply `func` to each element in `iterable`, collecting the results
        in a list that is returned.
        )Ú
_map_asyncr   r?   ©r   rS   ÚiterableÚ	chunksizer   r   r   r   j  s   zPool.mapc                 C   r¸   )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   q  s   zPool.starmapc                 C   ó   |   ||t|||¡S )z=
        Asynchronous version of `starmap()` method.
        )r¹   r   ©r   rS   r»   r¼   ÚcallbackÚerror_callbackr   r   r   Ústarmap_asyncy  s   ÿzPool.starmap_asyncc              
   c   sn   � zd}t |ƒD ]\}}||||fi fV  qW dS  ty6 } z||d t|fi fV  W Y d}~dS d}~ww )zšProvides a generator of tasks for imap and imap_unordered with
        appropriate handling for iterables which throw exceptions during
        iteration.éÿÿÿÿr   N)Ú	enumeraterF   rG   )r   Ú
result_jobrS   r»   rR   ÚxrV   r   r   r   Ú_guarded_task_generation�  s   €ÿ$€ÿzPool._guarded_task_generationr   c                 C   ó’   |   ¡  |dkrt| ƒ}| j |  |j||¡|jf¡ |S |dk r(td |¡ƒ‚t	 
|||¡}t| ƒ}| j |  |jt|¡|jf¡ dd„ |D ƒS )zP
        Equivalent of `map()` -- can be MUCH slower than `Pool.map()`.
        r   zChunksize must be 1+, not {0:n}c                 s   ó   � | ]
}|D ]}|V  qqd S r   r   ©r¡   Úchunkr_   r   r   r   Ú	<genexpr>§  ó   € zPool.imap.<locals>.<genexpr>)rµ   ÚIMapIteratorrm   r>   rÆ   Ú_jobÚ_set_lengthru   r=   r   Ú
_get_tasksr   ©r   rS   r»   r¼   rU   Útask_batchesr   r   r   ÚimapŒ  s4   þÿÿÿþüÿz	Pool.imapc                 C   rÇ   )zL
        Like `imap()` method but ordering of results is arbitrary.
        r   zChunksize must be 1+, not {0!r}c                 s   rÈ   r   r   rÉ   r   r   r   rË   Ã  rÌ   z&Pool.imap_unordered.<locals>.<genexpr>)rµ   ÚIMapUnorderedIteratorrm   r>   rÆ   rÎ   rÏ   ru   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µ   ÚApplyResultrm   r>   rÎ   )r   rS   r   rT   r¿   rÀ   rU   r   r   r   r¶   Å  s   zPool.apply_asyncc                 C   r½   )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u r%tt|ƒt| jƒd ƒ\}}|r%|d7 }t|ƒdkr-d}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œ   rg   r   rÐ   Ú	MapResultrm   r>   rÆ   rÎ   )
r   rS   r»   Úmapperr¼   r¿   rÀ   ÚextrarÒ   rU   r   r   r   r¹   ×  s,   
ÿþüÿzPool._map_asyncc                 C   s,   t | |d� | ¡ s| ¡  | ¡ r
d S d S )N)Útimeout)r   Úemptyr?   )r“   Úchange_notifierrß   r   r   r   Ú_wait_for_updatesô  s   ÿzPool._wait_for_updatesc                 C   sŠ   t  ¡ }|jtks|r9|jtkr9|  |||||||	|
||¡
 g |  |¡¢|¢}|  ||¡ |jtks|r9|jtks| d ¡ t	 
d¡ d S )Nzworker handler exiting)r}   Úcurrent_threadrh   r
   r   r²   r¥   râ   r>   r   rE   )r�   ÚcacheÚ	taskqueuerb   ra   r�   r©   rI   rJ   rK   rL   r�   rN   r“   rá   ÚthreadÚcurrent_sentinelsr   r   r   r   ú  s   þù
	zPool._handle_workersc                 C   st  t  ¡ }t| jd ƒD ]z\}}d }zm|D ]D}|jtkr!t d¡  nTz||ƒ W q tyW }	 z$|d d… \}
}z||
  	|d|	f¡ W n	 t
yL   Y nw W Y d }	~	qd }	~	ww |rmt d¡ |re|d nd}||d ƒ W d  } }}
q
W d  } }}
 nd  } }}
w t d¡ zt d¡ | d ¡ t d	¡ |D ]}|d ƒ qœW n ty²   t d
¡ Y nw 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)r}   rã   Úiterr?   rh   r
   r   rE   rF   Ú_setÚKeyErrorr>   rD   )rå   r>   rJ   r©   rä   ræ   ÚtaskseqÚ
set_lengthrP   rV   rQ   Úidxr’   r   r   r   r†     sN   

ÿ€ü
þ




ÿÿzPool._handle_tasksc              	   C   sº  t  ¡ }	 z|ƒ }W n ttfy   t d¡ Y d S w |jtkr0|jtks*J dƒ‚t d¡ n*|d u r:t d¡ n |\}}}z
||  	||¡ W n	 t
yR   Y nw d  } }}q|r¨|jtkr¨z|ƒ }W n ttfyw   t d¡ Y d S w |d u r‚t d¡ qZ|\}}}z
||  	||¡ W n	 t
yš   Y nw d  } }}|r¨|jtksat| dƒrÑt d¡ ztd	ƒD ]}| j ¡ sÀ n|ƒ  q·W n ttfyÐ   Y nw t d
t|ƒ|j¡ d S )Nr   z.result handler got EOFError/OSError -- exitingzThread not in TERMINATEz,result handler found thread._state=TERMINATEzresult handler got sentinelz&result handler ignoring extra sentinelrB   z"ensuring that outqueue is not fullé
   z7result handler exiting: len(cache)=%s, thread._state=%s)r}   rã   rD   rC   r   rE   rh   r
   r   rê   rë   r@   r¨   rB   Úpollrœ   )rJ   r?   rä   ræ   rP   rQ   rR   Úobjr   r   r   r‰   =  sn   

þ



ÿë

þ

ÿñ


€ÿ
ÿzPool._handle_resultsc                 c   s0   � t |ƒ}	 tt ||¡ƒ}|sd S | |fV  qr   )ré   Útupler   Úislice)rS   ÚitÚsizerÅ   r   r   r   rÐ   y  s   €
üzPool._get_tasksc                 C   s   t dƒ‚)Nz:pool objects cannot be passed between processes or pickled)ÚNotImplementedErrorr    r   r   r   r-   ‚  s   ÿzPool.__reduce__c                 C   s6   t  d¡ | jtkrt| _t| j_| j d ¡ d S d S )Nzclosing pool)r   rE   rh   r
   r   rƒ   rn   r>   r    r   r   r   rA   ‡  s   

ýz
Pool.closec                 C   s   t  d¡ t| _|  ¡  d S )Nzterminating pool)r   rE   r   rh   rŽ   r    r   r   r   r{   Ž  s   
zPool.terminatec                 C   sh   t  d¡ | jtkrtdƒ‚| jttfvrtdƒ‚| j ¡  | j	 ¡  | j
 ¡  | jD ]}| ¡  q+d S )Nzjoining poolzPool is still runningzIn unknown state)r   rE   rh   r
   ru   r   r   rƒ   r*   rˆ   r‹   rg   )r   r’   r   r   r   r*   “  s   






ÿz	Pool.joinc                 C   s\   t  d¡ | j ¡  | ¡ r(| j ¡ r,| j ¡  t 	d¡ | ¡ r*| j ¡ sd S d S d S d S )Nz7removing tasks from inqueue until task handler finishedr   )
r   rE   Ú_rlockÚacquireÚis_aliverB   rð   r´   ÚtimeÚsleep)rI   Útask_handlerrõ   r   r   r   Ú_help_stuff_finishŸ  s   



"þzPool._help_stuff_finishc
                 C   sV  t  d¡ t|_| d ¡ t|_t  d¡ |  ||t|ƒ¡ | ¡ s,t|	ƒdkr,tdƒ‚t|_| d ¡ | d ¡ t  d¡ t	 
¡ |urH| ¡  |rdt|d dƒrdt  d¡ |D ]}
|
jd u rc|
 ¡  qXt  d¡ t	 
¡ |urs| ¡  t  d	¡ t	 
¡ |ur‚| ¡  |r¥t|d dƒr§t  d
¡ |D ]}
|
 ¡ r¤t  d|
j ¡ |
 ¡  q’d S d S d S )Nzfinalizing poolz&helping task handler/workers to finishr   z.Cannot have cache with result_hander not alivezjoining worker handlerr{   zterminating workerszjoining task handlerzjoining result handlerzjoining pool workersr¦   )r   rE   r   rh   r>   rý   rœ   rù   r<   r}   rã   r*   r@   rz   r{   Úpid)r�   rå   rI   rJ   r©   rá   Úworker_handlerrü   Úresult_handlerrä   r’   r   r   r   r�   ¨  sJ   


ÿ




€


€úzPool._terminate_poolc                 C   s   |   ¡  | S r   )rµ   r    r   r   r   Ú	__enter__Þ  s   zPool.__enter__c                 C   s   |   ¡  d S r   )r{   )r   Úexc_typeÚexc_valÚexc_tbr   r   r   Ú__exit__â  r   zPool.__exit__)NNr   NNr   )NNN)r   )-r"   r#   r$   r7   r‚   Ústaticmethodra   r   ÚwarningsÚwarnr
   r˜   r6   r|   r¥   r«   ry   r¬   r²   rj   rµ   r·   r   r   rÁ   rÆ   rÓ   rÕ   r¶   r×   r¹   râ   Úclassmethodr   r†   r‰   rÐ   r-   rA   r{   r*   rý   r�   r  r  r   r   r   r   r   ­   sx    

ÿS

	




ÿ


ÿ

ÿ
ÿ

-
;


5c                   @   sJ   e Zd Zdd„ Zdd„ Zdd„ Zddd	„Zdd
d„Zdd„ Ze	e
jƒZdS )rÖ   c                 C   s>   || _ t ¡ | _ttƒ| _|j| _|| _|| _	| | j| j< d S r   )
rg   r}   ÚEventÚ_eventÚnextÚjob_counterrÎ   ro   Ú	_callbackÚ_error_callback)r   r©   r¿   rÀ   r   r   r   r   ë  s   

zApplyResult.__init__c                 C   s
   | j  ¡ S r   )r  Úis_setr    r   r   r   Úreadyô  r   zApplyResult.readyc                 C   s   |   ¡ std | ¡ƒ‚| jS )Nz{0!r} not ready)r  ru   r=   Ú_successr    r   r   r   Ú
successful÷  s   zApplyResult.successfulNc                 C   s   | j  |¡ d S r   )r  r   ©r   rß   r   r   r   r   ü  r.   zApplyResult.waitc                 C   s(   |   |¡ |  ¡ st‚| jr| jS | j‚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  Úsetro   rÎ   rg   ©r   rR   rñ   r   r   r   rê     s   


zApplyResult._setr   )r"   r#   r$   r   r  r  r   r?   rê   r	  ÚtypesÚGenericAliasÚ__class_getitem__r   r   r   r   rÖ   é  s    	

	
rÖ   c                   @   r   )rÜ   c                 C   sj   t j| |||d� d| _d g| | _|| _|dkr(d| _| j ¡  | j| j	= d S || t
|| ƒ | _d S )NrÚ   Tr   )rÖ   r   r  r  Ú
_chunksizeÚ_number_leftr  r  ro   rÎ   Úbool)r   r©   r¼   Úlengthr¿   rÀ   r   r   r   r     s   
ÿ
zMapResult.__init__c                 C   sÐ   |  j d8  _ |\}}|r>| jr>|| j|| j |d | j …< | j dkr<| jr-|  | j¡ | j| j= | j ¡  d | _	d S d S |sI| jrId| _|| _| j dkrf| j
rW|  
| j¡ | j| j= | j ¡  d | _	d S d S )Nr   r   F)r  r  r  r  r  ro   rÎ   r  r  rg   r  )r   rR   Úsuccess_resultÚsuccessrU   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   )rg   r}   Ú	ConditionÚLockÚ_condr  r  rÎ   ro   ÚcollectionsÚdequeÚ_itemsÚ_indexÚ_lengthÚ	_unsorted)r   r©   r   r   r   r   G  s   

zIMapIterator.__init__c                 C   s   | S r   r   r    r   r   r   Ú__iter__R  s   zIMapIterator.__iter__Nc                 C   s¼   | j �I z| j ¡ }W n9 tyD   | j| jkrd | _td ‚| j  |¡ z| j ¡ }W n tyA   | j| jkr>d | _td ‚t	d ‚w Y nw W d   ƒ n1 sOw   Y  |\}}|r\|S |‚r   )
r#  r&  ÚpopleftÚ
IndexErrorr'  r(  rg   ÚStopIterationr   r   )r   rß   r_   r   r2   r   r   r   r  U  s0   üÿú€ýzIMapIterator.nextc                 C   sÒ   | j �\ | j|kr<| j |¡ |  jd7  _| j| jv r6| j | j¡}| j |¡ |  jd7  _| j| jv s| j  ¡  n|| j|< | j| jkrW| j| j	= d | _
W d   ƒ d S W d   ƒ d S 1 sbw   Y  d S ©Nr   )r#  r'  r&  r°   r)  ÚpopÚnotifyr(  ro   rÎ   rg   r  r   r   r   rê   m  s"   
ý

ò"ôzIMapIterator._setc                 C   sh   | j �' || _| j| jkr"| j  ¡  | j| j= d | _W d   ƒ d S W d   ƒ d S 1 s-w   Y  d S r   )r#  r(  r'  r0  ro   rÎ   rg   )r   r  r   r   r   rÏ   ~  s   

û"þzIMapIterator._set_lengthr   )	r"   r#   r$   r   r*  r  Ú__next__rê   rÏ   r   r   r   r   rÍ   E  s    
rÍ   c                   @   s   e Zd Zdd„ ZdS )rÔ   c                 C   s|   | j �1 | j |¡ |  jd7  _| j  ¡  | j| jkr,| j| j= d | _W d   ƒ d S W d   ƒ d S 1 s7w   Y  d S r.  )	r#  r&  r°   r'  r0  r(  ro   rÎ   rg   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} ||i |¤ŽS )Nr   r`   )Údummyra   )rb   r   rT   ra   r   r   r   ra   œ  s   zThreadPool.ProcessNr   c                 C   s   t  | |||¡ d S r   )r   r   )r   r�   rK   rL   r   r   r   r   ¡  s   zThreadPool.__init__c                 C   s,   t  ¡ | _t  ¡ | _| jj| _| jj| _d S r   )rk   rl   r€   r�   r>   r‡   r?   rŠ   r    r   r   r   rj   ¤  s   


zThreadPool._setup_queuesc                 C   s
   | j jgS r   )rn   rB   r    r   r   r   r|   ª  r   zThreadPool._get_sentinelsc                 C   s   g S r   r   r£   r   r   r   r¥   ­  rZ   z ThreadPool._get_worker_sentinelsc                 C   sB   z	 | j dd� q tjy   Y nw t|ƒD ]}|  d ¡ qd S )NTF)Úblock)r?   rk   ÚEmptyr¨   r>   )rI   rü   rõ   rR   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$   r‚   r  ra   r   rj   r|   r¥   rý   râ   r   r   r   r   r   ™  s    




)Nr   NF))Ú__all__r$  r   rs   rk   r}   rú   r'   r  r  r&   r   r   r   Ú
connectionr   r	   r
   r   r   Úcountr  r   r   rF   r   r%   r,   r0   rX   rG   Údictr[   Úobjectr   rÖ   ÚAsyncResultrÜ   rÍ   rÔ   r   r   r   r   r   Ú<module>   sP   		
ÿ-    @++E