o
    wvXjÞ‘  ã                   @   s0  d Z ddlmZmZ ddlm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mZmZ ddlmZ ddlmZmZmZ ddlmZ ddlZ ddl m!Z!m"Z"m#Z#m$Z$ ddl%m&Z& ddlm'Z'm(Z(m)Z)m*Z*m+Z+m,Z,m-Z- ddl.m/Z/m0Z0 ddl1m2Z2m3Z3m4Z4m5Z5m6Z6 ddl7m8Z8 ddl9m:Z:m;Z; ddl<m=Z= ddl>m?Z?m@Z@mAZAmBZB ddlCmDZD dZEeFdhƒZGe=eHƒZIdZJeddƒZKdZLdZMdd „ ZNG d!d"„ d"eOƒZPG d#d$„ d$eQƒZRG d%d&„ d&eQƒZSG d'd(„ d(eReSƒZTeTZUG d)d*„ d*eRƒZVG d+d,„ d,eVeSƒZWG d-d.„ d.eTƒZXdS )/z°Result backend base classes.

- :class:`BaseBackend` defines the interface.

- :class:`KeyValueStoreBackend` is a common base class
    using K/V semantics like _get and _put.
é    )Úabsolute_importÚunicode_literals)Úraise_with_traceback)ÚdatetimeÚ	timedeltaN)Ú
namedtuple)Úpartial)ÚWeakValueDictionary)ÚExceptionInfo)ÚdumpsÚloadsÚprepare_accept_content)Úregistry)Úbytes_to_strÚensure_bytesÚ	from_utf8)Úmaybe_sanitize_url)Úcurrent_appÚgroupÚmaybe_signatureÚstates)Úget_current_task)Ú
ChordErrorÚImproperlyConfiguredÚNotRegisteredÚTaskRevokedErrorÚTimeoutErrorÚBackendGetMetaErrorÚBackendStoreError)ÚPY3Úitems)ÚGroupResultÚ
ResultBaseÚ	ResultSetÚallow_join_resultÚresult_from_tuple)Ú	BufferMap)ÚLRUCacheÚarity_greater)Ú
get_logger)Úcreate_exception_clsÚensure_serializableÚget_pickleable_exceptionÚget_pickled_exception)Ú get_exponential_backoff_interval)ÚBaseBackendÚKeyValueStoreBackendÚDisabledBackendÚpicklei    Úpending_results_t)ÚconcreteÚweakzU
No result backend is configured.
Please see the documentation for more information.
zû
Starting chords requires a result backend to be configured.

Note that a group chained with a task is also upgraded to be a chord,
as this pattern requires synchronization.

Result backends that supports chords: Redis, Database, Memcached, and more.
c                 C   s   | |dt  ¡ i|¤ŽS )zReturn an unpickled backend.Úapp)r   Ú_get_current_object)ÚclsÚargsÚkwargs© r;   úQ/var/www/html/myproject/venv/lib/python3.10/site-packages/celery/backends/base.pyÚunpickle_backendG   s   r=   c                   @   s    e Zd Zdd„ Ze Z ZZdS )Ú	_nulldictc                 O   ó   d S ©Nr;   )ÚselfÚaÚkwr;   r;   r<   ÚignoreM   ó   z_nulldict.ignoreN)Ú__name__Ú
__module__Ú__qualname__rD   Ú__setitem__ÚupdateÚ
setdefaultr;   r;   r;   r<   r>   L   s    r>   c                   @   s,  e Zd ZejZejZejZeZdZdZ	dZ
dZdddddœZ		dndd	„Zdod
d„Zdd„ Zddejfdd„Zddddejfdd„Zdd„ Zdddejfdd„Zdddejfdd„Zdpdd„Zdpdd„Zdpdd„Zdd „ Zd!d"„ Zd#d$„ Zd%d&„ Zd'd(„ Z d)d*„ Z!d+d,„ Z"dpd-d.„Z#dpd/d0„Z$d1d2„ Z%d3d4„ Z&		dqd5d6„Z'd7d8„ Z(	drd9d:„Z)d;d<„ Z*d=d>„ Z+d?d@„ Z,e,Z-dAdB„ Z.dCdD„ Z/dEdF„ Z0dGdH„ Z1dIdJ„ Z2dsdKdL„Z3dMdN„ Z4dOdP„ Z5dsdQdR„Z6dsdSdT„Z7dUdV„ Z8dWdX„ Z9dYdZ„ Z:d[d\„ Z;d]d^„ Z<d_d`„ Z=dadb„ Z>dtdcdd„Z?dedf„ Z@dgdh„ ZAdpdidj„ZBdudldm„ZCdS )vÚBackendNFTé   r   é   )Úmax_retriesÚinterval_startÚinterval_stepÚinterval_maxc                 K   sú   || _ | j j}	|p|	j| _tj| j \| _| _| _|p|	j	}
|
dkr%t
ƒ nt|
d�| _|  ||¡| _|d u r9|	jn|| _| jd u rD|	jn| j| _t| jƒ| _|	 dd¡| _|	 dd¡| _|	 dd¡| _|	 d	td
ƒ¡| _ti tƒ ƒ| _ttƒ| _|| _d S )Néÿÿÿÿ)ÚlimitÚresult_backend_always_retryFÚ+result_backend_max_sleep_between_retries_msi'  Ú,result_backend_base_sleep_between_retries_msé
   Úresult_backend_max_retriesÚinf) r6   ÚconfÚresult_serializerÚ
