o
    svXj¼ ã                   @   s0  d dl m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Zd dlmZ d dlmZ ddlmZmZ ddl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! ddl"m#Z#m$Z$m%Z%m&Z&m'Z'm(Z(m)Z) ddl*m+Z+m,Z,m-Z-m.Z.m/Z/m0Z0 ddlm1Z1m2Z2m3Z3 dZ4ej5d  dkZ6e 7¡ dkr°ddl8m9Z: eZ;n	d dlm<Z: ej;Z;ze	j=Z=W n e>yÉ   dZ=Y nw ej5dkrÓe	j?Z@ne	j@Z@d ZAdZBdZCd ZDdZEdZFdZGdZHd ZIdZJdZKeLeddƒZMdZNeLedd ƒZIdZOdZPe Q¡ ZRe	jSZSdd„ ZTd d!„ ZUd"d#„ ZVd$d%„ ZWdHd&d'„ZXG d(d)„ d)e@ƒZYG d*d+„ d+eZƒZ[G d,d-„ d-eZƒZ\d.d/„ Z]G d0d1„ d1e^ƒZ_G d2d3„ d3e!ƒZ`G d4d5„ d5e`ƒZaG d6d7„ d7e`ƒZbG d8d9„ d9e`ƒZcG d:d;„ d;e`ƒZdG d<d=„ d=e^ƒZeG d>d?„ d?e^ƒZfG d@dA„ dAefƒZgG dBdC„ dCe^ƒZhG dDdE„ dEehƒZiG dFdG„ dGeeƒZjdS )Ié    )Úabsolute_importN)Údeque)Úpartialé   )Ú	cpu_countÚget_context)Úutil)ÚTERM_SIGNALÚhuman_statusÚpickle_loadsÚreset_signalsÚrestart_state)Ú	get_errnoÚmem_rssÚsend_offset)ÚExceptionInfo)ÚDummyProcess)ÚCoroStopÚRestartFreqExceededÚSoftTimeLimitExceededÚ
TerminatedÚTimeLimitExceededÚTimeoutErrorÚWorkerLostError)ÚEmptyÚQueueÚrangeÚvaluesÚreraiseÚ	monotonic)ÚFinalizeÚdebugÚwarningzEchild process exiting after exceeding memory limit ({0}KiB / {1}KiB)
é   ÚWindows)Úkill_processtree)Úkillg    _ B)r#   r#   é   é   é›   ÚSIGUSR1g      $@ÚEX_OKi,  çš™™™™™¹?c                 C   s<   z| j }W n ty   d }Y nw |d u rtt |  ¡ ƒS |S ©N)r   ÚAttributeErrorr   Úfileno)Ú
connectionÚnative© r2   úJ/var/www/html/myproject/venv/lib/python3.10/site-packages/billiard/pool.pyÚ_get_send_offsety   s   
ÿr4   c                 C   s   t t| Ž ƒS r-   )ÚlistÚmap©Úargsr2   r2   r3   Úmapstarƒ   ó   r9   c                 C   s   t t | d | d ¡ƒS )Nr   r   )r5   Ú	itertoolsÚstarmapr7   r2   r2   r3   Ústarmapstar‡   s   r=   c                 O   s    t  ¡ j| g|¢R i |¤Ž d S r-   )r   Ú
get_loggerÚerror)Úmsgr8   Úkwargsr2   r2   r3   r?   ‹   s    r?   c                 C   s   | t  ¡ ur|  |¡ d S d S r-   )Ú	threadingÚcurrent_threadÚstop)ÚthreadÚtimeoutr2   r2   r3   Ústop_if_not_current�   s   ÿrG   c                   @   sd   e Zd ZdZdd„ Zerddd„Zdd	„ Zd
d„ Zdd„ Z	dS ddd„Zdd	„ Zdd„ Zdd„ Z	dS )ÚLaxBoundedSemaphorez^Semaphore that checks that # release is <= # acquires,
    but ignores if # releases >= value.c                 C   s   |  j d8  _ |  ¡  d S ©Nr   )Ú_initial_valueÚacquire©Úselfr2   r2   r3   Úshrink˜   s   zLaxBoundedSemaphore.shrinkr   Nc                 C   s   t  | |¡ || _d S r-   ©Ú
_SemaphoreÚ__init__rJ   ©rM   ÚvalueÚverboser2   r2   r3   rQ   ž   s   
zLaxBoundedSemaphore.__init__c                 C   sR   | j � |  jd7  _|  jd7  _| j  ¡  W d   ƒ d S 1 s"w   Y  d S rI   )Ú_condrJ   Ú_valueÚnotifyrL   r2   r2   r3   Úgrow¢   s
   "ýzLaxBoundedSemaphore.growc                 C   ób   | j }|�" | j| jk r|  jd7  _| ¡  W d   ƒ d S W d   ƒ d S 1 s*w   Y  d S rI   )rU   rV   rJ   Ú
