Ë
    UV.jÄ!  ã                  ó  — d Z ddlm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 dd
lmZ ddlmZ 	 ddlZ ed«      ZdZdZ G d„ dej6                  «      Z G d„ dej8                  «      Zy# e$ r dZY Œ@w xY w)aÍ  Etcd Transport module for Kombu.

It uses Etcd as a store to transport messages in Queues

It uses python-etcd for talking to Etcd's HTTP API

Features
========
* Type: Virtual
* Supports Direct: *Unreviewed*
* Supports Topic: *Unreviewed*
* Supports Fanout: *Unreviewed*
* Supports Priority: *Unreviewed*
* Supports TTL: *Unreviewed*

Connection String
=================

Connection string has the following format:

.. code-block::

    'etcd'://SERVER:PORT

é    )ÚannotationsN)Údefaultdict)Úcontextmanager)ÚEmpty)ÚChannelError)Ú
get_logger)ÚdumpsÚloads)Úcached_propertyé   )Úvirtualzkombu.transport.etcdiK	  Ú	localhostc                  óˆ   ‡ — e Zd ZdZdZdZdZdZdZˆ fd„Z	d„ Z
ed„ «       Zd	„ Zd
„ Zd„ Zd„ Zdd„Zd„ Zd„ Zed„ «       Zˆ xZS )ÚChannelz+Etcd Channel class which talks to the Etcd.ÚkombuNé
   é   c                ó¼  •— t         €t        d«      ‚t        ‰| �  |i |¤Ž | j                  j
                  j                  xs | j                  j                  }| j                  j
                  j                  xs t        }t        j                  d||| j                  «       t        t        «      | _        t        j                   |t#        |«      ¬«      | _        y )NúMissing python-etcd libraryzHost: %s Port: %s Timeout: %s©ÚhostÚport)ÚetcdÚImportErrorÚsuperÚ__init__Ú
connectionÚclientr   Údefault_portÚhostnameÚDEFAULT_HOSTÚloggerÚdebugÚtimeoutr   ÚdictÚqueuesÚClientÚint)ÚselfÚargsÚkwargsr   r   Ú	__class__s        €úVC:\xampp\htdocs\tradingbinance\backend\.venv\Lib\site-packages\kombu/transport/etcd.pyr   zChannel.__init__>   s�   ø€ Üˆ<ÜÐ;Ó<Ð<ä‰Ñ˜$Ð) &Ò)à�‰×%Ñ%×*Ñ*ÒJ¨d¯o©o×.JÑ.JˆØ�‰×%Ñ%×.Ñ.Ò>´,ˆä�‰Ð4°d¸DÀ$Ç,Á,ÔOä!¤$Ó'ˆŒä—k‘k t´#°d³)Ô<ˆ�ó    c                ó$   — | j                   › d|› �S )z”Create and return the `queue` with the proper prefix.

        Arguments:
        ---------
            queue (str): The name of the queue.
        Ú/)Úprefix)r)   Úqueues     r-   Ú_key_prefixzChannel._key_prefixM   s   € ð —+‘+�˜a ˜wÐ'Ð'r.   c              #  óÈ  K  — t        j                  | j                  |«      }| j                  |_        t
        j                  d|j                  › �«       |j                  d| j                  ¬«       	 d–— t
        j                  d|j                  › �«       |j                  «        y# t
        j                  d|j                  › �«       |j                  «        w xY w­w)ay  Try to acquire a lock on the Queue.

        It does so by creating a object called 'lock' which is locked by the
        current session..

        This way other nodes are not able to write to the lock object which
        means that they have to wait before the lock is released.

        Arguments:
        ---------
            queue (str): The name of the queue.
        zAcquiring lock T)ÚblockingÚlock_ttlNzReleasing lock )r   ÚLockr   Ú