serializerÚserializer_registryÚ	_encodersÚcontent_typeÚcontent_encodingÚencoderÚresult_cache_maxr>   r'   Ú_cacheÚprepare_expiresÚexpiresÚresult_accept_contentÚacceptÚaccept_contentr   ÚgetÚalways_retryÚmax_sleep_between_retries_msÚbase_sleep_between_retries_msÚfloatrO   r3   r	   Ú_pending_resultsr&   ÚMESSAGE_BUFFER_MAXÚ_pending_messagesÚurl)rA   r6   r]   Úmax_cached_resultsrh   rf   Úexpires_typerr   r:   r[   Úcmaxr;   r;   r<   Ú__init__q   s(   
þ


zBackend.__init__c                 C   s2   |r| j S t| j p
dƒ}| d¡r|dd… S |S )z=Return the backend as an URI, sanitizing the password or not.Ú z:///NrS   )rr   r   Úendswith)rA   Úinclude_passwordrr   r;   r;   r<   Úas_uri�   s   zBackend.as_uric                 K   s   |   ||tj¡S )zMark a task as started.)Ústore_resultr   ÚSTARTED©rA   Útask_idÚmetar;   r;   r<   Úmark_as_started–   ó   zBackend.mark_as_startedc                 C   s:   |r| j ||||d� |r|jr|  |||¡ dS dS dS )z#Mark task as successfully executed.)ÚrequestN)r{   ÚchordÚon_chord_part_return)rA   r~   Úresultr‚   r{   Ústater;   r;   r<   Úmark_as_doneš   s
   
ÿzBackend.mark_as_donec                 C   sX   |r| j |||||d� |r&|jr|  |||¡ |r(|jr*|  |||¡ dS dS dS dS )z#Mark task as executed with failure.©Ú	tracebackr‚   N)r{   rƒ   r„   ÚerrbacksÚ_call_task_errbacks)rA   r~   Úexcr‰   r‚   r{   Úcall_errbacksr†   r;   r;   r<   Úmark_as_failure¢   s   
ÿ
üzBackend.mark_as_failurec           	   	   C   sô   g }|j D ]?}| j |¡}|js| j|_z"t|jdƒr0t|jjtƒs0t	|jjdƒr0||||ƒ n| 
|¡ W q tyD   | 
|¡ Y qw |rx|j}|jpN|}t|| jd�}| jjjsb|j dd¡rm|j|f||d� d S |j|f||d� d S d S )NÚ
__header__rN   ©r6   Úis_eagerF)Ú	parent_idÚroot_id)rŠ   r6   Ú	signatureÚ_appÚhasattrÚtypeÚ
isinstancer�   r   r(   Úappendr   Úidr“   r   r[   Útask_always_eagerÚdelivery_inforj   ÚapplyÚapply_async)	rA   r‚   rŒ   r‰   Úold_signatureÚerrbackr~   r“   Úgr;   r;   r<   r‹   °   s<   

ú
öõ
€û

ÿ
ÿõzBackend._call_task_errbacksrw   c                 C   sD   t |ƒ}|r| j|||d |d� |r|jr |  |||¡ d S d S d S )Nrˆ   )r   r{   rƒ   r„   )rA   r~   Úreasonr‚   r{   r†   rŒ   r;   r;   r<   Úmark_as_revokedÞ   s   
ÿ
ÿzBackend.mark_as_revokedc                 C   s   | j |||||d�S )zfMark task as being retries.

        Note:
            Stores the current exception (if any).
        rˆ   )r{   )rA   r~   rŒ   r‰   r‚   r{   r†   r;   r;   r<   Úmark_as_retryç   s   
ÿzBackend.mark_as_retryc              
      s¶   ddl m} | j‰ z	ˆ j|j j}W n ty   | }Y nw z|‡ fdd„|j d¡p,g D ƒˆ d� 	|j
