o
    wvXjŒŠ  ã                   @   sÂ  d Z ddlmZm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mZ d	dlmZ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!m"Z"m#Z# d	dl$m%Z% d	dl&m'Z'm(Z( d	dl)m*Z* zddl+Z+W n e,y‰   dZ+Y nw dZ-dZ.dd„ Z/e
dd„ ƒZ0e
dd„ ƒZ1G dd„ de2ƒZ3ej4e!G dd„ de3ƒƒƒZ5ej4e!G dd„ de3ƒƒƒZ6ej4e!G d d!„ d!e6ƒƒƒZ7ej4e!G d"d#„ d#e5ƒƒƒZ8d&d$d%„Z9dS )'z3Task results/state and results for groups of tasks.é    )Úabsolute_importÚunicode_literalsN)ÚOrderedDictÚdeque)Úcontextmanager)Úcopy)Úcached_property)ÚThenableÚbarrierÚpromiseé   )Úcurrent_appÚstates)Ú_set_task_join_will_blockÚtask_join_will_block)Úapp_or_default)ÚImproperlyConfiguredÚIncompleteStreamÚTimeoutError)ÚitemsÚ	monotonicÚpython_2_unicode_compatibleÚrangeÚstring_t)Ú
deprecated)ÚDependencyGraphÚGraphFormatter)Úparse_iso8601)Ú
ResultBaseÚAsyncResultÚ	ResultSetÚGroupResultÚEagerResultÚresult_from_tuplez|Never call result.get() within a task!
See http://docs.celeryq.org/en/latest/userguide/tasks.html#task-synchronous-subtasks
c                   C   s   t ƒ rttƒ‚d S ©N)r   ÚRuntimeErrorÚE_WOULDBLOCK© r'   r'   úJ/var/www/html/myproject/venv/lib/python3.10/site-packages/celery/result.pyÚassert_will_not_block)   s   ÿr)   c                  c   ó0   � t ƒ } tdƒ z
d V  W t| ƒ d S t| ƒ w ©NF©r   r   ©Úreset_valuer'   r'   r(   Úallow_join_result.   ó   €r/   c                  c   r*   ©NTr,   r-   r'   r'   r(   Údenied_join_result8   r0   r2   c                   @   s   e Zd ZdZdZdS )r   zBase class for results.N)Ú__name__Ú
__module__Ú__qualname__Ú__doc__Úparentr'   r'   r'   r(   r   B   s    r   c                   @   s4  e Zd ZdZdZeZdZdZ			dfdd„Ze	dd„ ƒZ
e
jdd„ ƒZ
dgd	d
„Zdd„ Zdd„ Zdd„ Z		dhdd„Zdddddddddejejfdd„ZeZdd„ Zdd„ Zdidd„Zdd„ Zdidd „Zd!d"„ Zd#d$„ Zd%d&„ Zd'd(„ Zdjd)d*„ZeZ d+d,„ Z!dkd-d.„Z"d/d0„ Z#d1d2„ Z$d3d4„ Z%d5d6„ Z&d7d8„ Z'd9d:„ Z(d;d<„ Z)d=d>„ Z*d?d@„ Z+e,dAdB„ ƒZ-e	dCdD„ ƒZ.e	dEdF„ ƒZ/dGdH„ Z0dIdJ„ Z1dKdL„ Z2dMdN„ Z3e	dOdP„ ƒZ4e4Z5e	dQdR„ ƒZ6e	dSdT„ ƒZ7e7Z8e	dUdV„ ƒZ9e9jdWdV„ ƒZ9e	dXdY„ ƒZ:e	dZd[„ ƒZ;e	d\d]„ ƒZ<e	d^d_„ ƒZ=e	d`da„ ƒZ>e	dbdc„ ƒZ?e	ddde„ ƒZ@dS )lr   zxQuery task state.

    Arguments:
        id (str): See :attr:`id`.
        backend (Backend): See :attr:`backend`.
    Nc                 C   sd   |d u rt d t|ƒ¡ƒ‚t|p| jƒ| _|| _|p| jj| _|| _t| j	dd�| _
d | _d| _d S )Nz&AsyncResult requires valid id, not {0}T©ÚweakF)Ú
ValueErrorÚformatÚtyper   ÚappÚidÚbackendr7   r   Ú_on_fulfilledÚon_readyÚ_cacheÚ_ignored)Úselfr>   r?   Ú	task_namer=   r7   r'   r'   r(   Ú__init__^   s   ÿ
zAsyncResult.__init__c                 C   s   t | dƒr| jS dS )z+If True, task result retrieval is disabled.rC   F)ÚhasattrrC   ©rD   r'   r'   r(   Úignoredl   s   
zAsyncResult.ignoredc                 C   s
   || _ dS )z%Enable/disable task result retrieval.N)rC   )rD   Úvaluer'   r'   r(   rI   s   s   
Fc                 C   s   | j j| |d� | j ||¡S )Nr8   )r?   Úadd_pending_resultrA   Úthen©rD   ÚcallbackÚon_errorr9   r'   r'   r(   rL   x   s   zAsyncResult.thenc                 C   s   | j  | ¡ |S r$   ©r?   Úremove_pending_result©rD   Úresultr'   r'   r(   r@   |   s   zAsyncResult._on_fulfilledc                 C   s   | j }| j|o
| ¡ fd fS r$   )r7   r>   Úas_tuple)rD   r7   r'   r'   r(   rT   €   s   zAsyncResult.as_tuplec                 C   s(   d| _ | jr| j ¡  | j | j¡ dS )z/Forget the result of this task and its parents.N)rB   r7   Úforgetr?   r>   rH   r'   r'   r(   rU   „   s   
zAsyncResult.forgetc                 C   s    | j jj| j|||||d� dS )aŠ  Send revoke signal to all workers.

        Any worker receiving the task, or having reserved the
        task, *must* ignore it.

        Arguments:
            terminate (bool): Also terminate the process currently working
                on the task (if any).
            signal (str): Name of signal to send to process if terminate.
                Default is TERM.
            wait (bool): Wait for replies from workers.
                The ``timeout`` argument specifies the seconds to wait.
                Disabled by default.
            timeout (float): Time in seconds to wait for replies when
                ``wait`` is enabled.
        )Ú
