o
    wvXj‡É  ã                   @   s   d Z ddlmZ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mZmZ ddlmZ ddlm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 ddl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+ ddl,m-Z- ddl.m/Z/ ddl0m1Z1 ddl2m3Z3m4Z4m5Z5 ddl6m7Z7m8Z8m9Z9 ddl:m;Z; ddl<m=Z= ddl>m?Z@ zddlAmBZC dZDe	jEd dkrÓe	jEdk rÓe9fdd„Z9W n eFyì   ejBfdd„ZCd ZDe8fd!d„Z9Y nw d"ZGe=eHƒZIeIjJeIjKZJZKeLejMejNhƒZOd#ZPd$ZQd%ZRd&ZSeSeReReSd'œZTd(d)„ eT 4¡ D ƒZUed*d+ƒZVd,d-„ ZWd.d/„ ZXeYed0ƒ�r?ddddejZej[ej\ej]fd1d2„Z^nd>d3d2„Z^dddde^fd4d5„Z_ze` W n ea�y^   ebZ`Y nw d6d7„ ZcG d8d9„ d9ejdƒZdG d:d;„ d;ejeƒZeG d<d=„ d=ejfƒZgdS )?a‡  Version of multiprocessing.Pool using Async I/O.

.. note::

    This module will be moved soon, so don't use it directly.

This is a non-blocking version of :class:`multiprocessing.Pool`.

This code deals with three major challenges:

#. Starting up child processes and keeping them running.
#. Sending jobs to the processes and receiving results back.
#. Safely shutting down this system.
é    )Úabsolute_importÚunicode_literalsN)ÚdequeÚ
namedtuple)ÚBytesIO)ÚIntegral)ÚHIGHEST_PROTOCOL)Úsleep)ÚWeakValueDictionaryÚref)Úpool)Úbuf_tÚ
isblockingÚsetblocking)ÚACKÚNACKÚRUNÚ	TERMINATEÚWorkersJoined)Ú_SimpleQueue)ÚERRÚWRITE)Úpickle)ÚSELECT_BAD_FD)Úfxrange)Úpromise)ÚCounterÚitemsÚvalues)ÚpackÚunpackÚunpack_from)Únoop)Ú
get_logger)Ústate)ÚreadTé   )r&   é   é   c                 C   ó   || |  ¡ ƒS ©N)Útobytes)ÚfmtÚviewÚ_unpack_from© r/   úX/var/www/html/myproject/venv/lib/python3.10/site-packages/celery/concurrency/asynpool.pyr!   :   ó   r!   c                 C   s(   || |ƒ}t |ƒ}|dkr| |¡ |S ©Nr   )ÚlenÚwrite)ÚfdÚbufÚsizer%   ÚchunkÚnr/   r/   r0   Ú__read__?   s
   

