o
    €¨ÊhÝ ã                   @   s  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( d dl)m*Z*m+Z+ ddlm,Z,m-Z-m.Z. dZ/ej0d  dkZ1e 2¡ dkr¨ddl3m4Z5 eZ6n	d dlm7Z5 ej6Z6zej8Z8W n e9yÁ   dZ8Y nw ej0dkrËej:Z;nej;Z;d Z<dZ=dZ>d Z?dZ@dZAdZBdZCd ZDdZEdZFeGeddƒZHdZIeGedd ƒZDdZJdZKe L¡ ZMejNZNdd„ ZOd d!„ ZPd"d#„ ZQd$d%„ ZRdHd&d'„ZSG d(d)„ d)e;ƒZTG d*d+„ d+eUƒZVG d,d-„ d-eUƒZWd.d/„ ZXG d0d1„ d1ƒZYG d2d3„ d3eƒZZG d4d5„ d5eZƒZ[G d6d7„ d7eZƒZ\G d8d9„ d9eZƒZ]G d:d;„ d;eZƒZ^G d<d=„ d=ƒZ_G d>d?„ d?ƒZ`G d@dA„ dAe`ƒZaG dBdC„ dCƒZbG dDdE„ dEebƒZcG dFdG„ dGe_ƒZddS )Ié    N)Ú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©Ú	monotonic©ÚQueueÚEmpty)Ú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© r0   ú?/var/www/html/env/lib/python3.10/site-packages/billiard/pool.pyÚ_get_send_offsetx   s   
ÿr2   c                 C   s   t t| Ž ƒS r+   )ÚlistÚmap©Úargsr0   r0   r1   Úmapstar‚   ó   r7   c                 C   s   t t | d | d ¡ƒS )Nr   r   )r3   Ú	itertoolsÚstarmapr5   r0   r0   r1   Ústarmapstar†   s   r;   c                 O   s    t  ¡ j| g|¢R i |¤Ž d S r+   )r   Ú
get_loggerÚerror)Úmsgr6   Úkwargsr0   r0   r1   r=   Š   s    r=   c                 C   s   | t  ¡ ur|  |¡ d S d S r+   )Ú	threadingÚcurrent_threadÚstop)ÚthreadÚtimeoutr0   r0   r1   Ústop_if_not_currentŽ   s   ÿrE   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©Úselfr0   r0   r1   Úshrink—   s   zLaxBoundedSemaphore.shrinkr   Nc                 C   s   t  | |¡ || _d S r+   ©Ú
_SemaphoreÚ__init__rH   ©rK   ÚvalueÚverboser0   r0   r1   rO   �   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 rG   )Ú_condrH   Ú_valueÚnotifyrJ   r0   r0   r1   Ú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 rG   )rS   rT   rH   Ú
notify_all©rK   Úcondr0   r0   r1   Úrelease§   ó   
ý"ÿzLaxBoundedSemaphore.releasec                 C   ó*   | j | jk rt | ¡ | j | jk sd S d S r+   )rT   rH   rN   r[   rJ   r0   r0   r1   Úclear®   ó   
ÿzLaxBoundedSemaphore.clearc                 C   s   t  | ||¡ || _d S r+   rM   rP   r0   r0   r1   rO   ³   s   
c                 C   sT   | j }|� |  jd7  _|  jd7  _| ¡  W d   ƒ d S 1 s#w   Y  d S rG   )Ú_Semaphore__condrH   Ú_Semaphore__valuerU   rY   r0   r0   r1   rV   ·   s   
"ýc                 C   rW   rG   )r`   ra   rH   Ú	notifyAllrY   r0   r0   r1   r[   ¾   r\   c                 C   r]   r+   )ra   rH   rN   r[   rJ   r0   r0   r1   r^   Å   r_   ©r   N)
Ú__name__Ú
__module__Ú__qualname__Ú__doc__rL   ÚPY3rO   rV   r[   r^   r0   r0   r0   r1   rF   “   s    

rF   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ƒ  | j| j¡ d S r+   )ÚreprÚexcrQ   ÚsuperrO   )rK   rk   rQ   ©Ú	__class__r0   r1   rO   Ò   s   

