o
    ›¨Êh‹  ã                   @   st  d 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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 zddl Z W n e!yk   dZ Y nw dZ"dZ#dd„ Z$edd„ ƒZ%edd„ ƒZ&G dd„ dƒZ'ej(G dd„ de'ƒƒZ)ej(G dd„ de'ƒƒZ*ej(G dd„ de*ƒƒZ+ej(G dd „ d e)ƒƒZ,d#d!d"„Z-dS )$z3Task results/state and results for groups of tasks.é    N)Údeque)Úcontextmanager)Úproxy)Úisoparse)Úcached_property)ÚThenableÚbarrierÚpromiseé   )Úcurrent_appÚstates)Ú_set_task_join_will_blockÚtask_join_will_block)Úapp_or_default)ÚImproperlyConfiguredÚIncompleteStreamÚTimeoutError)ÚDependencyGraphÚGraphFormatter)Ú
ResultBaseÚAsyncResultÚ	ResultSetÚGroupResultÚEagerResultÚresult_from_tuplezˆNever call result.get() within a task!
See https://docs.celeryq.dev/en/latest/userguide/tasks.html#avoid-launching-synchronous-subtasks
c                   C   s   t ƒ rttƒ‚d S ©N)r   ÚRuntimeErrorÚE_WOULDBLOCK© r   r   ú?/var/www/html/env/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_result3   r'   r)   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   =   s    r   c                   @   sB  e Zd ZdZdZeZdZdZ			dhdd„Ze	dd„ ƒZ
e
jdd„ ƒZ
did	d
„Zdd„ Zdd„ Zdd„ Zdd„ Z		djdd„Z		djd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dkdd „Zd!d"„ Zdkd#d$„Zd%d&„ Zd'd(„ Zd)d*„ Zd+d,„ Z dld-d.„Z!e!Z"d/d0„ Z#dmd1d2„Z$d3d4„ Z%d5d6„ Z&d7d8„ Z'd9d:„ Z(d;d<„ Z)d=d>„ Z*d?d@„ Z+dAdB„ Z,e-dCdD„ ƒZ.e	dEdF„ ƒZ/e	dGdH„ ƒZ0dIdJ„ Z1dKdL„ Z2dMdN„ Z3dOdP„ Z4e	dQdR„ ƒZ5e5Z6e	dSdT„ ƒZ7e	dUdV„ ƒZ8e8Z9e	dWdX„ ƒZ:e:jdYdX„ ƒZ:e	dZd[„ ƒZ;e	d\d]„ ƒZ<e	d^d_„ ƒZ=e	d`da„ ƒZ>e	dbdc„ ƒZ?e	ddde„ ƒZ@e	dfdg„ ƒZAdS )nr   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 T©ÚweakF)Ú
ValueErrorÚtyper   ÚappÚidÚbackendr.   r	   Ú_on_fulfilledÚon_readyÚ_cacheÚ_ignored)Úselfr4   r5   Ú	task_namer3   r.   r   r   r   Ú__init__X   s   ÿ
zAsyncResult.__init__c                 C   s   t | dƒr| jS dS )z+If True, task result retrieval is disabled.r9   F)Úhasattrr9   ©r:   r   r   r   Úignoredf   s   
zAsyncResult.ignoredc                 C   s
   || _ dS )z%Enable/disable task result retrieval.N)r9   )r:   Úvaluer   r   r   r?   m   s   
Fc                 C   s   | j j| |d� | j ||¡S )Nr/   )r5   Úadd_pending_resultr7   Úthen©r:   ÚcallbackÚon_errorr0   r   r   r   rB   r   s   zAsyncResult.thenc                 C   s   | j  | ¡ |S r   ©r5   Úremove_pending_result©r:   Úresultr   r   r   r6   v   s   zAsyncResult._on_fulfilledc                 C   s   | j }| j|o
| ¡ fd fS r   )r.   r4   Úas_tuple)r:   r.   r   r   r   rJ   z   s   zAsyncResult.as_tuplec                 C   s0   g }| j }| | j¡ |dur| | ¡ ¡ |S )zReturn as a list of task IDs.N)r.   Úappendr4   ÚextendÚas_list)r:   Úresultsr.   r   r   r   rM   ~   s   zAsyncResult.as_listc                 C   s(   d| _ | jr| j ¡  | j | j¡ dS )z/Forget the result of this task and its parents.N)r8   r.   Úforgetr5   r4   r>   r   r   r   rO   ‡   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)r3   ÚcontrolÚrevoker4   ©r:   rQ   rR   rS   ÚwaitrU   r   r   r   rW   Ž   s   
þzAsyncResult.revokec                 C   s   | j jj||||||d� dS )a7  Send revoke signal to all workers only for tasks with matching headers values.

        Any worker receiving the task, or having reserved the
        task, *must* ignore it.
        All header fields *must* match.

        Arguments:
            headers (dict[str, Union(str, list)]): Headers to match when revoking tasks.
            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.
        rP   N)r3   rV   Úrevoke_by_stamped_headers)r:   ÚheadersrQ   rR   rS   rY   rU   r   r   r   rZ   ¤   s   