r:   Fc                 C   r)   r*   )Úgetvalue)r,   Úiobufr    r/   r/   r0   r!   G   r1   )ÚAsynPoolé   g      @é   é   )NÚfastÚfcfsÚfairc                 C   s   i | ]\}}||“qS r/   r/   )Ú.0ÚkÚvr/   r/   r0   Ú
<dictcomp>`   ó    rG   ÚAck)Úidr5   Úpayloadc                 C   s   | j o| j jdkS )z(Return true if generator is not started.éÿÿÿÿ)Úgi_frameÚf_lasti)Úgenr/   r/   r0   Úgen_not_startede   s   rP   c                 C   s$   z| j }W |ƒ S  ty   Y d S w r*   )Ú_writerÚAttributeError)ÚjobÚwriterr/   r/   r0   Ú_get_job_writerk   s   ýÿrU   Úpollc                    sè   |ƒ }|j ‰| r‡‡fdd„| D ƒ |r‡‡fdd„|D ƒ |r*‡ ‡fdd„|D ƒ tƒ tƒ }	}
|r9|dk r9dnt|d ƒ}| |¡}|D ](\}}t|tƒsS| ¡ }|ˆ@ r\|	 |¡ |ˆ@ re|
 |¡ |ˆ @ rn|	 |¡ qF|	|
dfS )Nc                    ó   g | ]}ˆ|ˆ ƒ‘qS r/   r/   ©rD   r5   )ÚPOLLINÚregisterr/   r0   Ú
<listcomp>|   rH   z_select_imp.<locals>.<listcomp>c                    rW   r/   r/   rX   )ÚPOLLOUTrZ   r/   r0   r[   ~   rH   c                    rW   r/   r/   rX   )ÚPOLLERRrZ   r/   r0   r[   €   rH   r   g     @�@)rZ   ÚsetÚroundrV   Ú
isinstancer   ÚfilenoÚadd)ÚreadersÚwritersÚerrÚtimeoutrV   rY   r\   r]   ÚpollerÚRÚWÚeventsr5   Úeventr/   )r]   rY   r\   rZ   r0   Ú_select_impu   s,   




€
rl   c                 C   s8   t   | |||¡\}}}|rtt|ƒt|ƒB ƒ}||dfS r2   )ÚselectÚlistr^   )rc   rd   re   rf   ÚrÚwÚer/   r/   r0   rl   �   s   
c                 C   s|  | du rt ƒ n| } |du rt ƒ n|}|du rt ƒ n|}z|| |||ƒW S  tjtjfy½ } zŠz|j}W n tyB   |jd }Y nw |tjkrUt ƒ t ƒ dfW  Y d}~S |tv r¸| |B |B D ]K}zt |gg g d¡ W q_ tjtjfyª } z.z|j}W n ty‹   |jd }Y nw |tvr‘‚ |  	|¡ | 	|¡ | 	|¡ W Y d}~q_d}~ww t ƒ t ƒ dfW  Y d}~S ‚ d}~ww )a<  Simple wrapper to :class:`~select.select`, using :`~select.poll`.

    Arguments:
        readers (Set[Fd]): Set of reader fds to test if readable.
        writers (Set[Fd]): Set of writer fds to test if writable.
        err (Set[Fd]): Set of fds to test for error condition.

    All fd sets passed must be mutable as this function
    will remove non-working fds from them, this also means
    the caller must make sure there are still fds in the sets
    before calling us again.

    Returns:
        Tuple[Set, Set, Set]: of ``(readable, writable, again)``, where
        ``readable`` is a set of fds that have data available for read,
        ``writable`` is a set of fds that's ready to be written to
        and ``again`` is a flag that if set means the caller must
        throw away the result and call us again.
    Nr   r?   )
r^   rm   ÚerrorÚsocketÚerrnorR   ÚargsÚEINTRr   Údiscard)rc   rd   re   rf   rV   ÚexcÚ_errnor5   r/   r/   r0   Ú_select—   sD   
ÿ

ÿ

€ö€ärz   c           	   
      sÎ   ‡ ‡fdd„}g }| D ]-‰|ƒ |}}z|ˆg|¢R i |¤Ž W q t tfy8   tjdˆdd� | ˆ¡ Y qw |rc|D ]'‰zt|dƒrK| ˆ¡ n| ˆd¡ W q= tyb   t dˆ|¡ Y q=w dS dS )	a˜  Apply hub method to fds in iter, remove from list if failure.

    Some file descriptors may become stale through OS reasons
    or possibly other reasons, so safely manage our lists of FDs.
    :param fds_iter: the file descriptors to iterate and apply hub_method
    :param source_data: data source to remove FD if it renders OSError
    :param hub_method: the method to call with with each fd and kwargs
    :*args to pass through to the hub_method;
    with a special syntax string '*fd*' represents a substitution
    for the current fd object in the iteration (for some callers).
    :**kwargs to pass through to the hub method (no substitutions needed)
    c                     s"   ˆ } d| v r‡fdd„ˆ D ƒ} | S )Nú*fd*c                    s   g | ]
}|d kr
ˆ n|‘qS )r{   r/   )rD   Úarg)r5   r/   r0   r[   è   s    zTiterate_file_descriptors_safely.<locals>._meta_fd_argument_maker.<locals>.<listcomp>r/   )Ú	call_args©ru   r5   r/   r0   Ú_meta_fd_argument_makerä   s   z@iterate_file_descriptors_safely.<locals>._meta_fd_argument_makerz)Encountered OSError when accessing fd %s T©Úexc_infoÚremoveNz*ValueError trying to invalidate %s from %s)	ÚOSErrorÚFileNotFoundErrorÚloggerÚwarningÚappendÚhasattrr‚   ÚpopÚ
ValueError)	Úfds_iterÚsource_dataÚ
hub_methodru   Úkwargsr   Ú	stale_fdsÚhub_argsÚ
hub_kwargsr/   r~   r0   Úiterate_file_descriptors_safelyÖ   s6   þü
€ÿÿùr’   c                   @   s   e Zd ZdZdd„ ZdS )ÚWorkerzPool worker process.c                 C   s   | j  t|ff¡ d S r*   )ÚoutqÚputÚ	WORKER_UP)ÚselfÚpidr/   r/   r0   Úon_loop_start  s   zWorker.on_loop_startN)Ú__name__Ú
__module__Ú__qualname__Ú__doc__r™   r/   r/   r/   r0   r“     s    r“   c                       s^   e Zd ZdZ‡ fdd„Zeeeee	j
fdd„Zdd„ Zdd	„ Zd
d„ Zdd„ Zdd„ Z‡  ZS )ÚResultHandlerz)Handles messages from the pool processes.c                    s>   |  d¡| _|  d¡| _tt| ƒj|i |¤Ž | j| jt< d S )NÚfileno_to_outqÚon_process_alive)r‰   rŸ   r    Úsuperrž   Ú__init__Ústate_handlersr–   )r—   ru   rŽ   ©Ú	__class__r/   r0   r¢     s   zResultHandler.__init__c	              
   c   s¸  � d }	}
|rt dƒ}t|ƒ}n|ƒ  }}|	dk r\z|||r$||	d … n|d|	 ƒ}W n tyF } z|jtvr9‚ d V  W Y d }~nd }~ww |dkrT|	rQtdƒ‚tƒ ‚|	|7 }	|	dk s|d|ƒ\}|rmt |ƒ}t|ƒ}n|ƒ  }}|
|k r¹z|||r�||
d … n|||
 ƒ}W n ty£ } z|jtvr–‚ d V  W Y d }~nd }~ww |dkr±|
r®tdƒ‚tƒ ‚|
|7 }
|
|k sv||| j|ƒ |rÉ|||ƒƒ}n	| d¡ ||ƒ}|rÚ||ƒ d S d S )Nr   r@   zEnd of file during messagez>i)Ú	bytearrayÚ
memoryviewrƒ   rt   ÚUNAVAILÚEOFErrorÚhandle_eventÚseek)r—   Ú
add_readerr5   Úcallbackr:   Ú
readcanbufr   r!   ÚloadÚHrÚBrr6   Úbufvr9   rx   Ú	body_sizeÚmessager/   r/   r0   Ú_recv_message  sj   €

ÿ
€ýÿó

ÿ
€ýÿó
ÿzResultHandler._recv_messagec                    s6   | j ‰| j‰|j‰ |j‰| j‰‡ ‡‡‡‡fdd„}|S )z3Coroutine reading messages from the pool processes.c              
      s„   zˆ|   W n t y   ˆ| ƒ Y S w ˆˆ | ˆƒ}zt|ƒ W n ty*   Y d S  tttfy:   ˆ| ƒ Y d S w ˆ | |ƒ d S r*   )ÚKeyErrorÚnextÚStopIterationÚIOErrorrƒ   r©   )ra   Úit©r¬   rŸ   Úon_state_changeÚrecv_messageÚremove_readerr/   r0   Úon_result_readableX  s   ÿÿz>ResultHandler._make_process_result.<locals>.on_result_readable)rŸ   r¼   r¬   r¾   rµ   )r—   Úhubr¿   r/   r»   r0   Ú_make_process_resultP  s   z"ResultHandler._make_process_resultc                 C   s   |   |¡| _d S r*   )rÁ   rª   )r—   rÀ   r/   r/   r0   Úregister_with_event_looph  s   z&ResultHandler.register_with_event_loopc                 G   s   t dƒ‚)NzNot registered with event loop)ÚRuntimeError)r—   ru   r/   r/   r0   rª   k  ó   zResultHandler.handle_eventc           	   	   C   sÐ   | j }| j}| j}| j}| j}t|ƒ}|r^|r`| jtkrb|d ur#|ƒ  tƒ }|D ]%}t|g| j| j	|j