connectionÚ	terminateÚsignalÚreplyÚtimeoutN)r=   ÚcontrolÚrevoker>   ©rD   rV   rW   rX   ÚwaitrZ   r'   r'   r(   r\   ‹   s   
þzAsyncResult.revokeTç      à?c              
   C   s�   | j rdS |	r
tƒ  tƒ }|r|r| jrt| jdd�}|  ¡  |r&| |¡ | jr4|r1| j|d� | jS | j	 
| ¡ | j	j| |||||||d�S )aæ  Wait until task is ready, and return its result.

        Warning:
           Waiting for tasks within a task may lead to deadlocks.
           Please read :ref:`task-synchronous-subtasks`.

        Warning:
           Backends use resources to store and transmit results. To ensure
           that resources are released, you must eventually call
           :meth:`~@AsyncResult.get` or :meth:`~@AsyncResult.forget` on
           EVERY :class:`~@AsyncResult` instance returned after calling
           a task.

        Arguments:
            timeout (float): How long to wait, in seconds, before the
                operation times out.
            propagate (bool): Re-raise exception if the task failed.
            interval (float): Time to wait (in seconds) before retrying to
                retrieve the result.  Note that this does not have any effect
                when using the RPC/redis result store backends, as they don't
                use polling.
            no_ack (bool): Enable amqp no ack (automatically acknowledge
                message).  If this is :const:`False` then the message will
                **not be acked**.
            follow_parents (bool): Re-raise any exception raised by
                parent tasks.
            disable_sync_subtasks (bool): Disable tasks to wait for sub tasks
                this is the default configuration. CAUTION do not enable this
                unless you must.

        Raises:
            celery.exceptions.TimeoutError: if `timeout` isn't
                :const:`None` and the result does not arrive within
                `timeout` seconds.
            Exception: If the remote call raised an exception then that
                exception will be re-raised in the caller process.
        NTr8   )rN   )rZ   ÚintervalÚon_intervalÚno_ackÚ	propagaterN   Ú