f¡ W n tyR } z|j|j
|d�W  Y d }~S d }~ww |j|j
|d�S )Nr   )r   c                    ó   g | ]}ˆ   |¡‘qS r;   )r”   )Ú.0r    r�   r;   r<   Ú
<listcomp>û   ó    ÿz2Backend.chord_error_from_stack.<locals>.<listcomp>Ú
link_errorr�   )rŒ   )Úceleryr   r6   Ú_tasksÚtaskÚbackendÚKeyErrorÚoptionsrj   rž   rš   Ú	ExceptionÚfail_from_current_stack)rA   ÚcallbackrŒ   r   r­   Úeb_excr;   r�   r<   Úchord_error_from_stackñ   s(   ÿ
ÿý€ÿzBackend.chord_error_from_stackc                 C   s.  t  ¡ \}}}zT|d u r|n|}t|||fƒ}|  |||j¡ |W t jdkrH|d urFz|j ¡  |jj W n	 t	y>   Y nw |j
}|d us*~S dt j  krSdk rZn ~S t  ¡  ~S t jdkrƒ|d ur�z|j ¡  |jj W n	 t	yy   Y nw |j
}|d use~w dt j  krŽdk r•n ~w t  ¡  ~w )N)é   é   r   )é   é   r   )rµ   r   r   )ÚsysÚexc_infor
   rŽ   r‰   Úversion_infoÚtb_frameÚclearÚf_localsÚRuntimeErrorÚtb_nextÚ	exc_clear)rA   r~   rŒ   Útype_Úreal_excÚtbÚexception_infor;   r;   r<   r±     sH   

þùýþ
ó
þùýþzBackend.fail_from_current_stackc                 C   sL   |du r| j n|}|tv rt|ƒS t|ƒ}t|d|jƒt|j| jƒ|j	dœS )z$Prepare exception for serialization.NrH   )Úexc_typeÚexc_messageÚ
exc_module)
r]   ÚEXCEPTION_ABLE_CODECSr,   r—   ÚgetattrrF   r+   r9   ÚencoderG   )rA   rŒ   r]   Úexctyper;   r;   r<   Úprepare_exception  s   þzBackend.prepare_exceptionc              
   C   s  |r…t |tƒs|| d¡}|du rtt|d ƒtƒ}n1t|ƒ}t|d ƒ}ztj| }| d¡D ]}t	||ƒ}q/W n t
tfyJ   t|tjjƒ}Y nw |d }zt |ttfƒr\||Ž }n||ƒ}W n ty{ } ztd ||¡ƒ}W Y d}~nd}~ww | jtv r…t|ƒ}|S )z1Convert serialized exception to Python exception.rÈ   NrÆ   Ú.rÇ   z{}({}))r˜   ÚBaseExceptionrj   r*   r   rF   r¹   ÚmodulesÚsplitrÊ   r®   ÚAttributeErrorrª   Ú
exceptionsÚtupleÚlistr°   Úformatr]   rÉ   r-   )rA   rŒ   rÈ   r8   rÆ   ÚnameÚexc_msgÚerrr;   r;   r<   Úexception_to_python%  s@   

ÿ
ÿÿÿ
€€ÿ
zBackend.exception_to_pythonc                 C   s    | j dkrt|tƒr| ¡ S |S )zPrepare value for storage.r2   )r]   r˜   r"   Úas_tuple©rA   r…   r;   r;   r<   Úprepare_valueE  s   zBackend.prepare_valuec                 C   s   |   |¡\}}}|S r@   )Ú_encode)rA   ÚdataÚ_Úpayloadr;   r;   r<   rË   K  s   zBackend.encodec                 C   s   t || jd�S )N)r]   )r   r]   )rA   rß   r;   r;   r<   rÞ   O  ó   zBackend._encodec                 C   s$   |d | j v r|  |d ¡|d< |S )NÚstatusr…   )ÚEXCEPTION_STATESrÚ   )rA   r   r;   r;   r<   Úmeta_from_decodedR  s   zBackend.meta_from_decodedc                 C   s   |   |  |¡¡S r@   )rå   Údecode©rA   rá   r;   r;   r<   Údecode_resultW  s   zBackend.decode_resultc                 C   s2   |d u r|S t r
|pt|ƒ}t|| j| j| jd�S )N)r`   ra   rh   )r   Ústrr   r`   ra   rh   rç   r;   r;   r<   ræ   Z  s   ýzBackend.decodec                 C   s<   |d u r	| j jj}t|tƒr| ¡ }|d ur|r||ƒS |S r@   )r6   r[   Úresult_expiresr˜   r   Útotal_seconds)rA   Úvaluer—   r;   r;   r<   re   c  s   

zBackend.prepare_expiresc                 C   s(   |d ur|S | j jj}|d u r| jS |S r@   )r6   r[   Úresult_persistentÚ
persistent)rA   Úenabledrî   r;   r;   r<   Úprepare_persistentl  s   
zBackend.prepare_persistentc                 C   s(   || j v rt|tƒr|  |¡S |  |¡S r@   )rä   r˜   r°   rÍ   rÝ   )rA   r…   r†   r;   r;   r<   Úencode_resultr  s   

zBackend.encode_resultc                 C   s
   || j v S r@   )rd   ©rA   r~   r;   r;   r<   Ú	is_cachedw  s   