||ƒ z|dd� W q( tyM   tdƒ Y  d S w | |¡ |rd|rf| jtksd S d S d S d S d S d S )NT)Úshutdownz&result handler: all workers terminated)ÚcacheÚcheck_timeoutsrŸ   r¼   Újoin_exited_workersr^   Ú_stater   r’   Ú_flush_outqueuerb   r   ÚdebugÚdifference_update)	r—   rÆ   rÇ   rŸ   r¼   rÈ   Ú	outqueuesÚpending_remove_fdr5   r/   r/   r0   Úon_stop_not_startedp  s.   þþ
*ðz!ResultHandler.on_stop_not_startedc                 C   sN  z|| }W n t y   ||ƒ Y S w |jj}zt|dƒ W n ttfy.   ||ƒ Y S w zZz| d¡r;| ¡ }nd }tdƒ W n( tt	fyj   ||ƒ Y W zt|dƒ W S  ttfyi   ||ƒ Y   S w w |rq||ƒ W zt|dƒ W d S  ttfy‰   ||ƒ Y S w zt|dƒ W w  ttfy¦   ||ƒ Y      Y S w )Nr?   r   ç      à?)
r¶   r”   Ú_readerr   rƒ   r¹   rV   Úrecvr	   r©   )r—   r5   r‚   Úprocess_indexr¼   ÚprocÚreaderÚtaskr/   r/   r0   rÊ   Œ  sL   üÿ

€ÿø€ÿþÿzResultHandler._flush_outqueue)rš   r›   rœ   r�   r¢   r:   r®   r   r!   Ú_pickler¯   rµ   rÁ   rÂ   rª   rÏ   rÊ   Ú__classcell__r/   r/   r¤   r0   rž     s    
ý9rž   c                       s~  e Zd ZdZeZeZ‡ fdd„Z		dN‡ fd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dd„ Zdd„ Zeeje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-„ Zd.d/„ Zd0d1„ Z d2d3„ Z!e"d4d5„ ƒZ#‡ fd6d7„Z$d8d9„ Z%d:d;„ Z&d<d=„ Z'd>d?„ Z(d@dA„ Z)dBdC„ Z*ejeefdDdE„Z+e,dFdG„ ƒZ-dHdI„ Z.e,dJdK„ ƒZ/e0dLdM„ ƒZ1‡  Z2S )Or=   zAsyncIO Pool (no threads).c                    s   t t| ƒ |¡}d|_|S )NF)r¡   r=   ÚWorkerProcessÚdead)r—   Úworkerr¤   r/   r0   rÙ   ²  s   zAsynPool.WorkerProcessNFc                    s  t  ||¡ˆ _|d u rˆ  ¡ n|}|ˆ _‡ fdd„t|ƒD ƒˆ _i ˆ _i ˆ _i ˆ _	|d u r/t
n|ˆ _tƒ ˆ _tƒ ˆ _tƒ ˆ _tƒ ˆ _tƒ ˆ _ˆ jjˆ _tƒ ˆ _tƒ ˆ _ttˆ ƒj|g|¢R i |¤Ž ˆ jD ]}|ˆ j|j< |ˆ j	|j< qetˆ jdt ƒˆ _!tˆ jdt ƒˆ _"d S )Nc                    ó   i | ]}ˆ   ¡ d “qS r*   ©Úcreate_process_queues©rD   Ú_©r—   r/   r0   rG   ¿  ó    
ÿz%AsynPool.__init__.<locals>.<dictcomp>Úon_soft_timeoutÚon_hard_timeout)#ÚSCHED_STRATEGIESÚgetÚsched_strategyÚ	cpu_countÚsynackÚrangeÚ_queuesÚ_fileno_to_inqÚ_fileno_to_outqÚ_fileno_to_synqÚPROC_ALIVE_TIMEOUTÚ_proc_alive_timeoutr^   Ú_waiting_to_startÚ_all_inqueuesÚ_active_writesÚ_active_writersÚ_busy_workersrw   Ú_mark_worker_as_availabler   Úoutbound_bufferr   Úwrite_statsr¡   r=   r¢   Ú_poolÚoutqR_fdÚsynqW_fdÚgetattrÚ_timeout_handlerr"   rã   rä   )r—   Ú	processesré   rç   Úproc_alive_timeoutru   rŽ   rÔ   r¤   rá   r0   r¢   ·  s@   ÿ
ÿþ

ÿ
ÿzAsynPool.__init__c                    s   t  ¡  tt| ƒ |¡S r*   )ÚgcÚcollectr¡   r=   Ú_create_worker_process)r—   Úir¤   r/   r0   r  õ  s   zAsynPool._create_worker_processc                 C   s   |   ||¡ |  ¡  d S r*   )Ú_untrack_child_processÚmaintain_pool)r—   rÀ   rÔ   r/   r/   r0   Ú_event_process_exitù  s   zAsynPool._event_process_exitc                 C   sN   z|j }W n ty   t |jj¡ }|_ Y nw t|gd|j| j||ƒ dS )z4Helper method determines appropriate fd for process.N)	Ú_sentinel_pollrR   ÚosÚdupÚ_popenÚsentinelr’   r¬   r  ©r—   rÔ   rÀ   r5   r/   r/   r0   Ú_track_child_processþ  s   
û
þzAsynPool._track_child_processc                 C   s4   |j d ur|j d }|_ | |¡ t |¡ d S d S r*   )r  r‚   r  Úcloser  r/   r/   r0   r    s
   

