o
    wvXj%f  ã                	   @   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 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dlmZ ddlmZmZmZ ddlmZm Z m!Z! ddl"m#Z# zddl$m%Z% W n e&y…   ddlm%Z% Y nw dZ'e(edƒZ)dZ*dZ+dZ,e#e-ƒZ.e.j/Z0dZ1dZ2dZ3ej4ej5ej6ej7ej8ej9ej:ej;dœZ<G dd„ deƒZ=e% >e=¡ e ddd„ d �d!d"„ ƒZ?d#e*ee@eAfd$d%„ZBd&d'„ ZCd(d)„ ZDeDd*ƒeG d+d,„ d,eEƒƒƒZFeDd-ƒeG d.d/„ d/eEƒƒƒZGG d0d1„ d1eEƒZHd2d3„ ZId4d5„ ZJdS )6aÙ  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.
é    )Úabsolute_importÚunicode_literalsN)Údefaultdict)Údatetime)ÚDecimal)Úislice)Ú
itemgetter)Útime)ÚWeakSetÚref©Ú	timetuple)Úcached_property)Ústates)ÚitemsÚpython_2_unicode_compatibleÚvalues)ÚLRUCacheÚmemoizeÚpass1)Ú
get_logger)ÚCallable)ÚWorkerÚTaskÚStateÚheartbeat_expiresÚpypy_version_infoéÈ   é   znSubstantial 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t| ƒj|i |¤Ž d S ©N)ÚfunÚsuperr'   Ú__init__)Úselfr)   ÚargsÚkwargs©Ú	__class__© úP/var/www/html/myproject/venv/lib/python3.10/site-packages/celery/events/state.pyr+   h   s   zCallableDefaultdict.__init__c                 O   s   | j |i |¤ŽS r(   )r)   )r,   r-   r.   r1   r1   r2   Ú__call__l   ó   zCallableDefaultdict.__call__)Ú__name__Ú
__module__Ú__qualname__Ú__doc__r+   r3   Ú__classcell__r1   r1   r/   r2   r'   X   s    r'   iè  c                 C   s   | d S ©Nr   r1   )ÚaÚ_r1   r1   r2   Ú<lambda>s   s    r=   )ÚmaxsizeÚkeyfunc                 C   s    t t| |t |¡t |¡ƒ d S r(   )ÚwarnÚDRIFT_WARNINGr   Úfromtimestamp)ÚhostnameÚdriftÚlocal_receivedÚ	timestampr1   r1   r2   Ú_warn_drifts   s   þrG   é<   c                 C   s8   |||ƒr	||ƒn|}|| |ƒr|| ƒ} | ||d   S )z#Return time when heartbeat expires.g      Y@r1   )rF   ÚfreqÚexpire_windowr   ÚfloatÚ
isinstancer1   r1   r2   r   {   s   
r   c                 C   s   | di |¤ŽS )Nr1   r1   )ÚclsÚfieldsr1   r1   r2   Ú_depickle_task‡   ó   rO   c                    s   ‡ fdd„}|S )Nc                    s6   ‡ fdd„}|| _ dd„ }|| _‡ fdd„}|| _| S )Nc                    s$   t || jƒrt| ˆ ƒt|ˆ ƒkS tS r(   )rL   r0   ÚgetattrÚNotImplemented)ÚthisÚother©Úattrr1   r2   Ú__eq__�   s   z8with_unique_field.<locals>._decorate_cls.<locals>.__eq__c                 S   s   |   |¡}|tu rdS | S ©NT)rW   rR   )rS   rT   Úresr1   r1   r2   Ú__ne__•   s   
z8with_unique_field.<locals>._decorate_cls.<locals>.__ne__c                    s   t t| ˆ ƒƒS r(   )ÚhashrQ   )rS   rU   r1   r2   Ú__hash__š   rP   z:with_unique_field.<locals>._decorate_cls.<locals>.__hash__)rW   rZ   r\   )rM   rW   rZ   r\   rU   r1   r2   Ú_decorate_cls�   s   z(with_unique_field.<locals>._decorate_clsr1   )rV   r]   r1   rU   r2   Úwith_unique_field‹   s   r^   rC   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.é   )rC   ÚpidrI   Ú
heartbeatsÚclockÚactiveÚ	processedÚloadavgÚsw_identÚsw_verÚsw_sys)ÚeventÚ__dict__Ú__weakref__NrH   r   c                 C   s`   || _ || _|| _|d u rg n|| _|pd| _|| _|| _|| _|	| _|
| _	|| _
|  ¡ | _d S r:   )rC   r`   rI   ra   rb   rc   rd   re   rf   rg   rh   Ú_create_event_handlerri   )r,   rC   r`   rI   ra   rb   rc   rd   re   rf   rg   rh   r1   r1   r2   r+   °   s   
zWorker.__init__c                 C   s6   | j | j| j| j| j| j| j| j| j| j	| j
| jffS r(   )r0   rC   r`   rI   ra   rb   rc   rd   re   rf   rg   rh   ©r,   r1   r1   r2   Ú
__reduce__À   s
   ýzWorker.__reduce__c              	      sR   t j‰ ˆj‰ˆj‰ˆjj‰ˆjj‰d d d t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   éÿÿÿÿ)rG   rC   )Útype_rF   rE   rN   Ú	max_driftr   ÚabsÚintÚinsortÚlenÚkÚvrD   Úhearts©Ú_setÚ	hb_appendÚhb_popÚhbmaxra   r,   r1   r2   ri   Í   s(   ÿùz+Worker._create_event_handler.<locals>.event)ÚobjectÚ__setattr__Úheartbeat_maxra   ÚpopÚappendÚHEARTBEAT_DRIFT_MAXr   rt   ru   Úbisectrv   rw   ©r,   ri   r1   r{   r2   rl   Æ   s   ýzWorker._create_event_handlerc                 K   s6   t |rt|fi |¤Žn|ƒD ]