zBackend.is_cachedc                 C   s  || j v rt ¡ }|r| ¡ }nd }||||  |¡|dœ}|r*t|dd ƒr*|j|d< |r7t|dd ƒr7|j|d< | jj	 
dd¡r‹|r‹t|dd ƒt|dd ƒt|d	d ƒt|d
d ƒt|dd ƒt|dƒrh|jrh|j d¡nd dœ}	|r†dd	h}
|
D ]}|	| }|  |¡}t|ƒ|	|< qt| |	¡ |S )N)rã   r…   r‰   ÚchildrenÚ	date_doner   Úgroup_idr’   Úextendedr…   r¬   r9   r:   ÚhostnameÚretriesrœ   Úrouting_key)r×   r9   r:   Úworkerrù   Úqueue)ÚREADY_STATESr   ÚutcnowÚ	isoformatÚcurrent_task_childrenrÊ   r   r’   r6   r[   Úfind_value_for_keyr–   rœ   rj   rË   r   rJ   )rA   r…   r†   r‰   r‚   Úformat_daterË   rõ   r   Úrequest_metaÚencode_needed_fieldsÚfieldrì   Úencoded_valuer;   r;   r<   Ú_get_result_metaz  sJ   
€û






ÿþø

zBackend._get_result_metac                 C   s   t  |¡ d S r@   )ÚtimeÚsleep)rA   Úamountr;   r;   r<   Ú_sleepª  râ   zBackend._sleepc           
   
   K   s¶   |   ||¡}d}	 z| j||||fd|i|¤Ž |W S  tyY } z3| jrN|  |¡rN|| jk rD|d7 }t| j|| jdƒd }	|  	|	¡ nt
td||d�ƒ n‚ W Y d}~nd}~ww q	)	záUpdate task state and result.

        if always_retry_backend_operation is activated, in the event of a recoverable exception,
        then retry operation with an exponential backoff until a limit has been reached.
        r   Tr‚   rN   éè  z%failed to store result on the backend)r~   r†   N)rñ   Ú_store_resultr°   rk   Úexception_safe_to_retryrO   r.   rm   rl   r  r   r   )
rA   r~   r…   r†   r‰   r‚   r:   rù   rŒ   Úsleep_amountr;   r;   r<   r{   ­  s4   ÿÿ
þþ€òûzBackend.store_resultc                 C   s   | j  |d ¡ |  |¡ d S r@   )rd   ÚpopÚ_forgetrò   r;   r;   r<   ÚforgetÍ  s   zBackend.forgetc                 C   ó   t dƒ‚)Nz"backend does not implement forget.©ÚNotImplementedErrorrò   r;   r;   r<   r  Ñ  ó   zBackend._forgetc                 C   s   |   |¡d S )zGet the state of a task.rã   )Úget_task_metarò   r;   r;   r<   Ú	get_stateÔ  s   zBackend.get_statec                 C   ó   |   |¡ d¡S )z$Get the traceback for a failed task.r‰   ©r  rj   rò   r;   r;   r<   Úget_tracebackÚ  r�   zBackend.get_tracebackc                 C   r  )zGet the result of a task.r…   r  rò   r;   r;   r<   Ú
get_resultÞ  r�   zBackend.get_resultc                 C   s&   z|   |¡d W S  ty   Y dS w )z(Get the list of subtasks sent by a task.rô   N)r  r®   rò   r;   r;   r<   Úget_childrenâ  s
   ÿzBackend.get_childrenc                 C   s   | j jjrt dt¡ d S d S )Nz9Shouldn't retrieve result with task_always_eager enabled.)r6   r[   r›   ÚwarningsÚwarnÚRuntimeWarning©rA   r;   r;   r<   Ú_ensure_not_eageré  s   
þÿzBackend._ensure_not_eagerc                 C   ó   dS )a  Check if an exception is safe to retry.

        Backends have to overload this method with correct predicates dealing with their exceptions.

        By default no exception is safe to retry, it's up to backend implementation
        to define which exceptions are safe.
        Fr;   )rA   rŒ   r;   r;   r<   r  ð  s   zBackend.exception_safe_to_retryc              
   C   sâ   |   ¡  |rz| j| W S  ty   Y nw d}	 z|  |¡}W n? ty^ } z2| jrS|  |¡rS|| jk rJ|d7 }t| j	|| j
dƒd }|  |¡ n
ttd|d�ƒ n‚ W Y d}~nd}~ww q|ro| d¡tjkro|| j|< |S )	zßGet task meta from backend.

        if always_retry_backend_operation is activated, in the event of a recoverable exception,
        then retry operation with an exponential backoff until a limit has been reached.
        r   TrN   r  zfailed to get meta)r~   Nrã   )r"  rd   r®   Ú_get_task_meta_forr°   rk   r  rO   r.   rm   rl   r  r   r   rj   r   ÚSUCCESS)rA   r~   Úcacherù   r   rŒ   r  r;   r;   r<   r  ú  s>   ÿ

þþ€òü
zBackend.get_task_metac                 C   ó   | j |dd�| j|< dS )z;Reload task result, even if it has been previously fetched.F©r&  N)r  rd   rò   r;   r;   r<   Úreload_task_result  ó   zBackend.reload_task_resultc                 C   r'  )z<Reload group result, even if it has been previously fetched.Fr(  N)Úget_group_metard   ©rA   rö   r;   r;   r<   Úreload_group_result#  r*  zBackend.reload_group_resultc                 C   sP   |   ¡  |rz| j| W S  ty   Y nw |  |¡}|r&|d ur&|| j|< |S r@   )r"  rd   r®   Ú_restore_group©rA   rö   r&  r   r;   r;   r<   r+  '  s   ÿ