ýzAsynPool._untrack_child_processc                    s”   ˆj  ˆ ¡ ˆj jˆ_ˆ ˆ ¡ ˆ ˆ ¡ ˆ ˆ ¡ ‡ ‡fdd„ˆjD ƒ tˆj	ˆj	ˆ j
ˆjdƒ tˆjƒD ]
\}}ˆ  ||¡ q6ˆ j ˆj¡ dS )z4Register the async pool with the current event loop.c                    s   g | ]}ˆ  |ˆ ¡‘qS r/   )r  ©rD   rp   ©rÀ   r—   r/   r0   r[     s    z5AsynPool.register_with_event_loop.<locals>.<listcomp>r{   N)Ú_result_handlerrÂ   rª   Úhandle_result_eventÚ_create_timelimit_handlersÚ_create_process_handlersÚ_create_write_handlersrù   r’   rí   r¬   r   ÚtimersÚcall_repeatedlyÚon_tickrb   Úon_poll_start)r—   rÀ   ÚhandlerÚintervalr/   r  r0   rÂ     s   



þz!AsynPool.register_with_event_loopc                    sR   ˆj ‰tƒ  ‰ˆ_‡‡‡‡fdd„}|ˆ_‡fdd„‰ ˆ ˆ_‡ fdd„}|ˆ_dS )z.Create handlers used to implement time limits.c                    sF   |rˆ |ˆj | j||ˆƒˆ| j< d S |r!ˆ |ˆj| jƒˆ| j< d S d S r*   )Ú_on_soft_timeoutÚ_jobÚ_on_hard_timeout)rh   ÚsoftÚhard)Ú
call_laterrÀ   r—   Útrefsr/   r0   Úon_timeout_set/  s   ÿ
ÿÿz;AsynPool._create_timelimit_handlers.<locals>.on_timeout_setc              	      s4   zˆ   | ¡}| ¡  ~W d S  ttfy   Y d S w r*   )r‰   Úcancelr¶   rR   )rS   Útref)r"  r/   r0   Ú_discard_tref:  s   
ÿz:AsynPool._create_timelimit_handlers.<locals>._discard_trefc                    s   ˆ | j ƒ d S r*   )r  )rh   )r&  r/   r0   Úon_timeout_cancelC  r1   z>AsynPool._create_timelimit_handlers.<locals>.on_timeout_cancelN)r!  r
   Ú_tref_for_idr#  r&  r'  )r—   rÀ   r#  r'  r/   )r&  r!  rÀ   r—   r"  r0   r  *  s   	
z#AsynPool._create_timelimit_handlersc              	   C   sv   |r|  || | j|¡| j|< z"z| j| }W n	 ty    Y nw |  |¡ W |s0|  |¡ d S d S |s:|  |¡ w w r*   )r!  r  r(  Ú_cacher¶   rã   r&  )r—   rS   r  r   rÀ   Úresultr/   r/   r0   r  G  s    
ÿÿ
€þþzAsynPool._on_soft_timeoutc              	   C   sZ   z&z| j | }W n	 ty   Y nw |  |¡ W |  |¡ d S W |  |¡ d S |  |¡ w r*   )r)  r¶   rä   r&  )r—   rS   r*  r/   r/   r0   r  X  s   ÿûzAsynPool._on_hard_timeoutc                 C   s   |   |¡ d S r*   )rö   )r—   rS   r  ÚobjÚinqW_fdr/   r/   r0   Úon_job_readyd  r1   zAsynPool.on_job_readyc                    s²   ˆ	j ˆ	jˆ	j‰‰‰ˆj‰ˆj‰ˆj‰ˆj‰ˆj‰ˆj‰ˆj	‰ˆj
‰
ˆj‰‡‡	‡fdd„‰‡‡‡‡‡	‡‡‡fdd„}|ˆ_d
dd„‰ ‡ ‡‡‡‡‡‡	‡
‡‡‡‡fdd	„}|ˆ_dS )z/Create handlers called on process up/down, etc.c                    sv   | ƒ } | d ur5|   ¡ r7| ˆv r9| jˆ v sJ ‚ˆ | j | u sJ ‚| jˆjv s'J ‚td| ƒ t | jd¡ d S d S d S d S )Nz(Timed out waiting for UP message from %ré	   )Ú	_is_aliverú   rc   rr   r  Úkillr˜   ©rÔ   )rŸ   rÀ   Úwaiting_to_startr/   r0   Úverify_process_alivev  s   
úz?AsynPool._create_process_handlers.<locals>.verify_process_alivec                    sœ   | j }tˆƒD ]}|jr|jj |kr| |_|jr!|jj |kr!| |_q| ˆ| j< ˆ | ˆ¡ t| jjƒr5J ‚ˆ | jˆ| jƒ ˆ 	| ¡ ˆ 
ˆjˆt| ƒ¡ dS )z"Called when a process has started.N)r,  r   Ú	_write_toÚ_scheduled_forrú   r  r   r”   rÑ   rb   r!  rð   r   )rÔ   ÚinfdrS   )r¬   rÆ   rŸ   r  rÀ   r—   r3  r2  r/   r0   Úon_process_up€  s   €

ÿz8AsynPool._create_process_handlers.<locals>.on_process_upNc              	   S   st   z|   ¡ }W n ttfy   Y d S w z|| |u r | |d ¡ W n
 ty+   Y |S w ||ƒ |d ur8||ƒ |S r*   )ra   r¹   rƒ   r‰   r¶   )r+  rÔ   ÚindexÚ
remove_funr­   r5   r/   r/   r0   Ú_remove_from_indexž  s"   ÿ€úz=AsynPool._create_process_handlers.<locals>._remove_from_indexc                    sÞ   t | ddƒrdS ˆ| ƒ ˆ | jj| ˆˆƒ | jr!ˆ | jj| ˆˆ	ƒ ˆ | jj| ˆˆ	ˆjd�}|r4ˆ |¡ ˆ
 | ˆ¡ ˆ | ¡ ˆ