notify_all©rM   Úcondr2   r2   r3   Úrelease¨   ó   
ý"ÿzLaxBoundedSemaphore.releasec                 C   ó*   | j | jk rt | ¡ | j | jk sd S d S r-   )rV   rJ   rP   r]   rL   r2   r2   r3   Úclear¯   ó   
ÿzLaxBoundedSemaphore.clearc                 C   s   t  | ||¡ || _d S r-   rO   rR   r2   r2   r3   rQ   ´   s   
c                 C   sT   | j }|� |  jd7  _|  jd7  _| ¡  W d   ƒ d S 1 s#w   Y  d S rI   )Ú_Semaphore__condrJ   Ú_Semaphore__valuerW   r[   r2   r2   r3   rX   ¸   s   
"ýc                 C   rY   rI   )rb   rc   rJ   Ú	notifyAllr[   r2   r2   r3   r]   ¿   r^   c                 C   r_   r-   )rc   rJ   rP   r]   rL   r2   r2   r3   r`   Æ   ra   ©r   N)
Ú__name__Ú
__module__Ú__qualname__Ú__doc__rN   ÚPY3rQ   rX   r]   r`   r2   r2   r2   r3   rH   ”   s    

rH   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-   )ÚreprÚexcrS   Úsuperrk   rQ   )rM   rm   rS   ©Ú	__class__r2   r3   rQ   Ó   s   

zMaybeEncodingError.__init__c                 C   s   d| j jt| ƒf S )Nz<%s: %s>)rp   rf   ÚstrrL   r2   r2   r3   Ú__repr__Ø   ó   zMaybeEncodingError.__repr__c                 C   s   d| j | jf S )Nz)Error sending result: '%r'. Reason: '%r'.)rS   rm   rL   r2   r2   r3   Ú__str__Û   s   ÿzMaybeEncodingError.__str__)rf   rg   rh   ri   rQ   rr   rt   Ú__classcell__r2   r2   ro   r3   rk   Ï   s
    rk   c                   @   s   e Zd ZdZdS )ÚWorkersJoinedzAll workers have terminated.N)rf   rg   rh   ri   r2   r2   r2   r3   rv   à   s    rv   c                 C   s   t ƒ ‚r-   )r   )ÚsignumÚframer2   r2   r3   Úsoft_timeout_sighandlerä   ó   ry   c                   @   sŒ   e Zd Z				ddd„Zdd„ Zdd	„ Zd
d„ Zddd„Zdd„ Zdd„ Z	e
edfdd„Zdd„ Zdd„ Zdd„ Zefdd„Zdd„ ZdS ) ÚWorkerNr2   Tc                 C   sz   |d u st |ƒtkr|dksJ ‚|| _|| _|| _|| _|| _|| _|	| _|||| _	| _
| _|
| _|| _|  | ¡ d S ©Nr   )ÚtypeÚintÚinitializerÚinitargsÚmaxtasksÚmax_memory_per_childÚ	_shutdownÚon_exitÚsigprotectionÚinqÚoutqÚsynqÚwrap_exceptionÚon_ready_counterÚcontribute_to_object)rM   r†   r‡   rˆ   r   r€   r�   Úsentinelr„   r…   r‰   r‚   rŠ   r2   r2   r3   rQ   î   s    zWorker.__init__c                 C   s¦   | j | j| j|_ |_|_| j j ¡ |_| jj ¡ |_| jr5| jj ¡ |_| jj ¡ |_	t
| jjƒ|_n	d  |_ |_	|_| j jj|_| jjj|_t
| j jƒ|_|S r-   )r†   r‡   rˆ   Ú_writerr/   ÚinqW_fdÚ_readerÚoutqR_fdÚsynqR_fdÚsynqW_fdr4   Úsend_syn_offsetÚ_send_syn_offsetÚsendÚ
_quick_putÚrecvÚ
_quick_getÚsend_job_offset)rM   Úobjr2   r2   r3   r‹   ÿ   s   zWorker.contribute_to_objectc                 C   s6   | j | j| j| j| j| j| j| j| j| j	| j
| jffS r-   )rp   r†   r‡   rˆ   r   r€   r�   rƒ   r„   r…   r‰   r‚   rL   r2   r2   r3   Ú
__reduce__  s
   ýzWorker.__reduce__c                    sê   t j‰ d g‰d‡ ‡fdd„	}|t _t ¡ }|  ¡  |  ¡  | j|d� zGzt  | j|d�¡ W n# tyR } zt	d| |dd� |  
|ˆd |¡ W Y d }~nd }~ww W |  
|ˆd d ¡ d S W |  
|ˆd d ¡ d S |  
|ˆd d ¡ w )	Nc                    s   | ˆd< ˆ | ƒS r|   r2   )Ústatus©Ú_exitÚ	_exitcoder2   r3   Úexit  s   zWorker.__call__.<locals>.exit©ÚpidzPool process %r error: %rr   ©Úexc_infor   r-   )Úsysr    ÚosÚgetpidÚ_make_child_methodsÚ
after_forkÚon_loop_startÚworkloopÚ	Exceptionr?   Ú_do_exit)rM   r    r¢   rm   r2   r�   r3   Ú__call__  s&   €þÿþ*zWorker.__call__c              	   C   s~   |d u r
|rt nt}| jd ur|  ||¡ tjdkr8z| j t||ff¡ t 	d¡ W t
 |¡ d S t
 |¡ w t
 |¡ d S )NÚwin32r   )Ú
EX_FAILUREr+   r„   r¥   Úplatformr‡   ÚputÚDEATHÚtimeÚsleepr¦   rž   )rM   r¢   Úexitcoderm   r2   r2   r3   r­   +  s   

zWorker._do_exitc                 C   ó   d S r-   r2   ©rM   r¢   r2   r2   r3   rª   ;  ó   zWorker.on_loop_startc                 C   s   |S r-   r2   )rM   Úresultr2   r2   r3   Úprepare_result>  r¹   zWorker.prepare_resultc              
      s@  |pt  ¡ }ˆjj}ˆj}ˆj}ˆj}ˆjpd}ˆj}	ˆj	}
ˆj
‰ ‡ ‡fdd„}d}zî|d u s5|rø||k rø|
ƒ }|rî|\}}|tksDJ ‚|\}}}}}|t|||ƒ ||ffƒ ˆ r`||ƒ}|s`q+zd|	||i |¤Žƒf}W n ty{   dtƒ f}Y nw z|t||||ffƒ W n9 tyÁ } z-t ¡ \}}}zt||d ƒ}tt||fƒ}|t||d|f|ffƒ W ~n~w W Y d }~nd }~ww |d7 }|dkrîtƒ }|dkrÕtdƒ |dkrî||krîtt ||¡ƒ tW ˆj|d� S |d u s5|rø||k s5|d	|ƒ |�r||k�rtntW ˆj|d� S tW ˆj|d� S ˆj|d� w )
Nr   c                    s^   d}	 |dkrt d| ˆjj ¡ dd� ˆ ƒ }|r*|\}}|tkr"dS |tks(J ‚dS |d7 }q)Nr   r   é<   z(!!!WAIT FOR ACK TIMEOUT: job:%r fd:%r!!!r£   FT)r?   rˆ   r�   r/   ÚNACKÚACK)ÚjidÚiÚreqÚtype_r8   ©Ú_wait_for_synrM   r2   r3   Úwait_for_synM  s   ÿõz%Worker.workloop.<locals>.wait_for_synTFr   z'worker unable to determine memory usage)Ú	completedzworker exiting after %d tasks)r¦   r§   r‡   r²   rŽ   r’   r�   r‚   r»   Úwait_for_jobrÅ   ÚTASKr¾   r¬   r   ÚREADYr¥   r¤   rk   r   r?   r"   ÚMAXMEM_USED_FMTÚformatÚ
EX_RECYCLEÚ_ensure_messages_consumedr°   r+   )rM   r!   Únowr¢   r²   rŽ   r’   r�   r‚   r»   rÇ   rÅ   rÆ   rÁ   rÂ   Úargs_ÚjobrÀ   Úfunr8   rA   Úconfirmrº   rm   Ú_ÚtbÚwrappedÚeinfoÚused_kbr2   rÃ   r3   r«   A  sv   
ÿÿ€÷
ÿÒ
%úzWorker.workloopc                 C   sJ   | j sdS ttƒD ]}| j j|krtd|ƒ  dS t t¡ q	tdƒ dS )zr Returns true if all messages sent out have been received and
        consumed within a reasonable amount of time Fz*ensured messages consumed after %d retriesTz<could not ensure all messages were consumed prior to exiting)	rŠ   r   Ú)GUARANTEE_MESSAGE_CONSUMPTION_RETRY_LIMITrS   r!   r´   rµ   Ú,GUARANTEE_MESSAGE_CONSUMPTION_RETRY_INTERVALr"   )rM   rÆ   Úretryr2   r2   r3   rÍ   Ž  s   
z Worker._ensure_messages_consumedc                 C   s’   t | jdƒr| jj ¡  t | jdƒr| jj ¡  | jd ur#| j| jŽ  t| j	d� t
d ur3t t
t¡ zt tjtj¡ W d S  tyH   Y d S w )Nr�   r�   )Úfull)Úhasattrr†   r�   Úcloser‡   r�   r   r€   r   r…   ÚSIG_SOFT_TIMEOUTÚsignalry   ÚSIGINTÚSIG_IGNr.   rL   r2   r2   r3   r©   ž  s   
ÿzWorker.after_forkc                    sd   |j ‰t|dƒr*|jj‰ t|dƒr!|jr!|j‰tf‡fdd„	}|S ‡ ‡fdd„}|S ‡fdd„}|S )Nr�   Úget_payloadc                    s   d|ˆ ƒ ƒfS ©NTr2   )rF   Úloads)râ   r2   r3   Ú_recv¼  ó   z'Worker._make_recv_method.<locals>._recvc                    s   ˆ | ƒr	dˆƒ fS dS ©NT©FNr2   ©rF   )Ú_pollÚgetr2   r3   rå   ¿  s   
c                    s(   zdˆ | d�fW S  t jy   Y dS w ©NTré   rè   )r   r   ré   )rë   r2   r3   rå   Ä  s
   ÿ)rë   rÜ   r�   Úpollrâ   r   )rM   Úconnrå   r2   )rê   rë   râ   r3   Ú_make_recv_method´  s   
ö
ûzWorker._make_recv_methodc                 C   s0   |   | j¡| _| jr|   | j¡| _d S d | _d S r-   )Ú_make_protected_receiver†   rÇ   rˆ   rÅ   )rM   rä   r2   r2   r3   r¨   Ë  s
   ÿÿzWorker._make_child_methodsc                    s2   |   |¡‰ | jr| jjnd ‰tf‡ ‡fdd„	}|S )Nc              
      s¢   ˆrˆƒ r| dƒ t tƒ‚zˆ dƒ\}}|sW d S W n( ttfyB } zt|ƒtjkr2W Y d }~d S | dt|ƒjƒ t t	ƒ‚d }~ww |d u rO| dƒ t t	ƒ‚|S )Nzworker got sentinel -- exitingç      ð?zworker got %s -- exiting)
Ú
SystemExitr+   ÚEOFErrorÚIOErrorr   ÚerrnoÚEINTRr}   rf   r°   )r!   ÚreadyrÁ   rm   ©Ú_receiveÚshould_shutdownr2   r3   ÚreceiveÔ  s&   
ÿ€üz/Worker._make_protected_receive.<locals>.receive)rï   rƒ   Úis_setr!   )rM   rî   rû   r2   rø   r3   rð   Ð  s   
zWorker._make_protected_receive)
NNr2   NNNTTNNr-   )rf   rg   rh   rQ   r‹   r›   r®   r­   rª   r»   r!   r   r«   rÍ   r©   rï   r   r¨   rð   r2   r2   r2   r3   r{   ì   s$    
ý
Mr{   c                       sN   e Zd Zdd„ Zdd„ Z‡ fdd„Zdd„ Zdd
d„Zdd„ Zdd„ Z	‡  Z
S )Ú
PoolThreadc                 O   s    t  | ¡ t| _d| _d| _d S ©NFT)r   rQ   ÚRUNÚ_stateÚ_was_startedÚdaemon©rM   r8   rA   r2   r2   r3   rQ   ð  s   

zPoolThread.__init__c              
   C   s¢   z|   ¡ W S  ty. } ztdt| ƒj|dd� tt ¡ tƒ t	 
¡  W Y d }~d S d }~w tyP } ztdt| ƒj|dd� t d¡ W Y d }~d S d }~ww )NzThread %r crashed: %rr   r£   )Úbodyr   r?   r}   rf   Ú_killr¦   r§   r	   r¥   r    r¬   rž   ©rM   rm   r2   r2   r3   Úrunö  s    
ÿ€ÿ€ýzPoolThread.runc                    s    d| _ tt| ƒj|i |¤Ž d S rã   )r  rn   rý   Ústartr  ro   r2   r3   r    s   zPoolThread.startc                 C   r·   r-   r2   rL   r2   r2   r3   Úon_stop_not_started  r¹   zPoolThread.on_stop_not_startedNc                 C   s    | j r
|  |¡ d S |  ¡  d S r-   )r  Újoinr	  ©rM   rF   r2   r2   r3   rD   
  s   
zPoolThread.stopc                 C   ó
   t | _d S r-   )Ú	TERMINATEr   rL   r2   r2   r3   Ú	terminate  ó   
zPoolThread.terminatec                 C   r  r-   )ÚCLOSEr   rL   r2   r2   r3   rÝ     r  zPoolThread.closer-   )rf   rg   rh   rQ   r  r  r	  rD   r  rÝ   ru   r2   r2   ro   r3   rý   î  s    
rý   c                       s$   e Zd Z‡ fdd„Zdd„ Z‡  ZS )Ú
Supervisorc                    s   || _ tt| ƒ ¡  d S r-   )Úpoolrn   r  rQ   )rM   r  ro   r2   r3   rQ     s   zSupervisor.__init__c                 C   sÖ   t dƒ t d¡ | j}zH|j}td|j dƒ|_tdƒD ]}| jtkr2|jtkr2| 	¡  t d¡ q||_| jtkrS|jtkrS| 	¡  t d¡ | jtkrS|jtks@W n t
yd   | ¡  | ¡  ‚ w t dƒ d S )Nzworker handler startinggš™™™™™é?é
   r   r,   zworker handler exiting)r!   r´   rµ   r  r   Ú
_processesr   r   rÿ   Ú_maintain_poolr   rÝ   r
  )rM   r  Ú
prev_staterÓ   r2   r2   r3   r    s.   

€
þ€ýzSupervisor.body)rf   rg   rh   rQ   r  ru   r2   r2   ro   r3   r    s    r  c                       s4   e Zd Z‡ fdd„Zdd„ Zdd„ Zdd„ Z‡  ZS )	ÚTaskHandlerc                    s0   || _ || _|| _|| _|| _tt| ƒ ¡  d S r-   )Ú	taskqueuer²   Úoutqueuer  Úcachern   r  rQ   )rM   r  r²   r  r  r  ro   r2   r3   rQ   >  ó   zTaskHandler.__init__c           
      C   sf  | j }| j}| j}t|jd ƒD ]™\}}d }d}z^t|ƒD ]H\}}| jr)tdƒ  nJz||ƒ W q ty=   tdƒ Y  n6 t	yd   |d d… \}}	z||  
|	dtƒ f¡ W n	 tya   Y nw Y qw |rqtdƒ ||d ƒ W qW  n7 t	y¨   |r„|d d… nd\}}	||v r™||  
|	d dtƒ f¡ |r¦t d¡ ||d ƒ Y qw td	ƒ |  ¡  d S )
Néÿÿÿÿz'task handler found thread._state != RUNzcould not put task on queuer'   Fzdoing set_length()r   )r   r   ztask handler got sentinel)r  r  r²   Úiterrë   Ú	enumerater   r!   rô   r¬   Ú_setr   ÚKeyErrorr   Útell_others)
rM   r  r  r²   ÚtaskseqÚ
set_lengthÚtaskrÀ   rÐ   Úindr2   r2   r3   r  F  sR   ÿ€ü
€úzTaskHandler.bodyc                 C   sj   | j }| j}| j}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 sending sentinel to result handlerz(task handler sending sentinel to workersz/task handler got IOError when sending sentinelsztask handler exiting)r  r²   r  r!   rô   )rM   r  r²   r  Úpr2   r2   r3   r!  p  s   

ÿÿzTaskHandler.tell_othersc                 C   s   |   ¡  d S r-   )r!  rL   r2   r2   r3   r	  ƒ  r:   zTaskHandler.on_stop_not_started)rf   rg   rh   rQ   r  r!  r	  ru   r2   r2   ro   r3   r  <  s
    *r  c                       sT   e Zd Z‡ fdd„Zdd„ Zdd„ Zdd„ Zd	d
„ Zdd„ Zdd„ Z	dd„ Z
‡  ZS )ÚTimeoutHandlerc                    s0   || _ || _|| _|| _d | _tt| ƒ ¡  d S r-   )Ú	processesr  Út_softÚt_hardÚ_itrn   r'  rQ   )rM   r(  r  r)  r*  ro   r2   r3   rQ   ‰  r  zTimeoutHandler.__init__c                    ó   t ‡ fdd„t| jƒD ƒdƒS )Nc                 3   ó&   � | ]\}}|j ˆ kr||fV  qd S r-   r¡   ©Ú.0rÀ   Úprocr¡   r2   r3   Ú	<genexpr>’  ó   € 
ÿþz1TimeoutHandler._process_by_pid.<locals>.<genexpr>©NN)Únextr  r(  r¸   r2   r¡   r3   Ú_process_by_pid‘  ó
   ÿýzTimeoutHandler._process_by_pidc              
   C   sx   t d|ƒ |  |j¡\}}|sd S |jdd� z	t|jtƒ W d S  ty; } zt|ƒtj	kr0‚ W Y d }~d S d }~ww )Nzsoft time limit exceeded for %rT©Úsoft)
r!   r5  Ú_worker_pidÚhandle_timeoutr  rÞ   ÚOSErrorr   rõ   ÚESRCH)rM   rÐ   ÚprocessÚ_indexrm   r2   r2   r3   Úon_soft_timeout—  s   
ÿ€ÿzTimeoutHandler.on_soft_timeoutc                 C   sz   |  ¡ rd S td|ƒ zt|jƒ‚ ty#   | |jdtƒ f¡ Y nw |  |j¡\}}|j	dd� |r;|  
|¡ d S d S )Nzhard time limit exceeded for %rFr7  )r÷   r!   r   Ú_timeoutr  Ú_jobr   r5  r9  r:  Ú_trywaitkill)rM   rÐ   r=  r>  r2   r2   r3   Úon_hard_timeout¦  s   

ÿÿzTimeoutHandler.on_hard_timeoutc                 C   sâ   t d|jƒ z!t |j¡|jkr"t d|jƒ t t |j¡tj¡ n| ¡  W n	 t	y0   Y n
w |j
jdd�r:d S t d|jƒ z&t |j¡|jkr^t d|jƒ t t |j¡tj¡ W d S t|jtƒ W d S  t	yp   Y d S w )Nztimeout: sending TERM to %szIworker %s is a group leader. It is safe to kill (SIGTERM) the whole groupr,   ré   z/timeout: TERM timed-out, now sending KILL to %szIworker %s is a group leader. It is safe to kill (SIGKILL) the whole group)r!   Ú_namer¦   Úgetpgidr¢   Úkillpgrß   ÚSIGTERMr  r;  Ú_popenÚwaitÚSIGKILLr  ©rM   Úworkerr2   r2   r3   rB  »  s*   €ÿÿzTimeoutHandler._trywaitkillc                 #   sæ   � | j | j}}tƒ }| j}| j}dd„ }| jtkrqt | j¡‰ |r-t‡ fdd„|D ƒƒ}ˆ  	¡ D ]5\}}|j
}	|j}
|
d u rA|}
|j}|d u rJ|}||	|ƒrT||ƒ q1||vrf||	|
ƒrf||ƒ | |¡ q1d V  | jtksd S d S )Nc                 S   s"   | r|sdS t ƒ | | krdS d S rþ   )r   )r  rF   r2   r2   r3   Ú
_timed_outØ  s
   ÿz2TimeoutHandler.handle_timeouts.<locals>._timed_outc                 3   s   � | ]	}|ˆ v r|V  qd S r-   r2   )r/  Úk©r  r2   r3   r1  ç  ó   € z1TimeoutHandler.handle_timeouts.<locals>.<genexpr>)r*  r)  Úsetr?  rC  r   rÿ   Úcopyr  ÚitemsÚ_time_acceptedÚ_soft_timeoutr@  Úadd)rM   r*  r)  Údirtyr?  rC  rM  rÀ   rÐ   Úack_timeÚsoft_timeoutÚhard_timeoutr2   rO  r3   Úhandle_timeoutsÒ  s4   €



€ézTimeoutHandler.handle_timeoutsc                 C   sP   | j tkr"z|  ¡ D ]}t d¡ q
W n	 ty   Y nw | j tkstdƒ d S )Nrñ   ztimeout handler exiting)r   rÿ   r[  r´   rµ   r   r!   ©rM   rÓ   r2   r2   r3   r  ø  s   
ÿÿ
üzTimeoutHandler.bodyc                 G   s@   | j d u r
|  ¡ | _ zt| j ƒ W d S  ty   d | _ Y d S w r-   )r+  r[  r4  ÚStopIteration©rM   r8   r2   r2   r3   Úhandle_event  s   

ÿzTimeoutHandler.handle_event)rf   rg   rh   rQ   r5  r?  rC  rB  r[  r  r_  ru   r2   r2   ro   r3   r'  ‡  s    &	r'  c                       sV   e Zd Z	d‡ fdd„	Zdd„ Zdd„ Zdd	d
„Zddd„Zdd„ Zddd„Z	‡  Z
S )ÚResultHandlerNc                    sb   || _ || _|| _|| _|| _|| _|| _d | _d| _|| _	|	| _
|
| _|  ¡  tt| ƒ ¡  d S )NF)r  rë   r  rí   Újoin_exited_workersÚputlockr   r+  Ú_shutdown_completeÚcheck_timeoutsÚon_job_readyÚon_ready_countersÚ_make_methodsrn   r`  rQ   )rM   r  rë   r  rí   ra  rb  r   rd  re  rf  ro   r2   r3   rQ     s   zResultHandler.__init__c                 C   s   | j dd� d S )NT)r[  )Úfinish_at_shutdownrL   r2   r2   r3   r	    ó   z!ResultHandler.on_stop_not_startedc                    sl   ˆj ‰ ˆj‰ˆj‰ˆj‰‡ ‡fdd„}‡ ‡‡‡fdd„}dd„ }t|t|t|i ‰ˆ_‡fdd„}|ˆ_d S )	Nc              	      s:   dˆ_ zˆ |   ||||¡ W d S  ttfy   Y d S w r|   )ÚRÚ_ackr   r.   )rÐ   rÀ   Útime_acceptedr¢   r’   )r  r   r2   r3   Úon_ack(  s   þz+ResultHandler._make_methods.<locals>.on_ackc                    sÞ   ˆd urˆ| |||ƒ zˆ |  }W n
 t y   Y d S w ˆjrOtt| ¡ ƒd ƒ}|rO|ˆjv rOˆj| }| ¡ � | jd7  _W d   ƒ n1 sJw   Y  | ¡ s[ˆd ur[ˆ ¡  z	| 	||¡ W d S  t yn   Y d S w rI   )
r   rf  r4  r  Úworker_pidsÚget_lockrS   r÷   r]   r  )rÐ   rÀ   rš   rŽ   ÚitemÚ
worker_pidrŠ   )r  re  rb  rM   r2   r3   Úon_ready0  s,   ÿ