lock_valueÚ_uuidr"   r#   ÚnameÚacquirer6   Úrelease)r)   r2   Úlocks      r-   Ú_queue_lockzChannel._queue_lockV   s�   è ø€ ô �y‰y˜Ÿ™ eÓ,ˆØ—_‘_ˆŒ
Ü�‰� t§y¡y kÐ2Ô3Ø�‰˜d¨T¯]©]ˆÔ;ð	Ûä�L‰L˜?¨4¯9©9¨+Ð6Ô7Ø�L‰L�Nøô �L‰L˜?¨4¯9©9¨+Ð6Ô7Ø�L‰L�Nüs   ‚A1C"Á4B+ Á83C"Â+4CÃC"c                ó–  — || j                   |<   | j                  |«      5  	 | j                  j                  | j	                  |«      dd¬«      cddd«       S # t
        j                  $ rP t        j                  d|› d�«       | j                  j                  | j	                  |«      ¬«      cY cddd«       S w xY w# 1 sw Y   yxY w)z™Create a new `queue` if the `queue` doesn't already exist.

        Arguments:
        ---------
            queue (str): The name of the queue.
        TN)ÚkeyÚdirÚvaluezQueue "z" already exists©r@   )
r&   r>   r   Úwriter3   r   ÚEtcdNotFiler"   r#   Úread)r)   r2   Ú_s      r-   Ú
_new_queuezChannel._new_queuen   s»   € ð #ˆ�‰�EÑØ×Ñ˜eÕ$ðEØ—{‘{×(Ñ(Ø×(Ñ(¨Ó/°TÀð )ó G÷ %Ñ$øô ×#Ñ#ò EÜ—‘˜w u gÐ-=Ð>Ô?Ø—{‘{×'Ñ'¨D×,<Ñ,<¸UÓ,CÐ'ÓDÑD÷ %Ñ$ðEú÷	 %Ð$ús)   ¡B?£,AÁAB<Â0B?Â;B<Â<B?Â?Cc                óŒ   — 	 | j                   j                  | j                  |«      «       y# t        j                  $ r Y yw xY w)z²Verify that queue exists.

        Returns
        -------
            bool: Should return :const:`True` if the queue exists
                or :const:`False` otherwise.
        TF)r   rF   r3   r   ÚEtcdKeyNotFound)r)   r2   r+   s      r-   Ú
_has_queuezChannel._has_queue~   s?   € ð	Ø�K‰K×Ñ˜T×-Ñ-¨eÓ4Ô5ØøÜ×#Ñ#ò 	Ùð	ús   ‚*- ­AÁAc                ó^   — | j                   j                  |d«       | j                  |«       y)zpDelete a `queue`.

        Arguments:
        ---------
            queue (str): The name of the queue.
        N)r&   ÚpopÚ_purge)r)   r2   r*   rG   s       r-   Ú_deletezChannel._deleteŒ   s"   € ð 	�‰�‰˜˜tÔ$Ø�‰�EÕr.   c                óà   — | j                  |«      5  | j                  |«      }| j                  j                  |t	        |«      d¬«      st        d|›d�«      ‚	 ddd«       y# 1 sw Y   yxY w)zõPut `message` onto `queue`.

        This simply writes a key to the Etcd store

        Arguments:
        ---------
            queue (str): The name of the queue.
            payload (dict): Message data which will be dumped to etcd.
        T)r@   rB   ÚappendzCannot add key z to etcdN)r>   r3   r   rD   r	   r   )r)   r2   ÚpayloadrG   r@   s        r-   Ú_putzChannel._put–   sm   € ð ×Ñ˜eÕ$Ø×"Ñ" 5Ó)ˆCØ—;‘;×$Ñ$ØÜ ›.Øð %ô !ô # _°S°G¸8Ð#DÓEÐEð	!÷ %×$Ñ$ús   ’AA$Á$A-c                ó®  — | j                  |«      5  | j                  |«      }t        j                  d|| j                  «       	 | j
                  j                  |d| j                  | j                  ¬«      }|€
t        «       ‚|j                  d   }t        j                  dj                  |d   «      «       t        |d   «      }| j
                  j                  |d   ¬	«       |cddd«       S # t        t        t        j                   f$ r7}t        j                  d