j | j	¡ ˆ	| jjƒ ˆ| jjƒ | j
r[ˆ| jjƒ | jrmˆ
j | j¡ ˆ| jjƒ dS dS )z#Called when a worker process exits.rÚ   N©r­   )rü   r”   rÑ   ÚsynqrQ   Úinqrw   r  ró   r,  ÚsynqR_fdrû   )rÔ   r=  )r:  Úall_inqueuesÚbusy_workersÚfileno_to_inqrŸ   Úfileno_to_synqrÀ   Úprocess_flush_queuesr¾   Úremove_writerr—   r2  r/   r0   Úon_process_down³  s6   ÿÿþ

þz:AsynPool._create_process_handlers.<locals>.on_process_downr*   )r¬   r¾   rD  r)  rò   rì   rí   rî   rõ   r  rC  rñ   r7  rE  )r—   rÀ   r7  rE  r/   )r:  r¬   r?  r@  rÆ   rA  rŸ   rB  r  rÀ   rC  r¾   rD  r—   r3  r2  r0   r  g  s"   
ÿ

"
z!AsynPool._create_process_handlersc                    sœ  ˆj ‰ˆj‰	ˆj‰ˆj‰ˆj‰ˆj‰ˆj‰ˆj}ˆj‰ˆj	‰ˆj
‰ˆjˆj‰‰ˆj‰|j‰ˆj‰|j‰ˆjj‰
ˆj‰ˆjtk‰tj‰tj‰tˆ td¡tˆ td¡i‰tjf‡‡‡fdd„	}|ˆ_‡‡‡‡‡‡‡‡fdd„}|ˆ_‡‡‡‡fdd„}|ˆ_ˆˆ_d‡‡‡‡‡‡‡‡‡‡‡‡‡‡fd	d
„	}	|	ˆ_‡‡
‡‡‡fdd„}
|
ˆ_ ‡‡fdd„‰‡‡‡‡‡fdd„‰‡ ‡‡‡‡‡fdd„}|ˆ_!d‡‡	fdd„	‰ dS )z6Create handlers used to write data to child processes.)r   c                    sX   | j d us
| jˆv r| js|  d |ƒ ˆ ƒ d ¡ |  | j ¡ d S | ˆvr*ˆ | ¡ d S d S r*   )Ú_terminatedÚcorrelation_idÚ	_acceptedÚ_ackÚ_set_terminatedÚ
appendleft)rS   Ú_time)ÚgetpidÚoutboundÚrevoked_tasksr/   r0   Ú	_put_backî  s   

ÿz2AsynPool._create_write_handlers.<locals>._put_backc                     sV   ˆˆ ƒ} ˆrˆot ˆƒt ˆƒk }nˆ}|r#t| ˆˆd ttB dd� d S t| ˆˆƒ d S )NT)Úconsolidate)r3   r’   r   r   )ÚinactiveÚadd_cond)Úactive_writesr?  r@  ÚdiffÚhub_addÚ
hub_removeÚis_fair_strategyrN  r/   r0   r    s   

þÿz6AsynPool._create_write_handlers.<locals>.on_poll_startc                    sX   ˆ  | ¡ zˆ|  |u rˆ | d ¡ ˆ   | ¡ ˆ  | ¡ W d S W d S  ty+   Y d S w r*   )rw   r‰   r¶   )r5   rÔ   )rT  r?  r@  rA  r/   r0   Úon_inqueue_close  s   

ýÿz9AsynPool._create_write_handlers.<locals>.on_inqueue_closeNc           
         sb  |sdg}t | ƒ}t|ƒD ]¡}| |d |  }|d  d7  < |ˆv r$qˆr+|ˆv r+q|ˆvr4ˆ|ƒ qzˆƒ }W n tyO   ˆˆƒD ]}ˆ|ƒ qDY  d S w |js®z	ˆ|  }|_W n tyi   ˆ|ƒ Y qw ˆ |||ƒ}t|ƒ|_ˆ|ƒ ˆ
|ƒ ˆ	|ƒ zt|ƒ W n! t	y�   Y q t
y¨ }	 z|	jtjkrž‚ W Y d }	~	qd }	~	ww ˆ||ƒ qd S )Nr   r?   )r3   rê   Ú
IndexErrorrH  r5  r¶   r   rQ   r·   r¸   rƒ   rt   ÚEBADF)
Ú	ready_fdsÚtotal_write_countÚ	num_readyrà   Úready_fdrS   ÚinqfdrÔ   Úcorrx   )Ú
_write_jobrT  Ú
add_writerr?  r@  rU  rA  rW  rX  Úmark_worker_as_busyÚmark_write_fd_as_activeÚmark_write_gen_as_activeÚpop_messageÚput_messager/   r0   Úschedule_writes)  sZ   

øû
ÿ€ÿ
€Íz8AsynPool._create_write_handlers.<locals>.schedule_writesc                    sN   ˆ | ˆd�}t |ƒ}ˆd|ƒ}ˆ| d d ƒ}t|ƒt|ƒ|f|_ˆ|ƒ d S )N©Úprotocolú>Ir?   r   )r3   r   Ú_payload)ÚtupÚbodyr³   ÚheaderrS   )ÚdumpsÚget_jobr   rk  rh  r/   r0   Úsend_jobr  s   
z1AsynPool._create_write_handlers.<locals>.send_jobc                    s:   t  d| | j|¡ |  ¡ r|  ¡  ˆ  |¡ ˆ |¡ d S )Nz"Process inqueue damaged: %r %r: %r)r…   Ú	exceptionÚexitcoder/  Ú	terminater‚   rP  )rÔ   r5   rS   rx   r  r/   r0   Úon_not_recovering~  s   
ÿ
z:AsynPool._create_write_handlers.<locals>.on_not_recoveringc              
   3   sÖ  � |j \}}}d}zÈ| |_| j}d }}	|dk rXz	||||ƒ7 }W n0 tyQ }
 z$t|
