Ë
    —V.j-/  ã                   óÚ   — 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mZ dd	lmZ dd
lmZmZ dZdZ G d„ de«      Zd„ Z G d„ de«      Z G d„ dej2                  e«      Zy)zqThe ``RPC`` result backend for AMQP brokers.

RPC-style result backend, using reply-to and one queue per client.
é    N)Úmaybe_declare)Úregister_after_fork)Úcached_property)Ústates)Úcurrent_taskÚtask_join_will_blocké   )Úbase)ÚAsyncBackendMixinÚBaseResultConsumer)ÚBacklogLimitExceededÚ
RPCBackendzñ
The "rpc" result backend does not support chords!

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                   ó   — e Zd ZdZy)r   z'Too much state history to fast-forward.N)Ú__name__Ú
__module__Ú__qualname__Ú__doc__© ó    úUC:\xampp\htdocs\tradingbinance\backend\.venv\Lib\site-packages\celery/backends/rpc.pyr   r      s   „ Ú1r   r   c                 ó$   — | j                  «        y ©N)Ú_after_fork)Úbackends    r   Ú_on_after_fork_cleanup_backendr   "   s   € Ø×ÑÕr   c                   óf   ‡ — e Zd Zej                  ZdZdZˆ fd„Zd	d„Zd
d„Z	d„ Z
d„ Zd„ Zd„ Zˆ xZS )ÚResultConsumerNc                 óZ   •— t        ‰| �  |i |¤Ž | j                  j                  | _        y r   )ÚsuperÚ__init__r   Ú_create_binding©ÚselfÚargsÚkwargsÚ	__class__s      €r   r    zResultConsumer.__init__,   s'   ø€ Ü‰Ñ˜$Ð) &Ò)Ø#Ÿ|™|×;Ñ;ˆÕr   c                 ó"  — | j                   j                  «       | _        | j                  |«      }| j	                  | j                  j
                  |g| j                  g|| j                  ¬«      | _        | j                  j                  «        y )N)Ú	callbacksÚno_ackÚaccept)
ÚappÚ
connectionÚ_connectionr!   ÚConsumerÚdefault_channelÚon_state_changer*   Ú	_consumerÚconsume)r#   Úinitial_task_idr)   r%   Úinitial_queues        r   ÚstartzResultConsumer.start0   sv   € ØŸ8™8×.Ñ.Ó0ˆÔØ×,Ñ,¨_Ó=ˆØŸ™Ø×Ñ×,Ñ,¨}¨oØ×+Ñ+Ð,°VØ—;‘;ð 'ó  ˆŒð 	�‰×ÑÕ r   c                 ó„   — | j                   r| j                   j                  |¬«      S |rt        j                  |«       y y )N)Útimeout)r-   Údrain_eventsÚtimeÚsleep)r#   r7   s     r   r8   zResultConsumer.drain_events9   s9   € Ø×ÒØ×#Ñ#×0Ñ0¸Ð0ÓAÐAÙÜ�J‰J�wÕð r   c                 ó¬   — 	 | j                   j                  «        | j                  j                  «        y # | j                  j                  «        w xY wr   )r1   Úcancelr-   Úclose©r#   s    r   ÚstopzResultConsumer.stop?   s<   € ð	%Ø�N‰N×!Ñ!Ô#à×Ñ×"Ñ"Õ$øˆD×Ñ×"Ñ"Õ$ús	   ‚7 ·Ac                 ón   — d | _         | j                  �"| j                  j                  «        d | _        y y r   )r1   r-   Úcollectr>   s    r   Úon_after_forkzResultConsumer.on_after_forkE   s4   € ØˆŒØ×ÑÐ'Ø×Ñ×$Ñ$Ô&Ø#ˆDÕð (r   c                 ó  — | j                   €| j                  |«      S | j                  |«      }| j                   j                  |«      s6| j                   j	                  |«       | j                   j                  «        y y r   )r1   r5   r!   Úconsuming_fromÚ	add_queuer2   )r#   Útask_idÚqueues      r   Úconsume_fromzResultConsumer.consume_fromK   sd   € Ø�>‰>Ð!Ø—:‘:˜gÓ&Ð&Ø×$Ñ$ WÓ-ˆØ�~‰~×,Ñ,¨UÔ3Ø�N‰N×$Ñ$ UÔ+Ø�N‰N×"Ñ"Õ$ð 4r   c                 ó†   — | j                   r5| j                   j                  | j                  |«      j                  «       y y r   )r1   Úcancel_by_queuer!   Úname©r#   rF   s     r   Ú
cancel_forzResultConsumer.cancel_forS   s1   € Ø�>Š>Ø�N‰N×*Ñ*¨4×+?Ñ+?ÀÓ+H×+MÑ+MÕNð r   ©Tr   )r   r   r   Úkombur.   r-   r1   r    r5   r8   r?   rB   rH   rM   Ú__classcell__©r&   s   @r   r   r   &   s:   ø„ Ø�~‰~€Hà€KØ€Iô<ó!ó ò%ò$ò%öOr   r   c                   ó’  ‡ — e Zd ZdZej
                  Zej                  ZeZeZdZ	dZ