t#        |«      › d|› �«       Y d}~t        «       ‚d}~ww xY w# 1 sw Y   yxY w)a[  Get the first available message from the queue.

        Before it does so it acquires a lock on the store so
        only one node reads at the same time. This is for read consistency

        Arguments:
        ---------
            queue (str): The name of the queue.
            timeout (int): Optional seconds to wait for a response.
        zFetching key %s with index %sT)r@   Ú	recursiveÚindexr$   NéÿÿÿÿzRemoving key {}r@   rB   rC   z_get failed: Ú:)r>   r3   r"   r#   rV   r   rF   r$   r   Ú	_childrenÚformatr
   ÚdeleteÚ	TypeErrorÚ
IndexErrorr   ÚEtcdExceptionÚtype)r)   r2   r$   r@   ÚresultÚitemÚmsg_contentÚerrors           r-   Ú_getzChannel._get¨   s%  € ð ×Ñ˜eÕ$Ø×"Ñ" 5Ó)ˆCÜ�L‰LÐ8¸#¸t¿z¹zÔJðDØŸ™×)Ñ)Ø tØŸ*™*¨d¯l©lð *ó <�ð �>Ü›'�Mà×'Ñ'¨Ñ+�Ü—‘Ð.×5Ñ5°d¸5±kÓBÔCä# D¨¡MÓ2�Ø—‘×"Ñ" t¨E¡{Ð"Ô3Ø"÷# %Ñ$øô$ œz¬4×+=Ñ+=Ð>ò DÜ—‘˜}¬T°%«[¨M¸¸5¸'ÐB×CÐCä“'ˆMûðDú÷% %Ð$ús0   ’3EÁB#C3Ã3EÄ$EÄ5EÅEÅEÅEc                óÜ   — | j                  |«      5  | j                  |«      }t        j                  d|› �«       | j                  j                  |d¬«      cddd«       S # 1 sw Y   yxY w)z„Remove all `message`s from a `queue`.

        Arguments:
        ---------
            queue (str): The name of the queue.
        zPurging queue at key T)r@   rU   N)r>   r3   r"   r#   r   r[   )r)   r2   r@   s      r-   rN   zChannel._purgeÊ   sY   € ð ×Ñ˜eÕ$Ø×"Ñ" 5Ó)ˆCÜ�L‰LÐ0°°Ð6Ô7Ø—;‘;×%Ñ%¨#¸Ð%Ó>÷ %×$Ò$ús   ’AA"Á"A+c                óš  — | j                  |«      5  d}	 | j                  |«      }t        j                  d|| j                  «       | j
                  j                  |d| j                  ¬«      }t        |j                  «      }t        j                  d|| j                  «       |cddd«       S # t        $ r Y Œ8w xY w# 1 sw Y   yxY w)z~Return the size of the `queue`.

        Arguments:
        ---------
            queue (str): The name of the queue.
        r   z)Fetching key recursively %s with index %sT)r@   rU   rV   z$Found %s keys under %s with index %sN)
