
    xaiP                     n   d Z ddlZddlZddlZddlmZ ddlmZ ddlm	Z	 ddl
mZ ddlmZ ddlmZmZmZ dd	lmZ dd
l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" ddlm#Z#m$Z$m%Z%m&Z&m'Z'm(Z(m)Z)m*Z* ddl+m,Z,m-Z-m.Z.m/Z/m0Z0 ddl1m2Z2 ddl3m4Z4m5Z5 ddl6m7Z7 ddl8m9Z9m:Z:m;Z;m<Z<m=Z= ddl>m?Z? dZ@ eAdh      ZB e7eC      ZDdZE edd      ZFdZGdZHd ZId0de%fdZJd1de"fd ZK G d! d"eL      ZMd# ZN G d$ d%      ZO G d& d'      ZP G d( d)eOeP      ZQeQZR G d* d+eO      ZS G d, d-eSeP      ZT G d. d/eQ      ZUy)2zResult backend base classes.

- :class:`BaseBackend` defines the interface.

- :class:`KeyValueStoreBackend` is a common base class
    using K/V semantics like _get and _put.
    N)
namedtuple)	timedelta)partial)WeakValueDictionary)ExceptionInfo)dumpsloadsprepare_accept_content)registry)bytes_to_strensure_bytes)maybe_sanitize_url)current_appgroupmaybe_signaturestates)get_current_task)Context)BackendGetMetaErrorBackendStoreError
ChordErrorImproperlyConfiguredNotRegisteredSecurityErrorTaskRevokedErrorTimeoutError)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raise_with_context) 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                 :     | |dt        j                         i|S )zReturn an unpickled backend.app)r   _get_current_object)clsargskwargss      f/var/www/html/BankruptcyAI-uat/bankruptcy-ai/venv/lib/python3.12/site-packages/celery/backends/base.pyunpickle_backendr:   ?   s     F+99;FvFF    returnc                 J    t        |       }t        |t              r||_        |S )zCreate a ChordError preserving the original exception as __cause__.

    This helper reduces code duplication across the codebase when creating
    ChordError instances that need to preserve the original exception.
    )r   
isinstance	Exception	__cause__)messageoriginal_excchord_errors      r9   _create_chord_error_with_causerD   D   s&     W%K,	* ,r;   c                 >    t        | |xs g t               |d|      S )zCreate a fake task request context for error callbacks.

    This helper reduces code duplication when creating fake request contexts
    for error callback handling.
    )iderrbacksdelivery_infotask)r   dict)task_idrG   	task_nameextras       r9   _create_fake_task_requestrN   P   s2     N	
   r;   c                       e Zd Zd ZexZxZZy)	_nulldictc                      y N )selfakws      r9   ignorez_nulldict.ignore`       r;   N)__name__
__module____qualname__rW   __setitem__update
setdefaultrS   r;   r9   rP   rP   _   s     )/.K.&:r;   rP   c                      | y| j                   S NF)ignore_resultrequests    r9   _is_request_ignore_resultrd   f   s       r;   c                   T   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	 	 d=dZd>d	Zd
 Zddej                   fdZddddej$                  fdZd Zdddej*                  fdZdddej.                  fdZd?dZd?dZd?dZd?dZd Zd Zd Zd Z d Z!d Z"d Z#d?dZ$d?dZ%d Z&d Z'	 	 d@d Z(d! Z)	 dAd"Z*d# Z+d$ Z,d% Z-e-Z.d& Z/d' Z0d( Z1d) Z2d* Z3dBd+Z4d, Z5d- Z6dBd.Z7dBd/Z8d0 Z9d1 Z:d2 Z;d3 Z<d4 Z=d5 Z>d6 Z?d7 Z@dCd8ZAd9 ZBd: ZCd?d;ZDdDd<ZEy)EBackendNFT   r      )max_retriesinterval_startinterval_stepinterval_maxc                 f   || _         | j                   j                  }	|xs |	j                  | _        t        j
                  | j                     \  | _        | _        | _        |xs |	j                  }
|
dk(  r
t               nt        |
      | _        | j                  ||      | _        ||	j                  n|| _        | j                   |	j"                  n| j                   | _        t%        | j                         | _        |	j'                  dd      | _        |	j'                  dd      | _        |	j'                  dd      | _        |	j'                  d	t/        d
            | _        |	j'                  dd      | _        t5        i t7                     | _        t;        t<              | _        || _         y )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result_backend_thread_safe)!r4   confresult_serializer
serializerserializer_registry	_encoderscontent_typecontent_encodingencoderresult_cache_maxrP   r#   _cacheprepare_expiresexpiresresult_accept_contentacceptaccept_contentr
   getalways_retrymax_sleep_between_retries_msbase_sleep_between_retries_msfloatri   thread_safer0   r   _pending_resultsr"   MESSAGE_BUFFER_MAX_pending_messagesurl)rT   r4   ry   max_cached_resultsr   r   expires_typer   r8   rw   cmaxs              r9   __init__zBackend.__init__   s`    xx}}$>(>(> -66tG					!:T%:%:%)RZikXD5I++G\B 5;Nd00-1[[-@d))dkk,T[[9 HH%BEJ,0HH5bdi,j)-1XX6dfh-i*88$@%,O88$@%H 1"6I6K L!*+=!>r;   c                     |r| j                   S t        | j                   xs d      }|j                  d      r|dd S |S )z=Return the backend as an URI, sanitizing the password or not. z:///Nrn   )r   r   endswith)rT   include_passwordr   s      r9   as_urizBackend.as_uri   s>     88O R0<</s3Bx8S8r;   c                 D    | j                  ||t        j                        S )zMark a task as started.)store_resultr   STARTEDrT   rK   metas      r9   mark_as_startedzBackend.mark_as_started   s      $??r;   c                     |r t        |      s| j                  ||||       |r!|j                  r| j                  |||       yyy)z#Mark task as successfully executed.rb   N)rd   r   chordon_chord_part_return)rT   rK   resultrc   r   states         r9   mark_as_donezBackend.mark_as_done   sH     !:7!CgvugFw}}%%guf= %7r;   c                    |r| j                  |||||       |r(|j                  r| j                  |||       	 t        |j                        }|D ]  }	t        |	      }
|
j                  |
j                         |
j                  j                  d      |
_        |
j                  j                  d      |
_        |r>|t        j                  v r,|
j                    | j                  |
j                   ||||
       d|
j                  v s| j                  |
||        |r!|j"                  r| j%                  |||       yyyy# t
        t        f$ r t               }Y w xY w)z#Mark task as executed with failure.	tracebackrc   rK   group_idNr   )r   r   r   iterchainAttributeError	TypeErrortupler   r]   optionsr   rF   r   r   PROPAGATE_STATESrK   rG   _call_task_errbacks)rT   rK   excr   rc   r   call_errbacksr   
