o
    ›¨Êh‡!  ã                   @   sB  d 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 ddl	m
Z
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	Zd
ee ¡ dœZeej dd¡ƒZeej dd¡ƒZeej dd¡ƒZeej dd¡ƒZi Z e !¡ Z"e !¡ Z#eeed�Z$eƒ Z%dgZ&eeed�Z'i Z(dZ)dZ*dd„ Z+dd„ Z,e j-e"j.fdd„Z/de j-e#j.e%j0fdd„Z1de j2e#j3e"j3fdd„Z4ej d¡pÆej d ¡Z5eej d!¡pÕej d"¡pÕdƒZ6e5�rddl7Z7dd#l8m9Z9 dd$l:m;Z; dd%l<m=Z=m>Z> da?da@daAdaBe6ZCg ZDe/ZEe4ZFe;ƒ jGd&k�re7jHd'd(„ ƒZId)d„ Z/d*d„ Z4G d+d,„ d,ƒZJdS )-zwInternal worker state (global).

This includes the currently active and reserved tasks,
statistics, and revoked tasks.
é    N)ÚCounter)ÚpickleÚpickle_protocol)Úcached_property)Ú__version__)ÚWorkerShutdownÚWorkerTerminate)Ú
LimitedSet)
ÚSOFTWARE_INFOÚreserved_requestsÚactive_requestsÚtotal_countÚrevokedÚtask_reservedÚmaybe_shutdownÚtask_acceptedÚ
task_readyÚ
Persistentz	py-celery)Úsw_identÚsw_verÚsw_sysÚCELERY_WORKER_REVOKES_MAXiPÃ  ÚCELERY_WORKER_SUCCESSFUL_MAXiè  ÚCELERY_WORKER_REVOKE_EXPIRESi0*  Ú CELERY_WORKER_SUCCESSFUL_EXPIRES)ÚmaxlenÚexpiresc                   C   sJ   t  ¡  t ¡  t ¡  t ¡  t ¡  dgtd d …< t ¡  t ¡  d S )Nr   )	ÚrequestsÚclearr   r   Úsuccessful_requestsr   Úall_total_countr   Úrevoked_stamps© r"   r"   úE/var/www/html/env/lib/python3.10/site-packages/celery/worker/state.pyÚreset_stateM   s   r$   c                   C   s8   t durt durtt ƒ‚tdurtdurttƒ‚dS dS )z Shutdown if flags have been set.NF)Úshould_terminater   Úshould_stopr   r"   r"   r"   r#   r   X   s
   ÿr   c                 C   s   || j | ƒ || ƒ dS )z2Update global state when a task has been reserved.N)Úid)ÚrequestÚadd_requestÚadd_reserved_requestr"   r"   r#   r   `   s   r   c                 C   s>   |st }|| j| ƒ || ƒ || jdiƒ t d  d7  < dS )z2Update global state when a task has been accepted.é   r   N)r    r'   Úname)r(   Ú_all_total_countr)   Úadd_active_requestÚadd_to_total_countr"   r"   r#   r   h   s   r   Fc                 C   s0   |rt  | j¡ || jdƒ || ƒ || ƒ dS )z)Update global state when a task is ready.N)r   Úaddr'   )r(   Ú
successfulÚremove_requestÚdiscard_active_requestÚdiscard_reserved_requestr"   r"   r#   r   v   s
   r   ÚC_BENCHÚCELERY_BENCHÚC_BENCH_EVERYÚCELERY_BENCH_EVERY)Ú	monotonic)Úcurrent_process)ÚmemdumpÚ
sample_memÚMainProcessc                   C   sN   t d ur#td ur%td tt  ¡ƒ td ttƒttƒ ¡ƒ tƒ  d S d S d S )Nz- Time spent in benchmark: {!r}z	- Avg: {})Úbench_firstÚ
bench_lastÚprintÚformatÚsumÚbench_sampleÚlenr;   r"   r"   r"   r#   Úon_shutdown™   s   ÿÿ
ûrE   c                 C   s*   d}t du rtƒ  a }tdu r|at| ƒS )z-Called when a task is reserved by the worker.N)Úbench_startr9   r>   Ú
__reserved)r(   Únowr"   r"   r#   r   ¢   s   
c                 C   sX   t d7 a t t s(tƒ }|t }td t|¡ƒ tj ¡  | aa	t
 |¡ tƒ  t| ƒS )z Called when a task is completed.r+   zG- Time spent processing {} tasks (since first task received): ~{:.4f}s
)Ú	all_countÚbench_everyr9   rF   r@   rA   ÚsysÚstdoutÚflushr?   rC   Úappendr<   Ú__ready)r(   rH   Údiffr"   r"   r#   r   ®   s   ÿ

c                   @   s²   e Zd ZdZeZeZej	Z	ej
Z
dZd$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dd„ Zdd„ Zdd„ Zdd„ Zed d!„ ƒZed"d#„ ƒZdS )%r   zÅStores worker state between restarts.

    This is the persistent data stored by the worker when
    :option:`celery worker --statedb` is enabled.

    Currently only stores revoked task id's.
    FNc                 C   s   || _ || _|| _|  ¡  d S ©N)ÚstateÚfilenameÚclockÚmerge)ÚselfrR   rS   rT   r"   r"   r#   Ú__init__Ï   s   zPersistent.__init__c                 C   s   | j j| j| jdd�S )NT)ÚprotocolÚ	writeback)ÚstorageÚopenrS   rX   ©rV   r"   r"   r#   r[   Õ   s   