on_message)rI   r)   r   r7   Ú_maybe_reraise_parent_errorrL   rB   Úmaybe_throwrS   r?   rK   Úwait_for_pending)rD   rZ   rc   r`   rb   Úfollow_parentsrN   rd   ra   Údisable_sync_subtasksÚEXCEPTION_STATESÚPROPAGATE_STATESÚ_on_intervalr'   r'   r(   Úget¡   s0   *
ùzAsyncResult.getc                 C   s"   t t|  ¡ ƒƒD ]}| ¡  qd S r$   )ÚreversedÚlistÚ_parentsrf   ©rD   Únoder'   r'   r(   re   è   s   
ÿz'AsyncResult._maybe_reraise_parent_errorc                 c   s$   � | j }|r|V  |j }|sd S d S r$   ©r7   rq   r'   r'   r(   rp   ì   s   €þzAsyncResult._parentsc                 k   s2   � | j |d�D ]\}}||jdi |¤ŽfV  qdS )a´  Collect results as they return.

        Iterator, like :meth:`get` will wait for the task to complete,
        but will also follow :class:`AsyncResult` and :class:`ResultSet`
        returned by the task, yielding ``(result, value)`` tuples for each
        result in the tree.

        An example would be having the following tasks:

        .. code-block:: python

            from celery import group
            from proj.celery import app

            @app.task(trail=True)
            def A(how_many):
                return group(B.s(i) for i in range(how_many))()

            @app.task(trail=True)
            def B(i):
                return pow2.delay(i)

            @app.task(trail=True)
            def pow2(i):
                return i ** 2

        .. code-block:: pycon

            >>> from celery.result import ResultBase
            >>> from proj.tasks import A

            >>> result = A.delay(10)
            >>> [v for v in result.collect()
            ...  if not isinstance(v, (ResultBase, tuple))]
            [0, 1, 4, 9, 16, 25, 36, 49, 64, 81]

        Note:
            The ``Task.trail`` option must be enabled
            so that the list of children is stored in ``result.children``.
            This is the default but enabled explicitly for illustration.

        Yields:
            Tuple[AsyncResult, Any]: tuples containing the result instance
            of the child task, and the return value of that task.
        ©ÚintermediateNr'   ©Úiterdepsrm   )rD   ru   ÚkwargsÚ_ÚRr'   r'   r(   Úcollectò   s   €.ÿzAsyncResult.collectc                 C   s"   d }|   ¡ D ]\}}| ¡ }q|S r$   rv   )rD   rJ   ry   rz   r'   r'   r(   Úget_leaf#  s   
zAsyncResult.get_leafc                 #   sh   � t d | fgƒ}|r2| ¡ \}‰ |ˆ fV  ˆ  ¡ r)| ‡ fdd„ˆ jp$g D ƒ¡ n|s.tƒ ‚|s
d S d S )Nc                 3   s   � | ]}ˆ |fV  qd S r$   r'   ©Ú.0Úchild©rr   r'   r(   Ú	<genexpr>0  ó   € z'AsyncResult.iterdeps.<locals>.<genexpr>)r   ÚpopleftÚreadyÚextendÚchildrenr   )rD   ru   Ústackr7   r'   r€   r(   rw   )  s   €
 ùzAsyncResult.iterdepsc                 C   s   | j | jjv S )z¨Return :const:`True` if the task has executed.

        If the task is still running, pending, or is waiting
        for retry then :const:`False` is returned.
        )Ústater?   ÚREADY_STATESrH   r'   r'   r(   r„   5  s   zAsyncResult.readyc                 C   ó   | j tjkS )z7Return :const:`True` if the task executed successfully.)rˆ   r   ÚSUCCESSrH   r'   r'   r(   Ú
successful=  ó   zAsyncResult.successfulc                 C   rŠ   )z(Return :const:`True` if the task failed.)rˆ   r   ÚFAILURErH   r'   r'   r(   ÚfailedA  r�   zAsyncResult.failedc                 O   s   | j j|i |¤Ž d S r$   )rA   Úthrow©rD   Úargsrx   r'   r'   r(   r�   E  s   zAsyncResult.throwc                 C   sn   | j d u r	|  ¡ n| j }|d |d | d¡}}}|tjv r+|r+|  ||  |¡¡ |d ur5|| j|ƒ |S )NÚstatusrS   Ú	traceback)rB   Ú_get_task_metarm   r   rk   r�   Ú_to_remote_tracebackr>   )rD   rc   rN   Úcacherˆ   rJ   Útbr'   r'   r(   rf   H  s   
ÿzAsyncResult.maybe_throwc                 C   s2   |rt d ur| jjjrt j |¡ ¡ S d S d S d S r$   )Útblibr=   ÚconfÚtask_remote_tracebacksÚ	TracebackÚfrom_stringÚas_traceback)rD   r˜   r'   r'   r(   r–   S  s   ÿz AsyncResult._to_remote_tracebackc                 C   sL   t |p	t| jdd�d�}| j|d�D ]\}}| |¡ |r#| ||¡ q|S )NÚoval)ÚrootÚshape)Ú	formatterrt   )r   r   r>   rw   Úadd_arcÚadd_edge)rD   ru   r¢   Úgraphr7   rr   r'   r'   r(   Úbuild_graphW  s   ÿ
€zAsyncResult.build_graphc                 C   ó
   t | jƒS ©z`str(self) -> self.id`.©Ústrr>   rH   r'   r'   r(   Ú__str__a  ó   
zAsyncResult.__str__c                 C   r§   ©z`hash(self) -> hash(self.id)`.©Úhashr>   rH   r'   r'   r(   Ú__hash__e  r¬   zAsyncResult.__hash__c                 C   s   d  t| ƒj| j¡S )Nz
<{0}: {1}>)r;   r<   r3   r>   rH   r'   r'   r(   Ú__repr__i  ó   zAsyncResult.__repr__c                 C   s.   t |tƒr|j| jkS t |tƒr|| jkS tS r$   )Ú
isinstancer   r>   r   ÚNotImplemented©rD   Úotherr'   r'   r(   Ú__eq__l  s
   


zAsyncResult.__eq__c                 C   ó   |   |¡}|tu rdS | S r1   ©r·   r´   ©rD   r¶   Úresr'   r'   r(   Ú__ne__s  ó   
zAsyncResult.__ne__c                 C   s   |   | j| jd | j| j¡S r$   )Ú	__class__r>   r?   r=   r7   rH   r'   r'   r(   Ú__copy__w  s   ÿzAsyncResult.__copy__c                 C   ó   | j |  ¡ fS r$   ©r¾   Ú__reduce_args__rH   r'   r'   r(   Ú
__reduce__|  ó   zAsyncResult.__reduce__c                 C   s   | j | jd d | jfS r$   )r>   r?   r7   rH   r'   r'   r(   rÂ     r²   zAsyncResult.__reduce_args__c                 C   s   | j dur| j  | ¡ dS dS )z9Cancel pending operations when the instance is destroyed.NrP   rH   r'   r'   r(   Ú__del__‚  s   
ÿzAsyncResult.__del__c                 C   s   |   ¡ S r$   )r¦   rH   r'   r'   r(   r¥   ‡  ó   zAsyncResult.graphc                 C   s   | j jS r$   )r?   Úsupports_native_joinrH   r'   r'   r(   rÇ   ‹  rÆ   z AsyncResult.supports_native_joinc                 C   ó   |   ¡  d¡S )Nr†   ©r•   rm   rH   r'   r'   r(   r†   �  ó   zAsyncResult.childrenc                 C   s:   |r|d }|t jv r|  | j |¡¡}|  | ¡ |S |S )Nr“   )r   r‰   Ú
_set_cacher?   Úmeta_from_decodedrA   )rD   Úmetarˆ   Údr'   r'   r(   Ú_maybe_set_cache“  s   

zAsyncResult._maybe_set_cachec                 C   s$   | j d u r|  | j | j¡¡S | j S r$   )rB   rÏ   r?   Úget_task_metar>   rH   r'   r'   r(   r•   œ  s   
zAsyncResult._get_task_metac                 K   s   t |  ¡ gƒS r$   )Úiterr•   ©rD   rx   r'   r'   r(   Ú
_iter_meta¡  rÄ   zAsyncResult._iter_metac                    s.   |  d¡}|r‡ fdd„|D ƒ|d< |ˆ _|S )Nr†   c                    s   g | ]}t |ˆ jƒ‘qS r'   )r#   r=   r}   rH   r'   r(   Ú
<listcomp>§  ó    ÿz*AsyncResult._set_cache.<locals>.<listcomp>)rm   rB   )rD   rÎ   r†   r'   rH   r(   rË   ¤  s   


ÿzAsyncResult._set_cachec                 C   ó   |   ¡ d S )zÕTask return value.

        Note:
            When the task has been executed, this contains the return value.
            If the task raised an exception, this will be the exception
            instance.
        rS   ©r•   rH   r'   r'   r(   rS   ­  s   	zAsyncResult.resultc                 C   rÈ   )z#Get the traceback of a failed task.r”   rÉ   rH   r'   r'   r(   r”   ¹  s   zAsyncResult.tracebackc                 C   rÖ   )a   The tasks current state.

        Possible values includes:

            *PENDING*

                The task is waiting for execution.

            *STARTED*

                The task has been started.

            *RETRY*

                The task is to be retried, possibly because of failure.

            *FAILURE*

                The task raised an exception, or has exceeded the retry limit.
                The :attr:`result` attribute then contains the
                exception raised by the task.

            *SUCCESS*

                The task executed successfully.  The :attr:`result` attribute
                then contains the tasks return value.
        r“   r×   rH   r'   r'   r(   rˆ   ¾  s   zAsyncResult.statec                 C   ó   | j S )zCompat. alias to :attr:`id`.©r>   rH   r'   r'   r(   Útask_idÞ  ó   zAsyncResult.task_idc                 C   ó
   || _ d S r$   rÙ   )rD   r>   r'   r'   r(   rÚ   ã  r¬   c                 C   rÈ   )NÚnamerÉ   rH   r'   r'   r(   rÝ   ç  rÊ   zAsyncResult.namec                 C   rÈ   )Nr’   rÉ   rH   r'   r'   r(   r’   ë  rÊ   zAsyncResult.argsc                 C   rÈ   )Nrx   rÉ   rH   r'   r'   r(   rx   ï  rÊ   zAsyncResult.kwargsc                 C   rÈ   )NÚworkerrÉ   rH   r'   r'   r(   rÞ   ó  rÊ   zAsyncResult.workerc                 C   s*   |   ¡  d¡}|rt|tjƒst|ƒS |S )zUTC date and time.Ú	date_done)r•   rm   r³   Údatetimer   )rD   rß   r'   r'   r(   rß   ÷  s   zAsyncResult.date_donec                 C   rÈ   )NÚretriesrÉ   rH   r'   r'   r(   rá   ÿ  rÊ   zAsyncResult.retriesc                 C   rÈ   )NÚqueuerÉ   rH   r'   r'   r(   râ     rÊ   zAsyncResult.queue)NNNNr+   ©NFNFN)F)TN)FN)Ar3   r4   r5   r6   r=   r   r>   r?   rF   ÚpropertyrI   ÚsetterrL   r@   rT   rU   r\   r   rj   rk   rm   r^   re   rp   r{   r|   rw   r„   rŒ   r�   r�   rf   Úmaybe_reraiser–   r¦   r«   r°   r±   r·   r¼   r¿   rÃ   rÂ   rÅ   r   r¥   rÇ   r†   rÏ   r•   rÓ   rË   rS   Úinfor”   rˆ   r“   rÚ   rÝ   r’   rx   rÞ   rß   rá   râ   r'   r'   r'   r(   r   I   s¬    
þ



ÿ
üE
1

	




		
	









r   c                   @   sp  e Zd ZdZdZdZdId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dJdd„ZeZdd„ Zdd„ Zdd„ Zdd„ Z		dKd!d"„Zd#d$„ Zd%d&„ Ze d'd(¡dLd*d+„ƒZ	)		dMd,d-„Z	)		dMd.d/„ZdNd0d1„Z		dOd2d3„Z				dPd4d5„Zd6d7„ Z d8d9„ Z!d:d;„ Z"d<d=„ Z#d>d?„ Z$d@dA„ Z%e&dBdC„ ƒZ'e&dDdE„ ƒZ(e(j)dFdE„ ƒZ(e&dGdH„ ƒZ*dS )Qr    zpA collection of results.

    Arguments:
        results (Sequence[AsyncResult]): List of result instances.
    Nc                 K   sL   || _ || _t| fd�| _|pt|ƒ| _| jr$| j t| jdd�¡ d S d S )N)r’   Tr8   )Ú_appÚresultsr   rA   r
   Ú_on_fullrL   Ú	_on_ready)rD   ré   r=   Úready_barrierrx   r'   r'   r(   rF     s   ÿzResultSet.__init__c                 C   s4   || j vr| j  |¡ | jr| j |¡ dS dS dS )zvAdd :class:`AsyncResult` as a new member of the set.

        Does nothing if the result is already a member.
        N)ré   Úappendrê   ÚaddrR   r'   r'   r(   rî     s   
ýzResultSet.addc                 C   s   | j jr
|  ¡  d S d S r$   )r?   Úis_asyncrA   rH   r'   r'   r(   rë   (  s   ÿzResultSet._on_readyc                 C   s@   t |tƒr| j |¡}z	| j |¡ W dS  ty   t|ƒ‚w )z~Remove result from the set; it must be a member.

        Raises:
            KeyError: if the result isn't a member.
        N)r³   r   r=   r   ré   Úremover:   ÚKeyErrorrR   r'   r'   r(   rð   ,  s   
ÿzResultSet.removec                 C   s&   z|   |¡ W dS  ty   Y dS w )zbRemove result from the set if it is a member.

        Does nothing if it's not a member.
        N)rð   rñ   rR   r'   r'   r(   Údiscard9  s
   ÿzResultSet.discardc                    s   ˆ j  ‡ fdd„|D ƒ¡ dS )z Extend from iterable of results.c                 3   s   � | ]
}|ˆ j vr|V  qd S r$   ©ré   ©r~   ÚrrH   r'   r(   r�   E  s   € z#ResultSet.update.<locals>.<genexpr>N)ré   r…   )rD   ré   r'   rH   r(   ÚupdateC  s   zResultSet.updatec                 C   s   g | j dd…< dS )z!Remove all results from this set.Nró   rH   r'   r'   r(   ÚclearG  s   zResultSet.clearc                 C   ó   t dd„ | jD ƒƒS )z²Return true if all tasks successful.

        Returns:
            bool: true if all of the tasks finished
                successfully (i.e. didn't raise an exception).
        c                 s   ó   � | ]}|  ¡ V  qd S r$   )rŒ   ©r~   rS   r'   r'   r(   r�   R  r‚   z'ResultSet.successful.<locals>.<genexpr>©Úallré   rH   r'   r'   r(   rŒ   K  ó   zResultSet.successfulc                 C   rø   )z¡Return true if any of the tasks failed.

        Returns:
            bool: true if one of the tasks failed.
                (i.e., raised an exception)
        c                 s   rù   r$   )r�   rú   r'   r'   r(   r�   [  r‚   z#ResultSet.failed.<locals>.<genexpr>©Úanyré   rH   r'   r'   r(   r�   T  rý   zResultSet.failedTc                 C   s   | j D ]	}|j||d� qd S )N)rN   rc   )ré   rf   )rD   rN   rc   rS   r'   r'   r(   rf   ]  s   
ÿzResultSet.maybe_throwc                 C   rø   )z¦Return true if any of the tasks are incomplete.

        Returns:
            bool: true if one of the tasks are still
                waiting for execution.
        c                 s   s   � | ]}|  ¡  V  qd S r$   ©r„   rú   r'   r'   r(   r�   i  s   € z$ResultSet.waiting.<locals>.<genexpr>rþ   rH   r'   r'   r(   Úwaitingb  rý   zResultSet.waitingc                 C   rø   )z˜Did all of the tasks complete? (either by success of failure).

        Returns:
            bool: true if all of the tasks have been executed.
        c                 s   rù   r$   r   rú   r'   r'   r(   r�   q  r‚   z"ResultSet.ready.<locals>.<genexpr>rû   rH   r'   r'   r(   r„   k  ó   zResultSet.readyc                 C   rø   )zaTask completion count.

        Returns:
            int: the number of tasks completed.
        c                 s   s   � | ]	}t | ¡ ƒV  qd S r$   )ÚintrŒ   rú   r'   r'   r(   r�   y  s   € z,ResultSet.completed_count.<locals>.<genexpr>)Úsumré   rH   r'   r'   r(   Úcompleted_counts  r  zResultSet.completed_countc                 C   s   | j D ]}| ¡  qdS )z?Forget about (and possible remove the result of) all the tasks.N)ré   rU   rR   r'   r'   r(   rU   {  s   

ÿzResultSet.forgetFc                 C   s*   | j jjdd„ | jD ƒ|||||d� dS )a[  Send revoke signal to all workers for all tasks in the set.

        Arguments:
            terminate (bool): Also terminate the process currently working
                on the task (if any).
            signal (str): Name of signal to send to process if terminate.
                Default is TERM.
            wait (bool): Wait for replies from worker.
                The ``timeout`` argument specifies the number of seconds
                to wait.  Disabled by default.
            timeout (float): Time in seconds to wait for replies when
                the ``wait`` argument is enabled.
        c                 S   s   g | ]}|j ‘qS r'   rÙ   rô   r'   r'   r(   rÔ   �  ó    z$ResultSet.revoke.<locals>.<listcomp>)rV   rZ   rW   rX   rY   N)r=   r[   r\   ré   r]   r'   r'   r(   r\   €  s   
þzResultSet.revokec                 C   r§   r$   )rÑ   ré   rH   r'   r'   r(   Ú__iter__“  ó   
zResultSet.__iter__c                 C   s
   | j | S )z`res[i] -> res.results[i]`.ró   )rD   Úindexr'   r'   r(   Ú__getitem__–  r¬   zResultSet.__getitem__z4.0z5.0r_   c           	      c   sÀ   � d}t dd„ | jD ƒƒ}|r^tƒ }t|ƒD ]%\}}| ¡ r0|j|o%|| |d�V  | |¡ q|jjr;t	 
|jj¡ q|D ]}| |d¡ q>t	 
|¡ ||7 }|rZ||krZtdƒ‚|sdS dS )z<Deprecated method, use :meth:`get` with a callback argument.ç        c                 s   s   � | ]
}|j t|ƒfV  qd S r$   )r>   r   rú   r'   r'   r(   r�   ž  s   € ÿz$ResultSet.iterate.<locals>.<genexpr>)rZ   rc   NzThe operation timed out)r   ré   Úsetr   r„   rm   rî   r?   Úsubpolling_intervalÚtimeÚsleepÚpopr   )	rD   rZ   rc   r`   Úelapsedré   ÚremovedrÚ   rS   r'   r'   r(   Úiterateš  s.   €ÿÿ€
ñzResultSet.iteratec	           	   
   C   s&   | j r| jn| j||||||||d�S )zÆSee :meth:`join`.

        This is here for API compatibility with :class:`AsyncResult`,
        in addition it uses :meth:`join_native` if available for the
        current result backend.
        )rZ   rc   r`   rN   rb   rd   ri   ra   )rÇ   Újoin_nativeÚjoin)	rD   rZ   rc   r`   rN   rb   rd   ri   ra   r'   r'   r(   rm   ²  s   	üzResultSet.getc	              	   C   s�   |rt ƒ  tƒ }	d}