dZdddddœZ G d„ d	ej                  «      Z G d
„ dej                  «      Z	 	 d&ˆ fd„	Zd„ Zd'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eZd„ Z	 d+d„Zd„ Z d„ Z!d„ Z"d„ Z#d „ Z$d(d!„Z%d"„ Z&d,ˆ fd#„	Z'e(d$„ «       Z)e*d%„ «       Z+ˆ xZ,S )-r   z&Base class for the RPC result backend.FTé   r   r	   )Úmax_retriesÚinterval_startÚinterval_stepÚinterval_maxc                   ó   — e Zd ZdZdZy)úRPCBackend.Consumerz4Consumer that requires manual declaration of queues.FN)r   r   r   r   Úauto_declarer   r   r   r.   rY   m   s
   „ ÙBà‰r   r.   c                   ó   — e Zd ZdZdZy)úRPCBackend.Queuez$Queue that never caches declaration.FN)r   r   r   r   Úcan_cache_declarationr   r   r   ÚQueuer\   r   s   „ Ù2à %Ñr   r^   c                 ó2  •— t        ‰
| �  |fi |¤Ž | j                  j                  }	|| _        i | _        | j                  |«      | _        | j                  rdnd| _        |xs |	j                  }|xs |	j                  }| j                  ||| j                  «      | _        |xs |	j                  | _        || _        | j!                  | | j                  | j"                  | j$                  | j&                  «      | _        t*        �t+        | t,        «       y y )Né   r	   )r   r    r+   Úconfr-   Ú_out_of_bandÚprepare_persistentÚ
persistentÚdelivery_modeÚresult_exchangeÚresult_exchange_typeÚ_create_exchangeÚexchangeÚresult_serializerÚ
serializerÚauto_deleter   r*   Ú_pending_resultsÚ_pending_messagesÚresult_consumerr   r   )r#   r+   r,   ri   Úexchange_typerd   rk   rl   r%   ra   r&   s             €r   r    zRPCBackend.__init__w   sú   ø€ ä‰Ñ˜Ñ' Ò'Ø�x‰x�}‰}ˆØ%ˆÔØˆÔØ×1Ñ1°*Ó=ˆŒØ"&§/¢/™Q°qˆÔØÒ3˜t×3Ñ3ˆØ%ÒB¨×)BÑ)BˆØ×-Ñ-Ø�m T×%7Ñ%7ó
ˆŒð %Ò>¨×(>Ñ(>ˆŒØ&ˆÔØ#×2Ñ2Ø�$—(‘(˜DŸK™KØ×!Ñ! 4×#9Ñ#9ó 
ˆÔô Ð*Ü Ô&DÕEð +r   c                 ól   — | j                   j                  «        | j                  j                  «        y r   )rm   Úclearro   r   r>   s    r   r   zRPCBackend._after_fork�   s&   € à×Ñ×#Ñ#Ô%Ø×Ñ×(Ñ(Õ*r   c                 ó$   — | j                  d «      S r   )ÚExchange)r#   rK   Útypere   s       r   rh   zRPCBackend._create_exchange’   s   € à�}‰}˜TÓ"Ð"r   c                 ó   — | j                   S )z$Create new binding for task with id.)ÚbindingrL   s     r   r!   zRPCBackend._create_binding–   s   € ð �|‰|Ðr   c                 ó<   — t        t        j                  «       «      ‚r   )ÚNotImplementedErrorÚE_NO_CHORD_SUPPORTÚstripr>   s    r   Úensure_chords_allowedz RPCBackend.ensure_chords_allowed›   s   € Ü!Ô"4×":Ñ":Ó"<Ó=Ð=r   c                 óf   — t        «       s't        | j                  |j                  «      d¬«       y y )NT)Úretry)r   r   rw   Úchannel)r#   ÚproducerrF   s      r   Úon_task_callzRPCBackend.on_task_callž   s(   € ô
 $Ô%Ü˜$Ÿ,™, x×'7Ñ'7Ó8ÀÖEð &r   c                 óš   — 	 |xs t         j                  }|j                  |j
                  xs |fS # t        $ r t        d|›�«      ‚w xY w)z‹Get the destination for result by task id.

        Returns:
            Tuple[str, str]: tuple of ``(reply_to, correlation_id)``.
        z%RPC backend missing task request for )r   ÚrequestÚAttributeErrorÚRuntimeErrorÚreply_toÚcorrelation_id)r#   rF   rƒ   s      r   Údestination_forzRPCBackend.destination_for¦   sb   € ð	EØÒ5¤×!5Ñ!5ˆGð ×Ñ ×!7Ñ!7Ò!B¸7ÐBÐBøô ò 	EÜØ7¸°{ÐCóEð Eð	Eús	   ‚2 ²A