chain_data
chain_elemchain_elem_ctxs              r9   mark_as_failurezBackend.mark_as_failure   sh   
 gsE(17  D}}))'5#>%!'--0
 ) J
 ")!4%%n&<&<=$2$:$:$>$>y$I!'5'='='A'A*'M$ !Uf.E.E%E"**6%%&..U"+^ &  n444--neSI;J> !1!1((#yA "2}[  #I. %"W
%s   E E"!E"c                    g }|j                   D ]  }| j                  j                  |      }|j                  s| j                  |_        	 t	        |j
                  d      rOt        |j
                  j                  t              s+t        |j
                  j                  d      r ||||       n|j                  |        |r|j                  }|j                  xs |}t        || j                        }| j                  j                  j                   s|j"                  j%                  dd      r|j'                  |f||       y |j)                  |f||       y y # t        $ r |j                  |       Y ow xY w)N
__header__rh   r4   is_eagerF)	parent_idroot_id)rG   r4   	signature_apphasattrtyper>   r   r   r$   appendr   rF   r   r   rw   task_always_eagerrH   r   applyapply_async)	rT   rc   r   r   old_signatureerrbackrK   r   gs	            r9   r   zBackend._call_task_errbacks   sM   '' 	.Ghh((1G<<#xx.  l; 'w||'>'>H%gll&=&=qAGS)4!((1+	.:  jjGoo0Gm2Axx}}..'2G2G2K2KJX]2^J'7   J'7    ! .
 $$W-.s   A6E!!E?>E?r   c                     t        |      }|r| j                  |||d |       |r!|j                  r| j                  |||       y y y )Nr   )r   r   r   r   )rT   rK   reasonrc   r   r   r   s          r9   mark_as_revokedzBackend.mark_as_revoked"  sO    v&gsE(,g  ?w}}%%guc: %7r;   c                 .    | j                  |||||      S )zfMark task as being retries.

        Note:
            Stores the current exception (if any).
        r   )r   )rT   rK   r   r   rc   r   r   s          r9   mark_as_retryzBackend.mark_as_retry+  s)       #u+4g ! G 	Gr;   c                    | j                   }	 |j                  |j                     j                  }t        |t              r| j                  |||      S t        d|j                  j                  d      |j                  j                  dg       d|}	 | j                  ||d        |j                  |j                  |      S # t        $ r | }Y w xY w# t        $ r'}|j                  |j                  |      cY d }~S d }~ww xY w)N)group_callbackbackendr   rK   