|durtdƒ‚g }| jD ].}d}
|r,|tƒ |	  }
|
dkr,tdƒ‚|j|
|||||d�}|r@||j|ƒ q| |¡ q|S )aÓ  Gather the results of all tasks as a list in order.

        Note:
            This can be an expensive operation for result store
            backends that must resort to polling (e.g., database).

            You should consider using :meth:`join_native` if your backend
            supports it.

        Warning:
            Waiting for tasks within a task may lead to deadlocks.
            Please see :ref:`task-synchronous-subtasks`.

        Arguments:
            timeout (float): The number of seconds to wait for results
                before the operation times out.
            propagate (bool): If any of the tasks raises an exception,
                the exception will be re-raised when this flag is set.
            interval (float): Time to wait (in seconds) before retrying to
                retrieve a result from the set.  Note that this does not have
                any effect when using the amqp result store backend,
                as it does not use polling.
            callback (Callable): Optional callback to be called for every
                result received.  Must have signature ``(task_id, value)``
                No results will be returned by this function if a callback
                is specified.  The order of results is also arbitrary when a
                callback is used.  To get access to the result object for
                a particular id you'll have to generate an index first:
                ``index = {r.id: r for r in gres.results.values()}``
                Or you can create new result objects on the fly:
                ``result = app.AsyncResult(task_id)`` (both will
                take advantage of the backend cache anyway).
            no_ack (bool): Automatic message acknowledgment (Note that if this
                is set to :const:`False` then the messages
                *will not be acknowledged*).
            disable_sync_subtasks (bool): Disable tasks to wait for sub tasks
                this is the default configuration. CAUTION do not enable this
                unless you must.

        Raises:
            celery.exceptions.TimeoutError: if ``timeout`` isn't
                :const:`None` and the operation takes longer than ``timeout``
                seconds.
        Nz,Backend does not support on_message callbackr  zjoin operation timed out)rZ   rc   r`   rb   ra   ri   )r)   r   r   ré   r   rm   r>   rí   )rD   rZ   rc   r`   rN   rb   rd   ri   ra   Ú