c                  ó   — y r   r   rL   s     r   Úon_reply_declarezRPCBackend.on_reply_declareµ   s   € ð 	r   c                  ó   — y r   r   )r#   Úresults     r   Úon_result_fulfilledzRPCBackend.on_result_fulfilled»   s   € ð 	r   c                  ó   — y)Nzrpc://r   )r#   Úinclude_passwords     r   Úas_urizRPCBackend.as_uriÀ   s   € Ør   c                 óŠ  — | j                  ||«      \  }}|sy| j                  j                  j                  j	                  d¬«      5 }	|	j                  | j                  |||||«      | j                  ||| j                  d| j                  | j                  |«      | j                  ¬«	       ddd«       |S # 1 sw Y   |S xY w)z!Send task return value and state.NT©Úblock)ri   Úrouting_keyr‡   rk   r~   Úretry_policyÚdeclarere   )rˆ   r+   ÚamqpÚproducer_poolÚacquireÚpublishÚ
_to_resultri   rk   r•   rŠ   re   )
r#   rF   rŒ   ÚstateÚ	tracebackrƒ   r%   r”   r‡   r€   s
             r   Ústore_resultzRPCBackend.store_resultÃ   s¸   € ð '+×&:Ñ&:¸7ÀGÓ&LÑ#ˆ�^ÙØØ�X‰X�]‰]×(Ñ(×0Ñ0°tÐ0Ô<ÀØ×ÑØ—‘ ¨°¸	À7ÓKØŸ™Ø'Ø-ØŸ?™?Ø¨×):Ñ):Ø×-Ñ-¨gÓ6Ø"×0Ñ0ð ô 	÷ =ð ˆ÷ =ð ˆús   Á	A%B8Â8Cc                 óP   — ||| j                  ||«      || j                  |«      dœS )N)rF   ÚstatusrŒ   r�   Úchildren)Úencode_resultÚcurrent_task_children)r#   rF   rœ   rŒ   r�   rƒ   s         r   r›   zRPCBackend._to_resultÖ   s3   € àØØ×(Ñ(¨°Ó7Ø"Ø×2Ñ2°7Ó;ñ