dd ƒtvr2‚ |d7 }|dkrDˆ| |||
ƒ tƒ ‚d V  W Y d }
~
nd }
~
ww d}|dk s|	|k r·z	|	|||	ƒ7 }	W n0 ty• }
 z$t|
dd ƒtvrv‚ |d7 }|dkrˆˆ| |||
ƒ tƒ ‚d V  W Y d }
~
nd }
~
ww d}|	|k s\W ˆ|ƒ ˆ| j  d7  < ˆ  |¡ ˆ| 	¡ ƒ d S W ˆ|ƒ ˆ| j  d7  < ˆ  |¡ ˆ| 	¡ ƒ d S ˆ|ƒ ˆ| j  d7  < ˆ  |¡ ˆ| 	¡ ƒ w )Nr   r@   rt   r?   éd   )
rm  r4  Úsend_job_offsetÚ	Exceptionrü   r¨   r¸   r8  rw   rQ   )rÔ   r5   rS   rp  ro  r³   ÚerrorsÚsendÚHwÚBwrx   )rT  rW  rw  Úwrite_generator_donerø   r/   r0   rb  †  sd   €€ø
ó€ø

ó
í
ü
z3AsynPool._create_write_handlers.<locals>._write_jobc                    sL   t ||ˆ|  ƒ}tˆƒ}ˆ |||d�}ˆ|ƒ ˆ|ƒ |f|_ˆ||ƒ d S )Nr;  )rI   r   ru   )Úresponser˜   rS   r5   Úmsgr­   ra  )Ú
_write_ackrc  re  rf  Úprecalcr  r/   r0   Úsend_ack¹  s   z1AsynPool._create_write_handlers.<locals>.send_ackc              
   3   s2  � |d \}}}z…zˆ|  }W n
 t y   tƒ ‚w |j}d }}	|dk rQz	||||ƒ7 }W n tyL }
 zt|
dd ƒtvr?‚ d V  W Y d }
~
nd }
~
ww |dk s%|	|k r�z	|	|||	ƒ7 }	W n ty| }
 zt|
dd ƒtvro‚ d V  W Y d }
~
nd }
~
ww |	|k sUW |r‡|ƒ  ˆ  | ¡ d S |r“|ƒ  ˆ  | ¡ w )Nr&   r   r@   rt   )r¶   r¸   Úsend_syn_offsetrz  rü   r¨   rw   )r5   Úackr­   rp  ro  r³   rÔ   r|  r}  r~  rx   )rT  rB  r/   r0   r‚  Å  sJ   €ý€ýý	€üý€	ýz3AsynPool._create_write_handlers.<locals>._write_ackr*   )"rì   rî   r÷   Úpopleftr‡   rò   ró   rô   rõ   Ú
differencerc  rb   r‚   rw   r)  Ú__getitem__rø   rç   ÚSCHED_STRATEGY_FAIRÚworker_stateÚrevokedr  rM  r   Ú_create_payloadr   ÚtimerP  r  rY  rW  Úconsolidate_callbackÚ
_quick_putr„  )r—   rÀ   r   rq  rk  Úactive_writersrP  r  rY  ri  rs  r„  r/   )r‚  rb  rT  rc  r?  r@  rU  rq  rA  rB  rr  rM  rÀ   rV  rW  rX  rd  re  rf  rw  rN  r   rg  rƒ  rk  rh  rO  r—   r  rø   r0   r  Ñ  sP   
ÿ(G
3
zAsynPool._create_write_handlersc              	   C   sö  | j tkrd S t| jƒD ]	}|js| ¡  q| jr| j ¡  |  ¡  zÃ| j t	kr¸t
ddddd�}i }t| jƒD ]}t|ƒ}|d urE|||< q7| jrÏt| jƒ}|D ]C}|jdkrvt|ƒrvz|| }W n	 tyj   Y nw | ¡  | j |¡ qPz|| }W n	 ty…   Y qPw |j}| ¡ r“|  ||¡ qP|  ¡  tt|ƒƒ | jsIW | j ¡  | j ¡  | j ¡  | j ¡  d S W | j ¡  | j ¡  | j ¡  | j ¡  d S W | j ¡  | j ¡  | j ¡  | j ¡  d S | j ¡  | j ¡  | j ¡  | j ¡  w )Nç{®Gáz„?gš™™™™™¹?T)Ú
repeatlastrb  )rÉ   r   r   r)  rH  Ú_cancelr÷   Úclearr  r   r   rU   rô   rn   rš   rP   r¶   rw   r4  r/  Ú_flush_writerr	   r·   ró   rõ   )r—   rS   Ú	intervalsÚowned_byrT   rd   rO   Újob_procr/   r/   r0   Úflushì  sz   
€

€

ÿÿÿ€å


×
&

à



ý

zAsynPool.flushc                 C   s¼   |j jh}zQ|r<| ¡ sn8t||dd�\}}}|s1|s|r1zt|ƒ W n ttttfy0   Y nw |sW | j	 
|¡ d S W | j	 
|¡ d S W | j	 
|¡ d S W | j	 
|¡ d S | j	 
|¡ w )NrÐ   )rd   re   rf   )r=  rQ   r/  rz   r·   r¸   rƒ   r¹   r©   rô   rw   )r—   rÔ   rT   ÚfdsÚreadableÚwritableÚagainr/   r/   r0   r–  *  s,   
ÿÿ÷ôö
þzAsynPool._flush_writerc                 C   s   t dd„ t| jƒD ƒƒS )zœGet queues for a new process.

        Here we'll find an unused slot, as there should always
        be one available when we start a new process.
        c                 s   s    � | ]\}}|d u r|V  qd S r*   r/   ©rD   ÚqÚownerr/   r/   r0   Ú	<genexpr>A  ó   €
 ÿÿz.AsynPool.get_process_queues.<locals>.<genexpr>)r·   r   rë   rá   r/   r/   r0   Úget_process_queues;  s   zAsynPool.get_process_queuesc                    s@   t ˆ jtˆ jƒ dƒ}|rˆ j ‡ fdd„t|ƒD ƒ¡ dS dS )z Grow the pool by ``n`` proceses.r   c                    rÜ   r*   rÝ   rß   rá   r/   r0   rG   H  râ   z$AsynPool.on_grow.<locals>.<dictcomp>N)ÚmaxÚ