ÿÿz-ResultHandler._make_methods.<locals>.on_readyc              
   S   sJ   z	t  | t¡ W d S  ty$ } zt|ƒtjkr‚ W Y d }~d S d }~ww r-   )r¦   r&   r	   r;  r   rõ   r<  )r¢   r¶   rm   r2   r2   r3   Úon_deathG  s   ÿ€ÿz-ResultHandler._make_methods.<locals>.on_deathc                    s<   | \}}z	ˆ | |Ž  W d S  t y   td||ƒ Y d S w )NzUnknown job state: %s (args=%s))r   r!   )r$  Ústater8   )Ústate_handlersr2   r3   Úon_state_changeR  s   ÿz4ResultHandler._make_methods.<locals>.on_state_change)	r  rb  r   re  r¾   rÉ   r³   ru  rv  )rM   rm  rr  rs  rv  r2   )r  re  rb  r   rM   ru  r3   rg  "  s   
ÿ
zResultHandler._make_methodsrñ   c              
   c   s¬   � | j }| j}	 z||ƒ\}}W n ttfy& } ztd|ƒ tƒ ‚d }~ww | jr8| jtks1J ‚tdƒ tƒ ‚|rP|d u rEtdƒ tƒ ‚||ƒ |dkrOd S nd S d V  q)Nr   ú result handler got %r -- exitingz,result handler found thread._state=TERMINATEzresult handler got sentinelr   )rí   rv  rô   ró   r!   r   r   r  )rM   rF   rí   rv  r÷   r$  rm   r2   r2   r3   Ú_process_resultZ  s4   €