time_startÚ	remainingré   rS   rJ   r'   r'   r(   r  Â  s0   /ÿ
ýzResultSet.joinc                 C   ó   | j  ||¡S r$   ©rA   rL   rM   r'   r'   r(   rL     rÄ   zResultSet.thenc                 C   s   | j j| |||||d�S )a0  Backend optimized version of :meth:`iterate`.

        .. versionadded:: 2.2

        Note that this does not support collecting the results
        for different task types using different backends.

        This is currently only supported by the amqp, Redis and cache
        result backends.
        )rZ   r`   rb   rd   ra   )r?   Úiter_native)rD   rZ   r`   rb   rd   ra   r'   r'   r(   r    s
   ýzResultSet.iter_nativec	                 C   sÆ   |rt ƒ  |r	dn	dd„ t| jƒD ƒ}	|rdn
dd„ tt| ƒƒD ƒ}
|  |||||¡D ]5\}}t|tƒrCg }|D ]	}| | 	¡ ¡ q8n|d }|rR|d t
jv rR|‚|rZ|||ƒ q+||
|	| < q+|
S )a-  Backend optimized version of :meth:`join`.

        .. versionadded:: 2.2

        Note that this does not support collecting the results
        for different task types using different backends.

        This is currently only supported by the amqp, Redis and cache
        result backends.
        Nc                 S   s   i | ]\}}|j |“qS r'   rÙ   )r~   ÚirS   r'   r'   r(   Ú