ð 	
r   c                 óp   — | j                   r| j                   j                  |«       || j                  |<   y r   )ro   Úon_out_of_band_resultrb   )r#   rF   Úmessages      r   r¥   z RPCBackend.on_out_of_band_resultß   s1   € ð
 ×ÒØ× Ñ ×6Ñ6°wÔ?Ø%,ˆ×Ñ˜'Ò"r   c                 óL  — | j                   j                  |d «      }|r| j                  ||«      S i }d }| j                  || j                  |«      D ]?  }| j                  |«      }|j                  |«      |c}||<   |sŒ.|j                  «        d }ŒA |j                  |d «      }|j                  «       D ]  \  }}	| j                  ||	«       Œ |r"|j                  «        | j                  ||«      S 	 | j                  |   S # t        $ r t        j                  d dœcY S w xY w)N)r    rŒ   )rb   ÚpopÚ_set_cache_by_messageÚ_slurp_from_queuer*   Ú_get_message_task_idÚgetÚackÚitemsr¥   ÚrequeueÚ_cacheÚKeyErrorr   ÚPENDING)
r#   rF   Úbacklog_limitÚbufferedÚlatest_by_idÚprevÚaccÚtidÚlatestÚmsgs
             r   Úget_task_metazRPCBackend.get_task_metaè   s+  € Ø×$Ñ$×(Ñ(¨°$Ó7ˆÙØ×-Ñ-¨g°xÓ@Ð@ð ˆØˆØ×)Ñ)¨'°4·;±;ÀÖNˆCØ×+Ñ+¨CÓ0ˆCØ&2×&6Ñ&6°sÓ&;¸SÐ#ˆD�,˜sÑ#Úð —‘”
Ø‘ð Oð ×!Ñ! '¨4Ó0ˆØ$×*Ñ*Ö,‰HˆC�Ø×&Ñ& s¨CÕ0ð -ñ Ø�N‰NÔØ×-Ñ-¨g°vÓ>Ð>ðBØ—{‘{ 7Ñ+Ð+øÜò Bä"(§.¡.¸DÑAÒAðBús   Ã5D ÄD#Ä"D#c                 óZ   — | j                  |j                  «      x}| j                  |<   |S r   )Úmeta_from_decodedÚpayloadr°   )r#   rF   r¦   r¾   s       r   r©   z RPCBackend._set_cache_by_message	  s.   € Ø)-×)?Ñ)?Ø�O‰Oó*ð 	ˆ�$—+‘+˜gÑ&àˆr   c              #   óP  K  — | j                   j                  j                  d¬«      5 \  }} | j                  |«      |«      }|j	                  «        t        |«      D ]  }|j                  ||¬«      }|s n|–— Œ | j                  |«      ‚	 d d d «       y # 1 sw Y   y xY w­w)NTr’   )r*   r)   )r+   ÚpoolÚacquire_channelr!   r–   Úranger¬   r   )	r#   rF   r*   Úlimitr)   Ú_r   rw   rº   s	            r   rª   zRPCBackend._slurp_from_queue  s”   è ø€ à�X‰X�]‰]×*Ñ*°Ð*Ô6¹,¸1¸gØ3�d×*Ñ*¨7Ó3°GÓ<ˆGØ�O‰OÔä˜5–\�Ø—k‘k¨¸�kÓ?�ÙÙØ“	ð	 "ð ×/Ñ/°Ó8Ð8ð ÷ 7×6Ñ6üs   ‚'B&©A'BÂ	B&ÂB#ÂB&c                 ój   — 	 |j                   d   S # t        t        f$ r |j                  d   cY S w xY w)Nr‡   rF   )Ú
propertiesr„   r±   r¾   )r#   r¦   s     r   r«   zRPCBackend._get_message_task_id  s>   € ð	.ð ×%Ñ%Ð&6Ñ7Ð7øÜ¤Ð)ò 	.à—?‘? 9Ñ-Ò-ð	.ús   ‚ ‘2±2c                  ó   — y r   r   )r#   r   s     r   ÚrevivezRPCBackend.revive%  s   € Ør   c                 ó   — t        d«      ‚)Nz4reload_task_result is not supported by this backend.©ry   rL   s     r   Úreload_task_resultzRPCBackend.reload_task_result(  s   € Ü!ØBóDð 	Dr   c                 ó   — t        d«      ‚)z<Reload group result, even if it has been previously fetched.z5reload_group_result is not supported by this backend.rÊ   rL   s     r   Úreload_group_resultzRPCBackend.reload_group_result,  s   € ä!ØCóEð 	Er   c                 ó   — t        d«      ‚)Nz,save_group is not supported by this backend.rÊ   )r#   Úgroup_idrŒ   s      r   Ú
save_groupzRPCBackend.save_group1  s   € Ü!Ø:ó<ð 	<r   c                 ó   — t        d«      ‚)Nz/restore_group is not supported by this backend.rÊ   )r#   rÏ   Úcaches      r   Úrestore_groupzRPCBackend.restore_group5  s   € Ü!Ø=ó?ð 	?r   c                 ó   — t        d«      ‚)Nz.delete_group is not supported by this backend.rÊ   )r#   rÏ   s     r   Údelete_groupzRPCBackend.delete_group9  s   € Ü!Ø<ó>ð 	>r   c                 ó  •— |si n|}t         ‰| �  |t        || j                  | j                  j
                  | j                  j                  | j                  | j                  | j                  | j                  ¬«      «      S )N)r,   ri   rp   rd   rk   rl   Úexpires)r   Ú
