o
    ›¨Êh0d  ã                	   @   sÒ  d 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
 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dlmZmZmZ ddl m!Z! dZ"e#edƒZ$dZ%dZ&dZ'e!e(ƒZ)e)j*Z+dZ,dZ-dZ.ej/ej0ej1ej2ej3ej4ej5ej6dœZ7G dd„ deƒZ8e 9e8¡ eddd„ d�d d!„ ƒZ:d"e%e
e;e<fd#d$„Z=d%d&„ Z>d'd(„ Z?e?d)ƒG d*d+„ d+ƒƒZ@e?d,ƒG d-d.„ d.ƒƒZAG d/d0„ d0ƒZBd1d2„ ZCd3d4„ ZDdS )5aÙ  In-memory representation of cluster state.

This module implements a data-structure used to keep
track of the state of a cluster of workers and the tasks
it is working on (by consuming events).

For every event consumed the state is updated,
so the state represents the state of the cluster
at the time of the last event.

Snapshots (:mod:`celery.events.snapshot`) can be used to
take "pictures" of this state at regular intervals
to for example, store that in a database.
é    N)Údefaultdict)ÚCallable)Údatetime)ÚDecimal)Úislice)Ú
itemgetter)Útime)ÚMappingÚOptional)ÚWeakSetÚref©Ú	timetuple)Úcached_property)Ústates)ÚLRUCacheÚmemoizeÚpass1)Ú
get_logger)ÚWorkerÚTaskÚStateÚheartbeat_expiresÚpypy_version_infoéÈ   é   zmSubstantial drift from %s may mean clocks are out of sync.  Current drift is %s seconds.  [orig: %s recv: %s]z4<State: events={0.event_count} tasks={0.task_count}>z9<Worker: {0.hostname} ({0.status_string} clock:{0.clock})z4<Task: {0.name}({0.uuid}) {0.state} clock:{0.clock}>)ÚsentÚreceivedÚstartedÚfailedÚretriedÚ	succeededÚrevokedÚrejectedc                       s(   e Zd ZdZ‡ fdd„Zdd„ Z‡  ZS )ÚCallableDefaultdicta¬  :class:`~collections.defaultdict` with configurable __call__.

    We use this for backwards compatibility in State.tasks_by_type
    etc, which used to be a method but is now an index instead.

    So you can do::

        >>> add_tasks = state.tasks_by_type['proj.tasks.add']

    while still supporting the method call::

        >>> add_tasks = list(state.tasks_by_type(
        ...     'proj.tasks.add', reverse=True))
    c                    s   || _ tƒ j|i |¤Ž d S ©N)ÚfunÚsuperÚ__init__)Úselfr&   ÚargsÚkwargs©Ú	__class__© úE/var/www/html/env/lib/python3.10/site-packages/celery/events/state.pyr(   _   s   zCallableDefaultdict.__init__c                 O   s   | j |i |¤ŽS r%   )r&   )r)   r*   r+   r.   r.   r/   Ú__call__c   ó   zCallableDefaultdict.__call__)Ú__name__Ú
__module__Ú__qualname__Ú__doc__r(   r0   Ú__classcell__r.   r.   r,   r/   r$   O   s    r$   iè  c                 C   s   | d S ©Nr   r.   )ÚaÚ_r.   r.   r/   Ú<lambda>j   s    r:   )ÚmaxsizeÚkeyfunc                 C   s    t t| |t |¡t |¡ƒ d S r%   )ÚwarnÚDRIFT_WARNINGr   Úfromtimestamp)ÚhostnameÚdriftÚlocal_receivedÚ	timestampr.   r.   r/   Ú_warn_driftj   s   þrD   é<   c                 C   s8   |||ƒr	||ƒn|}|| |ƒr|| ƒ} | ||d   S )z#Return time when heartbeat expires.g      Y@r.   )rC   ÚfreqÚexpire_windowr   ÚfloatÚ
isinstancer.   r.   r/   r   r   s   
r   c                 C   s   | di |¤ŽS )Nr.   r.   )ÚclsÚfieldsr.   r.   r/   Ú_depickle_task~   ó   rL   c                    s   ‡ fdd„}|S )Nc                    s(   ‡ fdd„}|| _ ‡ fdd„}|| _| S )Nc                    s$   t || jƒrt| ˆ ƒt|ˆ ƒkS tS r%   )rI   r-   ÚgetattrÚNotImplemented)ÚthisÚother©Úattrr.   r/   Ú__eq__†   s   z8with_unique_field.<locals>._decorate_cls.<locals>.__eq__c                    s   t t| ˆ ƒƒS r%   )ÚhashrN   )rP   rR   r.   r/   Ú__hash__Œ   rM   z:with_unique_field.<locals>._decorate_cls.<locals>.__hash__)rT   rV   )rJ   rT   rV   rR   r.   r/   Ú_decorate_cls„   s
   z(with_unique_field.<locals>._decorate_clsr.   )rS   rW   r.   rR   r/   Úwith_unique_field‚   s   rX   r@   c                   @   sŒ   e Zd ZdZdZeZdZesed Z				ddd	„Z