þz%AsyncResult.revoke_by_stamped_headersTç      à?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. This is the setting for the publisher
                (celery client) and is different from `timeout` parameter of
                `@app.task`, which is the setting for the worker. The task
                isn't terminated even if timeout occurs.
            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.
        NTr/   )rD   )rU   ÚintervalÚon_intervalÚno_ackÚ	propagaterD   Ú
on_message)r?   r    r	   r.   Ú_maybe_reraise_parent_errorrB   r8   Úmaybe_throwrI   r5   rA   Úwait_for_pending)r:   rU   r`   r]   r_   Úfollow_parentsrD   ra   r^   Ú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Ú_parentsrc   ©r:   Únoder   r   r   rb     s   
ÿz'AsyncResult._maybe_reraise_parent_errorc                 c   s$   � | j }|r|V  |j }|sd S d S r   ©r.   rn   r   r   r   rm   
  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   ©Úiterdepsrj   )r:   rr   ÚkwargsÚ_ÚRr   r   r   Úcollect  s   €.ÿzAsyncResult.collectc                 C   s"   d }|   ¡ D ]\}}| ¡ }q|S r   rs   )r:   r@   rv   rw   r   r   r   Úget_leafA  s   
zAsyncResult.get_leafc                 #   sn   � t d | fgƒ}| }|r5| ¡ \}‰ |ˆ fV  ˆ  ¡ r,| ‡ fdd„ˆ jp'g D ƒ¡ n|r1tƒ ‚|sd S d S )Nc                 3   s   � | ]}ˆ |fV  qd S r   r   ©Ú.0Úchild©ro   r   r   Ú	<genexpr>P  ó   € z'AsyncResult.iterdeps.<locals>.<genexpr>)r   ÚpopleftÚreadyrL   Úchildrenr   )r:   rr   ÚstackÚis_incomplete_streamr.   r   r}   r   rt   G  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.
        )Ústater5   ÚREADY_STATESr>   r   r   r   r�   U  s   zAsyncResult.readyc                 C   ó   | j tjkS )z7Return :const:`True` if the task executed successfully.)r…   r   ÚSUCCESSr>   r   r   r   Ú
successful]  ó   zAsyncResult.successfulc                 C   r‡   )z(Return :const:`True` if the task failed.)r…   r   ÚFAILUREr>   r   r   r   Úfaileda  rŠ   zAsyncResult.failedc                 O   s   | j j|i |¤Ž d S r   )r7   Úthrow©r:   Úargsru   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ÚstatusrI   Ú	traceback)r8   Ú_get_task_metarj   r   rh   r�   Ú_to_remote_tracebackr4   )r:   r`   rD   Úcacher…   r@   Útbr   r   r   rc   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   )Útblibr3   ÚconfÚtask_remote_tracebacksÚ	TracebackÚfrom_stringÚas_traceback)r:   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)Ú	formatterrq   )r   r   r4   rt   Úadd_arcÚadd_edge)r:   rr   rŸ   Úgraphr.   ro   r   r   r   Úbuild_graphw  s   ÿ
€zAsyncResult.build_graphc                 C   ó
   t | jƒS ©z`str(self) -> self.id`.©Ústrr4   r>   r   r   r   Ú__str__�  ó   
zAsyncResult.__str__c                 C   r¤   ©z`hash(self) -> hash(self.id)`.©Úhashr4   r>   r   r   r   Ú__hash__…  r©   zAsyncResult.__hash__c                 C   s   dt | ƒj› d| j› d�S )Nú<ú: ú>)r2   r*   r4   r>   r   r   r   Ú__repr__‰  s   zAsyncResult.__repr__c                 C   s.   t |tƒr|j| jkS t |tƒr|| jkS tS r   )Ú
isinstancer   r4   r§   ÚNotImplemented©r:   Úotherr   r   r   Ú__eq__Œ  s
   