€þÿëzResultHandler._process_resultc              	   C   sT   | j tkr(| jd u r|  d¡| _zt| jƒ W d S  ttfy'   d | _Y d S w d S r|   )r   rÿ   r+  rx  r4  r]  r   )rM   r/   Úeventsr2   r2   r3   r_  u  s   

ÿûzResultHandler.handle_eventc                 C   sz   t dƒ z3| jtkr*z
|  d¡D ]}qW n	 ty   Y nw | jtks
W |  ¡  d S W |  ¡  d S W |  ¡  d S |  ¡  w )Nzresult handler startingrñ   )r!   r   rÿ   rx  r   rh  r\  r2   r2   r3   r  ~  s    
ÿÿüùþzResultHandler.bodyFc              
   C   sŽ  d| _ | j}| j}| j}| j}| j}| j}| j}d }	|r”| jt	kr”|d ur(|ƒ  z|dƒ\}
}W n t
tfyJ } ztd|ƒ W Y d }~d S d }~ww |
rZ|d u rVtdƒ q||ƒ z|dd� W n+ tyŒ   tƒ }|	sp|}	n||	 dkr|tdƒ Y ntdtt||	 d d	ƒƒƒ Y nw |r”| jt	ks!t|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 )NTrñ   rw  z&result handler ignoring extra sentinel)Úshutdowng      @z!result handler exiting: timed outz6result handler: all workers terminated, timeout in %ssr   r�   z"ensuring that outqueue is not fullr  z7result handler exiting: len(cache)=%s, thread._state=%s)rc  rë   r  r  rí   ra  rd  rv  r   r  rô   ró   r!   rv   r   ÚabsÚminrÜ   r   r�   Úlen)rM   r[  rë   r  r  rí   ra  rd  rv  Útime_terminater÷   r$  rm   rÎ   rÀ   r2   r2   r3   rh  Š  sj   
€þþ€øï

€ÿ
ÿz ResultHandler.finish_at_shutdownr-   )rñ   r3  ©F)rf   rg   rh   rQ   r	  rg  rx  r_  r  rh  ru   r2   r2   ro   r3   r`  
  s    þ
8
	r`  c                   @   sn  e Zd ZdZdZeZeZeZeZe	Z	e
Z
																	dwd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„ Zdd„ Zdd„ Zdd „ Zdxd!d"„Zd#d$„ Zd%d&„ Zd'd(„ Zd)d*„ Zd+d,„ Zd-d.„ Zd/d0„ Zd1d2„ Z d3d4„ Z!dyd5d6„Z"dyd7d8„Z#d9d:„ Z$d;d<„ Z%d=d>„ Z&d?d@„ Z'dAdB„ Z(dCdD„ Z)dEdF„ Z*dGdH„ Z+dIdJ„ Z,di fdKdL„Z-dzdMdN„Z.		d{dOdP„Z/dzdQdR„Z0d|dSdT„Z1		d|dUdV„Z2di ddddddddddfdWdX„Z3dYdZ„ Z4dzd[d\„Z5		d{d]d^„Z6		d{d_d`„Z7e8dadb„ ƒZ9dcdd„ Z:dedf„ Z;dgdh„ Z<e8didj„ ƒZ=dkdl„ Z>dmdn„ Z?e8dodp„ ƒZ@eAdqdr„ ƒZBeAdsdt„ ƒZCeDdudv„ ƒZEdS )}ÚPoolzS
    Class which supports an async version of applying functions to arguments.
    TNr2   r   Fc                 K   sŒ  |pt ƒ | _|| _|  ¡  tƒ | _i | _t| _|| _	|| _
|| _|| _|| _|| _|| _|p/t| _|
| _|| _|| _|| _|| _i | _|| _t|pR| j	d upR| j
d uƒ| _|rdtd u rdt tdƒ¡ d }|d u rl|  ¡ n|| _ |pwt!| j d ƒ| _"t#||	p~dƒ| _#|d ur�t$|ƒs�t%dƒ‚|d ur™t$|ƒs™t%dƒ‚| jj&| _'g | _(i | _)i | _*|| _+|p°t,| j ƒ| _-t.| j ƒD ]}|  /|¡ q·|  0| ¡| _1|rÌ| j1 2¡  |  3| j| j4| j5| j(| j¡| _6|râ| j6 2¡  d | _7| j�r
|  8| j(| j| j
| j	¡| _9t:ƒ | _;d| _<|  =¡  |�s	| j9j>| _7n	d | _9d| _<d | _;|  ?¡ | _@| j@j>| _A|�r%| j@ 2¡  tB| | jC| j| jD| j5| j(| j1| j6| j@| j| j9|  E¡ f
dd�| _Fd S )	NúWSoft timeouts are not supported: on this platform: It does not have the SIGUSR1 signal.éd   r   zinitializer must be a callablez on_process_exit must be callableFé   )r8   Úexitpriority)Gr   Ú_ctxÚsynackÚ_setup_queuesr   Ú
_taskqueueÚ_cacherÿ   r   rF   rY  Ú_maxtasksperchildÚ_max_memory_per_childÚ_initializerÚ	_initargsÚ_on_process_exitÚLOST_WORKER_TIMEOUTÚlost_worker_timeoutÚon_process_upÚon_process_downÚon_timeout_setÚon_timeout_cancelÚthreadsÚreadersÚallow_restartÚboolÚenable_timeoutsrÞ   ÚwarningsÚwarnÚUserWarningr   r  ÚroundÚmax_restartsr   ÚcallableÚ	TypeErrorÚProcessÚ_ProcessÚ_poolÚ	_poolctrlÚ_on_ready_countersÚputlocksrH   Ú_putlockr   Ú_create_worker_processr  Ú_worker_handlerr  r  r–   Ú	_outqueueÚ_task_handlerrd  r'  Ú_timeout_handlerÚLockÚ_timeout_handler_mutexÚ_timeout_handler_startedÚ_start_timeout_handlerr_  Úcreate_result_handlerÚ_result_handlerÚhandle_result_eventr    Ú_terminate_poolÚ_inqueueÚ_help_stuff_finish_argsÚ
_terminate)rM   r(  r   r€   ÚmaxtasksperchildrF   rY  r�  rž  Úmax_restart_freqr‘  r’  r“  r”  r•  Ú	semaphorer¦  r—  r†  Úon_process_exitÚcontextr‚   r™  rA   rÀ   r2   r2   r3   rQ   Ï  s®   
ÿýÿ

ü
þ
€


üùzPool.__init__c                 O   s   | j |i |¤ŽS r-   )r¢  )rM   r8   Úkwdsr2   r2   r3   r¡  I  ó   zPool.Processc                 C   s   |  | j|d�¡S )N)Útarget)r‹   r¡  rK  r2   r2   r3   ÚWorkerProcessL  ó   zPool.WorkerProcessc              
   K   s:   | j | j| j| j| j| j| j| j| j| j	f	d| j