r>   r3   r"   r#   rV   r   rF   ÚlenrY   r\   )r)   r2   Úsizer@   r`   s        r-   Ú_sizezChannel._sizeÖ   s¿   € ð ×Ñ˜eÕ$ØˆDð	Ø×&Ñ& uÓ-�Ü—‘ÐHØ  $§*¡*ô.àŸ™×)Ñ)Ø tØŸ*™*ð *ó &�ô ˜6×+Ñ+Ó,�ô �L‰LÐ?Ø˜s D§J¡Jô0à÷ %Ñ$øô ò Ùðú÷ %Ð$ús/   ’C–A/B2Â#CÂ2	B>Â;CÂ=B>Â>CÃC
c                óX   — t        j                  «       › dt        j                  «       › �S )NÚ.)ÚsocketÚgethostnameÚosÚgetpid)r)   s    r-   r8   zChannel.lock_valueî   s#   € ä×$Ñ$Ó&Ð' q¬¯©«¨Ð6Ð6r.   )N)Ú__name__Ú
__module__Ú__qualname__Ú__doc__r1   rV   r$   Úsession_ttlr6   r   r3   r   r>   rH   rK   rO   rS   rd   rN   ri   r   r8   Ú__classcell__©r,   s   @r-   r   r   5   sw   ø„ Ù5à€FØ€EØ€GØ€KØ€Hô=ò(ð ñó ðò.Eò òòFó$ òD
?òð0 ñ7ó ô7r.   r   c                  ó0  ‡ — e Zd ZdZeZeZdZdZdZ	e
j                  j                  j                   edg«      ¬«      ZerHe
j                  j                   ej"                  fz   Ze
j                  j$                  ej"                  fz   Zˆ fd„Zd„ Zd	„ Zˆ xZS )
Ú	Transportz!Etcd storage Transport for Kombu.r   úpython-etcdé   Údirect)Úexchange_typec                óF   •— t         €t        d«      ‚t        ‰| �  |i |¤Ž y)z(Create a new instance of etcd.Transport.Nr   )r   r   r   r   )r)   r*   r+   r,   s      €r-   r   zTransport.__init__	  s&   ø€ äˆ<ÜÐ;Ó<Ð<ä‰Ñ˜$Ð) &Ó)r.   c                ó  — |j                   j                  xs | j                  }|j                   j                  xs t        }t
        j                  d||«       	 t        j                  |t        |«      ¬«       y# t        $ r Y yw xY w)zVerify the connection works.zVerify Etcd connection to %s:%sr   TF)r   r   r   r    r!   r"   r#   r   r'   r(   Ú
ValueError)r)   r   r   r   s       r-   Úverify_connectionzTransport.verify_connection  st   € à× Ñ ×%Ñ%Ò:¨×):Ñ):ˆØ× Ñ ×)Ñ)Ò9¬\ˆä�‰Ð6¸¸dÔCð	Ü�K‰K˜T¬¨D«	Õ2ØøÜò 	Øàð	ús   Á A< Á<	BÂBc                ó  — 	 ddl }|j                  j                  j                  «       D ])  }|j                  d«      sŒ|j	                  d«      d   c S  y# t
        t        f$ r t        j                  d«       Y yw xY w)z„Return the version of the etcd library.

        .. note::
           python-etcd has no __version__. This is a workaround.
        r   Nry   z==r   z'Unable to find the python-etcd version.ÚUnknown)	Úpip.commands.freezeÚcommandsÚfreezeÚ
startswithÚsplitr   r]   r"   Úwarning)r)   ÚpipÚxs      r-   Údriver_versionzTransport.driver_version  sl   € ð	Û&Ø—\‘\×(Ñ(×/Ñ/Ö1�Ø—<‘< Õ.ØŸ7™7 4›=¨Ñ+Ò+ñ 2øô œZÐ(ò 	Ü�N‰NÐDÔEÙð	ús   ‚<A ¿A ÁA Á$A>Á=A>)rp   rq   rr   rs   r   ÚDEFAULT_PORTr   Údriver_typeÚdriver_nameÚpolling_intervalr   rx   Ú
implementsÚextendÚ	frozensetr   Úconnection_errorsr^   Úchannel_errorsr   r€   r‹   ru   rv   s   @r-   rx   rx   ó   s¥   ø„ Ù+à€Gà€LØ€KØ€KØÐà×"Ñ"×-Ñ-×4Ñ4Ù  
Ó+ð 5ó -€Jñ à×Ñ×/Ñ/°4×3EÑ3EÐ2HÑHð 	ð
 ×Ñ×,Ñ,°×0BÑ0BÐ/EÑEð 	ô*òör.   rx   )rs   Ú
__future__r   rn   rl   Úcollectionsr   Ú
contextlibr   r2   r   Úkombu.exceptionsr   Ú	kombu.logr   Úkombu.utils.jsonr	   r
   Úkombu.utils.objectsr   Ú r   r   r   r"   rŒ   r!   r   rx   © r.   r-   Ú<module>rž      sˆ   ðñõ4 #ã 	Û Ý #Ý %Ý å )Ý  ß )Ý /å ðÛñ 
Ð*Ó	+€à€Ø€ô{7ˆg�o‰oô {7ô|9�×!Ñ!õ 9øðO ò Ø‚Dðús   ÁA? Á?B	ÂB	