link_error)rK   rG   r   rS   )r4   _tasksrI   r   KeyErrorr>   r   _handle_group_chord_errorrN   r   r   r   fail_from_current_stackrF   r?   )rT   callbackr   r4   r   fake_requesteb_excs          r9   chord_error_from_stackzBackend.chord_error_from_stack5  s
   hh	jj/77G
 h&11SZ`c1dd 1 
$$((3%%)),;
 

	I$$\3= 228;;C2HH+  	G	$  	L228;;F2KK	Ls/   #C C CC	D!D=DDc           
      6   t        |t              r%t        |d      r|j                  r|j                  }n|}	 |j	                         }t        |t
              r|j                          |j                  D ]q  }	 t        |j                  |j                  j                  dg       t        |dd            }	 |j                  ||d       |j                  |j                  |       s t        |d	d      }	|	r|j%                  |	|       y# t        $ r Y Lw xY w# t        $ r,}t         j#                  dt        |d	d      |       Y d}~d}~ww xY w# t        $ r=}
t         j#                  d
|
       |j                  |j                  |      cY d}
~
S d}
~
ww xY w)ah  Handle chord errors when the callback is a group.

        When a chord header fails and the body is a group, we need to:
        1. Revoke all pending tasks in the group body
        2. Mark them as failed with the chord error
        3. Call error callbacks for each task

        This prevents the group body tasks from hanging indefinitely (#8786)
        r@   r   rI   unknown)rK   rG   rL   Nr   z,Failed to handle chord error for task %s: %rrF   zLFailed to handle group chord error, falling back to single task handling: %r)r>   r   r   r@   freezer   revokeresultsrN   rF   r   r   getattrr   r?   r   logger	exceptionr   )rT   r   r   r   rB   frozen_groupr   r   task_excfrozen_group_idcleanup_excs              r9   r   z!Backend._handle_group_chord_errorQ  s    c:&73+D==LL1	O)002L,4##% +22 F'@$*II%3%;%;%?%?b%Q&-ffi&H(!#77lTXY  77		|7T#6 #*,d"C"++O\J'  ) ! ! % ((J#FD)<h   	O^
 22>3D3D#2NN	Osk   ?E 9=D7D
D'#E 	DDDD	E#"E
E 
EE 	F2FFFc                    t        j                         \  }}}	 ||n|}t        |||f      }| j                  |||j                         ||@	 |j
                  j                          |j
                  j                   |j                  }|@~S # t        $ r Y w xY w# |P	 |j
                  j                          |j
                  j                   n# t        $ r Y nw xY w|j                  }|P~w xY wrR   )
sysexc_infor   r   r   tb_frameclearf_localsRuntimeErrortb_next)rT   rK   r   type_real_exctbexception_infos          r9   r   zBackend.fail_from_current_stack  s    !llnx	!k(sC*E3+;<N  #~/G/GH!.KK%%'KK(( ZZ .  $ 	 .KK%%'KK((#  ZZ . sG   2B 0B	BBC4#0CC4	C C4C  C42C4c                     || j                   n|}|t        v rt        |      S t        |      }t	        |d|j
                        t        |j                  | j                        |j                  dS )z$Prepare exception for serialization.r[   )exc_typeexc_message
exc_module)
ry   EXCEPTION_ABLE_CODECSr(   r   r   rY   r'   r7   encoderZ   )rT   r   ry   exctypes       r9   prepare_exceptionzBackend.prepare_exception  sf    (2(:T__

..+C00s)#G^W=M=MN2388T[[I%002 	2r;   c                    |syt        |t              r| j                  t        v rt	        |      }|S t        |t
              s	 t        |      }|j                  d      }	 |d   }|t        |t              }n7	 t        j                  |   }|j                  d      D ]  }t        ||      } 	 |j                  dd      }t        |t&              rt)        |t              s||n| d| }t+        d	| d
|       	 t        |t,        t.        f      r || }|S  ||      }	 |S # t        $ r}t        d|       |d}~ww xY w# t        $ r}t        d      |d}~ww xY w# t        t         f$ r' t        |t"        j$                  j                        }Y w xY w# t0        $ r}	t1        | d| d      }Y d}	~	|S d}	~	ww xY w)z1Convert serialized exception to Python exception.NzbIf the stored exception isn't an instance of BaseException, it must be a dictionary.
Instead got: r   r   z5Exception information must include the exception type.r   r   z!Expected an exception class, got z with payload ())r>   BaseExceptionry   r   r)   rJ   r   r   r   
ValueErrorr&   rY   r   modulessplitr   r   celery
exceptionsr   
issubclassr   r   listr?   )
rT   r   er   r   r6   nameexc_msgfake_exc_typeerrs
             r9   exception_to_pythonzBackend.exception_to_python  s   ]+"77+C0JC&>3i WW\*
	::H &($CGkk*-$NN3/ -D!#t,C-
 ''-," #t$JsM,J(2(:H:,aPXz@ZM3M?.QXPYZ\ \	1'E4=17m 
	 'l 
u  > #0 14u!6 7 =>>>  	: 2 389:	: n- G*8+1+<+<+E+EGGB  	1se1WIQ/0C
	1s`   D' !E :5E$ ?F F '	E0D??E	E!EE!$3FF	G&F==Gc                 d    | j                   dk7  r t        |t              r|j                         S |S )zPrepare value for storage.r/   )ry   r>   r   as_tuplerT   r   s     r9   prepare_valuezBackend.prepare_value  s)    ??h&:fj+I??$$r;   c                 0    | j                  |      \  }}}|S rR   )_encode)rT   data_payloads       r9   r   zBackend.encode  s    T*1gr;   c                 0    t        || j                        S )N)ry   )r   ry   )rT   r  s     r9   r  zBackend._encode  s    Tdoo66r;   c                 V    |d   | j                   v r| j                  |d         |d<   |S )Nstatusr   )EXCEPTION_STATESr  )rT   r   s     r9   meta_from_decodedzBackend.meta_from_decoded  s1    >T222!55d8nEDNr;   c                 B    | j                  | j                  |            S rR   )r  decoderT   r  s     r9   decode_resultzBackend.decode_result  s    %%dkk'&:;;r;   c                     ||S |xs t        |      }t        || j                  | j                  | j                        S )N)r|   r}   r   )strr	   r|   r}   r   r  s     r9   r  zBackend.decode  sB    ?N)S\W"&"3"3&*&;&; KK) 	)r;   c                     | | j                   j                  j                  }t        |t              r|j                         }|