i|¤ŽS )Nrf  )r`  rª  r˜   r‰  Ú_poll_resultÚ_join_exited_workersr§  r   rd  re  r¥  )rM   Úextra_kwargsr2   r2   r3   r±  O  s   üüûzPool.create_result_handlerc                 C   r·   r-   r2   )rM   rÐ   rÀ   rš   rŽ   r2   r2   r3   re  X  r¹   zPool.on_job_readyc                 C   s   | j | j| jfS r-   )rµ  r«  r£  rL   r2   r2   r3   r¶  [  r¾  zPool._help_stuff_finish_argsc                 C   s   zt ƒ W S  ty   Y dS w rI   )r   ÚNotImplementedErrorrL   r2   r2   r3   r   ^  s
   ÿzPool.cpu_countc                 G   s   | j j|Ž S r-   )r²  r_  r^  r2   r2   r3   r³  d  r:   zPool.handle_result_eventc                 C   r·   r-   r2   )rM   rL  Úqueuesr2   r2   r3   Ú_process_register_queuesg  r¹   zPool._process_register_queuesc                    r,  )Nc                 3   r-  r-   r¡   r.  r¡   r2   r3   r1  k  r2  z'Pool._process_by_pid.<locals>.<genexpr>r3  )r4  r  r£  r¸   r2   r¡   r3   r5  j  r6  zPool._process_by_pidc                 C   s   | j | jd fS r-   )rµ  rª  rL   r2   r2   r3   Úget_process_queuesp  ræ   zPool.get_process_queuesc                 C   sÒ   | j r| j ¡ nd }|  ¡ \}}}| j d¡}|  | j|||| j| j| j	|| j
| j| j| j|d�¡}| j |¡ |  ||||f¡ |j dd¡|_d|_||_| ¡  || j|j< || j|j< | jrg|  |¡ |S )NrÀ   )r…   r‰   r‚   rŠ   r¡  Ú
PoolWorkerT)r—  r…  ÚEventrÈ  ÚValuerÀ  r{   rŒ  r�  rŠ  rŽ  r•  Ú_wrap_exceptionr‹  r£  ÚappendrÇ  ÚnameÚreplacer  Úindexr  r¤  r¢   r¥  r‘  )rM   rÀ   rŒ   r†   r‡   rˆ   rŠ   Úwr2   r2   r3   r¨  s  s,   
ø

zPool._create_worker_processc                 C   r·   r-   r2   rK  r2   r2   r3   Úprocess_flush_queues�  r¹   zPool.process_flush_queuesc                    sl  d}dd„ t | j ¡ ƒD ƒD ]}|ptƒ }|j\}}|| |jkr'|  ||¡ q|r2t| jƒs2t	ƒ ‚i i ‰}t
tt| jƒƒƒD ]]}| j| }|j}	|j}
|
du sU|	dur�td|ƒ |
durb| ¡  td|ƒ |ˆ|j< |	||j< |	ttfvrŠt|ddƒsŠtd|j|jt|	ƒd	d
� |  |¡ | j|= | j|j= | j|j= q@ˆ�r4dd„ | jD ƒ‰ t | j ¡ ƒD ]d}t‡ ‡fdd„| ¡ D ƒdƒ}|rï|  ||¡ | ¡ sî| |¡pÓd	}	ˆ |¡}|rçt|ddƒrç| |	¡ q°|   |||	¡ q°|j!}|j"}|�r| #¡ �s|  ||j¡ q°|�r| #¡ �s|  ||j¡ q°tˆƒD ]}| j$�r,|�s'|  %|¡ |  $|¡ �qt | ¡ ƒS g S )z¤Cleanup after any worker processes which have exited due to
        reaching their specified lifetime. Returns True if any workers were
        cleaned up.
        Nc                 S   s   g | ]}|  ¡ s|jr|‘qS r2   )r÷   Ú_worker_lost)r/  rÐ   r2   r2   r3   Ú
<listcomp>š  s
    ÿ
ÿz-Pool._join_exited_workers.<locals>.<listcomp>z!Supervisor: cleaning up worker %dzSupervisor: worked %d joinedÚ_controlled_terminationFz Process %r pid:%r exited with %rr   r£   c                 S   s   g | ]}|j ‘qS r2   r¡   ©r/  rÑ  r2   r2   r3   rÔ  ½  s    c                 3   s$   � | ]}|ˆv s|ˆ vr|V  qd S r-   r2   ©r/  r¢   ©Úall_pidsÚcleanedr2   r3   r1  À  s   € ÿÿz,Pool._join_exited_workers.<locals>.<genexpr>Ú_job_terminated)&r5   r‰  r   r   rÓ  Ú_lost_worker_timeoutÚmark_as_worker_lostr}  r£  rv   Úreversedr   r¶   rH  r!   r
  r¢   r+   rÌ   Úgetattrr?   rÎ  r
   rÒ  r¤  r¥  r4  rn  Úon_job_process_downr÷   rë   Ú_set_terminatedÚon_job_process_lostÚ	_write_toÚ_scheduled_forÚ	_is_aliver’  Ú_process_cleanup_queues)rM   rz  rÎ   rÐ   Ú	lost_timeÚlost_retÚ	exitcodesrÀ   rL  r¶   ÚpopenÚacked_by_goner0  Úwrite_toÚ	sched_forr2   rØ  r3   rÃ  �  s†   

€






ÿý


€ý
ÿ€€

€zPool._join_exited_workersc                 C   r·   r-   r2   )rM   rÐ   rL  r2   r2   r3   Úon_partial_readã  r¹   zPool.on_partial_readc                 C   r·   r-   r2   rK  r2   r2   r3   ræ  æ  r¹   zPool._process_cleanup_queuesc                 C   r·   r-   r2   )rM   rÐ   Úpid_goner2   r2   r3   rà  é  r¹   zPool.on_job_process_downc                 C   s   t ƒ |f|_d S r-   )r   rÓ  )rM   rÐ   r¢   r¶   r2   r2   r3   râ  ì  r¾  zPool.on_job_process_lostc                 C   s:   z	t d t|ƒ¡ƒ‚ t y   | d dtƒ f¡ Y d S w )NzWorker exited prematurely: {0}.F)r   rË   r
   r  r   )rM   rÐ   r¶   r2   r2   r3   rÝ  ï  s   ÿÿÿzPool.mark_as_worker_lostc                 C   ó   | S r-   r2   rL   r2   r2   r3   Ú	__enter__ú  r¹   zPool.__enter__c                 G   s   |   ¡ S r-   )r  )rM   r¤   r2   r2   r3   Ú__exit__ý  s   zPool.__exit__c                 C   r·   r-   r2   ©rM   Únr2   r2   r3   Úon_grow   r¹   zPool.on_growc                 C   r·   r-   r2   ró  r2   r2   r3   Ú	on_shrink  r¹   zPool.on_shrinkc                 C   s`   t |  ¡ ƒD ]%\}}|  jd8  _| jr| j ¡  | ¡  |  d¡ ||d kr+ d S qtdƒ‚)Nr   z&Can't shrink pool. All processes busy!)r  Ú_iterinactiver  r§  rN   Úterminate_controlledrö  Ú
ValueError)rM   rô  rÀ   rL  r2   r2   r3   rN     s   

ÿzPool.shrinkc                 C   s:   t |ƒD ]}|  jd7  _| jr| j ¡  q|  |¡ d S rI   )r   r  r§  rX   rõ  )rM   rô  rÀ   r2   r2   r3   rX     s   
€z	Pool.growc                 c   s"   � | j D ]
}|  |¡s|V  qd S r-   )r£  Ú_worker_activerK  r2   r2   r3   r÷    s   €