zAsyncResult.__eq__c                 C   s   |   | j| jd | j| j¡S r   )Ú	__class__r4   r5   r3   r.   r>   r   r   r   Ú__copy__“  s   ÿzAsyncResult.__copy__c                 C   ó   | j |  ¡ fS r   ©r·   Ú__reduce_args__r>   r   r   r   Ú
__reduce__˜  ó   zAsyncResult.__reduce__c                 C   s   | j | jd d | jfS r   )r4   r5   r.   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.NrF   r>   r   r   r   Ú__del__ž  s   
ÿzAsyncResult.__del__c                 C   s   |   ¡ S r   )r£   r>   r   r   r   r¢   £  ó   zAsyncResult.graphc                 C   s   | j jS r   )r5   Úsupports_native_joinr>   r   r   r   rÁ   §  rÀ   z AsyncResult.supports_native_joinc                 C   ó   |   ¡  d¡S )Nr‚   ©r’   rj   r>   r   r   r   r‚   «  ó   zAsyncResult.childrenc                 C   s:   |r|d }|t jv r|  | j |¡¡}|  | ¡ |S |S )Nr�   )r   r†   Ú
_set_cacher5   Úmeta_from_decodedr7   )r:   Ú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   )r8   rÉ   r5   Úget_task_metar4   r>   r   r   r   r’   ¸  s   
zAsyncResult._get_task_metac                 K   s   t |  ¡ gƒS r   )Úiterr’   ©r:   ru   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   r3   rz   r>   r   r   Ú
<listcomp>Ã  ó    ÿz*AsyncResult._set_cache.<locals>.<listcomp>)rj   r8   )r:   rÈ   r‚   r   r>   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.
        rI   ©r’   r>   r   r   r   rI   É  s   	zAsyncResult.resultc                 C   rÂ   )z#Get the traceback of a failed task.r‘   rÃ   r>   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Ñ   r>   r   r   r   r…   Ú  s   zAsyncResult.statec                 C   ó   | j S )zCompat. alias to :attr:`id`.©r4   r>   r   r   r   Útask_idú  ó   zAsyncResult.task_idc                 C   ó
   || _ d S r   rÓ   )r:   r4   r   r   r   rÔ   ÿ  r©   c                 C   rÂ   )NÚnamerÃ   r>   r   r   r   r×     rÄ   zAsyncResult.namec                 C   rÂ   )Nr�   rÃ   r>   r   r   r   r�     rÄ   zAsyncResult.argsc                 C   rÂ   )Nru   rÃ   r>   r   r   r   ru     rÄ   zAsyncResult.kwargsc                 C   rÂ   )NÚworkerrÃ   r>   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’   rj   r²   Údatetimer   )r:   rÙ   r   r   r   rÙ     s   zAsyncResult.date_donec                 C   rÂ   )NÚretriesrÃ   r>   r   r   r   rÛ     rÄ   zAsyncResult.retriesc                 C   rÂ   )NÚqueuerÃ   r>   r   r   r   rÜ     rÄ   zAsyncResult.queue)NNNNr"   ©NFNFN)F)TN)FN)Br*   r+   r,   r-   r3   r   r4   r5   r<   Úpropertyr?   ÚsetterrB   r6   rJ   rM   rO   rW   rZ   r   rg   rh   rj   rY   rb   rm   rx   ry   rt   r�   r‰   rŒ   r�   rc   Úmaybe_reraiser“   r£   r¨   r­   r±   r¶   r¸   r¼   r»   r¿   r   r¢   rÁ   r‚   rÉ   r’   rÍ   rÅ   rI   Úinfor‘   r…   r�   rÔ   r×   r�   ru   rØ   rÙ   rÛ   rÜ   r   r   r   r   r   D   s²    
þ


	
ÿ
ÿ
üH
1

	




		
	