|r ||      S |S rR   )r4   rw   result_expiresr>   r   total_seconds)rT   valuer   s      r9   r   zBackend.prepare_expires  sI    =HHMM00EeY''')E;r;   c                 j    ||S | j                   j                  j                  }|| j                  S |S rR   )r4   rw   result_persistent
persistent)rT   enabledr&  s      r9   prepare_persistentzBackend.prepare_persistent   s4    NXX]]44
","4tD*Dr;   c                     || j                   v r!t        |t              r| j                  |      S | j	                  |      S rR   )r  r>   r?   r   r  )rT   r   r   s      r9   encode_resultzBackend.encode_result&  s;    D)))j.K))&11!!&))r;   c                     || j                   v S rR   )r   rT   rK   s     r9   	is_cachedzBackend.is_cached+  s    $++%%r;   c           	      N   || j                   v r-| j                  j                         }|r|j                         }nd }|||| j	                  |      |d}|rt        |dd       r|j                  |d<   |rt        |dd       r|j                  |d<   | j                  j                  j                  dd      r|rt        |dd       t        |dd       t        |d	d       t        |d
d       t        |dd       t        |d      r'|j                  r|j                  j                  d      nd d}	t        |dd       r*|j                  |	d<   |	j                  |j                         |r/dd	h}
|
D ]&  }|	|   }| j!                  |      }t#        |      |	|<   ( |j                  |	       |S )N)r  r   r   children	date_doner   r   r   extendedr   rI   r7   r8   hostnameretriesrH   routing_key)r  r7   r8   workerr3  queuestampsstamped_headers)READY_STATESr4   now	isoformatcurrent_task_childrenr   r   r   rw   find_value_for_keyr   rH   r   r8  r]   r7  r   r   )rT   r   r   r   rc   format_dater   r0  r   request_metaencode_needed_fieldsfieldr#  encoded_values                 r9   _get_result_metazBackend._get_result_meta.  s    D%%%I%//1	I "227;"
 ww6&}}DwwT: ' 1 1D88==++JA#GVT:#GVT:%gx>%gz4@&w	4@w8)) %2266}E/3	  7Hd36=6M6ML!23 ''7,2H+=(!5 J ,U 3(,E(:.:=.IU+J
 L)r;   c                 .    t        j                  |       y rR   )timesleep)rT   amounts     r9   _sleepzBackend._sleepa  s    

6r;   c                    | j                  ||      }d}	 	  | j                  ||||fd|i| |S # t        $ r}| j                  rt| j	                  |      rc|| j
                  k  r<|dz  }t        | j                  || j                  d      dz  }	| j                  |	       nt        t        d||             n Y d}~nd}~ww xY w)	zUpdate 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   Trc   rh     z%failed to store result on the backend)rK   r   N)r*  _store_resultr?   r   exception_safe_to_retryri   r+   r   r   rH  r*   r   )
rT   rK   r   r   r   rc   r8   r3  r   sleep_amounts
             r9   r   zBackend.store_resultd  s     ##FE2"""7FE9 >+2>6<> $$)E)Ec)J!1!111 (H >> ==t(EGK(L L1*-.U_fnst ! s   1 	CBC  Cc                 ^    | j                   j                  |d        | j                  |       y rR   )r   pop_forgetr,  s     r9   forgetzBackend.forget  s     &Wr;   c                     t        d      )Nz"backend does not implement forget.NotImplementedErrorr,  s     r9   rP  zBackend._forget  s    !"FGGr;   c                 *    | j                  |      d   S )zGet the state of a task.r  )get_task_metar,  s     r9   	get_statezBackend.get_state  s    !!'*844r;   c                 B    | j                  |      j                  d      S )z$Get the traceback for a failed task.r   rV  r   r,  s     r9   get_tracebackzBackend.get_traceback  s    !!'*..{;;r;   c                 B    | j                  |      j                  d      S )zGet the result of a task.r   rY  r,  s     r9   
get_resultzBackend.get_result  s    !!'*..x88r;   c                 J    	 | j                  |      d   S # t        $ r Y yw xY w)z(Get the list of subtasks sent by a task.r/  N)rV  r   r,  s     r9   get_childrenzBackend.get_children  s/    	%%g.z:: 		s    	""c                     | j                   j                  j                  r<| j                   j                  j                  st	        j
                  dt               y y y )NzResults are not stored in backend and should not be retrieved when task_always_eager is enabled, unless task_store_eager_result is enabled.)r4   rw   r   task_store_eager_resultwarningswarnRuntimeWarningrT   s    r9   _ensure_not_eagerzBackend._ensure_not_eager  sA    88==**488==3X3XMM[ 4Y*r;   c                      y)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.
        FrS   )rT   r   s     r9   rL  zBackend.exception_safe_to_retry  s     r;   c                 (   | j                          |r	 | j                  |   S d}	 	 | j                  |      }	 |r1|j                  d      t        j                   k(  r|| j                  |<   |S # t        $ r Y Vw xY w# t        $ r}| j
                  rs| j                  |      rb|| j                  k  r<|dz  }t        | j                  || j                  d      dz  }| j                  |       nt        t        d|             n Y d}~nd}~ww xY w)	zGet 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   Trh   rJ  zfailed to get meta)rK   Nr  )re  r   r   _get_task_meta_forr?   r   rL  ri   r+   r   r   rH  r*   r   r   r   SUCCESS)rT   rK   cacher3  r   r   rM  s          r9   rV  zBackend.get_task_meta  s     	 {{7++ ..w7& TXXh'6>>9#'DKK 7    $$)E)Ec)J!1!111 (H >> ==t(EGK(L L1*/0DgV !	 s)   A. A= .	A:9A:=	DB DDc                 D    | j                  |d      | j                  |<   y)z;Reload task result, even if it has been previously fetched.Frj  N)rV  r   r,  s     r9   reload_task_resultzBackend.reload_task_result  s     #11'1GGr;   c                 D    | j                  |d      | j                  |<   y)z<Reload group result, even if it has been previously fetched.Frl  N)get_group_metar   rT   r   s     r9   reload_group_resultzBackend.reload_group_result  s      $ 3 3HE 3 JHr;   c                     | j                          |r	 | j                  |   S | j                  |      }|r||| j                  |<   |S # t        $ r Y 1w xY wrR   )re  r   r   _restore_grouprT   r   rj  r   s       r9   ro  zBackend.get_group_meta  sf     {{8,, ""8,T%$(DKK!  s   A	 		AAc                 8    | j                  ||      }|r|d   S y)zGet the result for a group.rl  r   N)ro  rt  s       r9   restore_groupzBackend.restore_group  s)    ""85"9>! r;   c                 &    | j                  ||      S )z&Store the result of an executed group.)_save_grouprT   r   r   s      r9   
save_groupzBackend.save_group  s    &11r;   c                 \    | j                   j                  |d        | j                  |      S rR   )r   rO  _delete_grouprp  s     r9   delete_groupzBackend.delete_group  s%    $'!!(++r;   c                      y)zBackend cleanup.NrS   rd  s    r9   cleanupzBackend.cleanup      r;   c                      y)z:Cleanup actions to do at the end of a task worker process.NrS   rd  s    r9   process_cleanupzBackend.process_cleanup  r  r;   c                     i S rR   rS   )rT   producerrK   s      r9   on_task_callzBackend.on_task_call  s    	r;   c                     t        d      )Nz%Backend does not support add_to_chordrS  )rT   chord_idr   s      r9   add_to_chordzBackend.add_to_chord  s    !"IJJr;   c                      y rR   rS   )rT   rc   r   r   r8   s        r9   r   zBackend.on_chord_part_return
  rX   r;   c                      y rR   rS   )rT   r   
chord_sizes      r9   set_chord_sizezBackend.set_chord_size  rX   r;   c                 .   |D cg c]  }|j                          c}|d<   	 t        |dd       }|j                  j	                  dt        |dd             }|G| j
                  j                  j                  j                  ||j                        d   j                  }|j                  j	                  dt        |dd            }| j
                  j                  d   j                  |j                  |f||||       y c c}w # t        $ r d }Y w xY w)Nr   r   r6  priorityr   zcelery.chord_unlock)	countdownr6  r  )r  r   r   r   r   r4   amqprouterrouter  tasksr   rF   )	rT   header_resultbodyr  r8   r	body_typer6  r  s	            r9   fallback_chord_unlockzBackend.fallback_chord_unlock  s    2?@QAJJL@x	fd3I   ')Wd*KL= HHMM((..vtyyA'JOOE<<##J	:q0QR,-99t%v	 	: 	
 A  	I	s   DD DDc                      y rR   rS   rd  s    r9   ensure_chords_allowedzBackend.ensure_chords_allowed'  rX   r;   c                 ~    | j                           | j                  j                  | } | j                  ||fi | y rR   )r  r4   r   r  rT   header_result_argsr  r8   r  s        r9   apply_chordzBackend.apply_chord*  s<    ""$,,,.@A"""=$A&Ar;   c                     |xs t        t               dd       }|r)t        |dg       D cg c]  }|j                          c}S y c c}w )Nrc   r/  )r   r   r  )rT   rc   r  s      r9   r<  zBackend.current_task_children/  sD    IW%5%7DI*1':r*JKQAJJLKK Ks   Ac                 8    |si n|}t         | j                  ||ffS rR   )r:   	__class__rT   r7   r8   s      r9   
__reduce__zBackend.__reduce__4  s!    !v 4>>4"@AAr;   )NNNNNNFrR   )TFNN)T)rh   )rS   N)FrY   rZ   r[   r   r9  UNREADY_STATESr  r   subpolling_intervalsupports_native_joinsupports_autoexpirer&  retry_policyr   r   r   ri  r   FAILUREr   r   REVOKEDr   RETRYr   r   r   r   r   r  r  r   r  r  r  r  r   r(  r*  r-  rC  rH  r   rQ  rP  rW  
get_statusrZ  r\  r^  re  rL  rV  rm  rq  ro  rv  rz  r}  r  r  r  r  r   r  r  r  r  r<  r  rS   r;   r9   rf   rf   l   s   &&L**N..L
  !
   J 	L CG6::9@
 "FNN> #'%)$nn6Bp,\ /1 $4v~~; 59"V\\GI8BOH&2EN7
<)E*
& AE %1f .2 DH5 J<9%NHK"2,IK
.B
L
Br;   rf   c                   N    e Zd Z	 	 ddZ	 	 	 d	dZ	 d
dZddZd Zed        Z	y)SyncBackendMixinNc              #   :  K   | j                          |j                  }|sy t               }|D ]H  }t        |t              r|j
                  |j                  f .|j                  |j
                         J | j                  ||||||      E d {    y 7 w)N)timeoutintervalno_ack
on_messageon_interval)re  r   setr>   r   rF   addget_many)	rT   r   r  r  r  r  r  r   task_idss	            r9   iter_nativezSyncBackendMixin.iter_native:  s      ..5 	(F&),ii//VYY'		( ==hv!{ ! 
 	
 	
s   BBBBc	                     | j                          |t        d      | j                  |j                  ||||      }	|	r$|j	                  |	       |j                  ||      S y )Nz,Backend does not support on_message callback)r  r  r  r  )	propagater   )re  r   wait_forrF   _maybe_set_cachemaybe_throw)
rT   r   r  r  r  r  r  r   r  r   s
             r9   wait_for_pendingz!SyncBackendMixin.wait_for_pendingN  s}     	 !&>@ @ }}IIw#	  
 ##D)%%	H%MM r;   c                     | j                          d}	 | j                  |      }|d   t        j                  v r|S |r |        t	        j
                  |       ||z  }|r||k\  rt        d      ^)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        r  zThe operation timed out.)re  rV  r   r9  rE  rF  r   )rT   rK   r  r  r  r  time_elapsedr   s           r9   r  zSyncBackendMixin.wait_for`  sx     	 %%g.DH~!4!44JJx H$L<72"#=>> r;   c                     |S rR   rS   )rT   r   r2   s      r9   add_pending_resultz#SyncBackendMixin.add_pending_result|      r;   c                     |S rR   rS   r  s     r9   remove_pending_resultz&SyncBackendMixin.remove_pending_result  r  r;   c                      yr`   rS   rd  s    r9   is_asynczSyncBackendMixin.is_async  s    r;   )N      ?TNN)Nr  TNNNT)Nr  TNr  )
rY   rZ   r[   r  r  r  r  r  propertyr  rS   r;   r9   r  r  9  sI    EI15
( ?BCG26N& GK?8  r;   r  c                       e Zd ZdZy)r,   z"Base (synchronous) result backend.NrY   rZ   r[   __doc__rS   r;   r9   r,   r,     s    ,r;   r,   c                   (    e Zd ZeZdZdZdZdZ fdZ	d Z
d Zd Zd	 Zd
 Zd Zd Zd Zd Zd"dZd"dZd"dZd"dZd Zej2                  fdZej2                  fdZddddddej2                  fdZd Z	 d#dZd Zd Z d Z!d Z"d  Z#d! Z$ xZ%S )$BaseKeyValueStoreBackendzcelery-task-meta-zcelery-taskset-meta-zchord-unlock-Fc                    t        | j                  d      r| j                  j                  | _        t        |   |i | | j                          | j                          | j                  r| j                  | _	        y y )N__func__)
r   key_tr  superr   _add_global_keyprefix_encode_prefixesimplements_incr_apply_chord_incrr  )rT   r7   r8   r  s      r9   r   z!BaseKeyValueStoreBackend.__init__  sh    4::z*,,DJ$)&)""$#55D  r;   c                    | j                   j                  j                  di       j                  dd      }|rL|d   dvr|dz  }| | j                   | _        | | j                   | _        | | j
                   | _        yy)a/  
        This method prepends the global keyprefix to the existing keyprefixes.

        This method checks if a global keyprefix is configured in `result_backend_transport_options` using the
        `global_keyprefix` key. If so, then it is prepended to the task, group and chord key prefixes.
         result_backend_transport_optionsglobal_keyprefixNrn   z:_-.r  )r4   rw   r   task_keyprefixgroup_keyprefixchord_keyprefix)rT   r  s     r9   r  z.BaseKeyValueStoreBackend._add_global_keyprefix  s      88==,,-OQSTXXYkmqr#61 C' %5$6t7J7J6K"LD&6%78L8L7M#ND &6%78L8L7M#ND  r;   c                     | j                  | j                        | _        | j                  | j                        | _        | j                  | j                        | _        y rR   )r  r  r  r  rd  s    r9   r  z)BaseKeyValueStoreBackend._encode_prefixes  sG    "jj)<)<=#zz$*>*>?#zz$*>*>?r;   c                     t        d      )NzMust implement the get method.rS  rT   keys     r9   r   zBaseKeyValueStoreBackend.get      !"BCCr;   c                     t        d      )NzDoes not support get_manyrS  )rT   keyss     r9   mgetzBaseKeyValueStoreBackend.mget  s    !"=>>r;   c                 &    | j                  ||      S rR   )r  )rT   r  r#  r   s       r9   _set_with_statez(BaseKeyValueStoreBackend._set_with_state  s    xxU##r;   c                     t        d      )NzMust implement the set method.rS  rT   r  r#  s      r9   r  zBaseKeyValueStoreBackend.set  r  r;   c                     t        d      )Nz Must implement the delete methodrS  r  s     r9   deletezBaseKeyValueStoreBackend.delete  s    !"DEEr;   c                     t        d      )NzDoes not implement incrrS  r  s     r9   incrzBaseKeyValueStoreBackend.incr  s    !";<<r;   c                      y rR   rS   r  s      r9   expirezBaseKeyValueStoreBackend.expire  rX   r;   c                 ^    |st        d| d      | j                  | j                  ||      S )z#Get the cache key for a task by id.ztask_id must not be empty. Got 	 instead.)r   _get_key_forr  )rT   rK   r  s      r9   get_key_for_taskz)BaseKeyValueStoreBackend.get_key_for_task  s5    >wiyQRR  !4!4gsCCr;   c                 ^    |st        d| d      | j                  | j                  ||      S )z$Get the cache key for a group by id. group_id must not be empty. Got r  )r   r  r  rT   r   r  s      r9   get_key_for_groupz*BaseKeyValueStoreBackend.get_key_for_group  5    ?zSTT  !5!5xEEr;   c                 ^    |st        d| d      | j                  | j                  ||      S )z?Get the cache key for the chord waiting on group with given id.r  r  )r   r  r  r  s      r9   get_key_for_chordz*BaseKeyValueStoreBackend.get_key_for_chord  r  r;   c                 f    | j                   } |d      j                  | ||       ||      g      S )Nr   )r  join)rT   prefixrF   r  r  s        r9   r  z%BaseKeyValueStoreBackend._get_key_for  s4    

Ry~~E"IuSz
  	r;   c                     | j                  |      }| j                  | j                  fD ],  }|j                  |      st	        |t        |      d       c S  t	        |      S )zTake bytes: emit string.N)r  r  r  
startswithr   len)rT   r  r  s      r9   _strip_prefixz&BaseKeyValueStoreBackend._strip_prefix  s^    jjo))4+?+?? 	7F~~f%#CF$566	7 C  r;   c              #   d   K   |D ]'  \  }}|	| j                  |      }|d   |v s"||f ) y w)Nr  )r  )rT   valuesr9  kr#  s        r9   _filter_readyz&BaseKeyValueStoreBackend._filter_ready  sD      	#HAu **51?l2U(N		#s   00	0c                 .   t        |d      rC| j                  |j                         |      D ci c]  \  }}| j                  |      | c}}S | j                  t	        |      |      D ci c]  \  }}t        ||         | c}}S c c}}w c c}}w )Nitems)r   r  r	  r  	enumerater   )rT   r  r  r9  r  vis          r9   _mget_to_resultsz)BaseKeyValueStoreBackend._mget_to_results  s    67# !..v||~|LAq ""1%q(  !..y/@,OAq T!W%q( s   B.BNr  Tc	           
   #   >  K   |dn|}t        |t              r|n
t        |      }	t               }
| j                  }|	D ]0  }	 ||   }|d   |v st        |      |f |
j	                  |       2 |	j                  |
       d}|	rt        |	      }| j                  | j                  |D cg c]  }| j                  |       c}      ||      }|j                  |       |	j                  |D ch c]  }t        |       c}       |j                         D ]  \  }}| ||       t        |      |f   |r||z  |k\  rt        d| d      |r |        t        j                  |       |dz  }|r||k\  ry |	ry y # t
        $ r Y Qw xY wc c}w c c}w w)Nr  r  r   zOperation timed out (r   rh   )r>   r  r   r   r  r   difference_updater  r  r  r  r]   r	  r   rE  rF  )rT   r  r  r  r  r  r  max_iterationsr9  ids
cached_idsrj  rK   cached
iterationsr  r  r  r  r  r#  s                        r9   r  z!BaseKeyValueStoreBackend.get_many  s     #*3$Xs3hXU
 	,G,w (#|3&w/77NN7+	, 	j)
9D%%dii:>1@56 261F1Fq1I 1@ 'ABFVALLO!!A">q<?">?ggi /
U)u%"3'../ :0G;"%:7)1#EFFJJx !OJ*">#   1@ #?sO   ?FFFAF+F-F0FA>FF	FFFFc                 D    | j                  | j                  |             y rR   )r  r  r,  s     r9   rP  z BaseKeyValueStoreBackend._forget#  s    D))'23r;   c                 T   | j                  ||||      }t        |      |d<   | j                  |      }|d   t        j                  k(  r|S 	 | j                  | j                  |      | j                  |      |       |S # t        $ r}	t        t        |	      ||      |	d }	~	ww xY w)N)r   r   r   rc   rK   r  )r   rK   )
rC  r   rh  r   ri  r  r  r   r   r  )
rT   rK   r   r   r   rc   r8   r   current_metaexs
             r9   rK  z&BaseKeyValueStoreBackend._store_result&  s    $$F%/8' % K&w/Y ..w7!V^^3M	S  !6!6w!?TARTYZ  ! 	S#CG5'JPRR	Ss   1B 	B'
B""B'c                     | j                  | j                  |      | j                  d|j                         i      t        j
                         |S )Nr   )r  r  r   r  r   ri  ry  s      r9   rx  z$BaseKeyValueStoreBackend._save_group>  sA    T33H=![[(FOO4E)FG	Yr;   c                 D    | j                  | j                  |             y rR   )r  r  rp  s     r9   r|  z&BaseKeyValueStoreBackend._delete_groupC  s    D**845r;   c                     | j                  | j                  |            }|st        j                  ddS | j	                  |      S )$Get task meta-data for a task by id.N)r  r   )r   r  r   PENDINGr  r   s      r9   rh  z+BaseKeyValueStoreBackend._get_task_meta_forF  s>    xx--g67$nn==!!$''r;   c                     | j                  | j                  |            }|r1| j                  |      }|d   }t        || j                        |d<   |S y)r  r   N)r   r  r  r!   r4   )rT   r   r   r   s       r9   rs  z'BaseKeyValueStoreBackend._restore_groupM  sU    xx..x89 ;;t$D(^F.vtxx@DNK	 r;   c                 z    | j                           | j                  j                  | }|j                  |        y )Nr   )r  r4   r   saver  s        r9   r  z*BaseKeyValueStoreBackend._apply_chord_incrY  s6    ""$,,,.@A4(r;   c           	         | j                   sy | j                  }|j                  }|sy | j                  |      }	 t	        j
                  ||       }|	 t        |      | j                  |      }|j                  j                  d      }|t!        |      }||kD  rt        j#                  d	|       y ||k(  rt        |j                  |      }
|j$                  r|j&                  n|j(                  }	 t+               5   ||j,                  j.                  d
      }d d d        	 |
j1                         |j?                          | j?                  |       y | jA                  || jB                         y # t        $ rV}	t        |j                  |      }
t        j                  d||	       | j                  |
t        d|	            cY d }	~	S d }	~	ww xY w# t        $ rW}	t        |j                  |      }
t        j                  d||	       | j                  |
t        d| d            cY d }	~	S d }	~	ww xY w# 1 sw Y    xY w# t        $ rE}	t        j                  d||	       t3        d|	|	      }| j                  |
|       Y d }	~	Zd }	~	ww xY w# t        $ r}		 t5        |j7                               }dj9                  ||	      }n# t:        $ r t=        |	      }Y nw xY wt        j                  d||       t3        ||	      }| j                  |
|       Y d }	~	d }	~	ww xY w# |j?                          | j?                  |       w xY w)Nr   r   zChord %r raised: %rzCannot restore group: zChord callback %r raised: %rzGroupResult z no longer existsr  z/Chord counter incremented too many times for %rT)r  r  zCallback error: )rA   rB   )r   r   zDependency {0.id} raised {1!r})"r  r4   r   r  r   restorer?   r   r   r   r   r   r   r   r  r   r  warningr  join_nativer  r    rw   result_chord_join_timeoutdelayrD   next_failed_join_reportformatStopIterationreprr  r  r   )rT   rc   r   r   r8   r4   gidr  depsr   r   valsizejretrC   culpritr   s                     r9   r   z-BaseKeyValueStoreBackend.on_chord_part_return^  s   ##hhmm$$S)	&&sD9D < o% iin }}  .<t9D:NNL D[&w}}#>H$($=$=  499A!&( ( # B B"&(C( TNN3' C KKT\\*u  	&w}}#>H2C=..3C7;< 	  *7==cB  !?cJ22cU2CDE (( ($ ! T$$%:CE"@"23' :#K //{/SST  
P'"4#;#;#=>G=DDF % '!#YF'  !6VD<VZ]^++X;+OO
P* C s   E8 G 6
J  H=J (I
 8	GAGGG	H:#AH5/H:5H:=IJ 
	J:JL0 JL0 	L-%+KL(K(%L('K((:L("L0 (L--L0 0#M)r   r  )&rY   rZ   r[   r   r  r  r  r  r  r   r  r  r   r  r  r  r  r  r  r  r  r  r  r  r   r9  r  r  r  rP  rK  rx  r|  rh  rs  r  r   __classcell__)r  s   @r9   r  r    s    E(N,O%OO6O@
D?$DF=DFF! 281D1D # ;A:M:M  *.D d4$11$L4 /30
6(
)
D+r;   r  c                       e Zd ZdZy)r-   z/Result backend base class for key/value stores.Nr  rS   r;   r9   r-   r-     s    9r;   r-   c                   H    e Zd ZdZi Zd Zd Zd Zd ZexZ	xZ
xZZexZxZZy)r.   zDummy result backend.c                      y rR   rS   r  s      r9   r   zDisabledBackend.store_result  rX   r;   c                 <    t        t        j                               rR   )rT  E_CHORD_NO_BACKENDstriprd  s    r9   r  z%DisabledBackend.ensure_chords_allowed  s    !"4":":"<==r;   c                 <    t        t        j                               rR   )rT  E_NO_BACKENDr:  r  s      r9   _is_disabledzDisabledBackend._is_disabled  s    !,"4"4"677r;   c                      y)Nzdisabled://rS   r  s      r9   r   zDisabledBackend.as_uri  s    r;   N)rY   rZ   r[   r  r   r   r  r=  r   rW  r  r\  rZ  get_task_meta_forr  r  rS   r;   r9   r.   r.     sE    F>8 ;GFIF
FZ-.:::8r;   r.   rR   )Nr   )Vr  r   rE  ra  collectionsr   datetimer   	functoolsr   weakrefr   billiard.einfor   kombu.serializationr   r	   r
   r   rz   kombu.utils.encodingr   r   kombu.utils.urlr   celery.exceptionsr  r   r   r   r   celery._stater   celery.app.taskr   r   r   r   r   r   r   r   r   celery.resultr   r   r   r    r!   celery.utils.collectionsr"   celery.utils.functionalr#   r$   celery.utils.logr%   celery.utils.serializationr&   r'   r(   r)   r*   celery.utils.timer+   __all__	frozensetr   rY   r   r   r0   r<  r9  r:   rD   rN   rJ   rP   rd   rf   r  r,   BaseDictBackendr  r-   r.   rS   r;   r9   <module>rT     sC      "   ' ( D D ? ; .  > > * #] ] ] b b . ; 'S S >
D!8*- 	H	 2 5  
 G
	* 	W^ / /!JB JBZK K\-'+ - T+w T+n:35E :;k ;r;   