d
d„ Zdd„ Zdd„ Zdd„ Zedd„ ƒZedd„ ƒZeefdd„ƒZedd„ ƒZdS )r   zWorker State.é   )r@   ÚpidrF   Ú
heartbeatsÚclockÚactiveÚ	processedÚloadavgÚsw_identÚsw_verÚsw_sys)ÚeventÚ__dict__Ú__weakref__NrE   r   c                 C   s`   || _ || _|| _|d u rg n|| _|pd| _|| _|| _|| _|	| _|
| _	|| _
|  ¡ | _d S r7   )r@   rZ   rF   r[   r\   r]   r^   r_   r`   ra   rb   Ú_create_event_handlerrc   )r)   r@   rZ   rF   r[   r\   r]   r^   r_   r`   ra   rb   r.   r.   r/   r(   ¡   s   
zWorker.__init__c                 C   s6   | j | j| j| j| j| j| j| j| j| j	| j
| jffS r%   )r-   r@   rZ   rF   r[   r\   r]   r^   r_   r`   ra   rb   ©r)   r.   r.   r/   Ú
__reduce__±   s
   ýzWorker.__reduce__c                    sP   t j‰ ˆj‰ˆj‰ˆjj‰ˆjj‰d d d tttt	j
tf‡ ‡‡‡‡‡fdd„	}|S )Nc	                    sÄ   |pi }|  ¡ D ]
\}	}
ˆ ˆ|	|
ƒ q| dkrg ˆd d …< d S |r#|s%d S |||ƒ||ƒ ƒ}||kr;tˆj|||ƒ |r`|ˆƒ}|ˆd krKˆdƒ |rY|ˆd krYˆ|ƒ d S |ˆ|ƒ d S d S )NÚofflineé   r   éÿÿÿÿ)ÚitemsrD   r@   )Útype_rC   rB   rK   Ú	max_driftÚabsÚintÚinsortÚlenÚkÚvrA   Úhearts©Ú_setÚ	hb_appendÚhb_popÚhbmaxr[   r)   r.   r/   rc   ¾   s(   ÿùz+Worker._create_event_handler.<locals>.event)ÚobjectÚ__setattr__Úheartbeat_maxr[   ÚpopÚappendÚHEARTBEAT_DRIFT_MAXro   rp   Úbisectrq   rr   ©r)   rc   r.   rv   r/   rf   ·   s   ýzWorker._create_event_handlerc                 K   s:   |r
t |fi |¤Žn|}| ¡ D ]
\}}t| ||ƒ qd S r%   )Údictrl   Úsetattr)r)   ÚfÚkwÚdrs   rt   r.   r.   r/   ÚupdateØ   s   ÿzWorker.updatec                 C   ó
   t  | ¡S r%   )ÚR_WORKERÚformatrg   r.   r.   r/   Ú__repr__Ý   ó   
zWorker.__repr__c                 C   s   | j rdS dS )NÚONLINEÚOFFLINE©Úaliverg   r.   r.   r/   Ústatus_stringà   s   zWorker.status_stringc                 C   s   t | jd | j| jƒS )Nrk   )r   r[   rF   rG   rg   r.   r.   r/   r   ä   s   
ÿzWorker.heartbeat_expiresc                 C   s   t | jo	|ƒ | jk ƒS r%   )Úboolr[   r   )r)   Únowfunr.   r.   r/   r‘   é   s   zWorker.alivec                 C   s
   d  | ¡S )Nz{0.hostname}.{0.pid})r‹   rg   r.   r.   r/   Úidí   ó   