r   c                   @   sR  e Zd ZdZdZdZdCd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dDdd„ZeZdd„ Zdd„ Zdd„ Zdd„ Z		dEd!d"„Zd#d$„ Zd%d&„ Z	'		dFd(d)„Z	'		dFd*d+„ZdGd,d-„Z		dHd.d/„Z				dId0d1„Zd2d3„ Zd4d5„ Zd6d7„ Zd8d9„ Z d:d;„ Z!e"d<d=„ ƒZ#e"d>d?„ ƒZ$e$j%d@d?„ ƒZ$e"dAdB„ ƒZ&dS )Jr   zpA collection of results.

    Arguments:
        results (Sequence[AsyncResult]): List of result instances.
    Nc                 K   sP   || _ || _tt| ƒfd�| _|pt|ƒ| _| jr&| j t| jdd�¡ d S d S )N)r�   Tr/   )	Ú_apprN   r	   r   r7   r   Ú_on_fullrB   Ú	_on_ready)r:   rN   r3   Úready_barrierru   r   r   r   r<   1  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)rN   rK   rã   ÚaddrH   r   r   r   ræ   9  s   
ýzResultSet.addc                 C   s   | j jr
|  ¡  d S d S r   )r5   Úis_asyncr7   r>   r   r   r   rä   C  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§   r3   r   rN   Úremover1   ÚKeyErrorrH   r   r   r   rè   G  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é   rH   r   r   r   ÚdiscardT  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   ©rN   ©r{   Úrr>   r   r   r~   `  s   € z#ResultSet.update.<locals>.<genexpr>N)rN   rL   )r:   rN   r   r>   r   Úupdate^  s   zResultSet.updatec                 C   s   g | j dd…< dS )z!Remove all results from this set.Nrë   r>   r   r   r   Úclearb  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{   rI   r   r   r   r~   m  r   z'ResultSet.successful.<locals>.<genexpr>©ÚallrN   r>   r   r   r   r‰   f  ó   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~   v  r   z#ResultSet.failed.<locals>.<genexpr>©ÚanyrN   r>   r   r   r   rŒ   o  rõ   zResultSet.failedTc                 C   s   | j D ]	}|j||d� qd S )N)rD   r`   )rN   rc   )r:   rD   r`   rI   r   r   r   rc   x  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~   „  s   € z$ResultSet.waiting.<locals>.<genexpr>rö   r>   r   r   r   Úwaiting}  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~   Œ  r   z"ResultSet.ready.<locals>.<genexpr>ró   r>   r   r   r   r�   †  s   zResultSet.readyc                 C   rð   )a  Task completion count.

        Note that `complete` means `successful` in this context. In other words, the
        return value of this method is the number of ``successful`` tasks.

        Returns:
            int: the number of complete (i.e. successful) tasks.
        c                 s   s   � | ]	}t | ¡ ƒV  qd S r   )Úintr‰   rò   r   r   r   r~   —  s   € z,ResultSet.completed_count.<locals>.<genexpr>)ÚsumrN   r>   r   r   r   Úcompleted_countŽ  s   	zResultSet.completed_countc                 C   s   | j D ]}| ¡  qdS )z?Forget about (and possible remove the result of) all the tasks.N)rN   rO   rH   r   r   r   rO   ™  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>)rQ   rU   rR   rS   rT   N)r3   rV   rW   rN   rX   r   r   r   rW   ž  s   
þzResultSet.revokec                 C   r¤   r   )rË   rN   r>   r   r   r   Ú__iter__±  ó   
zResultSet.__iter__c                 C   s
   | j | S )z`res[i] -> res.results[i]`.rë   )r:   Úindexr   r   r   Ú__getitem__´  r©   zResultSet.__getitem__r\   c	           	   
   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.
        )rU   r`   r]   rD   r_   ra   rf   r^   )rÁ   Újoin_nativeÚjoin)	r:   rU   r`   r]   rD   r_   ra   rf   r^   r   r   r   rj   ¸  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�}|rB||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 callbackg        zjoin operation timed out)rU   r`   r]   r_   r^   rf   )	r    ÚtimeÚ	monotonicr   rN   r   rj   r4   rK   )r:   rU   r`   r]   rD   r_   ra   rf   r^   Ú
time_startÚ	remainingrN   rI   r@   r   r   r   r  È  s0   /ÿ
ýzResultSet.joinc                 C   ó   | j  ||¡S r   ©r7   rB   rC   r   r   r   rB     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.
        )rU   r]   r_   ra   r^   )r5   Úiter_native)r:   rU   r]   r_   ra   r^   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{   ÚirI   r   r   r   Ú
<dictcomp>7  rÏ   z)ResultSet.join_native.<locals>.<dictcomp>c                 S   s   g | ]}d ‘qS r   r   )r{   rv   r   r   r   rÎ   :  s    z)ResultSet.join_native.<locals>.<listcomp>rI   r�   )r    Ú	enumeraterN   ÚrangeÚlenr
  r²   rl   rK   rj   r   rh   )r:   rU   r`   r]   rD   r_   ra   r^   rf   Úorder_indexÚaccrÔ   rÇ   r@   Ú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{   rv   rÇ   r   r   r   r~   L  r   z'ResultSet._iter_meta.<locals>.<genexpr>c                 S   s   h | ]}|j ’qS r   rÓ   rì   r   r   r   Ú	<setcomp>M  rý   z'ResultSet._iter_meta.<locals>.<setcomp>Úmax_iterationsr
   )r5   Úget_manyrN   rÌ   r   r   r   rÍ   K  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   )r5   Ú	is_cachedr4   r…   r   rh   )r{   Úresr   r   r   r~   Q  s   € ÿþþz0ResultSet._failed_join_report.<locals>.<genexpr>rë   r>   r   r   r   Ú_failed_join_reportP  ó   zResultSet._failed_join_reportc                 C   r¤   r   )r  rN   r>   r   r   r   Ú__len__U  rÿ   zResultSet.__len__c                 C   s   t |tƒr|j| jkS tS r   )r²   r   rN   r³   r´   r   r   r   r¶   X  s   