zMaybeEncodingError.__init__c                 C   s   d| j jt| ƒf S )Nz<%s: %s>)rn   rd   ÚstrrJ   r0   r0   r1   Ú__repr__×   ó   zMaybeEncodingError.__repr__c                 C   s   d| j | jf S )Nz)Error sending result: '%r'. Reason: '%r'.)rQ   rk   rJ   r0   r0   r1   Ú__str__Ú   s   ÿzMaybeEncodingError.__str__)rd   re   rf   rg   rO   rp   rr   Ú__classcell__r0   r0   rm   r1   ri   Î   s
    ri   c                   @   s   e Zd ZdZdS )ÚWorkersJoinedzAll workers have terminated.N)rd   re   rf   rg   r0   r0   r0   r1   rt   ß   s    rt   c                 C   s   t ƒ ‚r+   )r   )ÚsignumÚframer0   r0   r1   Úsoft_timeout_sighandlerã   ó   rw   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 ) ÚWorkerNr0   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)rK   r„   r…   r†   r}   r~   r   Úsentinelr‚   rƒ   r‡   r€   rˆ   r0   r0   r1   rO   í   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_fdr2   Úsend_syn_offsetÚ_send_syn_offsetÚsendÚ
_quick_putÚrecvÚ
_quick_getÚsend_job_offset)rK   Úobjr0   r0   r1   r‰   þ   s   zWorker.contribute_to_objectc                 C   s:   | j | j| j| j| j| j| j| j| j| j	| j
| j| jffS r+   )rn   r„   r…   r†   r}   r~   r   r�   r‚   rƒ   r‡   r€   rˆ   rJ   r0   r0   r1   Ú
__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 rz   r0   )Ústatus©Ú_exitÚ	_exitcoder0   r1   Ú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)rK   rž   r    rk   r0   r›   r1   Ú__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œ   )rK   r    Úexitcoderk   r0   r0   r1   r«   +  s   

zWorker._do_exitc                 C   ó   d S r+   r0   ©rK   r    r0   r0   r1   r¨   ;  ó   zWorker.on_loop_startc                 C   s   |S r+   r0   )rK   Úresultr0   r0   r1   Ú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_r6   ©Ú_wait_for_synrK   r0   r1   Ú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¢   ri   r   r=   r    ÚMAXMEM_USED_FMTÚformatÚ
EX_RECYCLEÚ_ensure_messages_consumedr®   r)   )rK   r   Únowr    r°   rŒ   r�   r   r€   r¹   rÅ   rÃ   rÄ   r¿   rÀ   Úargs_Újobr¾   Úfunr6   r?   Úconfirmr¸   rk   Ú_ÚtbÚwrappedÚeinfoÚused_kbr0   rÁ   r1   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ˆ   ÚrangeÚ)GUARANTEE_MESSAGE_CONSUMPTION_RETRY_LIMITrQ   r   r²   r³   Ú,GUARANTEE_MESSAGE_CONSUMPTION_RETRY_INTERVALr    )rK   rÄ   Úretryr0   r0   r1   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Úsignalrw   ÚSIGINTÚSIG_IGNr,   rJ   r0   r0   r1   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 ©NTr0   )rD   Úloads)rá   r0   r1   Ú_recv¼  ó   z'Worker._make_recv_method.<locals>._recvc                    s   ˆ | ƒr	dˆƒ fS dS ©NT©FNr0   ©rD   )Ú_pollÚgetr0   r1   rä   ¿  s   
c                    s(   zdˆ | d�fW S  t jy   Y dS w ©NTrè   rç   r   rè   )rê   r0   r1   rä   Ä  s
   ÿ)rê   rÛ   r�   Úpollrá   r
   )rK   Úconnrä   r0   )ré   rê   rá   r1   Ú_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Ã   )rK   rã   r0   r0   r1   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{   rd   r®   )r   Úreadyr¿   rk   ©Ú_receiveÚshould_shutdownr0   r1   ÚreceiveÔ  s&   
ÿ€üz/Worker._make_protected_receive.<locals>.receive)rî   r�   Úis_setr   )rK   rí   rú   r0   r÷   r1   rï   Ð  s   
zWorker._make_protected_receive)
NNr0   NNNTTNNr+   )rd   re   rf   rO   r‰   r™   r¬   r«   r¨   r¹   r   r   r©   rË   r§   rî   r
   r¦   rï   r0   r0   r0   r1   ry   ë   s$    
ý
Mry   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   rO   ÚRUNÚ_stateÚ_was_startedÚdaemon©rK   r6   r?   r0   r0   r1   rO   ð  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{   rd   Ú_killr¤   r¥   r   r£   rž   rª   rœ   ©rK   rk   r0   r0   r1   Úrunö  s    
ÿ€ÿ€ýzPoolThread.runc                    s    d| _ tt| ƒj|i |¤Ž d S râ   )r   rl   rü   Ústartr  rm   r0   r1   r    s   zPoolThread.startc                 C   rµ   r+   r0   rJ   r0   r0   r1   Úon_stop_not_started  r·   zPoolThread.on_stop_not_startedNc                 C   s    | j r
|  |¡ d S |  ¡  d S r+   )r   Újoinr  ©rK   rD   r0   r0   r1   rB   
  s   