\}}t| ||ƒ qd S r(   )r   ÚdictÚsetattr)r,   ÚfÚkwrx   ry   r1   r1   r2   Úupdateç   s   $ÿzWorker.updatec                 C   ó
   t  | ¡S r(   )ÚR_WORKERÚformatrm   r1   r1   r2   Ú__repr__ë   ó   
zWorker.__repr__c                 C   s   | j rdS dS )NÚONLINEÚOFFLINE©Úaliverm   r1   r1   r2   Ústatus_stringî   s   zWorker.status_stringc                 C   s   t | jd | j| jƒS )Nrq   )r   ra   rI   rJ   rm   r1   r1   r2   r   ò   s   
ÿzWorker.heartbeat_expiresc                 C   s   t | jo	|ƒ | jk ƒS r(   )Úboolra   r   )r,   Únowfunr1   r1   r2   r•   ÷   s   zWorker.alivec                 C   s
   d  | ¡S )Nz{0.hostname}.{0.pid})r�   rm   r1   r1   r2   Úidû   ó   
z	Worker.id)NNrH   Nr   NNNNNN)r5   r6   r7   r8   r‚   ÚHEARTBEAT_EXPIRE_WINDOWrJ   Ú_fieldsÚPYPYÚ	__slots__r+   rn   rl   rŒ   r�   Úpropertyr–   r   r	   r•   r™   r1   r1   r1   r2   r   ¢   s.    
þ!