zBackend.get_group_metac                 C   s   | j ||d�}|r|d S dS )zGet the result for a group.r(  r…   N)r+  r/  r;   r;   r<   Úrestore_group4  s   ÿzBackend.restore_groupc                 C   s   |   ||¡S )z&Store the result of an executed group.)Ú_save_group©rA   rö   r…   r;   r;   r<   Ú
save_group:  s   zBackend.save_groupc                 C   s   | j  |d ¡ |  |¡S r@   )rd   r  Ú_delete_groupr,  r;   r;   r<   Údelete_group>  s   
zBackend.delete_groupc                 C   r#  )zsBackend cleanup.

        Note:
            This is run by :class:`celery.task.DeleteExpiredTaskMetaTask`.
        Nr;   r!  r;   r;   r<   ÚcleanupB  ó    zBackend.cleanupc                 C   r#  )z:Cleanup actions to do at the end of a task worker process.Nr;   r!  r;   r;   r<   Úprocess_cleanupI  r7  zBackend.process_cleanupc                 C   s   i S r@   r;   )rA   Úproducerr~   r;   r;   r<   Úon_task_callL  rE   zBackend.on_task_callc                 C   r  )Nz%Backend does not support add_to_chordr  )rA   Úchord_idr…   r;   r;   r<   Úadd_to_chordO  r  zBackend.add_to_chordc                 K   r?   r@   r;   )rA   r‚   r†   r…   r:   r;   r;   r<   r„   R  rE   zBackend.on_chord_part_returnc                 K   sh   dd„ |D ƒ|d< |j  dt|jdd ƒ¡}|j  dt|jddƒ¡}| jjd j|j|f||||d� d S )	Nc                 S   ó   g | ]}|  ¡ ‘qS r;   ©rÛ   ©r¦   Úrr;   r;   r<   r§   W  ó    z1Backend.fallback_chord_unlock.<locals>.<listcomp>r…   rü   Úpriorityr   zcelery.chord_unlock)Ú	countdownrü   rB  )r¯   rj   rÊ   r—   r6   Útasksrž   rš   )rA   Úheader_resultÚbodyrC  r:   rü   rB  r;   r;   r<   Úfallback_chord_unlockU  s   

üzBackend.fallback_chord_unlockc                 C   r?   r@   r;   r!  r;   r;   r<   Úensure_chords_alloweda  rE   zBackend.ensure_chords_allowedc                 K   s    |   ¡  | j||fi |¤Ž d S r@   )rH  rG  ©rA   rE  rF  r:   r;   r;   r<   Úapply_chordd  s   zBackend.apply_chordc                 C   s0   |pt tƒ dd ƒ}|rdd„ t |dg ƒD ƒS d S )Nr‚   c                 S   r=  r;   r>  r?  r;   r;   r<   r§   k  rA  z1Backend.current_task_children.<locals>.<listcomp>rô   )rÊ   r   )rA   r‚   r;   r;   r<   r   h  s   ÿzBackend.current_task_childrenr;   c                 C   s   |si n|}t | j||ffS r@   )r=   Ú	__class__©rA   r9   r:   r;   r;   r<   Ú
__reduce__m  s   zBackend.__reduce__)NNNNNN©Fr@   )TF©NN)T)rN   )r;   N)DrF   rG   rH   r   rý   ÚUNREADY_STATESrä   r   Úsubpolling_intervalÚsupports_native_joinÚsupports_autoexpirerî   Úretry_policyrv   rz   r€   r%  r‡   ÚFAILURErŽ   r‹   ÚREVOKEDr£   ÚRETRYr¤   r´   r±   rÍ   rÚ   rÝ   rË   rÞ   rå   rè   ræ   re   rð   rñ   ró   r  r  r{   r  r  r  Ú
get_statusr  r  r  r"  r  r  r)  r-  r+  r0  r3  r5  r6  r8  r:  r<  r„   rG  rH  rJ  r   rM  r;   r;   r;   r<   rL   S   sœ    ü
þ
	
ÿ	
ý.
ÿ	
ÿ




 
	
	
þ0
ÿ 

%



rL   c                   @   sT   e Zd Z		ddd„Z			ddd„Z	ddd	„Zddd„Zdd„ Zedd„ ƒZ	dS )ÚSyncBackendMixinNç      à?Tc                 c   s|   � |   ¡  |j}|sd S tƒ }|D ]}t|tƒr |j|jfV  q| |j¡ q| j||||||d�D ]	\}	}
|	|
fV  q2d S )N)ÚtimeoutÚintervalÚno_ackÚ
on_messageÚon_interval)r"  ÚresultsÚsetr˜   r#   rš   ÚaddÚget_many)rA   r…   r[  r\  r]  r^  r_  r`  Útask_idsr~   r   r;   r;   r<   Úiter_natives  s"   €
ýûzSyncBackendMixin.iter_nativec	           
      C   sN   |   ¡  |d urtdƒ‚| j|j||||d�}	|	r%| |	¡ |j||d�S d S )Nz,Backend does not support on_message callback)r[  r\  r_  r]  )Ú	propagater²   )r"  r   Úwait_forrš   Ú_maybe_set_cacheÚmaybe_throw)