z	Worker.id)NNrE   Nr   NNNNNN)r2   r3   r4   r5   r}   ÚHEARTBEAT_EXPIRE_WINDOWrG   Ú_fieldsÚPYPYÚ	__slots__r(   rh   rf   rˆ   rŒ   Úpropertyr’   r   r   r‘   r•   r.   r.   r.   r/   r   ”   s.    
þ!

r   Úuuidc                   @   s6  e Zd ZdZd Z Z Z Z Z Z	 Z
 Z Z Z Z Z Z Z Z Z Z Z Z Z Z Z Z ZZejZdZ dZ!e"sCdZ#ej$diZ%dZ&d$dd	„Z'dddej(e)e*j+ej,fd
d„Z-d%dd„Z.dd„ Z/dd„ Z0dd„ Z1dd„ Z2dd„ Z3dd„ Z4e5dd„ ƒZ6e5dd„ ƒZ7e5dd„ ƒZ8e9d d!„ ƒZ:e9d"d#„ ƒZ;dS )&r   zTask State.Nr   )rœ   ÚnameÚstater   r   r   r#   r!   r   r    r"   r*   r+   ÚetaÚexpiresÚretriesÚworkerÚresultÚ	exceptionrC   ÚruntimeÚ	tracebackÚexchangeÚrouting_keyr\   ÚclientÚrootÚroot_idÚparentÚ	parent_idÚchildren)rd   re   )r�   r*   r+   r­   r«   r¡   rŸ   r    )r*   r+   r¡   r£   rŸ   r¥   r    r¤   r§   r¨   r«   r­   c                    sh   |ˆ _ |ˆ _ˆ jd urt‡ fdd„|pdD ƒƒˆ _ntƒ ˆ _ˆ jˆ jˆ jdœˆ _|r2ˆ j 	|¡ d S d S )Nc                 3   s*   � | ]}|ˆ j jv rˆ j j |¡V  qd S r%   )Úcluster_stateÚtasksÚget)Ú.0Útask_idrg   r.   r/   Ú	<genexpr>"  s   € þýz Task.__init__.<locals>.<genexpr>r.   )r®   rª   r¬   )
rœ   r¯   r   r®   Ú_serializable_childrenÚ_serializable_rootÚ_serializable_parentÚ_serializer_handlersrd   rˆ   )r)   rœ   r¯   r®   r+   r.   rg   r/   r(     s   
þýÿzTask.__init__c	           
         sœ   |pi }||ƒ}	|	d ur|| ||ƒ n|  ¡ }	|	|kr?| j|kr?||	ƒ|| jƒkr?| j |	¡‰ ˆ d ur>‡ fdd„| ¡ D ƒ}n|j|	|d� | j |¡ d S )Nc                    s   i | ]\}}|ˆ v r||“qS r.   r.   )r²   rs   rt   ©Úkeepr.   r/   Ú
<dictcomp>E  s    zTask.event.<locals>.<dictcomp>)rž   rC   )Úupperrž   Úmerge_rulesr±   rl   rˆ   rd   )
r)   rm   rC   rB   rK   Ú
precedencer„   Útask_event_to_stateÚRETRYrž   r.   r¹   r/   rc   1  s   
ÿ€z
Task.eventc                    s8   ˆ sg nˆ ‰ ˆdu rˆj nˆ‰‡ ‡‡fdd„}t|ƒ ƒS )z;Information about this task suitable for on-screen display.Nc                  3   s:   � t ˆƒt ˆ ƒ D ]} tˆ| d ƒ}|d ur| |fV  q	d S r%   )ÚlistrN   )ÚkeyÚvalue©ÚextrarK   r)   r.   r/   Ú_keysS  s   €
€ýzTask.info.<locals>._keys)Ú_info_fieldsrƒ   )r)   rK   rÅ   rÆ   r.   rÄ   r/   ÚinfoN  s   
z	Task.infoc                 C   r‰   r%   )ÚR_TASKr‹   rg   r.   r.   r/   rŒ   [  r�   zTask.__repr__c                    s&   t j‰ ˆjj‰‡ ‡‡fdd„ˆjD ƒS )Nc                    s"   i | ]}|ˆ|t ƒˆ ˆ|ƒƒ“qS r.   )r   )r²   rs   ©r±   Úhandlerr)   r.   r/   r»   a  s    ÿz Task.as_dict.<locals>.<dictcomp>)r{   Ú__getattribute__r¸   r±   r˜   rg   r.   rÊ   r/   Úas_dict^  s
   ÿzTask.as_dictc                 C   s   dd„ | j D ƒS )Nc                 S   ó   g | ]}|j ‘qS r.   ©r•   )r²   Útaskr.   r.   r/   Ú
