o
    wvXjn  ã                   @   sþ  d Z ddlmZmZm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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d	lmZ d
Zdee ¡ dœZdZdZi Ze	 ¡ Z e	 ¡ Z!eƒ Z"dgZ#eeed�Z$dZ%dZ&dd„ Z'dd„ Z(ej)e j*fdd„Z+de!j*e"j,fdd„Z-ej.e!j/e j/fdd„Z0ej1 2d¡p¥ej1 2d¡Z3e4ej1 2d¡p´ej1 2d¡p´dƒZ5e3rõddl6Z6ddl7m8Z8 dd lm9Z9 dd!l:m;Z;m<Z< da=da>da?da@e5ZAg ZBe+ZCe0ZDe8ƒ jEd"kríe6jFd#d$„ ƒZGd%d„ Z+d&d„ Z0G d'd(„ d(eHƒZIdS ))zwInternal worker state (global).

This includes the currently active and reserved tasks,
statistics, and revoked tasks.
é    )Úabsolute_importÚprint_functionÚunicode_literalsN)ÚpickleÚpickle_protocol)Úcached_property)Ú__version__)ÚWorkerShutdownÚWorkerTerminate)ÚCounter)Ú
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_sysiPÃ  i0*  )ÚmaxlenÚexpiresc                   C   s:   t  ¡  t ¡  t ¡  t ¡  dgtd d …< t ¡  d S )Nr   )ÚrequestsÚclearr   r   r   Úall_total_countr   © r   r   úP/var/www/html/myproject/venv/lib/python3.10/site-packages/celery/worker/state.pyÚreset_stateB   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   K   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   s   r   c                 C   s2   |st }|| ƒ || jdiƒ t d  d7  < dS )z2Update global state when a task has been accepted.é   r   N)r   Úname)r&   Ú_all_total_countÚadd_active_requestÚadd_to_total_countr   r   r    r   [   s
   r   c                 C   s    || j dƒ || ƒ || ƒ dS )z)Update global state when a task is ready.Nr$   )r&   Úremove_requestÚdiscard_active_requestÚdiscard_reserved_requestr   r   r    r   g   s   r   ÚC_BENCHÚCELERY_BENCHÚC_BENCH_EVERYÚCELERY_BENCH_EVERYiè  )Úcurrent_process)Ú	monotonic)Ú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: {0!r}z
- Avg: {0})Úbench_firstÚ
bench_lastÚprintÚformatÚsumÚbench_sampleÚlenr7   r   r   r   r    Úon_shutdown…   s   ÿÿ
ûrA   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_startr6   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)   zI- Time spent processing {0} tasks (since first task received): ~{1:.4f}s
)Ú	all_countÚbench_everyr6   rB   r<   r=   ÚsysÚstdoutÚflushr;   r?   Úappendr8   Ú__ready)r&   rD   Ú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)ÚselfrN   rO   rP   r   r   r    Ú__init__»   s   zPersistent.__init__c                 C   s   | j j| j| jdd�S )NT)ÚprotocolÚ	writeback)ÚstorageÚopenrO   rT   ©rR   r   r   r    rW   Á   s   
ÿzPersistent.openc                 C   s   |   | j¡ d S rM   )Ú_merge_withÚdbrX   r   r   r    rQ   Æ   ó   zPersistent.mergec                 C   s   |   | j¡ | j ¡  d S rM   )Ú
_sync_withrZ   ÚsyncrX   r   r   r    r]   É   s   zPersistent.syncc                 C   s   | j r| j ¡  d| _ d S d S )NF)Ú_is_openrZ   ÚcloserX   r   r   r    r_   Í   s   

þzPersistent.closec                 C   s   |   ¡  |  ¡  d S rM   )r]   r_   rX   r   r   r    ÚsaveÒ   s   zPersistent.savec                 C   s   |   |¡ |  |¡ |S rM   )Ú_merge_revokedÚ_merge_clock©rR   Údr   r   r    rY   Ö   s   

zPersistent._merge_withc              
   C   sN   | j  ¡  | tdƒdtdƒ|  |  | j ¡¡tdƒ| jr!| j ¡ ndi¡ |S )NÚ	__proto__é   ÚzrevokedrP   r   )Ú_revoked_tasksÚpurgeÚupdateÚstrÚcompressÚ_dumpsrP   Úforwardrc   r   r   r    r\   Û   s   
ýzPersistent._sync_withc                 C   s0   | j r| j  | tdƒ¡pd¡|tdƒ< d S d S )NrP   r   )rP   ÚadjustÚgetrk   rc   r   r   r    rb   ä   s   &ÿzPersistent._merge_clockc                 C   sd   z|   |tdƒ ¡ W n ty*   z|  | tdƒ¡¡ W n	 ty'   Y nw Y nw | j ¡  d S )Nrg   r   )Ú_merge_revoked_v3rk   ÚKeyErrorÚ_merge_revoked_v2Úpoprh   ri   rc   r   r   r    ra   è   s   ÿ€ýzPersistent._merge_revokedc                 C   s$   |r| j  t |  |¡¡¡ d S d S rM   )rh   rj   r   ÚloadsÚ
decompress)rR   rg   r   r   r    rq   ó   s   ÿzPersistent._merge_revoked_v3c                 C   s$   t |tƒs
|  |¡S | j |¡ d S rM   )Ú
isinstancer   Ú_merge_revoked_v1rh   rj   )rR   Úsavedr   r   r    rs   ÷   s   

zPersistent._merge_revoked_v2c                 C   s   | j j}|D ]}||ƒ qd S rM   )rh   Úadd)rR   ry   rz   Úitemr   r   r    rx   ý   s   
ÿzPersistent._merge_revoked_v1c                 C   s   t j|| jd�S )N)rT   )r   ÚdumpsrT   )rR   Úobjr   r   r    rm     r[   zPersistent._dumpsc                 C   s   | j jS rM   )rN   r   rX   r   r   r    rh     s   zPersistent._revoked_tasksc                 C   s   d| _ |  ¡ S )NT)r^   rW   rX   r   r   r    rZ   	  s   zPersistent.dbrM   )Ú__name__Ú
__module__Ú__qualname__Ú__doc__ÚshelverV   r   rT   Úzlibrl   rv   r^   rS   rW   rQ   r]   r_   r`   rY   r\   rb   ra   rq   rs   rx   rm   Úpropertyrh   r   rZ   r   r   r   r    r   ¬   s2    
	
r   )Jr�   Ú
__future__r   r   r   ÚosÚplatformr‚   rG   Úweakrefrƒ   Úkombu.serializationr   r   Úkombu.utils.objectsr   Úceleryr   Úcelery.exceptionsr	   r
   Úcelery.fiver   Úcelery.utils.collectionsr   Ú__all__Úsystemr   ÚREVOKES_MAXÚREVOKE_EXPIRESr   ÚWeakSetr   r   r   r   r   r#   r"   r!   r   Ú__setitem__rz   r   rj   r   rt   Údiscardr   Úenvironrp   r1   Úintr3   ÚatexitÚbilliard.processr5   r6   Úcelery.utils.debugr7   r8   rE   r:   rB   r;   rF   r?   rC   rK   Ú_nameÚregisterrA   Úobjectr   r   r   r   r    Ú<module>   s†   ý		
þ	
ý
ý

ÿÿ