r   Úuuidc                   @   s8  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*e+j,ej-fd
d„Z.d%dd„Z/dd„ Z0dd„ Z1dd„ Z2dd„ Z3dd„ Z4dd„ Z5e6dd„ ƒZ7e6dd„ ƒZ8e6dd„ ƒZ9e:d d!„ ƒZ;e:d"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Ú	exceptionrF   ÚruntimeÚ	tracebackÚexchangeÚrouting_keyrb   ÚclientÚrootÚroot_idÚparentÚ	parent_idÚchildren)rj   rk   )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_idrm   r1   r2   Ú	<genexpr>1  s   € þýz Task.__init__.<locals>.<genexpr>r1   )r²   r®   r°   )
r    r³   r
   r²   Ú_serializable_childrenÚ_serializable_rootÚ_serializable_parentÚ_serializer_handlersrj   rŒ   )r,   r    r³   r²   r.   r1   rm   r2   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 r1   r1   )r¶   rx   ry   ©Úkeepr1   r2   Ú
<dictcomp>U  s    zTask.event.<locals>.<dictcomp>)r¢   rF   )Úupperr¢   Úmerge_rulesrµ   rŒ   rj   )r,   rr   rF   rE   rN   Ú
precedencer   r‰   Útask_event_to_stateÚRETRYr¢   r1   r½   r2   ri   @  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(   )ÚlistrQ   )ÚkeyÚvalue©ÚextrarN   r,   r1   r2   Ú_keysc  s   €
€ýzTask.info.<locals>._keys)Ú_info_fieldsrˆ   )r,   rN   rÉ   rÊ   r1   rÈ   r2   Úinfo^  s   
z	Task.infoc                 C   r�   r(   )ÚR_TASKr�   rm   r1   r1   r2   r�   k  r‘   zTask.__repr__c                    s&   t j‰ ˆjj‰‡ ‡‡fdd„ˆjD ƒS )Nc                    s"   i | ]}|ˆ|t ƒˆ ˆ|ƒƒ“qS r1   )r   )r¶   rx   ©rµ   Úhandlerr,   r1   r2   r¿   q  s    ÿz Task.as_dict.<locals>.<dictcomp>)r€   Ú__getattribute__r¼   rµ   rœ   rm   r1   rÎ   r2   Úas_dictn  s
   ÿzTask.as_dictc                 C   s   dd„ | j D ƒS )Nc                 S   ó   g | ]}|j ‘qS r1   ©r™   )r¶   Útaskr1   r1   r2   Ú
<listcomp>v  ó    z/Task._serializable_children.<locals>.<listcomp>)r²   ©r,   rÇ   r1   r1   r2   r¹   u  r4   zTask._serializable_childrenc                 C   ó   | j S r(   )r¯   r×   r1   r1   r2   rº   x  ó   zTask._serializable_rootc                 C   rØ   r(   )r±   r×   r1   r1   r2   r»   {  rÙ   zTask._serializable_parentc                 C   s   t | j|  ¡ ffS r(   )rO   r0   rÑ   rm   r1   r1   r2   rn   ~  ó   zTask.__reduce__c                 C   rØ   r(   )r    rm   r1   r1   r2   r™   �  s   zTask.idc                 C   s   | j d u r| jS | j jS r(   )r¦   r­   r™   rm   r1   r1   r2   Úorigin…  s   zTask.originc                 C   s   | j tjv S r(   ©r¢   r   ÚREADY_STATESrm   r1   r1   r2   Úready‰  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ÚKeyErrorrm   r1   r1   r2   r°   �  ó
   ÿzTask.parentc                 C   rß   r(   )r¯   r³   r´   rà   rá   rm   r1   r1   r2   r®   •  râ   z	Task.root)NNN)NN)=r5   r6   r7   r8   r¡   r    r   r!   r$   r"   r#   r%   r&   r-   r.   r£   r¤   r¥   r¦   r§   r¨   rF   r©   rª   r«   r¬   r¯   r±   r­   r   ÚPENDINGr¢   rb   rœ   r�   rž   ÚRECEIVEDrÁ   rË   r+   rÂ   r   r‰   ÚTASK_EVENT_TO_STATErµ   rÄ   ri   rÌ   r�   rÑ   r¹   rº   r»   rn   rŸ   r™   rÛ   rÞ   r   r°   r®   r1   r1   r1   r2   r      sˆ    ýÿÿÿÿÿÿÿþþþþþþýýýÿ