<listcomp>f  ó    z/Task._serializable_children.<locals>.<listcomp>)r®   ©r)   rÃ   r.   r.   r/   rµ   e  r1   zTask._serializable_childrenc                 C   ó   | j S r%   )r«   rÓ   r.   r.   r/   r¶   h  ó   zTask._serializable_rootc                 C   rÔ   r%   )r­   rÓ   r.   r.   r/   r·   k  rÕ   zTask._serializable_parentc                 C   s   t | j|  ¡ ffS r%   )rL   r-   rÍ   rg   r.   r.   r/   rh   n  ó   zTask.__reduce__c                 C   rÔ   r%   )rœ   rg   r.   r.   r/   r•   q  s   zTask.idc                 C   s   | j d u r| jS | j jS r%   )r¢   r©   r•   rg   r.   r.   r/   Úoriginu  s   zTask.originc                 C   s   | j tjv S r%   ©rž   r   ÚREADY_STATESrg   r.   r.   r/   Úreadyy  s   z
Task.readyc                 C   ó.   z| j o| jjj| j  W S  ty   Y d S w r%   )r­   r¯   r°   ÚdataÚKeyErrorrg   r.   r.   r/   r¬   }  ó
   ÿzTask.parentc                 C   rÛ   r%   )r«   r¯   r°   rÜ   rÝ   rg   r.   r.   r/   rª   …  rÞ   z	Task.root)NNN)NN)<r2   r3   r4   r5   r�   r   r   r   r!   r   r    r"   r#   r*   r+   rŸ   r    r¡   r¢   r£   r¤   rC   r¥   r¦   r§   r¨   r«   r­   r©   r   ÚPENDINGrž   r\   r˜   r™   rš   ÚRECEIVEDr½   rÇ   r(   r¾   r„   ÚTASK_EVENT_TO_STATEr±   rÀ   rc   rÈ   rŒ   rÍ   rµ   r¶   r·   rh   r›   r•   r×   rÚ   r   r¬   rª   r.   r.   r.   r/   r   ò   s†    ýÿÿÿÿÿÿÿþþþþþþýýýÿ

þ




r   c                   @   s   e Zd ZdZeZeZdZdZdZ					d9dd„Z	e
d	d
„ ƒZdd„ Zd:dd„Zd:defdd„Zd:dd„Zd:def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fd%d&„Zd;d'ee fd(d)„Zd<d*efd+d,„ZeZd<d-d.„Z d<d/d0„Z!d1d2„ Z"d3d4„ Z#d5d6„ Z$d7d8„ Z%dS )=r   zRecords clusters state.r   rY   Néˆ  é'  c                 C   sÊ   || _ |d u rt|ƒn|| _|d u rt|ƒn|| _|d u rg n|| _|| _|| _|| _|| _t	 
¡ | _i | _tƒ | _i | _|  ¡  t| jtƒ| _| j t|	| jƒ¡ t| jtƒ| _| j t|
| jƒ¡ d S r%   )Úevent_callbackr   Úworkersr°   Ú	_taskheapÚmax_workers_in_memoryÚmax_tasks_in_memoryÚon_node_joinÚon_node_leaveÚ	threadingÚLockÚ_mutexÚhandlersÚsetÚ_seen_typesÚ_tasks_to_resolveÚrebuild_taskheapr$   Ú_tasks_by_typer   Útasks_by_typerˆ   Ú!_deserialize_Task_WeakSet_MappingÚ_tasks_by_workerÚtasks_by_worker)r)   Úcallbackrå   r°   Útaskheaprç   rè   ré   rê   rô   r÷   r.   r.   r/   r(   —  s>   ÿÿÿÿ
ÿ
ÿÿ
ÿzState.__init__c                 C   s   |   ¡ S r%   )Ú_create_dispatcherrg   r.   r.   r/   Ú_event¶  s   zState._eventc              	   O   sd   |  dd¡}| j� z||i |¤ŽW |r|  ¡  W  d   ƒ S |r'|  ¡  w w 1 s+w   Y  d S )NÚclear_afterF)r~   rí   Ú_clear)r)   r&   r*   r+   rü   r.   r.   r/   Úfreeze_whileº  s   û
ÿüzState.freeze_whileTc                 C   ó4   | j � |  |¡W  d   ƒ S 1 sw   Y  d S r%   )rí   Ú_clear_tasks©r)   rÚ   r.   r.   r/   Úclear_tasksÃ  ó   $ÿzState.clear_tasksrÚ   c                 C   sJ   |rdd„ |   ¡ D ƒ}| j ¡  | j |¡ n| j ¡  g | jd d …< d S )Nc                 S   s"   i | ]\}}|j tjvr||“qS r.   rØ   ©r²   rœ   rÐ   r.   r.   r/   r»   É  s
    ÿz&State._clear_tasks.<locals>.<dictcomp>)Ú	itertasksr°   Úclearrˆ   ræ   )r)   rÚ   Úin_progressr.   r.   r/   r   Ç  s   ÿ