<dictcomp>1  rÕ   z)ResultSet.join_native.<locals>.<dictcomp>c                 S   s   g | ]}d ‘qS r$   r'   )r~   ry   r'   r'   r(   rÔ   4  s    z)ResultSet.join_native.<locals>.<listcomp>rS   r“   )r)   Ú	enumerateré   r   Úlenr  r³   ro   rí   rm   r   rk   )rD   rZ   rc   r`   rN   rb   rd   ra   ri   Úorder_indexÚaccrÚ   rÍ   rJ   Úchildren_resultr'   r'   r(   r  !  s*   ÿ
ÿ
ÿzResultSet.join_nativec                 K   s.   dd„ | j jdd„ | jD ƒfddi|¤ŽD ƒS )Nc                 s   s   � | ]\}}|V  qd S r$   r'   )r~   ry   rÍ   r'   r'   r(   r�   F  r‚   z'ResultSet._iter_meta.<locals>.<genexpr>c                 S   s   h | ]}|j ’qS r'   rÙ   rô   r'   r'   r(   Ú	<setcomp>G  r  z'ResultSet._iter_meta.<locals>.<setcomp>Úmax_iterationsr   )r?   Úget_manyré   rÒ   r'   r'   r(   rÓ   E  s   ÿÿ
ÿzResultSet._iter_metac                 C   s   dd„ | j D ƒS )Nc                 s   s.   � | ]}|j  |j¡r|jtjv r|V  qd S r$   )r?   Ú	is_cachedr>   rˆ   r   rk   )r~   r»   r'   r'   r(   r�   K  s   € ÿþþz0ResultSet._failed_join_report.<locals>.<genexpr>ró   rH   r'   r'   r(   Ú_failed_join_reportJ  ó   zResultSet._failed_join_reportc                 C   r§   r$   )r  ré   rH   r'   r'   r(   Ú__len__O  r  zResultSet.__len__c                 C   s   t |tƒr|j| jkS tS r$   )r³   r    ré   r´   rµ   r'   r'   r(   r·   R  s   
zResultSet.__eq__c                 C   r¸   r1   r¹   rº   r'   r'   r(   r¼   W  r½   zResultSet.__ne__c                 C   s$   d  t| ƒjd dd„ | jD ƒ¡¡S )Nz<{0}: [{1}]>ú, c                 s   ó   � | ]}|j V  qd S r$   rÙ   rô   r'   r'   r(   r�   ]  ó   € z%ResultSet.__repr__.<locals>.<genexpr>)r;   r<   r3   r  ré   rH   r'   r'   r(   r±   [  s   ÿzResultSet.__repr__c                 C   s$   z| j d jW S  ty   Y d S w ©Nr   )ré   rÇ   Ú
IndexErrorrH   r'   r'   r(   rÇ   _  s
   ÿzResultSet.supports_native_joinc                 C   s,   | j d u r| jr| jd jnt ¡ | _ | j S r,  )rè   ré   r=   r   Ú_get_current_objectrH   r'   r'   r(   r=   f  s
   
ÿzResultSet.appc                 C   rÜ   r$   )rè   )rD   r=   r'   r'   r(   r=   m  r¬   c                 C   s   | j r| j jS | jd jS r,  )r=   r?   ré   rH   r'   r'   r(   r?   q  s   zResultSet.backend©NNr1   rã   )NTr_   )NTr_   NTNTNr+   )Nr_   TNN)NTr_   NTNNT)+r3   r4   r5   r6   rè   ré   rF   rî   rë   rð   rò   rö   r÷   rŒ   r�   rf   ræ   r  r„   r  rU   r\   r  r
  r   ÚCallabler  rm   r  rL   r  r  rÓ   r&  r(  r·   r¼   r±   rä   rÇ   r=   rå   r?   r'   r'   r'   r(   r      sr    


	
		
ÿ

þ
þ
J
ÿ
ý$


r    c                   @   s¤   e Zd ZdZdZdZd!dd„Zdd„ Zd"dd„Zd"d	d
„Z	dd„ Z
dd„ Zdd„ Ze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d „ƒZdS )$r!   az  Like :class:`ResultSet`, but with an associated id.

    This type is returned by :class:`~celery.group`.

    It enables inspection of the tasks state and return values as
    a single entity.

    Arguments:
        id (str): The id of the group.
        results (Sequence[AsyncResult]): List of result instances.
        parent (ResultBase): Parent result of this group.
    Nc                 K   s$   || _ || _tj| |fi |¤Ž d S r$   )r>   r7   r    rF   )rD   r>   ré   r7   rx   r'   r'   r(   rF   Œ  s   zGroupResult.__init__c                 C   s   | j  | ¡ t | ¡ d S r$   )r?   rQ   r    rë   rH   r'   r'   r(   rë   ‘  s   zGroupResult._on_readyc                 C   s   |p| j j | j| ¡S )zãSave group-result for later retrieval using :meth:`restore`.

        Example:
            >>> def save_and_restore(result):
            ...     result.save()
            ...     result = GroupResult.restore(result.id)
        )r=   r?   Ú