zPoolThread.stopc                 C   ó
   t | _d S r+   )Ú	TERMINATErÿ   rJ   r0   r0   r1   Ú	terminate  ó   
zPoolThread.terminatec                 C   r  r+   )ÚCLOSErÿ   rJ   r0   r0   r1   rÜ     r  zPoolThread.closer+   )rd   re   rf   rO   r  r  r  rB   r  rÜ   rs   r0   r0   rm   r1   rü   î  s    
rü   c                       s$   e Zd Z‡ fdd„Zdd„ Z‡  ZS )Ú
Supervisorc                    s   || _ tƒ  ¡  d S r+   )Úpoolrl   rO   )rK   r  rm   r0   r1   rO     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	  )rK   r  Ú
prev_staterÑ   r0   r0   r1   r    s.   

€
þ€ýzSupervisor.body)rd   re   rf   rO   r  rs   r0   r0   rm   r1   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                    s,   || _ || _|| _|| _|| _tƒ  ¡  d S r+   )Ú	taskqueuer°   Úoutqueuer  Úcacherl   rO   )rK   r  r°   r  r  r  rm   r0   r1   rO   >  ó   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)
rK   r  r  r°   ÚtaskseqÚ
set_lengthÚtaskr¾   rÎ   Úindr0   r0   r1   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ó   )rK   r  r°   r  Úpr0   r0   r1   r   p  s   

