o
    ›¨Êhb  ã                   @   sÌ   d 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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jfdd„ZG dd„ dƒZdS )zBase Execution Pool.é    N)ÚAnyÚDict)ÚExceptionInfo)ÚWorkerLostError)Ú	safe_repr)ÚWorkerShutdownÚWorkerTerminateÚ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Ú	propagateÚ	monotonicÚ_ÚretÚexcr   r   úI/var/www/html/env/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		d8d	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d9d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/e$e%e&f fd0d1„Z'e(d2d3„ ƒZ)e(d4d5„ ƒZ*e(d6d7„ ƒ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__I   s   
zBasePool.__init__c                 C   ó   d S r%   r   ©r,   r   r   r"   Úon_startR   ó   zBasePool.on_startc                 C   s   dS )NTr   r/   r   r   r"   Údid_start_okU   r1   zBasePool.did_start_okc                 C   r.   r%   r   r/   r   r   r"   ÚflushX   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_applya   r1   zBasePool.on_applyc                 C   r.   r%   r   r/   r   r   r"   Úon_terminated   r1   zBasePool.on_terminatec                 C   r.   r%   r   ©r,   Újobr   r   r"   Úon_soft_timeoutg   r1   zBasePool.on_soft_timeoutc                 C   r.   r%   r   r:   r   r   r"   Úon_hard_timeoutj   r1   zBasePool.on_hard_timeoutc                 O   r.   r%   r   r7   r   r   r"   Úmaintain_poolm   r1   zBasePool.maintain_poolc                 C   ó   t t| ƒ› d�ƒ‚)Nz does not implement kill_job©ÚNotImplementedErrorÚtype)r,   r   Úsignalr   r   r"   Úterminate_jobp   ó   ÿzBasePool.terminate_jobc                 C   r?   )Nz does not implement restartr@   r/   r   r   r"   Úrestartt   rE   zBasePool.restartc                 C   s   |   ¡  | j| _d S r%   )r4   Ú	TERMINATEÚ_stater/   r   r   r"   Ústopx   ó   zBasePool.stopc                 C   ó   | j | _|  ¡  d S r%   )rG   rH   r9   r/   r   r   r"   Ú	terminate|   rJ   zBasePool.terminatec                 C   s"   t  tj¡| _|  ¡  | j| _d S r%   )ÚloggerÚisEnabledForÚloggingÚDEBUGÚ_does_debugr0   ÚRUNrH   r/   r   r   r"   Ústart€   s   zBasePool.startc                 C   rK   r%   )ÚCLOSErH   Úon_closer/   r   r   r"   Úclose…   rJ   zBasePool.closec                 C   r.   r%   r   r/   r   r   r"   rU   ‰   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*   )rQ   rM   Údebugr   r   r8   r'   r*   )r,   r   r   r   r(   r   r   r"   Úapply_asyncŒ   s   þþýzBasePool.apply_asyncÚreturnc                 C   s   | j jd | j j | jdœS )z¶
        Return configuration and statistics information. Subclasses should
        augment the data as required.

        :return: The returned value must be JSON-friendly.
        ú:)Úimplementationzmax-concurrency)Ú	__class__Ú
__module__Ú__name__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%   )rH   rR   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),r_   r^   Ú__qualname__Ú__doc__rR   rT   rG   r
   ÚTimerÚsignal_safeÚis_greenrH   Ú_poolrQ   Úuses_semaphoreÚtask_join_will_blockÚbody_can_be_bufferr-   r0   r2   r3   r4   r6   r8   r9   r<   r=   r>   rD   rF   rI   rL   rS   rV   rU   rY   r   Ústrr   r`   Úpropertyra   rb   rc   r   r   r   r"   r   /   sT    
ÿ	



r   )re   rO   Úosr   ÚtimeÚtypingr   r   Úbilliard.einfor   Úbilliard.exceptionsr   Úkombu.utils.encodingr   Úcelery.exceptionsr   r   r	   Úcelery.utilsr
   Úcelery.utils.logr   Úcelery.utils.textr   Ú__all__rM   r   r   r   r   r   r   r   r"   Ú<module>   s(    
þ