ÿzPersistent.openc                 C   s   |   | j¡ d S rQ   )Ú_merge_withÚdbr\   r"   r"   r#   rU   Ú   ó   zPersistent.mergec                 C   s   |   | j¡ | j ¡  d S rQ   )Ú
_sync_withr^   Úsyncr\   r"   r"   r#   ra   Ý   s   zPersistent.syncc                 C   s   | j r| j ¡  d| _ d S d S )NF)Ú_is_openr^   Úcloser\   r"   r"   r#   rc   á   s   

þzPersistent.closec                 C   s   |   ¡  |  ¡  d S rQ   )ra   rc   r\   r"   r"   r#   Úsaveæ   s   zPersistent.savec                 C   s   |   |¡ |  |¡ |S rQ   )Ú_merge_revokedÚ_merge_clock©rV   Údr"   r"   r#   r]   ê   s   

zPersistent._merge_withc                 C   s>   | j  ¡  | d|  |  | j ¡¡| jr| j ¡ nddœ¡ |S )Né   r   )Ú	__proto__ÚzrevokedrT   )Ú_revoked_tasksÚpurgeÚupdateÚcompressÚ_dumpsrT   Úforwardrg   r"   r"   r#   r`   ï   s   
ýzPersistent._sync_withc                 C   s(   | j r| j  | d¡pd¡|d< d S d S )NrT   r   )rT   ÚadjustÚgetrg   r"   r"   r#   rf   ø   s   ÿzPersistent._merge_clockc                 C   s\   z	|   |d ¡ W n ty&   z
|  | d¡¡ W n	 ty#   Y nw Y nw | j ¡  d S )Nrk   r   )Ú_merge_revoked_v3ÚKeyErrorÚ_merge_revoked_v2Úpoprl   rm   rg   r"   r"   r#   re   ü   s   ÿ€ýzPersistent._merge_revokedc                 C   s$   |r| j  t |  |¡¡¡ d S d S rQ   )rl   rn   r   ÚloadsÚ
decompress)rV   rk   r"   r"   r#   rt     s   ÿzPersistent._merge_revoked_v3c                 C   s$   t |tƒs
|  |¡S | j |¡ d S rQ   )Ú
isinstancer	   Ú_merge_revoked_v1rl   rn   )rV   Úsavedr"   r"   r#   rv     s   

zPersistent._merge_revoked_v2c                 C   s   | j j}|D ]}||ƒ qd S rQ   )rl   r0   )rV   r|   r0   Úitemr"   r"   r#   r{     s   
ÿzPersistent._merge_revoked_v1c                 C   s   t j|| jd�S )N)rX   )r   ÚdumpsrX   )rV   Úobjr"   r"   r#   rp     r_   zPersistent._dumpsc                 C   s   | j jS rQ   )rR   r   r\   r"   r"   r#   rl     s   zPersistent._revoked_tasksc                 C   s   d| _ |  ¡ S )NT)rb   r[   r\   r"   r"   r#   r^     s   zPersistent.dbrQ   )Ú__name__Ú
__module__Ú__qualname__Ú__doc__ÚshelverZ   r   rX   Úzlibro   ry   rb   rW   r[   rU   ra   rc   rd   r]   r`   rf   re   rt   rv   r{   rp   Úpropertyrl   r   r^   r"   r"   r"   r#   r   À   s2    
	
r   )Krƒ   ÚosÚplatformr„   rK   Úweakrefr…   Úcollectionsr   Úkombu.serializationr   r   Úkombu.utils.objectsr   Úceleryr   Úcelery.exceptionsr   r   Úcelery.utils.collectionsr	   Ú__all__Úsystemr
   ÚintÚenvironrs   ÚREVOKES_MAXÚSUCCESSFUL_MAXÚfloatÚREVOKE_EXPIRESÚSUCCESSFUL_EXPIRESr   ÚWeakSetr   r   r   r   r    r   r!   r&   r%   r$   r   Ú__setitem__r0   r   rn   r   rw   Údiscardr   r5   r7   ÚatexitÚtimer9   Úbilliard.processr:   Úcelery.utils.debugr;   r<   rI   r>   rF   r?   rJ   rC   rG   rO   Ú_nameÚregisterrE   r   r"   r"   r"   r#   Ú<module>   s”    ýÿ	
þ	
ü
ü
ÿÿ