zResultSet.__eq__c                 C   s*   dt | ƒj› dd dd„ | jD ƒ¡› d�S )Nr®   z: [ú, c                 s   ó   � | ]}|j V  qd S r   rÓ   rì   r   r   r   r~   ^  ó   € z%ResultSet.__repr__.<locals>.<genexpr>ú]>)r2   r*   r  rN   r>   r   r   r   r±   ]  s   *zResultSet.__repr__c                 C   s$   z| j d jW S  ty   Y d S w ©Nr   )rN   rÁ   Ú
IndexErrorr>   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â   rN   r3   r   Ú_get_current_objectr>   r   r   r   r3   g  s
   
ÿzResultSet.appc                 C   rÖ   r   )râ   )r:   r3   r   r   r   r3   n  r©   c                 C   s   | j r| j jS | jd jS r  )r3   r5   rN   r>   r   r   r   r5   r  s   zResultSet.backend©NNr(   rÝ   )NTr\   NTNTNr"   )Nr\   TNN)NTr\   NTNNT)'r*   r+   r,   r-   râ   rN   r<   ræ   rä   rè   rê   rî   rï   r‰   rŒ   rc   rà   rù   r�   rü   rO   rW   rþ   r  rj   r  rB   r
  r  rÍ   r  r  r¶   r±   rÞ   rÁ   r3   rß   r5   r   r   r   r   r   $  sl    


	
		
ÿ
þ
þ
J
ÿ
ý$


r   c                       s¨   e Zd ZdZdZdZd‡ fdd„	Z‡ f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edd„ ƒZed!dd„ƒZ‡  Z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                    s$   || _ || _tƒ j|fi |¤Ž d S r   )r4   r.   Úsuperr<   )r:   r4   rN   r.   ru   ©r·   r   r   r<   Œ  s   zGroupResult.__init__c                    s   | j  | ¡ tƒ  ¡  d S r   )r5   rG   r#  rä   r>   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)
        )r3   r5   Ú
save_groupr4   ©r:   r5   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)r3   r5   Údelete_groupr4   r&  r   r   r   ÚdeleteŸ  s   zGroupResult.deletec                 C   r¹   r   rº   r>   r   r   r   r¼   £  r½   zGroupResult.__reduce__c                 C   s   | j | jfS r   )r4   rN   r>   r   r   r   r»   ¦  ó   zGroupResult.__reduce_args__c                 C   s   t | jp| jƒS r   )Úboolr4   rN   r>   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   r4   rN   r.   r§   r³   r´   r   r   r   r¶   ­  s   

ÿ
ý

