o
    ›¨ÊhÚ  ã                   @   sú   d Z ddlZddlmZ ddlmZmZ ddlmZm	Z	 ddlm
Z ddlmZ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lmZ ddlmZ dZh d£Zer]dehZndhZeeƒZ e j!e j"Z!Z"dd„ Z#dd„ Z$G dd„ deƒZ%dS )zKPrefork execution pool.

Pool implementation using :mod:`multiprocessing`.
é    N)Úforking_enable)ÚREMAP_SIGTERMÚTERM_SIGNAME)ÚCLOSEÚRUN)ÚPool)Ú	platformsÚsignals)Ú_set_task_join_will_blockÚset_default_app)Útrace)ÚBasePool)Únoop)Ú
get_loggeré   )ÚAsynPool)ÚTaskPoolÚprocess_initializerÚprocess_destructor>   ÚSIGHUPÚSIGTERMÚSIGTTINÚSIGTTOUÚSIGUSR1ÚSIGINTc                 C   sL  t  d¡ tdƒ t jjtŽ  t jjtŽ  t jd|d� | j	 
¡  | j	 ¡  tj d¡p-d}|r:d| ¡ v r:d| j_| jjttj d	d
¡pFd
ƒ|ttj dd¡ƒttj d¡ƒ|d� tj d¡rht | |¡ n|  ¡  t| ƒ |  ¡  | jt_d
dlm} | j ¡ D ]\}}|||| j	|| d�|_ qƒd
dl!m"} | #¡  tj$j%dd� dS )z™Pool child process initializer.

    Initialize the child pool process to ensure the correct
    app instance is used and things like logging works.
    ÚSIGKILLTÚceleryd)ÚhostnameÚCELERY_LOG_FILENz%iFÚCELERY_LOG_LEVELr   ÚCELERY_LOG_REDIRECTÚCELERY_LOG_REDIRECT_LEVELÚFORKED_BY_MULTIPROCESSING)Úbuild_tracer)Úapp)Ústate)Úsender)&r   Úset_pdeathsigr
   r	   ÚresetÚWORKER_SIGRESETÚignoreÚWORKER_SIGIGNOREÚset_mp_process_titleÚloaderÚinit_workerÚinit_worker_processÚosÚenvironÚgetÚlowerÚlogÚalready_setupÚsetupÚintÚboolÚstrr   Úsetup_worker_optimizationsÚset_currentr   ÚfinalizeÚ_tasksÚcelery.app.tracer#   ÚtasksÚitemsÚ	__trace__Úcelery.workerr%   Úreset_stateÚworker_process_initÚsend)r$   r   Úlogfiler#   ÚnameÚtaskÚworker_state© rJ   úL/var/www/html/env/lib/python3.10/site-packages/celery/concurrency/prefork.pyr   &   s<   


ü
ÿr   c                 C   s   t jjd| |d� dS )z_Pool child process destructor.

    Dispatch the :signal:`worker_process_shutdown` signal.
    N)r&   ÚpidÚexitcode)r	   Úworker_process_shutdownrE   )rL   rM   rJ   rJ   rK   r   R   s   
ÿr   c                       st   e Zd ZdZeZeZdZdZdd„ Z	dd„ Z
dd	„ Zd
d„ Zdd„ Zdd„ Zdd„ Z‡ fdd„Zedd„ ƒZ‡  ZS )r   z$Multiprocessing Pool implementation.TNc              	   C   s˜   t | j ƒ | j dd¡r| jn| j}| jr| jjjnd }|d| jt	t
dd|dœ| j¤Ž }| _|j| _|j| _|j| _|j| _|j| _t|dd ƒ| _d S )NÚthreadsTF)Ú	processesÚinitializerÚon_process_exitÚenable_timeoutsÚsynackÚproc_alive_timeoutÚflushrJ   )r   Úoptionsr2   ÚBlockingPoolr   r$   ÚconfÚworker_proc_alive_timeoutÚlimitr   r   Ú_poolÚapply_asyncÚon_applyÚmaintain_poolÚterminate_jobÚgrowÚshrinkÚgetattrrV   )Úselfr   rU   ÚPrJ   rJ   rK   Úon_starte   s,   
ÿþûú	zTaskPool.on_startc                 C   s   | j  ¡  | j  t¡ d S ©N)r\   Úrestartr]   r   ©rd   rJ   rJ   rK   rh   }   s   
zTaskPool.restartc                 C   s
   | j  ¡ S rg   )r\   Údid_start_okri   rJ   rJ   rK   rj   �   s   
zTaskPool.did_start_okc                 C   s(   z	| j j}W ||ƒS  ty   Y d S w rg   )r\   Úregister_with_event_loopÚAttributeError)rd   ÚloopÚregrJ   rJ   rK   rk   „   s   
þÿz!TaskPool.register_with_event_loopc                 C   s@   | j dur| j jttfv r| j  ¡  | j  ¡  d| _ dS dS dS )zGracefully stop the pool.N)r\   Ú_stater   r   ÚcloseÚjoinri   rJ   rJ   rK   Úon_stop‹   s
   


ýzTaskPool.on_stopc                 C   s"   | j dur| j  ¡  d| _ dS dS )zForce terminate the pool.N)r\   Ú	terminateri   rJ   rJ   rK   Úon_terminate’   s   


þzTaskPool.on_terminatec                 C   s,   | j d ur| j jtkr| j  ¡  d S d S d S rg   )r\   ro   r   rp   ri   rJ   rJ   rK   Úon_close˜   s   ÿzTaskPool.on_closec              	      sp   t | jdd ƒ}tƒ  ¡ }| | jdd„ | jjD ƒ| jjpd| j| jjp$d| jj	p)df|d ur1|ƒ nddœ¡ |S )NÚhuman_write_statsc                 S   s   g | ]}|j ‘qS rJ   )rL   )Ú.0ÚprJ   rJ   rK   Ú
<listcomp>¡   s    z&TaskPool._get_info.<locals>.<listcomp>zN/Ar   )zmax-concurrencyrP   zmax-tasks-per-childzput-guarded-by-semaphoreÚtimeoutsÚwrites)
rc   r\   ÚsuperÚ	_get_infoÚupdater[   Ú_maxtasksperchildÚputlocksÚsoft_timeoutÚtimeout)rd   Úwrite_statsÚinfo©Ú	__class__rJ   rK   r}   œ   s   



ÿù	zTaskPool._get_infoc                 C   s   | j jS rg   )r\   Ú
_processesri   rJ   rJ   rK   Únum_processesª   s   zTaskPool.num_processes)Ú__name__Ú
__module__Ú__qualname__Ú__doc__r   r   rX   Úuses_semaphorerƒ   rf   rh   rj   rk   rr   rt   ru   r}   Úpropertyrˆ   Ú__classcell__rJ   rJ   r…   rK   r   \   s     r   )&rŒ   r0   Úbilliardr   Úbilliard.commonr   r   Úbilliard.poolr   r   r   rX   Úceleryr   r	   Úcelery._stater
   r   Ú
celery.appr   Úcelery.concurrency.baser   Úcelery.utils.functionalr   Úcelery.utils.logr   Úasynpoolr   Ú__all__r)   r+   r‰   ÚloggerÚwarningÚdebugr   r   r   rJ   rJ   rJ   rK   Ú<module>   s.    
,