€þzPool._iterinactivec                 C   s(   t | jƒD ]}|j| ¡ v r dS qdS )NTF)r   r‰  r¢   rn  )rM   rL  rÐ   r2   r2   r3   rú    s
   ÿzPool._worker_activec              	   C   s„   t | jt| jƒ ƒD ]5}| jtkr dS z|r$|| ttfvr$| j 	¡  W n t
y3   | j 	¡  Y nw |  |  ¡ ¡ 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.
        Nzadded worker)r   r  r}  r£  r   rÿ   r+   rÌ   r   ÚstepÚ
IndexErrorr¨  Ú_avail_indexr!   )rM   ré  rÀ   r2   r2   r3   Ú_repopulate_pool$  s   

€ÿ
÷zPool._repopulate_poolc                    sD   t | jƒ| jk s
J ‚tdd„ | jD ƒƒ‰ t‡ fdd„t| jƒD ƒƒS )Nc                 s   s   � | ]}|j V  qd S r-   )rÐ  )r/  r&  r2   r2   r3   r1  5  s   € z$Pool._avail_index.<locals>.<genexpr>c                 3   s   � | ]	}|ˆ vr|V  qd S r-   r2   )r/  rÀ   ©Úindicesr2   r3   r1  6  rP  )r}  r£  r  rQ  r4  r   rL   r2   rÿ  r3   rý  3  s   zPool._avail_indexc                 C   s
   |   ¡  S r-   )rÃ  rL   r2   r2   r3   Údid_start_ok8  r  zPool.did_start_okc                 C   s<   |   ¡ }|  |¡ tt|ƒƒD ]}| jdur| j ¡  qdS )zF"Clean up any exited workers and start replacements for them.
        N)rÃ  rþ  r   r}  r§  r]   )rM   ÚjoinedrÀ   r2   r2   r3   r  ;  s   


€þzPool._maintain_poolc              
   C   s�   | j jtkrD| jtkrFz|  ¡  W d S  ty"   |  ¡  |  ¡  ‚  tyC } zt|ƒt	j
kr>tttt|ƒƒt ¡ d ƒ ‚ d }~ww d S d S )Nr'   )r©  r   rÿ   r  r   rÝ   r
  r;  r   rõ   ÚENOMEMr   ÚMemoryErrorrq   r¥   r¤   r  r2   r2   r3   Úmaintain_poolD  s"   

þ€ûùzPool.maintain_poolc                    sF   ˆ j  ¡ ˆ _ˆ j  ¡ ˆ _ˆ jjjˆ _ˆ jjjˆ _	‡ fdd„}|ˆ _
d S )Nc                    s   ˆ j j | ¡rdˆ  ¡ fS dS rç   )rª  r�   rí   r˜   ré   rL   r2   r3   rÂ  Y  s   z(Pool._setup_queues.<locals>._poll_result)r…  ÚSimpleQueuerµ  rª  r�   r•   r–   r�   r—   r˜   rÂ  ©rM   rÂ  r2   rL   r3   r‡  S  s   
zPool._setup_queuesc                 C   sj   | j r1| jd ur3| j� | jsd| _| j ¡  W d   ƒ d S W d   ƒ d S 1 s*w   Y  d S d S d S rã   )r•  r¬  r®  r¯  r  rL   r2   r2   r3   r°  _  s   ý"ÿÿzPool._start_timeout_handlerc                 C   ó    | j tkr|  |||¡ ¡ S dS )z8
        Equivalent of `func(*args, **kwargs)`.
        N)r   rÿ   Úapply_asyncrë   )rM   Úfuncr8   r½  r2   r2   r3   Úapplyh  s   
ÿz
Pool.applyc                 C   s"   | j tkr|  ||t|¡ ¡ S d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).
        N)r   rÿ   Ú
_map_asyncr=   rë   ©rM   r
  ÚiterableÚ	chunksizer2   r2   r3   r<   o  s   
ÿÿÿzPool.starmapc                 C   s"   | j tkr|  ||t|||¡S dS )z=
        Asynchronous version of `starmap()` method.
        N)r   rÿ   r  r=   ©rM   r
  r  r  ÚcallbackÚerror_callbackr2   r2   r3   Ústarmap_asyncy  s
   
ÿÿzPool.starmap_asyncc                 C   r  )zx
        Apply `func` to each element in `iterable`, collecting the results
        in a list that is returned.
        N)r   rÿ   Ú	map_asyncrë   r  r2   r2   r3   r6   ‚  s   
ÿzPool.mapc                    ó²   | j tkrdS |p| j}|dkr,t| j|d�‰| j ‡ ‡fdd„t|ƒD ƒˆjf¡ ˆS |dks2J ‚t	 
ˆ ||¡}t| j|d�‰| j ‡fdd„t|ƒD ƒˆjf¡ dd„ ˆD ƒS )zP
        Equivalent of `map()` -- can be MUCH slower than `Pool.map()`.
        Nr   ©r�  c                 3   ó*   � | ]\}}t ˆj|ˆ |fi ffV  qd S r-   ©rÈ   rA  ©r/  rÀ   Úx©r
  rº   r2   r3   r1  •  ó   € ÿzPool.imap.<locals>.<genexpr>c                 3   ó*   � | ]\}}t ˆ j|t|fi ffV  qd S r-   ©rÈ   rA  r9   r  ©rº   r2   r3   r1     r  c                 s   ó   � | ]
}|D ]}|V  qqd S r-   r2   ©r/  Úchunkrp  r2   r2   r3   r1  ¤  ó   € )r   rÿ   r�  ÚIMapIteratorr‰  rˆ  r²   r  Ú_set_lengthr€  Ú
_get_tasks©rM   r
  r  r  r�  Útask_batchesr2   r  r3   ÚimapŠ  s4   

ÿÿýÿ
ÿýz	Pool.imapc                    r  )zL
        Like `imap()` method but ordering of results is arbitrary.
        Nr   r  c                 3   r  r-   r  r  r  r2   r3   r1  ³  r  z&Pool.imap_unordered.<locals>.<genexpr>c                 3   r  r-   r  r  r  r2   r3   r1  ¿  r  c                 s   r   r-   r2   r!  r2   r2   r3   r1  Ã  r#  )r   rÿ   r�  ÚIMapUnorderedIteratorr‰  rˆ  r²   r  r%  r€  r&  r'  r2   r  r3   Úimap_unordered¦  s4   

ÿÿýÿ
ÿýzPool.imap_unorderedc                 C   s  | j tkrdS |	p| j}	|
p| j}
|p| j}|	r%tdu r%t tdƒ¡ d}	| j tkr†|du r1| j	n|}|r?| j
dur?| j
 ¡  t| j|||||	|
|| j| j|| jrT| jnd|d�}|
s]|	ra|  ¡  | jrw| j t|jd|||ffgdf¡ |S |  t|jd|||ff¡ |S dS )a  
        Asynchronous equivalent of `apply()` method.

        Callback is called when the functions return value is ready.
        The accept callback is called when the job is accepted to be executed.

        Simplified the flow is like this:

            >>> def apply_async(func, args, kwds, callback, accept_callback):
            ...     if accept_callback:
            ...         accept_callback()
            ...     retval = func(*args, **kwds)
            ...     if callback:
            ...         callback(retval)

        Nr�  )r“  r”  Úcallbacks_propagateÚsend_ackÚcorrelation_id)r   rÿ   rY  rF   r�  rÞ   rš  r›  rœ  r¦  r§  rK   ÚApplyResultr‰  r“  r”  r†  r-  r°  r•  rˆ  r²   rÈ   rA  r–   )rM   r
  r8   r½  r  r  Úaccept_callbackÚtimeout_callbackÚwaitforslotrY  rF   r�  r,  r.  rº   r2   r2   r3   r	  Å  sF   



ÿ


ù	ÿÿÿëzPool.apply_asyncc                 C   r·   r-   r2   )rM   ÚresponserÐ   rÀ   Úfdr2   r2   r3   r-  ý  r¹   zPool.send_ackc              
   C   st   |   |¡\}}|d ur8z	t||ptƒ W n ty/ } zt|ƒtjkr$‚ W Y d }~d S d }~ww d|_d|_d S d S rã   )	r5  r  r	   r;  r   rõ   r<  rÕ  rÛ  )rM   r¢   Úsigr0  rÓ   rm   r2   r2   r3   Úterminate_job   s   ÿ€ÿ
øzPool.terminate_jobc                 C   s   |   ||t|||¡S )z<
        Asynchronous equivalent of `map()` method.
        )r  r9   r  r2   r2   r3   r    s   ÿzPool.map_asyncc           	         s®   | j tkrdS t|dƒst|ƒ}|du r(tt|ƒt| jƒd ƒ\}}|r(|d7 }t|ƒdkr0d}t |||¡}t	| j