ý




r   c                   @   s  e Zd ZdZeZeZdZdZdZ					d6dd„Z	e
d	d
„ ƒZdd„ Zd7dd„Zd7dd„Zd7d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efd$d%„Zd8d&d'„Zd9d(d)„ZeZd9d*d+„Zd9d,d-„Zd.d/„ Zd0d1„ Z d2d3„ Z!d4d5„ Z"dS ):r   zRecords clusters state.r   r_   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û   r1   r1   r2   r+   §  s>   ÿÿÿÿ
ÿ
ÿÿ
ÿzState.__init__c                 C   s   |   ¡ S r(   )Ú_create_dispatcherrm   r1   r1   r2   Ú_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   r1   r1   r2   Ú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Þ   r1   r1   r2   Úclear_tasksÕ  ó   $ÿzState.clear_tasksc                 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 r1   rÜ   ©r¶   r    rÔ   r1   r1   r2   r¿   Û  s
    ÿz&State._clear_tasks.<locals>.<dictcomp>)Ú	itertasksr´   ÚclearrŒ   rê   )r,   rÞ   Úin_progressr1   r1   r2   r  Ù  s   ÿ

zState._clear_tasksc                 C   s$   | j  ¡  |  |¡ d| _d| _d S r:   )ré   r
  r  Úevent_countÚ
task_countr  r1   r1   r2   r  å  s   


zState._clearc                 C   r  r(   )rñ   r  r  r1   r1   r2   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,   rC   r.   r¦   r1   r1   r2   Ú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Ô   r1   r1   r2   Úget_or_create_taskÿ  s   þzState.get_or_create_taskc                 C   r  r(   )rñ   rÿ   r‡   r1   r1   r2   ri     r  zState.eventc                 C   ó    |   t|d d|g¡d�¡d S )úDeprecated, use :meth:`event`.ú-rÔ   ©Útyper   ©rÿ   rˆ   Újoin©r,   rr   rN   r1   r1   r2   Ú
task_event  ó    zState.task_eventc                 C   r  )r  r  r¦   r  r   r  r  r1   r1   r2   Ú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 )	NrC   rF   rE   r    rb   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 )Nrp   r  r  r¦   ro   FÚonlinerÔ   r   r  Tr   rq   r    )r  Ú	partitionri   rƒ   r­   r¦   r™   rw   r   r  r¡   Úaddr±   r´   Ú_add_pending_task_childr²   rõ   rŒ   )ri   r   rá   rv   ÚcreatedÚgroupr<   ÚsubjectrÏ   rC   rF   rE   Ú
is_offliner¦   r    rb   Ú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é   r1   r2   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†   rv   )r,   rÿ   r1   r+  r2   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Ô   Úchr1   r1   r2   r  Ž  s   ÿzState._add_pending_task_childc                    s2   ‡ fdd„t | jƒD ƒ }| jd d …< | ¡  d S )Nc                    s$   g | ]}ˆ |j |j|jt|ƒƒ‘qS r1   )rb   rF   rÛ   r   ©r¶   Útr   r1   r2   rÕ   –  s    ÿÿz*State.rebuild_taskheap.<locals>.<listcomp>)r   r´   rê   Úsort)r,   r   Úheapr1   r   r2   rö   •  s   
þzState.rebuild_taskheapc                 c   s:   � t t| jƒƒD ]\}}|V  |r|d |kr d S qd S )Nrp   )Ú	enumerater   r´   )r,   ÚlimitÚindexÚrowr1   r1   r2   r	  œ  s   €€ýzState.itertasksc                 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?  ÚreverseÚ_heapÚseenÚevtuprÔ   r    r1   r1   r2   Ú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  rJ  r1   r2   r¸   À  s   €
 