__reduce__Údictr-   ri   rK   ru   rd   rk   rl   r×   r"   s      €r   rØ   zRPCBackend.__reduce__=  sk   ø€ Ù!‘ vˆÜ‰wÑ! $¬ØØ×'Ñ'Ø—]‘]×'Ñ'ØŸ-™-×,Ñ,Ø—‘Ø—‘Ø×(Ñ(Ø—L‘Lô	)
ó 	ð 		r   c                 ó€   — | j                  | j                  | j                  | j                  dd| j                  ¬«      S )NFT)Údurablerl   r×   )r^   Úoidri   r×   r>   s    r   rw   zRPCBackend.bindingJ  s8   € à�z‰zØ�H‰H�d—m‘m T§X¡XØØØ—L‘Lð	 ó 
ð 	
r   c                 ó.   — | j                   j                  S r   )r+   Ú
thread_oidr>   s    r   rÜ   zRPCBackend.oidS  s   € ð �x‰x×"Ñ"Ð"r   )NNNNNT)Údirectr`   rN   )NN)éè  )rà   F)r   N)-r   r   r   r   rO   rt   ÚProducerr   r   rd   Úsupports_autoexpireÚsupports_native_joinr•   r.   r^   r    r   rh   r!   r|   r�   rˆ   rŠ   r�   r�   rž   r›   r¥   r»   Úpollr©   rª   r«   rÈ   rË   rÍ   rÐ   rÓ   rÕ   rØ   Úpropertyrw   r   rÜ   rP   rQ   s   @r   r   r   X   s,  ø„ Ù0à�~‰~€HØ�~‰~€HØ#€Nð 0Ðà€JØÐØÐð ØØØñ	€Lô�5—>‘>ô ô
&�—‘ô &ð
 KOØ?CõFò,+ó
#òò
>òFòCòòó
ð .2óò&
ò-óBð> €Dòð .3ó9ò.òòDòEò
<ó?ò>õð ñ
ó ð
ð ñ#ó ô#r   r   )r   r9   rO   Úkombu.commonr   Úkombu.utils.compatr   Úkombu.utils.objectsr   Úceleryr   Úcelery._stater   r   Ú r
   Úasynchronousr   r   Ú__all__rz   Ú	Exceptionr   r   r   ÚBackendr   r   r   r   Ú<module>rð      sj   ðñó ã Ý &Ý 2Ý /å ß <å ß ?à
0€ðÐ ô2˜9ô 2òô/OÐ'ô /Oôd~#�—‘Ð0õ ~#r   