ÿÿzTaskHandler.tell_othersc                 C   s   |   ¡  d S r+   )r   rJ   r0   r0   r1   r  ƒ  r8   zTaskHandler.on_stop_not_started)rd   re   rf   rO   r  r   r  rs   r0   r0   rm   r1   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                    s,   || _ || _|| _|| _d | _tƒ  ¡  d S r+   )Ú	processesr  Út_softÚt_hardÚ_itrl   rO   )rK   r'  r  r(  r)  rm   r0   r1   rO   ‰  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Ÿ   r0   r1   Ú	<genexpr>’  ó   € 
ÿþz1TimeoutHandler._process_by_pid.<locals>.<genexpr>©NN)Únextr  r'  r¶   r0   rŸ   r1   Ú_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   r4  Ú_worker_pidÚhandle_timeoutr  rÝ   ÚOSErrorr   rô   ÚESRCH)rK   rÎ   ÚprocessÚ_indexrk   r0   r0   r1   Ú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 %rFr6  )rö   r   r   Ú_timeoutr  Ú_jobr   r4  r8  r9  Ú_trywaitkill)rK   rÎ   r<  r=  r0   r0   r1   Ú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  ©rK   Úworkerr0   r0   r1   rA  »  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  rD   r0   r0   r1   Ú
_timed_outØ  s
   ÿz2TimeoutHandler.handle_timeouts.<locals>._timed_outc                 3   s   � | ]	}|ˆ v r|V  qd S r+   r0   )r.  Úk©r  r0   r1   r0  ç  ó   € z1TimeoutHandler.handle_timeouts.<locals>.<genexpr>)r)  r(  Úsetr>  rB  rÿ   rþ   Úcopyr  ÚitemsÚ_time_acceptedÚ_soft_timeoutr?  Úadd)rK   r)  r(  Údirtyr>  rB  rL  r¾   rÎ   Úack_timeÚsoft_timeoutÚhard_timeoutr0   rN  r1   Ú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þ   rZ  r²   r³   r   r   ©rK   rÑ   r0   r0   r1   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*  rZ  r3  ÚStopIteration©rK   r6   r0   r0   r1   Úhandle_event  s   

ÿzTimeoutHandler.handle_event)rd   re   rf   rO   r4  r>  rB  rA  rZ  r  r^  rs   r0   r0   rm   r1   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                    s^   || _ || _|| _|| _|| _|| _|| _d | _d| _|| _	|	| _
|
| _|  ¡  tƒ  ¡  d S )NF)r  rê   r  rì   Újoin_exited_workersÚputlockr   r*  Ú_shutdown_completeÚcheck_timeoutsÚon_job_readyÚon_ready_countersÚ_make_methodsrl   rO   )rK   r  rê   r  rì   r`  ra  r   rc  rd  re  rm   r0   r1   rO     s   zResultHandler.__init__c                 C   s   | j dd� d S )NT)rZ  )Úfinish_at_shutdownrJ   r0   r0   r1   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 rz   )ÚRÚ_ackr  r,   )rÎ   r¾   Útime_acceptedr    r�   )r  r   r0   r1   Ú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 rG   )
r  re  r3  r  Úworker_pidsÚget_lockrQ   rö   r[   r  )rÎ   r¾   r˜   rŒ   ÚitemÚ
worker_pidrˆ   )r  rd  ra  rK   r0   r1   Ú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´   rk   r0   r0   r1   Ú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#  Ústater6   )Ústate_handlersr0   r1   Úon_state_changeR  s   ÿz4ResultHandler._make_methods.<locals>.on_state_change)	r  ra  r   rd  r¼   rÇ   r±   rt  ru  )rK   rl  rq  rr  ru  r0   )r  rd  ra  r   rK   rt  r1   rf  "  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ì   ru  ró   rò   r   r   rÿ   r  )rK   rD   rì   ru  rö   r#  rk   r0   r0   r1   Ú_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 rz   )rÿ   rþ   r*  rw  r3  r\  r   )rK   r-   Úeventsr0   r0   r1   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þ   rw  r   rg  r[  r0   r0   r1   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ð   rv  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)rb  rê   r  r  rì   r`  rc  ru  rÿ   r  ró   rò   r   rt   r   ÚabsÚminrÛ   rÖ   r�   Úlen)rK   rZ  rê   r  r  rì   r`  rc  ru  Útime_terminaterö   r#  rk   rÌ   r¾   r0   r0   r1   rg  Š  sj   
€þþ€øï

€ÿ
ÿz ResultHandler.finish_at_shutdownr+   )rð   r2  ©F)rd   re   rf   rO   r  rf  rw  r^  r  rg  rs   r0   r0   rm   r1   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.
    TNr0   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é   )r6   Úexitpriority)Gr   Ú_ctxÚsynackÚ_setup_queuesr   Ú
_taskqueueÚ_cacherþ   rÿ   rD   rX  Ú_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ÚputlocksrF   Ú_putlockrÖ   Ú_create_worker_processr  Ú_worker_handlerr  r  r”   Ú	_outqueueÚ_task_handlerrc  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)rK   r'  r}   r~   ÚmaxtasksperchildrD   rX  r�  r�  Úmax_restart_freqr�  r‘  r’  r“  r”  Ú	semaphorer¥  r–  r…  Úon_process_exitÚcontextr€   r˜  r?   r¾   r0   r0   r1   rO   Ï  s®   
ÿýÿ

ü
þ
€


üùzPool.__init__c                 O   s   | j |i |¤ŽS r+   )r¡  )rK   r6   Úkwdsr0   r0   r1   r   I  ó   zPool.Processc                 C   s   |  | j|d�¡S )N)Útarget)r‰   r   rJ  r0   r0   r1   ÚWorkerProcessL  ó   zPool.WorkerProcessc              
   K   s:   | j | j| j| j| j| j| j| j| j| j	f	d| j
i|¤ŽS )Nre  )r_  r©  r–   rˆ  Ú_poll_resultÚ_join_exited_workersr¦  r   rc  rd  r¤  )rK   Úextra_kwargsr0   r0   r1   r°  O  s   üüûzPool.create_result_handlerc                 C   rµ   r+   r0   )rK   rÎ   r¾   r˜   rŒ   r0   r0   r1   rd  X  r·   zPool.on_job_readyc                 C   s   | j | j| jfS r+   )r´  rª  r¢  rJ   r0   r0   r1   rµ  [  r½  zPool._help_stuff_finish_argsc                 C   s   zt ƒ W S  ty   Y dS w rG   )r   ÚNotImplementedErrorrJ   r0   r0   r1   r   ^  s
   ÿzPool.cpu_countc                 G   s   | j j|Ž S r+   )r±  r^  r]  r0   r0   r1   r²  d  r8   zPool.handle_result_eventc                 C   rµ   r+   r0   )rK   rK  Úqueuesr0   r0   r1   Ú_process_register_queuesg  r·   zPool._process_register_queuesc                    r+  )Nc                 3   r,  r+   rŸ   r-  rŸ   r0   r1   r0  k  r1  z'Pool._process_by_pid.<locals>.<genexpr>r2  )r3  r  r¢  r¶   r0   rŸ   r1   r4  j  r5  zPool._process_by_pidc                 C   s   | j | jd fS r+   )r´  r©  rJ   r0   r0   r1   Ú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¿  ry   r‹  rŒ  r‰  r�  r”  Ú_wrap_exceptionrŠ  r¢  ÚappendrÆ  ÚnameÚreplacer  Úindexr  r£  r    r¤  r�  )rK   r¾   rŠ   r„   r…   r†   rˆ   Úwr0   r0   r1   r§  s  s,   
ø

zPool._create_worker_processc                 C   rµ   r+   r0   rJ  r0   r0   r1   Ú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°ˆ ¡ 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 r0   )rö   Ú_worker_lost)r.  rÎ   r0   r0   r1   Ú
<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 r0   rŸ   ©r.  rÐ  r0   r0   r1   rÓ  ½  s    c                 3   s$   � | ]}|ˆv s|ˆ vr|V  qd S r+   r0   ©r.  r    ©Úall_pidsÚcleanedr0   r1   r0  À  s   € ÿÿz,Pool._join_exited_workers.<locals>.<genexpr>Ú_job_terminated)&r3   rˆ  Úvaluesr   rÒ  Ú_lost_worker_timeoutÚmark_as_worker_lostr|  r¢  rt   ÚreversedrÖ   r´   rG  r   r	  r    r)   rÊ   Úgetattrr=   rÍ  r	   rÑ  r£  r¤  r3  rm  Úon_job_process_downrö   rê   Ú_set_terminatedÚon_job_process_lostÚ	_write_toÚ_scheduled_forÚ	_is_aliver‘  Ú_process_cleanup_queues)rK   ry  rÌ   rÎ   Ú	lost_timeÚlost_retÚ	exitcodesr¾   rK  r´   ÚpopenÚacked_by_goner/  Úwrite_toÚ	sched_forr0   r×  r1   rÂ  �  s†   

€






ÿý


€ý
ÿ€€

€zPool._join_exited_workersc                 C   rµ   r+   r0   )rK   rÎ   rK  r0   r0   r1   Úon_partial_readã  r·   zPool.on_partial_readc                 C   rµ   r+   r0   rJ  r0   r0   r1   ræ  æ  r·   zPool._process_cleanup_queuesc                 C   rµ   r+   r0   )rK   rÎ   Úpid_goner0   r0   r1   rà  é  r·   zPool.on_job_process_downc                 C   s   t ƒ |f|_d S r+   )r   rÒ  )rK   rÎ   r    r´   r0   r0   r1   râ  ì  r½  zPool.on_job_process_lostc                 C   s>   zt d t|ƒ|j¡ƒ‚ t y   | d dtƒ f¡ Y d S w )Nz(Worker exited prematurely: {0} Job: {1}.F)r   rÉ   r	   r@  r  r   )rK   rÎ   r´   r0   r0   r1   rÝ  ï  s   
ÿÿÿzPool.mark_as_worker_lostc                 C   ó   | S r+   r0   rJ   r0   r0   r1   Ú	__enter__ú  r·   zPool.__enter__c                 G   s   |   ¡ S r+   )r  )rK   r¢   r0   r0   r1   Ú__exit__ý  s   zPool.__exit__c                 C   rµ   r+   r0   ©rK   Únr0   r0   r1   Úon_grow   r·   zPool.on_growc                 C   rµ   r+   r0   ró  r0   r0   r1   Ú	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¦  rL   Úterminate_controlledrö  Ú
ValueError)rK   rô  r¾   rK  r0   r0   r1   rL     s   

ÿzPool.shrinkc                 C   s:   t |ƒD ]}|  jd7  _| jr| j ¡  q|  |¡ d S rG   )rÖ   r  r¦  rV   rõ  )rK   rô  r¾   r0   r0   r1   rV     s   
€z	Pool.growc                 c   s"   � | j D ]
}|  |¡s|V  qd S r+   )r¢  Ú_worker_activerJ  r0   r0   r1   r÷    s   €

€þzPool._iterinactivec                 C   s(   | j  ¡ D ]}|j| ¡ v r dS qdS )NTF)rˆ  rÛ  r    rm  )rK   rK  rÎ   r0   r0   r1   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   )rK   ré  r¾   r0   r0   r1   Ú_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%  r0   r0   r1   r0  5  s   € z$Pool._avail_index.<locals>.<genexpr>c                 3   s   � | ]	}|ˆ vr|V  qd S r+   r0   )r.  r¾   ©Úindicesr0   r1   r0  6  rO  )r|  r¢  r  rP  r3  rÖ   rJ   r0   rÿ  r1   rý  3  s   zPool._avail_indexc                 C   s
   |   ¡  S r+   )rÂ  rJ   r0   r0   r1   Ú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[   )rK   Újoinedr¾   r0   r0   r1   r  ;  s   


€þzPool._maintain_poolc              
   C   sz   | j jtkr9| jtkr;z|  ¡  W d S  ty"   |  ¡  |  ¡  ‚  ty8 } zt|ƒt	j
kr3t|‚‚ d }~ww d S d S r+   )r¨  rÿ   rþ   r  r   rÜ   r	  r:  r   rô   ÚENOMEMÚMemoryErrorr  r0   r0   r1   Ú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è   rJ   r0   r1   rÁ  W  s   z(Pool._setup_queues.<locals>._poll_result)r„  ÚSimpleQueuer´  r©  r‹   r“   r”   r�   r•   r–   rÁ  ©rK   rÁ  r0   rJ   r1   r†  Q  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  rJ   r0   r0   r1   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ê   )rK   Úfuncr6   r¼  r0   r0   r1   Úapplyf  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ê   ©rK   r
  ÚiterableÚ	chunksizer0   r0   r1   r:   m  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;   ©rK   r
  r  r  ÚcallbackÚerror_callbackr0   r0   r1   Ústarmap_asyncw  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  r0   r0   r1   r4   €  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Æ   r@  ©r.  r¾   Úx©r
  r¸   r0   r1   r0  “  ó   € ÿzPool.imap.<locals>.<genexpr>c                 3   ó*   � | ]\}}t ˆ j|t|fi ffV  qd S r+   ©rÆ   r@  r7   r  ©r¸   r0   r1   r0  ž  r  c                 s   ó   � | ]
}|D ]}|V  qqd S r+   r0   ©r.  Úchunkro  r0   r0   r1   r0  ¢  ó   € )rÿ   rþ   r�  ÚIMapIteratorrˆ  r‡  r°   r  Ú_set_lengthr  Ú
_get_tasks©rK   r
  r  r  r�  Útask_batchesr0   r  r1   Ú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  r0   r1   r0  ±  r  z&Pool.imap_unordered.<locals>.<genexpr>c                 3   r  r+   r  r  r  r0   r1   r0  ½  r  c                 s   r   r+   r0   r!  r0   r0   r1   r0  Á  r#  )rÿ   rþ   r�  ÚIMapUnorderedIteratorrˆ  r‡  r°   r  r%  r  r&  r'  r0   r  r1   Ú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þ   rX  rD   r�  rÝ   r™  rš  r›  r¥  r¦  rI   ÚApplyResultrˆ  r’  r“  r…  r-  r¯  r”  r‡  r°   rÆ   r@  r”   )rK   r
  r6   r¼  r  r  Úaccept_callbackÚtimeout_callbackÚwaitforslotrX  rD   r�  r,  r.  r¸   r0   r0   r1   r	  Ã  sF   



ÿ


ù	ÿÿÿëzPool.apply_asyncc                 C   rµ   r+   r0   )rK   ÚresponserÎ   r¾   Úfdr0   r0   r1   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â   )	r4  r  r   r:  r   rô   r;  rÔ  rÚ  )rK   r    Úsigr/  rÑ   rk   r0   r0   r1   Úterminate_jobþ  s   ÿ€ÿ
øzPool.terminate_jobc                 C   s   |   ||t|||¡S )z<
        Asynchronous equivalent of `map()` method.
        )r  r7   r  r0   r0   r1   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¸   r0   r1   r0  '  r  z"Pool._map_async.<locals>.<genexpr>)rÿ   rþ   rÛ   r3   Údivmodr|  r¢  r  r&  Ú	MapResultrˆ  r‡  r°   r  )	rK   r
  r  r:  r  r  r  Úextrar(  r0   r9  r1   r    s(   

ÿÿÿzPool._map_asyncc                 c   s0   � t |ƒ}	 tt ||¡ƒ}|sd S | |fV  qr+   )r  Útupler9   Úislice)r
  ÚitÚsizer  r0   r0   r1   r&  +  s   €
üzPool._get_tasksc                 C   s   t dƒ‚)Nz:pool objects cannot be passed between processes or pickled)rÄ  rJ   r0   r0   r1   r™   4  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°   rE   rJ   r0   r0   r1   rÜ   9  s   


úz
Pool.closec                 C   s$   t dƒ t| _| j ¡  |  ¡  d S )Nzterminating pool)r   r  rÿ   r¨  r  r¶  rJ   r0   r0   r1   r  C  s   
zPool.terminatec                 C   s   t | ƒ d S r+   )rE   )Útask_handlerr0   r0   r1   Ú_stop_task_handlerI  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   rE   r¨  rC  rª  r±  r  r¢  r|  rG  r	  )rK   r¾   r%  r0   r0   r1   r	  M  s   


€z	Pool.joinc                 C   s   | j  ¡ D ]}| ¡  qd S r+   )r£  rÛ  rP  )rK   Úer0   r0   r1   Ú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   Ú_rlockrI   Úis_aliver�   rì   r•   r²   r³   )ÚinqueuerB  r¢  r0   r0   r1   Ú_help_stuff_finish`  s   