|t|ƒ||d�‰| j ‡ ‡fdd„t|ƒD ƒdf¡ ˆS )	zY
        Helper function to implement map, starmap and their async counterparts.
        NÚ__len__r(   r   r   ©r  c                 3   r  r-   r  r  ©Úmapperrº   r2   r3   r1  )  r  z"Pool._map_async.<locals>.<genexpr>)r   rÿ   rÜ   r5   Údivmodr}  r£  r€  r&  Ú	MapResultr‰  rˆ  r²   r  )	rM   r
  r  r:  r  r  r  Úextrar(  r2   r9  r3   r    s(   

ÿÿÿzPool._map_asyncc                 c   s0   � t |ƒ}	 tt ||¡ƒ}|sd S | |fV  qr-   )r  Útupler;   Úislice)r
  ÚitÚsizer  r2   r2   r3   r&  -  s   €
üzPool._get_tasksc                 C   s   t dƒ‚)Nz:pool objects cannot be passed between processes or pickled)rÅ  rL   r2   r2   r3   r›   6  s   ÿzPool.__reduce__c                 C   sP   t dƒ | jtkr&t| _| jr| j ¡  | j ¡  | j 	d ¡ t
| jƒ d S d S )Nzclosing pool)r!   r   rÿ   r  r§  r`   r©  rÝ   rˆ  r²   rG   rL   r2   r2   r3   rÝ   ;  s   


úz
Pool.closec                 C   s$   t dƒ t| _| j ¡  |  ¡  d S )Nzterminating pool)r!   r  r   r©  r  r·  rL   r2   r2   r3   r  E  s   
zPool.terminatec                 C   s   t | ƒ d S r-   )rG   )Útask_handlerr2   r2   r3   Ú_stop_task_handlerK  s   zPool._stop_task_handlerc                 C   sœ   | j ttfv s	J ‚tdƒ t| jƒ tdƒ |  | j¡ tdƒ t| jƒ tdƒ t	| j
ƒD ]\}}td|d t| j
ƒ|ƒ |jd urG| ¡  q.tdƒ d S )Nzjoining worker handlerújoining task handlerújoining result handlerzresult handler joinedzjoining worker %s/%s (%r)r   zpool join complete)r   r  r  r!   rG   r©  rC  r«  r²  r  r£  r}  rH  r
  )rM   rÀ   r&  r2   r2   r3   r
  O  s   


€z	Pool.joinc                 C   s   t | jƒD ]}| ¡  qd S r-   )r   r¤  rQ  )rM   Úer2   r2   r3   Úrestart^  s   
ÿzPool.restartc                 C   sZ   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!   Ú_rlockrK   Úis_aliver�   rí   r—   r´   rµ   )ÚinqueuerB  r£  r2   r2   r3   Ú_help_stuff_finishb  s   


"þzPool._help_stuff_finishc                 C   s   |  d ¡ d S r-   )r²   )Úclsr  r  r2   r2   r3   Ú_set_result_sentinelk  s   zPool._set_result_sentinelc                 C   s:  t dƒ | ¡  | ¡  | d ¡ t dƒ | j|
Ž  | ¡  |  ||¡ |	d ur,|	 ¡  |rFt|d dƒrFt dƒ |D ]
}| ¡ rE| ¡  q;t dƒ |  |¡ t dƒ | ¡  |	d urdt dƒ |	 t	¡ |r�t|d dƒr�t d	ƒ |D ]}| 
¡ rˆt d
|jƒ |jd urˆ| ¡  qst dƒ |r“| ¡  |r›| ¡  d S d S )Nzfinalizing poolz&helping task handler/workers to finishr   r  zterminating workersrD  rE  zjoining timeout handlerzjoining pool workerszcleaning up worker %dzpool workers joined)r!   r  r²   rK  rM  rÜ   rå  rC  rD   ÚTIMEOUT_MAXrI  r¢   rH  r
  rÝ   )rL  r  rJ  r  r  Úworker_handlerrB  Úresult_handlerr  Útimeout_handlerÚhelp_stuff_finish_argsr&  r2   r2   r3   r´  o  sJ   

€


€ÿzPool._terminate_poolc                 C   ó   dd„ | j D ƒS )Nc                 S   s   g | ]}|j j‘qS r2   )rH  rŒ   rÖ  r2   r2   r3   rÔ  ¨  ó    z*Pool.process_sentinels.<locals>.<listcomp>)r£  rL   r2   r2   r3   Úprocess_sentinels¦  ri  zPool.process_sentinels)NNr2   NNNNNr   NNNNTNFFFNNNFr  )r   r-   )NNNre   )Frf   rg   rh   ri   rÌ  r{   r  r  r'  r`  r   rQ   r¡  rÀ  r±  re  r¶  r   r³  rÇ  r5  rÈ  r¨  rÒ  rÃ  rî  ræ  rà  râ  rÝ  rñ  rò  rõ  rö  rN   rX   r÷  rú  rþ  rý  r  r  r  r‡  r°  r  r<   r  r6   r)  r+  r	  r-  r6  r  r  Ústaticmethodr&  r›   rÝ   r  rC  r
  rG  rK  ÚclassmethodrM  r´  ÚpropertyrU  r2   r2   r2   r3   r€  Ã  sÌ    
ðz	
S

		


ÿ
	

ÿ
û8

ÿ	
ÿ





6r€  c                   @   s¸   e Zd ZdZdZdZdddddedddddf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d$dd„Zd$dd„Zdd„ Zd%dd„Zd d!„ Zd"d#„ ZdS )&r/  Nr2   c                 C   sš   || _ tƒ | _t ¡ | _ttƒ| _|| _	|| _
|| _|| _|| _|| _|| _|| _|	| _|
| _|p2d| _|| _d| _d| _d | _d | _d | _| || j< d S )Nr2   F)r.  r­  Ú_mutexrB   rÊ  Ú_eventr4  Újob_counterrA  r‰  Ú	_callbackÚ_accept_callbackÚ_error_callbackÚ_timeout_callbackr@  rU  rÜ  Ú_on_timeout_setÚ_on_timeout_cancelÚ_callbacks_propagateÚ	_send_ackÚ	_acceptedÚ
_cancelledr9  rT  Ú_terminated)rM   r  r  r0  r1  r  rY  rF   r�  r“  r”  r,  r-  r.  r2   r2   r3   rQ   ´  s,   


zApplyResult.__init__c                 C   s   dj | jj| j| j|  ¡ d�S )Nz"<%s: {id} ack:{ack} ready:{ready}>)ÚidÚackr÷   )rË   rp   rf   rA  rd  r÷   rL   r2   r2   r3   rr   Ò  s   þzApplyResult.__repr__c                 C   s
   | j  ¡ S r-   )rZ  ÚisSetrL   r2   r2   r3   r÷   Ø  r  zApplyResult.readyc                 C   ó   | j S r-   )rd  rL   r2   r2   r3   ÚacceptedÛ  rz   zApplyResult.acceptedc                 C   s   |   ¡ sJ ‚| jS r-   )r÷   Ú_successrL   r2   r2   r3   Ú
successfulÞ  s   zApplyResult.successfulc                 C   s
   d| _ dS )zOnly works if synack is used.TN)re  rL   r2   r2   r3   Ú_cancelâ  s   
zApplyResult._cancelc                 C   s   | j  | jd ¡ d S r-   )r‰  ÚpoprA  rL   r2   r2   r3   Údiscardæ  rs   zApplyResult.discardc                 C   s
   || _ d S r-   )rf  ©rM   rw   r2   r2   r3   r  é  r  zApplyResult.terminatec                 C   s6   zt |pd ƒ‚ t y   |  d dtƒ f¡ Y d S w ©Nr   F)r   r  r   rq  r2   r2   r3   rá  ì  s
   ÿzApplyResult._set_terminatedc                 C   s   | j r| j gS g S r-   ©r9  rL   r2   r2   r3   rn  ò  rÁ  zApplyResult.worker_pidsc                 C   s   | j  |¡ d S r-   )rZ  rI  r  r2   r2   r3   rI  õ  r¾  zApplyResult.waitc                 C   s*   |   |¡ |  ¡ st‚| jr| jS | jj‚r-   )rI  r÷   r   rl  rV   Ú	exceptionr  r2   r2   r3   rë   ø  s   
zApplyResult.getc              
   O   sb   |r/z
||i |¤Ž W d S  | j y   ‚  ty. } ztd|dd� W Y d }~d S d }~ww d S )Nz"Pool callback raised exception: %rr   r£   )rb  r¬   r?   )rM   rÑ   r8   rA   rm   r2   r2   r3   Úsafe_apply_callback  s   ÿ€ÿûzApplyResult.safe_apply_callbackFc                 C   s0   | j d ur| j| j ||r| jn| jd� d S d S )N)r8  rF   )r_  ru  rU  r@  )rM   r8  r2   r2   r3   r:    s   