zGroupResult.__eq__c              	   C   s2   dt | ƒj› d| j› dd dd„ | jD ƒ¡› d�S )Nr®   r¯   z [r  c                 s   r  r   rÓ   rì   r   r   r   r~   ¹  r  z'GroupResult.__repr__.<locals>.<genexpr>r  )r2   r*   r4   r  rN   r>   r   r   r   r±   ¸  s   2zGroupResult.__repr__c                 C   r¤   r¥   r¦   r>   r   r   r   r¨   »  r©   zGroupResult.__str__c                 C   r¤   rª   r«   r>   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   )rJ   rì   r   r   r   rÎ   Æ  s    z(GroupResult.as_tuple.<locals>.<listcomp>)r4   r.   rJ   rN   r>   r   r   r   rJ   Ã  s   þzGroupResult.as_tuplec                 C   rÒ   r   rë   r>   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²   r3   rÞ   r   r5   Úrestore_group)Úclsr4   r5   r3   r   r   r   ÚrestoreÍ  s
   ÿ

zGroupResult.restore)NNNr   r"  )r*   r+   r,   r-   r4   rN   r<   rä   r'  r)  r¼   r»   r,  Ú__nonzero__r¶   r±   r¨   r­   rJ   rÞ   r‚   Úclassmethodr/  Ú__classcell__r   r   r$  r   r   w  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   s4   || _ || _|| _|| _|| _tƒ | _|  | ¡ d S r   )r4   Ú_resultÚ_stateÚ
_tracebackÚ_namer	   r7   )r:   r4   Ú	ret_valuer…   r‘   r×   r   r   r   r<   Û  s   zEagerResult.__init__Fc                 C   r  r   r	  rC   r   r   r   rB   æ  r½   zEagerResult.thenc                 C   rÒ   r   )r8   r>   r   r   r   r’   é  s   zEagerResult._get_task_metac                 C   r¹   r   rº   r>   r   r   r   r¼   ì  r½   zEagerResult.__reduce__c                 C   s   | j | j| j| jfS r   )r4   r3  r4  r5  r>   r   r   r   r»   ï  r¾   zEagerResult.__reduce_args__c                 C   s   |   ¡ \}}||Ž S r   )r¼   )r:   r.  r�   r   r   r   r¸   ò  s   zEagerResult.__copy__c                 C   ó   dS r(   r   r>   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‰   rI   r…   r   rh   r²   Ú	Exception)r:   rU   r`   rf   ru   r   r   r   rj   ù  s   
ÿÿüzEagerResult.getc                 C   s   d S r   r   r>   r   r   r   rO     r9  zEagerResult.forgetc                 O   s   t j| _d S r   )r   ÚREVOKEDr4  rŽ   r   r   r   rW   
  r*  zEagerResult.revokec                 C   s   d| j › d�S )Nz<EagerResult: r°   rÓ   r>   r   r   r   r±     r½   zEagerResult.__repr__c                 C   s   | j | j| j| j| jdœS )N)rÔ   rI   r�   r‘   r×   )r4   r3  r4  r5  r6  r>   r   r   r   r8     s   ûzEagerResult._cachec                 C   rÒ   )zThe tasks return value.)r3  r>   r   r   r   rI     rÕ   zEagerResult.resultc                 C   rÒ   )zThe tasks state.)r4  r>   r   r   r   r…     rÕ   zEagerResult.statec                 C   rÒ   )z!The traceback if the task failed.)r5  r>   r   r   r   r‘   %  rÕ   zEagerResult.tracebackc                 C   r8  r"   r   r>   r   r   r   rÁ   *  s   z EagerResult.supports_native_joinr"  r"   )NTT)r*   r+   r,   r-   r<   rB   r’   r¼   r»   r¸   r�   rj   rY   rO   rW   r±   rÞ   r8   rI   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   rz   ©r3   r   r   rÎ   =  s    z%result_from_tuple.<locals>.<listcomp>rp   )r   r   r²   r   rl   Útupler   r   )rí   r3   ÚResultr  Únodesr4   r.   r   r<  r   r   /  s   

þr   r   ).r-   rÚ   r  Úcollectionsr   Ú
contextlibr   Úweakrefr   Údateutil.parserr   Úkombu.utils.objectsr   Úviner   r   r	   Ú r   r   r4  r   r   r3   r   Ú
exceptionsr   r   r   Úutils.graphr   r   r–   ÚImportErrorÚ__all__r   r    r&   r)   r   Úregisterr   r   r   r   r   r   r   r   r   Ú<module>   sR    ÿ
	
	   b  T_W