save_groupr>   ©rD   r?   r'   r'   r(   Úsave•  s   zGroupResult.savec                 C   s   |p| j j | j¡ dS )z.Remove this result if it was previously saved.N)r=   r?   Údelete_groupr>   r2  r'   r'   r(   ÚdeleteŸ  s   zGroupResult.deletec                 C   rÀ   r$   rÁ   rH   r'   r'   r(   rÃ   £  rÄ   zGroupResult.__reduce__c                 C   s   | j | jfS r$   )r>   ré   rH   r'   r'   r(   rÂ   ¦  ó   zGroupResult.__reduce_args__c                 C   s   t | jp| jƒS r$   )Úboolr>   ré   rH   r'   r'   r(   Ú__bool__©  r'  zGroupResult.__bool__c                 C   sF   t |tƒr|j| jko|j| jko|j| jkS t |tƒr!|| jkS tS r$   )r³   r!   r>   ré   r7   r   r´   rµ   r'   r'   r(   r·   ­  s   

ÿ
ý

zGroupResult.__eq__c                 C   r¸   r1   r¹   rº   r'   r'   r(   r¼   ¸  r½   zGroupResult.__ne__c                 C   s(   d  t| ƒj| jd dd„ | jD ƒ¡¡S )Nz<{0}: {1} [{2}]>r)  c                 s   r*  r$   rÙ   rô   r'   r'   r(   r�   ¿  r+  z'GroupResult.__repr__.<locals>.<genexpr>)r;   r<   r3   r>   r  ré   rH   r'   r'   r(   r±   ¼  s   þzGroupResult.__repr__c                 C   r§   r¨   r©   rH   r'   r'   r(   r«   Â  r¬   zGroupResult.__str__c                 C   r§   r­   r®   rH   r'   r'   r(   r°   Æ  r¬   zGroupResult.__hash__c                 C   s&   | j | jo	| j ¡ fdd„ | jD ƒfS )Nc                 S   s   g | ]}|  ¡ ‘qS r'   )rT   rô   r'   r'   r(   rÔ   Í  s    z(GroupResult.as_tuple.<locals>.<listcomp>)r>   r7   rT   ré   rH   r'   r'   r(   rT   Ê  s   þzGroupResult.as_tuplec                 C   rØ   r$   ró   rH   r'   r'   r(   r†   Ð  s   zGroupResult.childrenc                 C   s.   |pt | jtƒs| jnt}|p|j}| |¡S )z&Restore previously saved group result.)r³   r=   rä   r   r?   Úrestore_group)Úclsr>   r?   r=   r'   r'   r(   ÚrestoreÔ  s
   ÿ

zGroupResult.restore)NNNr$   r/  )r3   r4   r5   r6   r>   ré   rF   rë   r3  r5  rÃ   rÂ   r8  Ú__nonzero__r·   r¼   r±   r«   r°   rT   rä   r†   Úclassmethodr;  r'   r'   r'   r(   r!   v  s,    