þÿzApplyResult.handle_timeoutc                 C   sÚ   | j �` | jr|  | ¡ |\| _| _| j ¡  | jr"| j | j	d ¡ | j
r0| jr0|  | j
| j¡ | jd urK| jrS| js[|  | j| j¡ W d   ƒ d S W d   ƒ d S W d   ƒ d S W d   ƒ d S 1 sfw   Y  d S r-   )rY  ra  rl  rV   rZ  rQ  rd  r‰  ro  rA  r\  ru  r^  ©rM   rÀ   rš   r2   r2   r3   r    s4   

ÿ
ÿÿÿïññ"ñzApplyResult._setc                 C   s¢  | j �Ä | jr(| jr(d| _|r|  t|| j|¡W  d   ƒ S 	 W d   ƒ d S d| _|| _|| _|  ¡ r=| j	 
| jd ¡ | jrI|  | | j| j¡ t}| jr¡z5z|  ||¡ W n | jyb   t}‚  tyl   t}Y nw W | jrƒ|rƒ|  ||| j|¡W  d   ƒ S n| jrŸ|r |  ||| j|¡     Y W  d   ƒ S w w | jr·|r¿|  ||| j|¡ W d   ƒ d S W d   ƒ d S W d   ƒ d S 1 sÊw   Y  d S rã   )rY  re  rc  rd  r½   rA  rT  r9  r÷   r‰  ro  r`  rU  r@  r¾   r]  Ú_propagate_errorsr¬   )rM   rÀ   rl  r¢   r’   r3  r2   r2   r3   rk  '  sZ   üûÿ€

ÿæ€

ÿæ
âã"ãzApplyResult._ackr-   r  )rf   rg   rh   rÓ  rã  rä  r�  rQ   rr   r÷   rk  rm  rn  rp  r  rá  rn  rI  rë   ru  r:  r  rk  r2   r2   r2   r3   r/  ¯  s4    
û


	

r/  c                   @   s4   e Zd Zdd„ Zdd„ Zdd„ Zdd„ Zd	d
„ ZdS )r<  c                 C   s’   t j| |||d� d| _|| _d g| | _dg| | _d g| | _d g| | _|| _|dkr<d| _	| j
 ¡  || j= d S || t|| ƒ | _	d S )Nr8  TFr   )r/  rQ   rl  Ú_lengthrV   rd  r9  rT  Ú
_chunksizeÚ_number_leftrZ  rQ  rA  r˜  )rM   r  r  Úlengthr  r  r2   r2   r3   rQ   O  s   ÿ
zMapResult.__init__c                 C   s¾   |\}}|r>|| j || j |d | j …< |  jd8  _| jdkr<| jr*|  | j ¡ | jr5| j | jd ¡ | j 	¡  d S d S d| _
|| _ | jrM|  | j ¡ | jrX| j | jd ¡ | j 	¡  d S )Nr   r   F)rV   ry  rz  r\  rd  r‰  ro  rA  rZ  rQ  rl  r^  )rM   rÀ   Úsuccess_resultÚsuccessrº   r2   r2   r3   r  a  s$   
ûzMapResult._setc                 G   sn   || j  }t|d | j  | jƒ}t||ƒD ]}d| j|< || j|< || j|< q|  ¡ r5| j 	| j
d ¡ d S d S ©Nr   T)ry  r|  rx  r   rd  r9  rT  r÷   r‰  ro  rA  )rM   rÀ   rl  r¢   r8   r  rD   Újr2   r2   r3   rk  u  s   


ÿzMapResult._ackc                 C   s
   t | jƒS r-   )Úallrd  rL   r2   r2   r3   rk    r  zMapResult.acceptedc                 C   rS  )Nc                 S   s   g | ]}|r|‘qS r2   r2   r×  r2   r2   r3   rÔ  ƒ  rT  z)MapResult.worker_pids.<locals>.<listcomp>rs  rL   r2   r2   r3   rn  ‚  r¾  zMapResult.worker_pidsN)rf   rg   rh   rQ   r  rk  rk  rn  r2   r2   r2   r3   r<  M  s    
r<  c                   @   sZ   e Zd ZdZefdd„Zdd„ Zddd„ZeZdd	„ Z	d
d„ Z
dd„ Zdd„ Zdd„ ZdS )r$  Nc                 C   sZ   t  t  ¡ ¡| _ttƒ| _|| _tƒ | _	d| _
d | _d| _i | _g | _|| _| || j< d S rr  )rB   Ú	Conditionr­  rU   r4  r[  rA  r‰  r   Ú_itemsr>  rx  Ú_readyÚ	_unsortedÚ_worker_pidsrÜ  )rM   r  r�  r2   r2   r3   rQ   �  s   
zIMapIterator.__init__c                 C   rð  r-   r2   rL   r2   r2   r3   Ú__iter__š  r¹   zIMapIterator.__iter__c                 C   sº   | j �F z| j ¡ }W n6 tyA   | j| jkrd| _t‚| j  |¡ z| j ¡ }W n ty>   | j| jkr<d| _t‚t	‚w Y nw W d   ƒ n1 sLw   Y  |\}}|rY|S t
|ƒ‚rã   )rU   r‚  Úpopleftrü  r>  rx  rƒ  r]  rI  r   r¬   )rM   rF   rp  r}  rS   r2   r2   r3   r4  �  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rWd| _| j	| j
= W d   ƒ d S W d   ƒ d S 1 sbw   Y  d S r~  )rU   r>  r‚  rÍ  r„  ro  rW   rx  rƒ  r‰  rA  rv  r2   r2   r3   r  µ  s"   
ý
ò"ôzIMapIterator._setc                 C   sh   | j �' || _| j| jkr"d| _| j  ¡  | j| j= W d   ƒ d S W d   ƒ d S 1 s-w   Y  d S rã   )rU   rx  r>  rƒ  rW   r‰  rA  )rM   r{  r2   r2   r3   r%  Æ  s   
û"þzIMapIterator._set_lengthc                 G   s   | j  |¡ d S r-   )r…  rÍ  )rM   rÀ   rl  r¢   r8   r2   r2   r3   rk  Î  r¾  zIMapIterator._ackc                 C   rj  r-   )rƒ  rL   r2   r2   r3   r÷   Ñ  rz   zIMapIterator.readyc                 C   rj  r-   )r…  rL   r2   r2   r3   rn  Ô  rz   zIMapIterator.worker_pidsr-   )rf   rg   rh   rÓ  r�  rQ   r†  r4  Ú__next__r  r%  rk  r÷   rn  r2   r2   r2   r3   r$  Š  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,d| _| j| j= W d   ƒ d S W d   ƒ d S 1 s7w   Y  d S r~  )	rU   r‚  rÍ  r>  rW   rx  rƒ  r‰  rA  rv  r2   r2   r3   r  Þ  s   
ú"üzIMapUnorderedIterator._setN)rf   rg   rh   r  r2   r2   r2   r3   r*  Ü  s    r*  c                   @   s:   e Zd ZddlmZ eZddd„Zdd„ Zed	d
„ ƒZ	dS )Ú
ThreadPoolr   )r¡  Nr2   c                 C   s   t  | |||¡ d S r-   )r€  rQ   )rM   r(  r   r€   r2   r2   r3   rQ   ñ  rs   zThreadPool.__init__c                    s:   t ƒ ˆ _t ƒ ˆ _ˆ jjˆ _ˆ jjˆ _‡ fdd„}|ˆ _d S )Nc                    s(   z	dˆ j | d�fW S  ty   Y dS w rì   )r˜   r   ré   rL   r2   r3   rÂ  ú  s
   ÿz.ThreadPool._setup_queues.<locals>._poll_result)r   rµ  rª  r²   r–   rë   r˜   rÂ  r  r2   rL   r3   r‡  ô  s   


zThreadPool._setup_queuesc                 C   sV   | j � | j ¡  | j d gt|ƒ ¡ | j  ¡  W d   ƒ d S 1 s$w   Y  d S r-   )Ú	not_emptyÚqueuer`   Úextendr}  rZ   )rJ  rB  r  r2   r2   r3   rK    s
   
"ýzThreadPool._help_stuff_finish)NNr2   )
rf   rg   rh   Údummyr¡  r   rQ   r‡  rV  rK  r2   r2   r2   r3   r‰  ì  s    
r‰  r-   )kÚ
__future__r   rR  rõ   r;   r¦   r±   rß   r¥   rB   r´   rš  Úcollectionsr   Ú	functoolsr   Ú r   r   r   Úcommonr	   r
   r   r   r   Úcompatr   r   r   rÖ   r   r�  r   Ú
exceptionsr   r   r   r   r   r   r   Úfiver   r   r   r   r   r   r    r!   r"   rÊ   Úversion_inforj   ÚsystemÚ_winr%   r  rJ  r&   rN  r.   Ú	SemaphorerP   rÿ   r  r  r¾   rÉ   rÈ   r½   r³   r+   r°   rÌ   rß  rÞ   r�  rØ   rÙ   Úcountr[  r­  r4   r9   r=   r?   rG   rH   r¬   rk   rv   ry   Úobjectr{   rý   r  r  r'  r`  r€  r/  r<  r$  r*  r‰  r2   r2   r2   r3   Ú<module>   s¬   	$ 	
ÿ


;  )%K  :     q =R