rA   r…   r[  r\  r]  r^  r_  r²   rf  r   r;   r;   r<   Úwait_for_pendingˆ  s   ÿü
þz!SyncBackendMixin.wait_for_pendingc                 C   s\   |   ¡  d}	 |  |¡}|d tjv r|S |r|ƒ  t |¡ ||7 }|r-||kr-tdƒ‚q)aL  Wait for task and return its result.

        If the task raises an exception, this exception
        will be re-raised by :func:`wait_for`.

        Raises:
            celery.exceptions.TimeoutError:
                If `timeout` is not :const:`None`, and the operation
                takes longer than `timeout` seconds.
        g        rN   rã   zThe operation timed out.)r"  r  r   rý   r  r	  r   )rA   r~   r[  r\  r]  r_  Útime_elapsedr   r;   r;   r<   rg  š  s   

özSyncBackendMixin.wait_forFc                 C   ó   |S r@   r;   )rA   r…   r5   r;   r;   r<   Úadd_pending_result¶  rE   z#SyncBackendMixin.add_pending_resultc                 C   rl  r@   r;   rÜ   r;   r;   r<   Úremove_pending_result¹  rE   z&SyncBackendMixin.remove_pending_resultc                 C   r#  )NFr;   r!  r;   r;   r<   Úis_async¼  s   zSyncBackendMixin.is_async)NrZ  TNN)NrZ  TNNNT)NrZ  TNrN  )
rF   rG   rH   re  rj  rg  rm  rn  Úpropertyro  r;   r;   r;   r<   rY  r  s    
ÿ
þ
ÿ
rY  c                   @   ó   e Zd ZdZdS )r/   z"Base (synchronous) result backend.N©rF   rG   rH   Ú__doc__r;   r;   r;   r<   r/   Á  ó    r/   c                       s  e Zd ZeZdZdZdZdZ‡ fdd„Z	dd„ Z
d	d
„ Zdd„ Zdd„ Zdd„ Zdd„ Zdd„ Zd7dd„Zd7dd„Zd7dd„Zdd„ Zejfdd„Zejfd d!„Zd"d#d$d"d"d"ejfd%d&„Zd'd(„ Z	"d8d)d*„Zd+d,„ Zd-d.„ Zd/d0„ Zd1d2„ Zd3d4„ Z d5d6„ Z!‡  Z"S )9ÚBaseKeyValueStoreBackendzcelery-task-meta-zcelery-taskset-meta-zchord-unlock-Fc                    sJ   t | jdƒr| jj| _|  ¡  tt| ƒj|i |¤Ž | jr#| j| _	d S d S )NÚ__func__)
r–   Úkey_trv  Ú_encode_prefixesÚsuperru  rv   Úimplements_incrÚ_apply_chord_incrrJ  rL  ©rK  r;   r<   rv   Ï  s   
ÿz!BaseKeyValueStoreBackend.__init__c                 C   s.   |   | j¡| _|   | j¡| _|   | j¡| _d S r@   )rw  Útask_keyprefixÚgroup_keyprefixÚchord_keyprefixr!  r;   r;   r<   rx  ×  s   z)BaseKeyValueStoreBackend._encode_prefixesc                 C   r  )NzMust implement the get method.r  ©rA   Úkeyr;   r;   r<   rj   Ü  r  zBaseKeyValueStoreBackend.getc                 C   r  )NzDoes not support get_manyr  )rA   Úkeysr;   r;   r<   Úmgetß  r  zBaseKeyValueStoreBackend.mgetc                 C   r  )NzMust implement the set method.r  )rA   r�  rì   r†   r;   r;   r<   ra  â  r  zBaseKeyValueStoreBackend.setc                 C   r  )Nz Must implement the delete methodr  r€  r;   r;   r<   Údeleteå  r  zBaseKeyValueStoreBackend.deletec                 C   r  )NzDoes not implement incrr  r€  r;   r;   r<   Úincrè  r  zBaseKeyValueStoreBackend.incrc                 C   r?   r@   r;   )rA   r�  rì   r;   r;   r<   Úexpireë  rE   zBaseKeyValueStoreBackend.expirerw   c                 C   ó$   | j }|dƒ | j||ƒ||ƒg¡S )z#Get the cache key for a task by id.rw   )rw  Újoinr}  )rA   r~   r�  rw  r;   r;   r<   Úget_key_for_taskî  ó   ÿz)BaseKeyValueStoreBackend.get_key_for_taskc                 C   r‡  )z$Get the cache key for a group by id.rw   )rw  rˆ  r~  ©rA   rö   r�  rw  r;   r;   r<   Úget_key_for_groupõ  rŠ  z*BaseKeyValueStoreBackend.get_key_for_groupc                 C   r‡  )z?Get the cache key for the chord waiting on group with given id.rw   )rw  rˆ  r  r‹  r;   r;   r<   Úget_key_for_chordü  rŠ  z*BaseKeyValueStoreBackend.get_key_for_chordc                 C   sF   |   |¡}| j| jfD ]}| |¡rt|t|ƒd… ƒ  S qt|ƒS )zTake bytes: emit string.N)rw  r}  r~  Ú
startswithr   Úlen)rA   r�  Úprefixr;   r;   r<   Ú_strip_prefix  s   

