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mZ ddl	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ZedƒZdddddejdefdd„ZG dd„ deƒZdS )zBase Execution Pool.é    )Úabsolute_importÚunicode_literalsN)ÚExceptionInfo)ÚWorkerLostError)Ú	safe_repr)ÚWorkerShutdownÚWorkerTerminate)Ú	monotonicÚreraise)Útimer2)Ú
get_logger)Útruncate)ÚBasePoolÚapply_targetzcelery.pool© c	                 K   sâ   |si n|}|r||p|ƒ |ƒ ƒ z	| |i |¤Ž}
W nP |y"   ‚  t y)   ‚  ttfy2   ‚  tyj } z-ztttt|ƒƒt ¡ d ƒ W n tyW   |t	ƒ ƒ Y nw W Y d}~dS W Y d}~dS d}~ww ||
ƒ dS )z#Apply function within pool context.é   N)
Ú	Exceptionr   r   ÚBaseExceptionr
   r   ÚreprÚsysÚexc_infor   )ÚtargetÚargsÚkwargsÚcallbackÚaccept_callbackÚpidÚgetpidÚ	propagater	   Ú_ÚretÚexcr   r   úT/var/www/html/myproject/venv/lib/python3.10/site-packages/celery/concurrency/base.pyr      s0   
ÿÿþ€ûr   c                   @   s  e Zd ZdZdZdZdZejZdZ	dZ
dZdZdZdZdZdZ		d7d	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d8dd „Zd!d"„ Zd#d$„ Zd%d&„ Zd'd(„ Z d)d*„ Z!d+d,„ Z"d9d-d.„Z#d/d0„ Z$e%d1d2„ ƒZ&e%d3d4„ ƒZ'e%d5d6„ ƒZ(dS ):r   z
Task pool.é   r   é   TFNr   c                 K   s(   || _ || _|| _|| _|| _|| _d S ©N)ÚlimitÚputlocksÚoptionsÚforking_enableÚcallbacks_propagateÚapp)Úselfr&   r'   r)   r*   r+   r(   r   r   r"   Ú__init__K   s   
zBasePool.__init__c                 C   ó   d S r%   r   ©r,   r   r   r"   Úon_startT   ó   zBasePool.on_startc                 C   s   dS )NTr   r/   r   r   r"   Údid_start_okW   r1   zBasePool.did_start_okc                 C   r.   r%   r   r/   r   r   r"   ÚflushZ   r1   zBasePool.flushc                 C   r.   r%   r   r/   r   r   r"   Úon_stop]   r1   zBasePool.on_stopc                 C   r.   r%   r   )r,   Úloopr   r   r"   Úregister_with_event_loop`   r1   z!BasePool.register_with_event_loopc                 O   r.   r%   r   ©r,   r   r   r   r   r"   Úon_applyc   r1   zBasePool.on_applyc                 C   r.   r%   r   r/   r   r   r"   Úon_terminatef   r1   zBasePool.on_terminatec                 C   r.   r%   r   ©r,   Újobr   r   r"   Úon_soft_timeouti   r1   zBasePool.on_soft_timeoutc                 C   r.   r%   r   r:   r   r   r"   Úon_hard_timeoutl   r1   zBasePool.on_hard_timeoutc                 O   r.   r%   r   r7   r   r   r"   Úmaintain_poolo   r1   zBasePool.maintain_poolc                 C   ó   t d t| ƒ¡ƒ‚)Nz{0} does not implement kill_job©ÚNotImplementedErrorÚformatÚtype)r,   r   Úsignalr   r   r"   Úterminate_jobr   ó   ÿzBasePool.terminate_jobc                 C   r?   )Nz{0} does not implement restartr@   r/   r   r   r"   Úrestartv   rF   zBasePool.restartc                 C   s   |   ¡  | j| _d S r%   )r4   Ú	TERMINATEÚ_stater/   r   r   r"   Ústopz   ó   zBasePool.stopc                 C   ó   | j | _|  ¡  d S r%   )rH   rI   r9   r/   r   r   r"   Ú	terminate~   rK   zBasePool.terminatec                 C   s"   t  tj¡| _|  ¡  | j| _d S r%   )ÚloggerÚisEnabledForÚloggingÚDEBUGÚ_does_debugr0   ÚRUNrI   r/   r   r   r"   Ústart‚   s   zBasePool.startc                 C   rL   r%   )ÚCLOSErI   Úon_closer/   r   r   r"   Úclose‡   rK   zBasePool.closec                 C   r.   r%   r   r/   r   r   r"   rV   ‹   r1   zBasePool.on_closec                 K   sb   |si n|}|s
g n|}| j r!t d|tt|ƒdƒtt|ƒdƒ¡ | j|||f| j| jdœ|¤ŽS )zÈEquivalent of the :func:`apply` built-in function.

        Callbacks should optimally return as soon as possible since
        otherwise the thread which handles the result will get blocked.
        z&TaskPool: Apply %s (args:%s kwargs:%s)i   )Úwaitforslotr*   )rR   rN   Údebugr   r   r8   r'   r*   )r,   r   r   r   r(   r   r   r"   Úapply_asyncŽ   s   þþýzBasePool.apply_asyncc                 C   s
   d| j iS )Nzmax-concurrency©r&   r/   r   r   r"   Ú	_get_info    s   ÿzBasePool._get_infoc                 C   s   |   ¡ S r%   )r\   r/   r   r   r"   Úinfo¥   s   zBasePool.infoc                 C   s   | j | jkS r%   )rI   rS   r/   r   r   r"   Úactive©   s   zBasePool.activec                 C   s   | j S r%   r[   r/   r   r   r"   Únum_processes­   s   zBasePool.num_processes)NTTr   Nr%   )NN))Ú__name__Ú
__module__Ú__qualname__Ú__doc__rS   rU   rH   r   ÚTimerÚsignal_safeÚis_greenrI   Ú_poolrR   Úuses_semaphoreÚtask_join_will_blockÚbody_can_be_bufferr-   r0   r2   r3   r4   r6   r8   r9   r<   r=   r>   rE   rG   rJ   rM   rT   rW   rV   rZ   r\   Úpropertyr]   r^   r_   r   r   r   r"   r   1   sT    
ÿ	



r   )rc   Ú
__future__r   r   rP   Úosr   Úbilliard.einfor   Úbilliard.exceptionsr   Úkombu.utils.encodingr   Úcelery.exceptionsr   r   Úcelery.fiver	   r
   Úcelery.utilsr   Úcelery.utils.logr   Úcelery.utils.textr   Ú__all__rN   r   r   Úobjectr   r   r   r   r"   Ú<module>   s(   
þ