_processesr3   rë   Úupdaterê   )r—   r9   rU  r/   rá   r0   Úon_growD  s   ÿÿzAsynPool.on_growc                 C   s   dS )z#Shrink the pool by ``n`` processes.Nr/   )r—   r9   r/   r/   r0   Ú	on_shrinkL  s    zAsynPool.on_shrinkc                 C   s†   t dd�}t dd�}d}t|jƒsJ ‚t|jƒrJ ‚t|jƒr!J ‚t|jƒs(J ‚| jr>t dd�}t|jƒs7J ‚t|jƒr>J ‚|||fS )z5Create new in, out, etc. queues, returned as a tuple.T)Ú	wnonblock)Ú	rnonblockN)r   r   rÑ   rQ   ré   )r—   r=  r”   r<  r/   r/   r0   rÞ   O  s   



zAsynPool.create_process_queuesc                    s’   zt ‡ fdd„| jD ƒƒ}W n ty   t dˆ ¡ Y S w |j| jvs&J ‚|j| jvs.J ‚| j 	|¡ || j|j< || j
|j< | j |j¡ dS )zsCalled when receiving the :const:`WORKER_UP` message.

        Marks the process as ready to receive work.
        c                 3   s   � | ]
}|j ˆ kr|V  qd S r*   ©r˜   r  r¬  r/   r0   r¢  g  s   € z,AsynPool.on_process_alive.<locals>.<genexpr>z"process with pid=%s already exitedN)r·   rù   r¸   r…   r†   r,  rì   rò   rñ   rw   rî   rû   rb   )r—   r˜   rÔ   r/   r¬  r0   r    a  s   ÿzAsynPool.on_process_alivec                 C   sH   |j r|j  ¡ s|  ||j ¡ dS |jr |j ¡ s"|  |¡ dS dS dS )z:Called for each job when the process assigned to it exits.N)r4  r/  Úon_partial_readr5  rP  )r—   rS   Úpid_goner/   r/   r0   Úon_job_process_downq  s
   ýzAsynPool.on_job_process_downc                 C   s   |   ||¡ dS )z¼Called when the process executing job' exits.

        This happens when the process job'
        was assigned to exited by mysterious means (error exitcodes and
        signals).
        N)Úmark_as_worker_lost)r—   rS   r˜   ru  r/   r/   r0   Úon_job_process_lost{  s   zAsynPool.on_job_process_lostc                    s–   | j d u rdS tt| j ƒƒ}t|ƒ‰dd„ ‰ ˆˆ ˆr!ˆt| j ƒ ndˆƒd ‡ ‡fdd„|D ƒ¡d tt|ƒ¡t 	| j
| j
¡t| jƒt| jƒdœd	œS )
NzN/Ac                 S   s   d  | rt| ƒ| ¡S d¡S )Nz{0:.2%}r   )ÚformatÚfloat)rF   Útotalr/   r/   r0   ÚperŠ  s   z'AsynPool.human_write_stats.<locals>.perr   z, c                 3   s   � | ]}ˆ |ˆƒV  qd S r*   r/   )rD   rF   ©rµ  r´  r/   r0   r¢  �  s   € z-AsynPool.human_write_stats.<locals>.<genexpr>)r´  Úactive)r´  ÚavgÚallÚrawÚstrategyÚinqueues)rø   rn   r   Úsumr3   ÚjoinÚmapÚstrÚSCHED_STRATEGY_TO_NAMEræ   rç   rò   ró   )r—   Úvalsr/   r¶  r0   Úhuman_write_stats„  s    
ÿþøzAsynPool.human_write_statsc              	   C   s:   |j szd| j|  |¡< W dS  ttfy   Y dS w dS )z-Called to clean up queues after process exit.N)rÚ   rë   Ú_find_worker_queuesr¶   rŠ   ©r—   rÔ   r/   r/   r0   Ú_process_cleanup_queues›  s   ÿýz AsynPool._process_cleanup_queuesc                 C   s|   | j D ]8}z	t|jjdƒ W n ttfy   Y qw z|j d¡ W q ty; } z|jtjkr1‚ W Y d}~qd}~ww dS )z>Called at shutdown to tell processes that we're shutting down.r?   N)	r   r   r=  rQ   rƒ   r¹   r•   rt   r[  )Útask_handlerrÔ   rx   r/   r/   r0   Ú_stop_task_handler£  s   
ÿÿ€ÿøzAsynPool._stop_task_handlerc                    s   t t| ƒj| j| jd�S )N)rŸ   r    )r¡   r=   Úcreate_result_handlerrí   r    rá   r¤   r/   r0   rÉ  ²  s   
þzAsynPool.create_result_handlerc                 C   s8   || j v sJ ‚t| j ƒ}|| j |< |t| j ƒksJ ‚dS )z;Mark new ownership for ``queues`` to update fileno indices.N)rë   r3   )r—   rÔ   ÚqueuesÚbr/   r/   r0   Ú_process_register_queues¸  s   

z!AsynPool._process_register_queuesc                    s6   zt ‡ fdd„t| jƒD ƒƒW S  ty   tˆ ƒ‚w )z"Find the queues owned by ``proc``.c                 3   s    � | ]\}}|ˆ kr|V  qd S r*   r/   rŸ  r1  r/   r0   r¢  Â  r£  z/AsynPool._find_worker_queues.<locals>.<genexpr>)r·   r   rë   r¸   rŠ   rÅ  r/   r1  r0   rÄ  ¿  s
   ÿzAsynPool._find_worker_queuesc                 C   s"   d | _ d  | _ | _ | _| _d S r*   )r�  Ú_inqueueÚ	_outqueueÚ