ÿz&BaseKeyValueStoreBackend._strip_prefixc                 c   s<   � |D ]\}}|d ur|   |¡}|d |v r||fV  qd S )Nrã   )rè   )rA   Úvaluesrý   Úkrì   r;   r;   r<   Ú_filter_ready  s   €

€üz&BaseKeyValueStoreBackend._filter_readyc                    sF   t |dƒr‡fdd„ˆ t|ƒ|¡D ƒS ‡ fdd„ˆ t|ƒ|¡D ƒS )Nr    c                    s   i | ]
\}}ˆ   |¡|“qS r;   )r‘  )r¦   r“  Úvr!  r;   r<   Ú
<dictcomp>  s    
ÿÿz=BaseKeyValueStoreBackend._mget_to_results.<locals>.<dictcomp>c                    s   i | ]\}}t ˆ | ƒ|“qS r;   ©r   )r¦   Úir•  )r‚  r;   r<   r–    s    ÿÿ)r–   r”  r    Ú	enumerate)rA   r’  r‚  rý   r;   )r‚  rA   r<   Ú_mget_to_results  s   

þ
þz)BaseKeyValueStoreBackend._mget_to_resultsNrZ  Tc	              	   #   sb  � |d u rdn|}t |tƒr|nt|ƒ}	tƒ }
ˆ j}|	D ]$}z|| }W n	 ty-   Y qw |d |v r@t|ƒ|fV  |
 |¡ q|	 |
¡ d}|	r¯t|	ƒ}ˆ  ˆ  	‡ fdd„|D ƒ¡||¡}| 
|¡ |	 dd„ |D ƒ¡ t|ƒD ]\}}|d ur~||ƒ t|ƒ|fV  qr|r•|| |kr•td |¡ƒ‚|rš|ƒ  t |¡ |d	7 }|r«||kr«d S |	sJd S d S )
NrZ  rã   r   c                    r¥   r;   )r‰  )r¦   r“  r!  r;   r<   r§   5  r¨   z5BaseKeyValueStoreBackend.get_many.<locals>.<listcomp>c                 S   s   h | ]}t |ƒ’qS r;   r—  )r¦   r•  r;   r;   r<   Ú	<setcomp>8  rA  z4BaseKeyValueStoreBackend.get_many.<locals>.<setcomp>zOperation timed out ({0})rN   )r˜   ra  rd   r®   r   rb  Údifference_updaterÕ   rš  rƒ  rJ   r    r   rÖ   r  r	  )rA   rd  r[  r\  r]  r^  r_  Úmax_iterationsrý   ÚidsÚ
cached_idsr&  r~   ÚcachedÚ
iterationsr‚  r@  r�  rì   r;   r!  r<   rc     sN   €ÿ
€
ÿÿ

ïz!BaseKeyValueStoreBackend.get_manyc                 C   ó   |   |  |¡¡ d S r@   )r„  r‰  rò   r;   r;   r<   r  F  ó   z BaseKeyValueStoreBackend._forgetc           	      K   sX   | j ||||d�}t|ƒ|d< |  |¡}|d tjkr|S |  |  |¡|  |¡|¡ |S )N)r…   r†   r‰   r‚   r~   rã   )r  r   r$  r   r%  ra  r‰  rË   )	rA   r~   r…   r†   r‰   r‚   r:   r   Úcurrent_metar;   r;   r<   r  I  s   ÿ
z&BaseKeyValueStoreBackend._store_resultc                 C   s(   |   |  |¡|  d| ¡ i¡tj¡ |S )Nr…   )ra  rŒ  rË   rÛ   r   r%  r2  r;   r;   r<   r1  ]  s   ÿz$BaseKeyValueStoreBackend._save_groupc                 C   r¢  r@   )r„  rŒ  r,  r;   r;   r<   r4  b  r£  z&BaseKeyValueStoreBackend._delete_groupc                 C   s*   |   |  |¡¡}|stjddœS |  |¡S )ú$Get task meta-data for a task by id.N)rã   r…   )rj   r‰  r   ÚPENDINGrè   r}   r;   r;   r<   r$  e  s   
z+BaseKeyValueStoreBackend._get_task_meta_forc                 C   s>   |   |  |¡¡}|r|  |¡}|d }t|| jƒ|d< |S dS )r¥  r…   N)rj   rŒ  ræ   r%   r6   )rA   rö   r   r…   r;   r;   r<   r.  l  s   
üz'BaseKeyValueStoreBackend._restore_groupc                 K   s   |   ¡  |j| d� d S )N©r­   )rH  ÚsaverI  r;   r;   r<   r{  x  s   z*BaseKeyValueStoreBackend._apply_chord_incrc                 K   sÔ  | j sd S | j}|j}|sd S |  |¡}z	tj|| d�}W n+ tyH }	 zt|j|d�}
t	 
d||	¡ |  |
td |	¡ƒ¡W  Y d }	~	S d }	~	ww |d u r}zt|ƒ‚ ty| }	 zt|j|d�}
t	 
d||	¡ |  |
td |¡ƒ¡W  Y d }	~	S d }	~	ww |  |¡}t|ƒ}||kr’t	 d|¡ d S ||k�rat|j|d�}
|jr¤|jn|j}z®ztƒ � |dd	d
�}W d   ƒ n1 s½w   Y  W n> t�y }	 z1zt| ¡ ƒ}d ||	¡}W n tyç   t|	ƒ}Y nw t	 
d||¡ |  |
t|ƒ¡ W Y d }	~	n2d }	~	ww z|
 |¡ W n2 t�y. }	 zt	 