r!   c                   @   s¶   e Zd ZdZd%dd„Zd&dd„Zdd	„ Zd
d„ Zdd„ Zdd„ Z	dd„ Z
		d'dd„ZeZdd„ Zdd„ Zdd„ Zedd„ ƒZedd„ ƒZedd „ ƒZeZed!d"„ ƒZed#d$„ ƒZdS )(r"   z.Result that we know has already been executed.Nc                 C   s.   || _ || _|| _|| _tƒ | _|  | ¡ d S r$   )r>   Ú_resultÚ_stateÚ
_tracebackr   rA   )rD   r>   Ú	ret_valuerˆ   r”   r'   r'   r(   rF   ã  s   zEagerResult.__init__Fc                 C   r  r$   r  rM   r'   r'   r(   rL   í  rÄ   zEagerResult.thenc                 C   rØ   r$   )rB   rH   r'   r'   r(   r•   ð  s   zEagerResult._get_task_metac                 C   rÀ   r$   rÁ   rH   r'   r'   r(   rÃ   ó  rÄ   zEagerResult.__reduce__c                 C   s   | j | j| j| jfS r$   ©r>   r>  r?  r@  rH   r'   r'   r(   rÂ   ö  r²   zEagerResult.__reduce_args__c                 C   s   |   ¡ \}}||Ž S r$   )rÃ   )rD   r:  r’   r'   r'   r(   r¿   ù  s   zEagerResult.__copy__c                 C   ó   dS r1   r'   rH   r'   r'   r(   r„   ý  ó   zEagerResult.readyTc                 K   sN   |rt ƒ  |  ¡ r| jS | jtjv r%|r"t| jtƒr| j‚t| jƒ‚| jS d S r$   )r)   rŒ   rS   rˆ   r   rk   r³   Ú	Exception)rD   rZ   rc   ri   rx   r'   r'   r(   rm      s   
ÿÿüzEagerResult.getc                 C   s   d S r$   r'   rH   r'   r'   r(   rU     rD  zEagerResult.forgetc                 O   s   t j| _d S r$   )r   ÚREVOKEDr?  r‘   r'   r'   r(   r\     r6  zEagerResult.revokec                 C   s
   d  | ¡S )Nz<EagerResult: {0.id}>)r;   rH   r'   r'   r(   r±     r  zEagerResult.__repr__c                 C   s   | j | j| j| jdœS )N)rÚ   rS   r“   r”   rB  rH   r'   r'   r(   rB     s
   üzEagerResult._cachec                 C   rØ   )zThe tasks return value.)r>  rH   r'   r'   r(   rS      rÛ   zEagerResult.resultc                 C   rØ   )zThe tasks state.)r?  rH   r'   r'   r(   rˆ   %  rÛ   zEagerResult.statec                 C   rØ   )z!The traceback if the task failed.)r@  rH   r'   r'   r(   r”   +  rÛ   zEagerResult.tracebackc                 C   rC  r+   r'   rH   r'   r'   r(   rÇ   0  s   z EagerResult.supports_native_joinr$   r+   )NTT)r3   r4   r5   r6   rF   rL   r•   rÃ   rÂ   r¿   r„   rm   r^   rU   r\   r±   rä   rB   rS   rˆ   r“   r”   rÇ   r'   r'   r'   r(   r"   Þ  s6    



ÿ



r"   c                    s‚   t ˆ ƒ‰ ˆ j}t| tƒs?| \}}t|ttfƒr|n|df\}}|r&t|ˆ ƒ}|dur9ˆ j|‡ fdd„|D ƒ|d�S |||d�S | S )zDeserialize result from tuple.Nc                    s   g | ]}t |ˆ ƒ‘qS r'   )r#   r}   ©r=   r'   r(   rÔ   C  s    z%result_from_tuple.<locals>.<listcomp>rs   )r   r   r³   r   ro   Útupler#   r!   )rõ   r=   ÚResultr»   Únodesr>   r7   r'   rG  r(   r#   5  s   

þr#   r$   ):r6   Ú
__future__r   r   rà   r  Úcollectionsr   r   Ú
contextlibr   r   Úkombu.utils.objectsr   Úviner	   r
   r   Ú r   r   r?  r   r   r=   r   Ú
exceptionsr   r   r   Úfiver   r   r   r   r   Úutilsr   Úutils.graphr   r   Úutils.iso8601r   r™   ÚImportErrorÚ__all__r&   r)   r/   r2   Úobjectr   Úregisterr   r    r!   r"   r#   r'   r'   r'   r(   Ú<module>   s`   ÿ
	
	   @  nfU