"þzPool._help_stuff_finishc                 C   s   |  d ¡ d S r+   )r°   )Úclsr  r  r0   r0   r1   Ú_set_result_sentineli  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  rB   ÚTIMEOUT_MAXrI  r    rG  r	  rÜ   )rL  r  rJ  r  r  Úworker_handlerrB  Úresult_handlerr  Útimeout_handlerÚhelp_stuff_finish_argsr%  r0   r0   r1   r³  m  sJ   

€


€ÿzPool._terminate_poolc                 C   ó   dd„ | j D ƒS )Nc                 S   s   g | ]}|j j‘qS r0   )rG  rŠ   rÕ  r0   r0   r1   rÓ  ¦  ó    z*Pool.process_sentinels.<locals>.<listcomp>)r¢  rJ   r0   r0   r1   Úprocess_sentinels¤  rh  zPool.process_sentinels)NNr0   NNNNNr   NNNNTNFFFNNNFr~  )r   r+   )NNNrc   )Frd   re   rf   rg   rË  ry   r  r  r&  r_  r   rO   r   r¿  r°  rd  rµ  r   r²  rÆ  r4  rÇ  r§  rÑ  rÂ  rî  ræ  rà  râ  rÝ  rñ  rò  rõ  rö  rL   rV   r÷  rú  rþ  rý  r  r  r  r†  r¯  r  r:   r  r4   r)  r+  r	  r-  r6  r  r  Ústaticmethodr&  r™   rÜ   r  rC  r	  rG  rK  ÚclassmethodrM  r³  ÚpropertyrU  r0   r0   r0   r1   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/  Nr0   c                 C   sš   || _ tƒ | _t ¡ | _ttƒ| _|| _	|| _
|| _|| _|| _|| _|| _|| _|	| _|
| _|p2d| _|| _d| _d| _d | _d | _d | _| || j< d S )Nr0   F)r.  r¬  Ú_mutexr@   rÉ  Ú_eventr3  Újob_counterr@  rˆ  Ú	_callbackÚ_accept_callbackÚ_error_callbackÚ_timeout_callbackr?  rT  rÜ  Ú_on_timeout_setÚ_on_timeout_cancelÚ_callbacks_propagateÚ	_send_ackÚ	_acceptedÚ
_cancelledr8  rS  Ú_terminated)rK   r  r  r0  r1  r  rX  rD   r�  r’  r“  r,  r-  r.  r0   r0   r1   rO   ²  s,   