zState._clear_tasksc                 C   s$   | j  ¡  |  |¡ d| _d| _d S r7   )rå   r  r   Úevent_countÚ
task_countr  r.   r.   r/   rý   Ó  s   


zState._clearc                 C   rÿ   r%   )rí   rý   r  r.   r.   r/   r  Ù  r  zState.clearc                 K   sZ   z| j | }|r| |¡ |dfW S  ty,   | j|fi |¤Ž }| j |< |df Y S w )zsGet or create worker by hostname.

        Returns:
            Tuple: of ``(worker, was_created)`` pairs.
        FT)rå   rˆ   rÝ   r   )r)   r@   r+   r¢   r.   r.   r/   Úget_or_create_workerÝ  s   


ÿÿýzState.get_or_create_workerc                 C   sD   z| j | dfW S  ty!   | j|| d� }| j |< |df Y S w )zGet or create task by uuid.F©r¯   T)r°   rÝ   r   )r)   rœ   rÐ   r.   r.   r/   Úget_or_create_taskí  s   þzState.get_or_create_taskc                 C   rÿ   r%   )rí   rû   r‚   r.   r.   r/   rc   õ  r  zState.eventc                 C   ó    |   t|d d|g¡d�¡d S )úDeprecated, use :meth:`event`.ú-rÐ   ©Útyper   ©rû   rƒ   Újoin©r)   rm   rK   r.   r.   r/   Ú
task_eventù  ó    zState.task_eventc                 C   r  )r  r  r¢   r  r   r  r  r.   r.   r/   Úworker_eventý  r  zState.worker_eventc                    sÞ   ˆj j‰ˆj‰tdddƒ‰tdddddƒ‰ˆj‰ˆj‰ˆj‰ˆjˆj ‰	ˆj	j
‰ˆjˆj‰
‰ˆjˆj‰‰ ˆjˆj‰‰ˆjjˆjj‰‰ˆjj‰ˆjj‰tttjdf‡ ‡‡‡‡‡‡‡‡‡	‡
‡‡‡‡‡‡‡‡‡fdd„	}|S )	Nr@   rC   rB   rœ   r\   Tc                    s<  ˆ j d7  _ ˆrˆˆ| ƒ | d  d¡\}}}zˆ|ƒ}W n	 |y'   Y nw ||| ƒ|fS |dkr˜z	ˆ| ƒ\}	}
}W n
 |yF   Y d S w |dk}z	ˆ|	ƒd}}W n |yo   |reˆ|	ƒd}}nˆ|	ƒ }ˆ|	< Y nw | ||
