o
    wvXjŸ  ã                   @   s  d Z ddlmZm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dlm Z  dZ!h d£Z"erkd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`.
é    )Úabsolute_importÚunicode_literalsN)Úforking_enable)ÚREMAP_SIGTERMÚTERM_SIGNAME)ÚCLOSEÚRUN)ÚPool)Ú	platformsÚsignals)Ú_set_task_join_will_blockÚset_default_app)Útrace)ÚBasePool)Úitems)Únoop)Ú
get_loggeré   )ÚAsynPool)ÚTaskPoolÚprocess_initializerÚprocess_destructor>   ÚSIGHUPÚSIGTERMÚSIGTTINÚSIGTTOUÚSIGUSR1ÚSIGINTc                 C   sB  t dƒ tjjtŽ  tjjtŽ  tjd|d� | j 	¡  | j 
¡  tj d¡p(d}|r5d| ¡ v r5d| j_| jjttj dd	¡pAd	ƒ|ttj d
d¡ƒttj d¡ƒ|d� tj d¡rct | |¡ n|  ¡  t| ƒ |  ¡  | jt_d	dlm} t| 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.
    TÚ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   r
   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%   r   ÚtasksÚ	__trace__Úcelery.workerr'   Úreset_stateÚworker_process_initÚsend)r&   r   Úlogfiler%   ÚnameÚtaskÚworker_state© rJ   úW/var/www/html/myproject/venv/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   T   s   
ÿr   c                   @   sl   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dd„ Zedd„ ƒZd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   Úoptionsr3   Ú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_startg   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                 C   s`   t | jdd ƒ}| jdd„ | jjD ƒ| jjpd| j| jjpd| jjp"df|d ur,|ƒ dœS 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\   r[   Ú_maxtasksperchildÚputlocksÚsoft_timeoutÚtimeout)rd   Úwrite_statsrJ   rJ   rK   Ú	_get_infož   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ƒ   rJ   rJ   rJ   rK   r   ^   s     r   )+r‡   Ú
__future__r   r   r1   Úbilliardr   Úbilliard.commonr   r   Úbilliard.poolr   r   r	   rX   Úceleryr
   r   Úcelery._stater   r   Ú
celery.appr   Úcelery.concurrency.baser   Úcelery.fiver   Úcelery.utils.functionalr   Úcelery.utils.logr   Úasynpoolr   Ú__all__r*   r,   r„   ÚloggerÚwarningÚdebugr   r   r   rJ   rJ   rJ   rK   Ú<module>   s2   
*