_quick_getÚ_poll_resultrá   r/   r/   r0   Ú_setup_queuesÇ  s   ÿzAsynPool._setup_queuesc           
   
   C   s   |j j}| jj}|h}|r†|jsˆ| jtkrŠt|d|dd�\}}}|rxz| ¡ }W n? t	t
tfyg } z0t|ddƒ}	|	tjkrDW Y d}~q|	tjkrPW Y d}~dS |	tvr\td||dd� W Y d}~dS d}~ww |du rstd|ƒ dS ||ƒ ndS |rŒ|jsŽ| jtksdS dS dS dS dS dS )	a  Flush all queues.

        Including the outbound buffer, so that
        all tasks that haven't been started will be discarded.

        In Celery this is called whenever the transport connection is lost
        (consumer restart), and when a process is terminated.
        Nr’  ©rf   rt   z got %r while flushing process %rr?   r€   z&got sentinel while flushing process %r)r”   rÑ   r  r¼   ÚclosedrÉ   r   rz   rÒ   rƒ   r¹   r©   rü   rt   rv   ÚEAGAINr¨   rË   )
r—   rÔ   Úresqr¼   r›  rœ  rà   rÖ   rx   ry   r/   r/   r0   rC  Ñ  s6   	

ÿ€÷

,êzAsynPool.process_flush_queuesc                 C   s–   |j s|  |¡ t|ƒ}|r| j |¡ ~|jsGd|_t| jƒ}z|  |¡}|  	||¡r3d| j|  
¡ < W n	 ty=   Y nw t| jƒ|ksIJ ‚dS dS )z8Called when a job was partially written to exited child.TN)rH  rP  rU   rô   rw   rÚ   r3   rë   rÄ  Údestroy_queuesrÞ   rŠ   )r—   rS   rÔ   rT   ÚbeforerÊ  r/   r/   r0   r­  õ  s(   


€ÿö
zAsynPool.on_partial_readc                 C   sÊ   |  ¡ rJ ‚| j |¡ d}z| j |¡ W n ty!   d}Y nw z|  |d j ¡ |¡ W n	 t	y8   Y nw |D ]'}|rb|j
|jfD ]}|jsa|  |¡ z| ¡  W qE t	tfy`   Y qEw qEq;|S )zqDestroy queues that can no longer be used.

        This way they can be replaced by new usable sockets.
        r?   r   )r/  rñ   rw   rë   r‰   r¶   rY  rQ   ra   r¹   rÑ   rÓ  rW  r  rƒ   )r—   rÊ  rÔ   ÚremovedÚqueueÚsockr/   r/   r0   rÖ    s4   ÿÿ
ÿü€zAsynPool.destroy_queuesc           	      C   s,   |||f|d�}t |ƒ}|d|ƒ}|||fS )Nrj  rl  )r3   )	r—   Útype_ru   rq  r   rk  ro  r7   rp  r/   r/   r0   r�  *  s   

zAsynPool._create_payloadc                 C   s   d S r*   r/   )ÚclsrÎ  rù   r/   r/   r0   Ú_set_result_sentinel2  s   zAsynPool._set_result_sentinelc                 C   s   | j fS r*   )rù   rá   r/   r/   r0   Ú_help_stuff_finish_args7  rÄ   z AsynPool._help_stuff_finish_argsc           	   	   C   s¢   t dƒ i }tƒ }|D ]}z|jj ¡ }| |¡ |||< W q ty'   Y qw |rOt|dd�\}}}|r6q(|s:d S |D ]
}|| jj ¡  q<t	dƒ |s*d S d S )Nz7removing tasks from inqueue until task handler finishedrÐ   rÒ  r   )
rË   r^   r=  rÑ   ra   rb   r¹   rz   rÒ   r	   )	rÜ  r   Úfileno_to_procÚinqRrp   r5   rœ  rà   rž  r/   r/   r0   Ú_help_stuff_finish<  s.   ÿ
ÿøzAsynPool._help_stuff_finishc                 C   s
   | j diS )Ng      @)r  rá   r/   r/   r0   r  U  s   
zAsynPool.timers)NFNN)3rš   r›   rœ   r�   rž   r“   rÙ   r¢   r  r  r  r  rÂ   r  r  r  r-  r  r   r×   rq  r   r  rš  r–  r¤  r¨  r©  rÞ   r    r¯  r±  rÃ  rÆ  ÚstaticmethodrÈ  rÉ  rÌ  rÄ  rÑ  rC  r­  rÖ  r�  ÚclassmethodrÝ  rÞ  rá  Úpropertyr  rØ   r/   r/   r¤   r0   r=   ¬  sj    ÿ>k
þ  >	
	

$
þ

r=   )NNNr   )hr�   Ú
__future__r   r   rt   r   r  rm   rs   ÚsysrŽ  Úcollectionsr   r   Úior   Únumbersr   r   r   r	   Úweakrefr
   r   Úbilliardr   rù   Úbilliard.compatr   r   r   Úbilliard.poolr   r   r   r   r   Úbilliard.queuesr   Úkombu.asynchronousr   r   Úkombu.serializationr×   Úkombu.utils.eventior   Úkombu.utils.functionalr   Úviner   Úcelery.fiver   r   r   Úcelery.platformsr   r    r!   Úcelery.utils.functionalr"   Úcelery.utils.logr#   Úcelery.workerr$   r‹  Ú	_billiardr%   r:   r®   Úversion_infoÚImportErrorÚ__all__rš   r…   rr   rË   Ú	frozensetrÔ  rv   r¨   r–   rï   ÚSCHED_STRATEGY_FCFSrŠ  rå   rÁ  rI   rP   rU   rˆ   rV   rY   r\   r]   rl   rz   r„   Ú	NameErrorr¹   r’   r“   rž   ÚPoolr=   r/   r/   r/   r0   Ú<module>   s˜   €öü
	þ

ÿ9ÿ-
  