|| ¡ ˆ
r„|s€|dkr„ˆ
|ƒ ˆr’|r’ˆ|ƒ ˆ |	d ¡ ||f|fS |dk�rœˆ| ƒ\}}	}
}}|d	k}z	ˆ|ƒd}}W n |yÈ   ˆ |ˆd
� }ˆ|< d}Y nw |rÏ|	|_n(zˆ|	ƒ}W n |yæ   ˆ|	ƒ }ˆ|	< Y nw ||_|d ur÷|r÷| d ||
¡ |rû|	n|j}tˆƒ}|d ˆ	k�rˆdƒ |||
|t|ƒƒ}|�r%|ˆd k�r%ˆ|ƒ n|ˆ|ƒ |dk�r6ˆ j	d7  _	| ||
|| ¡ |j
}|d u�r[ˆ|ƒ |�r[ˆ|ƒ |¡ ˆ|	ƒ |¡ |j�r}zˆj|j }W n |�yv   ˆ |¡ Y nw |j |¡ zˆj |¡}W n
 |�y�   Y nw |j |¡ ||f|fS d S )Nrj   r  r  r¢   ri   FÚonlinerÐ   r   r  Tr   rk   r   )r  Ú	partitionrc   r~   r©   r¢   r•   rr   r   r	  r�   Úaddr­   r°   Ú_add_pending_task_childr®   rñ   rˆ   )rc   r   rÝ   rq   ÚcreatedÚgroupr9   ÚsubjectrË   r@   rC   rB   Ú
is_offliner¢   rœ   r\   Úis_client_eventrÐ   Útask_createdr×   ÚheapsÚtimetupÚ	task_nameÚparent_taskÚ	_children©r   r   Úadd_typerä   Úget_handlerÚget_taskÚget_task_by_type_setÚget_task_by_worker_setÚ
get_workerÚmax_events_in_heapré   rê   r)   rù   r°   ÚtfieldsÚ	th_appendÚth_popÚwfieldsrå   r.   r/   rû     sª   
ÿÿ€ü
ÿþÿ



ÿÿÆz(State._create_dispatcher.<locals>._event)rî   Ú__getitem__rä   r   ræ   r   r~   rè   Úheap_multiplierrð   r  ré   rê   r°   r   rå   r   rÜ   rô   r÷   r   rÝ   r�   rq   )r)   rû   r.   r'  r/   rú     s*   ÿ4þ^zState._create_dispatcherc                 C   sD   z| j |j }W n ty   tƒ  }| j |j< Y nw | |¡ d S r%   )rñ   r­   rÝ   r   r  )r)   rÐ   Úchr.   r.   r/   r  |  s   ÿzState._add_pending_task_childc                    s2   ‡ fdd„| j  ¡ D ƒ }| jd d …< | ¡  d S )Nc                    s$   g | ]}ˆ |j |j|jt|ƒƒ‘qS r.   )r\   rC   r×   r   ©r²   Útr   r.   r/   rÑ   „  s    ÿÿz*State.rebuild_taskheap.<locals>.<listcomp>)r°   Úvaluesræ   Úsort)r)   r   Úheapr.   r   r/   rò   ƒ  s   
þzState.rebuild_taskheapÚlimitc                 c   s:   � t | j ¡ ƒD ]\}}|V  |r|d |kr d S qd S )Nrj   )Ú	enumerater°   rl   )r)   r;  ÚindexÚrowr.   r.   r/   r  Š  s   €€ýzState.itertasksÚreversec                 c   sd   � | j }|r
t|ƒ}tƒ }t|d|ƒD ]}|d ƒ }|dur/|j}||vr/||fV  | |¡ qdS )zkGenerator yielding tasks ordered by time.

        Yields:
            Tuples of ``(uuid, Task)``.
        r   é   N)ræ   Úreversedrï   r   rœ   r  )r)   r;  r?  Ú_heapÚseenÚevtuprÐ   rœ   r.   r.   r/   Útasks_by_time�  s   €


€úzState.tasks_by_timec                    ó"   t ‡ fdd„| j|d�D ƒd|ƒS )zÊGet all tasks by type.

        This is slower than accessing :attr:`tasks_by_type`,
        but will be ordered by time.

        Returns:
            Generator: giving ``(uuid, Task)`` pairs.
        c                 3   s&   � | ]\}}|j ˆ kr||fV  qd S r%   ©r�   r  rG  r.   r/   r´   ®  s   €
 