zApplyResult.__init__c                 C   s   dj | jj| j| j|  ¡ d�S )Nz&<{name}: {id} ack:{ack} ready:{ready}>)rÍ  ÚidÚackrö   )rÉ   rn   rd   r@  rd  rö   rJ   r0   r0   r1   rp   Ð  s   þzApplyResult.__repr__c                 C   s
   | j  ¡ S r+   )rZ  rû   rJ   r0   r0   r1   rö   Ö  r  zApplyResult.readyc                 C   ó   | j S r+   )rd  rJ   r0   r0   r1   ÚacceptedÙ  rx   zApplyResult.acceptedc                 C   s   |   ¡ sJ ‚| jS r+   )rö   Ú_successrJ   r0   r0   r1   Ú
successfulÜ  s   zApplyResult.successfulc                 C   s
   d| _ dS )zOnly works if synack is used.TN)re  rJ   r0   r0   r1   Ú_cancelà  s   
zApplyResult._cancelc                 C   s   | j  | jd ¡ d S r+   )rˆ  Úpopr@  rJ   r0   r0   r1   Údiscardä  rq   zApplyResult.discardc                 C   s
   || _ d S r+   )rf  ©rK   ru   r0   r0   r1   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   rp  r0   r0   r1   rá  ê  s
   ÿzApplyResult._set_terminatedc                 C   s   | j r| j gS g S r+   ©r8  rJ   r0   r0   r1   rm  ð  rÀ  zApplyResult.worker_pidsc                 C   s   | j  |¡ d S r+   )rZ  rH  r
  r0   r0   r1   rH  ó  r½  zApplyResult.waitc                 C   s*   |   |¡ |  ¡ st‚| jr| jS | jj‚r+   )rH  rö   r   rk  rT   Ú	exceptionr
  r0   r0   r1   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=   )rK   rÏ   r6   r?   rk   r0   r0   r1   Ú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)r7  rD   )r_  rt  rT  r?  )rK   r7  r0   r0   r1   r9  	  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  rk  rT   rZ  rP  rd  rˆ  rn  r@  r\  rt  r^  ©rK   r¾   r˜   r0   r0   r1   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»   r@  rS  r8  rö   rˆ  rn  r`  rT  r?  r¼   r]  Ú_propagate_errorsrª   )rK   r¾   rk  r    r�   r3  r0   r0   r1   rj  %  sZ   üûÿ€

ÿæ€

ÿæ
âã"ãzApplyResult._ackr+   r~  )rd   re   rf   rÒ  rã  rä  rŽ  rO   rp   rö   rj  rl  rm  ro  r  rá  rm  rH  rê   rt  r9  r  rj  r0   r0   r0   r1   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/  rO   rk  Ú_lengthrT   rd  r8  rS  Ú
_chunksizeÚ_number_leftrZ  rP  r@  r—  )rK   r  r  Úlengthr  r  r0   r0   r1   rO   M  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)rT   rx  ry  r\  rd  rˆ  rn  r@  rZ  rP  rk  r^  )rK   r¾   Úsuccess_resultÚsuccessr¸   r0   r0   r1   r  _  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)rx  r{  rw  rÖ   rd  r8  rS  rö   rˆ  rn  r@  )rK   r¾   rk  r    r6   r  rB   Újr0   r0   r1   rj  s  s   


ÿzMapResult._ackc                 C   s
   t | jƒS r+   )Úallrd  rJ   r0   r0   r1   rj  }  r  zMapResult.acceptedc                 C   rS  )Nc                 S   s   g | ]}|r|‘qS r0   r0   rÖ  r0   r0   r1   rÓ  �  rT  z)MapResult.worker_pids.<locals>.<listcomp>rr  rJ   r0   r0   r1   rm  €  r½  zMapResult.worker_pidsN)rd   re   rf   rO   r  rj  rj  rm  r0   r0   r0   r1   r<  K  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 rq  )r@   Ú	Conditionr¬  rS   r3  r[  r@  rˆ  r   Ú_itemsr=  rw  Ú_readyÚ	_unsortedÚ_worker_pidsrÜ  )rK   r  r�  r0   r0   r1   rO   ‹  s   
zIMapIterator.__init__c                 C   rð  r+   r0   rJ   r0   r0   r1   Ú__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â   )rS   r�  Úpopleftrü  r=  rw  r‚  r\  rH  r   rª   )rK   rD   ro  r|  rQ   r0   r0   r1   r3  ›  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}  )rS   r=  r�  rÌ  rƒ  rn  rU   rw  r‚  rˆ  r@  ru  r0   r0   r1   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â   )rS   rw  r=  r‚  rU   rˆ  r@  )rK   rz  r0   r0   r1   r%  Ä  s   
û"þzIMapIterator._set_lengthc                 G   s   | j  |¡ d S r+   )r„  rÌ  )rK   r¾   rk  r    r6   r0   r0   r1   rj  Ì  r½  zIMapIterator._ackc                 C   ri  r+   )r‚  rJ   r0   r0   r1   rö   Ï  rx   zIMapIterator.readyc                 C   ri  r+   )r„  rJ   r0   r0   r1   rm  Ò  rx   zIMapIterator.worker_pidsr+   )rd   re   rf   rÒ  rŽ  rO   r…  r3  Ú__next__r  r%  rj  rö   rm  r0   r0   r0   r1   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}  )	rS   r�  rÌ  r=  rU   rw  r‚  rˆ  r@  ru  r0   r0   r1   r  Ü  s   
ú"üzIMapUnorderedIterator._setN)rd   re   rf   r  r0   r0   r0   r1   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   Nr0   c                 C   s   t  | |||¡ d S r+   )r  rO   )rK   r'  r}   r~   r0   r0   r1   rO   ï  rq   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è   rJ   r0   r1   rÁ  ø  s
   ÿz.ThreadPool._setup_queues.<locals>._poll_result)r   r´  r©  r°   r”   rê   r–   rÁ  r  r0   rJ   r1   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|  rX   )rJ  rB  r  r0   r0   r1   rK  ÿ  s
   
"ýzThreadPool._help_stuff_finish)NNr0   )
rd   re   rf   Údummyr   r   rO   r†  rV  rK  r0   r0   r0   r1   rˆ  ê  s    
rˆ  r+   )erQ  rô   r9   r¤   r¯   rÞ   r£   r@   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   r   rŠ  r   r   r   r   r    rÈ   Úversion_inforh   ÚsystemÚ_winr#   r  rI  r$   rN  r,   Ú	SemaphorerN   rþ   r  r  r¼   rÇ   rÆ   r»   r±   r)   r®   rÊ   rß  rÝ   rŽ  r×   rØ   Úcountr[  r¬  r2   r7   r;   r=   rE   rF   rª   ri   rt   rw   ry   rü   r  r  r&  r_  r  r/  r<  r$  r*  rˆ  r0   r0   r0   r1   Ú<module>   s¬   $	
ÿ


;  )%K  :     o =R