d||	¡ |  |
td |	¡ƒ¡ W Y d }	~	nd }	~	ww W | ¡  | j |¡ d S W | ¡  | j |¡ d S W | ¡  | j |¡ d S | ¡  | j |¡ w |  || j¡ d S )Nr§  r�   zChord %r raised: %rzCannot restore group: {0!r}zChord callback %r raised: %rz GroupResult {0} no longer existsz/Chord counter incremented too many times for %rg      @T)r[  rf  zDependency {0.id} raised {1!r}zCallback error: {0!r})rz  r6   r   r�  r!   Úrestorer°   r   rƒ   ÚloggerÚ	exceptionr´   r   rÖ   Ú
ValueErrorr…  r�  ÚwarningrR  Újoin_nativerˆ  r$   ÚnextÚ_failed_join_reportÚStopIterationÚreprÚdelayr„  Úclientr†  rf   )rA   r‚   r†   r…   r:   r6   Úgidr�  ÚdepsrŒ   r²   ÚvalÚsizeÚjÚretÚculpritr¢   r;   r;   r<   r„   |  sž   
þ€ýþ€ý
ÿ
ÿ€ÿÿ€öþ€þü÷úÿz-BaseKeyValueStoreBackend.on_chord_part_return)rw   rO  )#rF   rG   rH   r   rw  r}  r~  r  rz  rv   rx  rj   rƒ  ra  r„  r…  r†  r‰  rŒ  r�  r‘  r   rý   r”  rš  rc  r  r  r1  r4  r$  r.  r{  r„   Ú__classcell__r;   r;   r|  r<   ru  È  sB    



þ&
ÿru  c                   @   rq  )r0   z/Result backend base class for key/value stores.Nrr  r;   r;   r;   r<   r0   ½  rt  r0   c                   @   sP   e Zd ZdZi Zdd„ Zdd„ Zdd„ Zdd	„ Ze Z	 Z
 ZZe Z ZZd
S )r1   zDummy result backend.c                 O   r?   r@   r;   rL  r;   r;   r<   r{   Æ  rE   zDisabledBackend.store_resultc                 C   ó   t t ¡ ƒ‚r@   )r  ÚE_CHORD_NO_BACKENDÚstripr!  r;   r;   r<   rH  É  ó   z%DisabledBackend.ensure_chords_allowedc                 O   r½  r@   )r  ÚE_NO_BACKENDr¿  rL  r;   r;   r<   Ú_is_disabledÌ  rÀ  zDisabledBackend._is_disabledc                 O   r#  )Nzdisabled://r;   rL  r;   r;   r<   rz   Ï  rE   zDisabledBackend.as_uriN)rF   rG   rH   rs  rd   r{   rH  rÂ  rz   r  rX  r  r  Úget_task_meta_forrg  rc  r;   r;   r;   r<   r1   Á  s    r1   )Yrs  Ú
__future__r   r   Úfuture.utilsr   r   r   r¹   r  r  Úcollectionsr   Ú	functoolsr   Úweakrefr	   Úbilliard.einfor
   Úkombu.serializationr   r   r   r   r^   Úkombu.utils.encodingr   r   r   Úkombu.utils.urlr   Úcelery.exceptionsrª   r   r   r   r   Úcelery._stater   r   r   r   r   r   r   r   Úcelery.fiver   r    Úcelery.resultr!   r"   r#   r$   r%   Úcelery.utils.collectionsr&   Úcelery.utils.functionalr'   r(   Úcelery.utils.logr)   Úcelery.utils.serializationr*   r+   r,   r-   Úcelery.utils.timer.   Ú__all__Ú	frozensetrÉ   rF   rª  rp   r3   rÁ  r¾  r=   Údictr>   ÚobjectrL   rY  r/   ÚBaseDictBackendru  r0   r1   r;   r;   r;   r<   Ú<module>   s^   $


    #O v