ÿÿz'State._tasks_by_type.<locals>.<genexpr>©r?  r   ©r   rE  )r)   r�   r;  r?  r.   rG  r/   ró   ¤  s   	ýzState._tasks_by_typec                    rF  )znGet all tasks by worker.

        Slower than accessing :attr:`tasks_by_worker`, but ordered by time.
        c                 3   s(   � | ]\}}|j jˆ kr||fV  qd S r%   )r¢   r@   r  ©r@   r.   r/   r´   ¹  s   €
 ÿÿz)State._tasks_by_worker.<locals>.<genexpr>rH  r   rI  )r)   r@   r;  r?  r.   rJ  r/   rö   ³  s   ýzState._tasks_by_workerc                 C   s
   t | jƒS )z%Return a list of all seen task types.)Úsortedrð   rg   r.   r.   r/   Ú
task_types¾  r–   zState.task_typesc                 C   s   dd„ | j  ¡ D ƒS )z+Return a list of (seemingly) alive workers.c                 s   s   � | ]}|j r|V  qd S r%   r�   )r²   Úwr.   r.   r/   r´   Ä  s   € z&State.alive_workers.<locals>.<genexpr>)rå   r8  rg   r.   r.   r/   Úalive_workersÂ  s   zState.alive_workersc                 C   r‰   r%   )ÚR_STATEr‹   rg   r.   r.   r/   rŒ   Æ  r�   zState.__repr__c                 C   s8   | j | j| j| jd | j| j| j| jt| j	ƒt| j
ƒf
fS r%   )r-   rä   rå   r°   rç   rè   ré   rê   Ú_serialize_Task_WeakSet_Mappingrô   r÷   rg   r.   r.   r/   rh   É  s   ûzState.__reduce__)
NNNNrâ   rã   NNNN)Tr%   )NT)&r2   r3   r4   r5   r   r   r  r	  r4  r(   r   rû   rþ   r  r“   r   rý   r  r
  r  rc   r  r  rú   r  r   rò   r
   rp   r  rE  Útasks_by_timestampró   rö   rL  rN  rŒ   rh   r.   r.   r.   r/   r   Ž  sJ    
ü

	
{

r   c                 C   s   dd„ |   ¡ D ƒS )Nc                 S   s    i | ]\}}|d d„ |D ƒ“qS )c                 S   rÎ   r.   rÏ   r6  r.   r.   r/   rÑ   Ô  rÒ   z>_serialize_Task_WeakSet_Mapping.<locals>.<dictcomp>.<listcomp>r.   )r²   r�   r°   r.   r.   r/   r»   Ô  s     z3_serialize_Task_WeakSet_Mapping.<locals>.<dictcomp>©rl   )Úmappingr.   r.   r/   rP  Ó  rÖ   rP  c                    s   | pi } ‡ fdd„|   ¡ D ƒS )Nc                    s(   i | ]\}}|t ‡ fd d„|D ƒƒ“qS )c                 3   s    � | ]}|ˆ v rˆ | V  qd S r%   r.   )r²   Úi©r°   r.   r/   r´   Ù  s   € z?_deserialize_Task_WeakSet_Mapping.<locals>.<dictcomp>.<genexpr>)r   )r²   r�   ÚidsrU  r.   r/   r»   Ù  s    ÿz5_deserialize_Task_WeakSet_Mapping.<locals>.<dictcomp>rR  )rS  r°   r.   rU  r/   rõ   ×  s   
ÿrõ   )Er5   r�   Úsysrë   Úcollectionsr   Úcollections.abcr   r   Údecimalr   Ú	itertoolsr   Úoperatorr   r   Útypingr	   r
   Úweakrefr   r   Úkombu.clocksr   Úkombu.utils.objectsr   Úceleryr   Úcelery.utils.functionalr   r   r   Úcelery.utils.logr   Ú__all__Úhasattrr™   r—   r€   r>   r2   ÚloggerÚwarningr=   rO  rŠ   rÉ   rß   rà   ÚSTARTEDÚFAILURErÀ   ÚSUCCESSÚREVOKEDÚREJECTEDrá   r$   ÚregisterrD   rH   rI   r   rL   rX   r   r   r   rP  rõ   r.   r.   r.   r/   Ú<module>   st    
ÿø


þ]   G