ÿÿz'State._tasks_by_type.<locals>.<genexpr>©rD  r   ©r   rH  )r,   r¡   r?  rD  r1   rJ  r2   r÷   ¶  s   	ýzState._tasks_by_typec                    rI  )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¦   rC   r  ©rC   r1   r2   r¸   Ë  s   €
 ÿÿz)State._tasks_by_worker.<locals>.<genexpr>rK  r   rL  )r,   rC   r?  rD  r1   rM  r2   rú   Å  s   ýzState._tasks_by_workerc                 C   s
   t | jƒS )z%Return a list of all seen task types.)Úsortedrô   rm   r1   r1   r2   Ú
task_typesÐ  rš   zState.task_typesc                 C   s   dd„ t | jƒD ƒS )z+Return a list of (seemingly) alive workers.c                 s   s   � | ]}|j r|V  qd S r(   r”   )r¶   Úwr1   r1   r2   r¸   Ö  s   € z&State.alive_workers.<locals>.<genexpr>)r   ré   rm   r1   r1   r2   Úalive_workersÔ  s   zState.alive_workersc                 C   r�   r(   )ÚR_STATEr�   rm   r1   r1   r2   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(   )r0   rè   ré   r´   rë   rì   rí   rî   Ú_serialize_Task_WeakSet_Mappingrø   rû   rm   r1   r1   r2   rn   Û  s   ûzState.__reduce__)
NNNNræ   rç   NNNN)Tr(   rX   )#r5   r6   r7   r8   r   r   r  r  r8  r+   r   rÿ   r  r  r  r  r
  r  r  ri   r  r  rþ   r  r   rö   r	  rH  Útasks_by_timestampr÷   rú   rO  rQ  r�   rn   r1   r1   r1   r2   r   ž  sJ    
ü!

	


{



r   c                 C   s   dd„ t | ƒD ƒS )Nc                 S   s    i | ]\}}|d d„ |D ƒ“qS )c                 S   rÒ   r1   rÓ   r:  r1   r1   r2   rÕ   æ  rÖ   z>_serialize_Task_WeakSet_Mapping.<locals>.<dictcomp>.<listcomp>r1   )r¶   r¡   r´   r1   r1   r2   r¿   æ  s     z3_serialize_Task_WeakSet_Mapping.<locals>.<dictcomp>©r   )Úmappingr1   r1   r2   rS  å  rÚ   rS  c                    s   ‡ fdd„t | p	i ƒD ƒS )Nc                    s(   i | ]\}}|t ‡ fd d„|D ƒƒ“qS )c                 3   s    � | ]}|ˆ v rˆ | V  qd S r(   r1   )r¶   Úi©r´   r1   r2   r¸   ê  s   € z?_deserialize_Task_WeakSet_Mapping.<locals>.<dictcomp>.<genexpr>)r
   )r¶   r¡   ÚidsrX  r1   r2   r¿   ê  s    ÿz5_deserialize_Task_WeakSet_Mapping.<locals>.<dictcomp>rU  )rV  r´   r1   rX  r2   rù   é  s   

ÿrù   )Kr8   Ú
__future__r   r   r†   Úsysrï   Úcollectionsr   r   Údecimalr   Ú	itertoolsr   Úoperatorr   r	   Úweakrefr
   r   Úkombu.clocksr   Úkombu.utils.objectsr   Úceleryr   Úcelery.fiver   r   r   Úcelery.utils.functionalr   r   r   Úcelery.utils.logr   Úcollections.abcr   ÚImportErrorÚ__all__Úhasattrr�   r›   r…   rA   r5   ÚloggerÚwarningr@   rR  rŽ   rÍ   rã   rä   ÚSTARTEDÚFAILURErÄ   ÚSUCCESSÚREVOKEDÚREJECTEDrå   r'   ÚregisterrG   rK   rL   r   rO   r^   r€   r   r   r   rS  rù   r1   r1   r1   r2   Ú<module>   s€   þ
ø


þ\   I