o
    wvXj[   ã                   @   sä   d Z ddlmZ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 zddlZW n eyU   dZY nw edƒZdZdZG dd„ dejƒZG dd„ dejƒZdS )z}Etcd Transport.

It uses Etcd as a store to transport messages in Queues

It uses python-etcd for talking to Etcd's HTTP API
é    )Úabsolute_importÚunicode_literalsN)Údefaultdict)Úcontextmanager)ÚChannelError)ÚEmpty)Ú
get_logger)ÚloadsÚdumps)Úcached_propertyé   )Úvirtualzkombu.transport.etcdiK	  Ú	localhostc                       sŽ   e Zd ZdZdZdZdZdZdZ‡ fdd„Z	dd	„ Z
ed
d„ ƒZdd„ Zdd„ Zdd„ Zdd„ Zddd„Zdd„ Zdd„ Zedd„ ƒZ‡  ZS )ÚChannelz+Etcd Channel class which talks to the Etcd.ÚkombuNé
   é   c                    sz   t d u rtdƒ‚tt| ƒj|i |¤Ž | jjjp| jj}| jjj	p"t
}t d||| j¡ ttƒ| _t j|t|ƒd�| _d S )NúMissing python-etcd libraryzHost: %s Port: %s Timeout: %s©ÚhostÚport)ÚetcdÚImportErrorÚsuperr   Ú__init__Ú
connectionÚclientr   Údefault_portÚhostnameÚDEFAULT_HOSTÚloggerÚdebugÚtimeoutr   ÚdictÚqueuesÚClientÚint)ÚselfÚargsÚkwargsr   r   ©Ú	__class__© úQ/var/www/html/myproject/venv/lib/python3.10/site-packages/kombu/transport/etcd.pyr   +   s   
zChannel.__init__c                 C   s   d  | j|¡S )z‚Create and return the `queue` with the proper prefix.

        Arguments:
            queue (str): The name of the queue.
        z{0}/{1})ÚformatÚprefix)r'   Úqueuer,   r,   r-   Ú_key_prefix:   s   zChannel._key_prefixc                 c   s~   � t  | j|¡}| j|_t d |j¡¡ |j	d| j
d� zdV  W t d |j¡¡ | ¡  dS t d |j¡¡ | ¡  w )ag  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 {0}T)ÚblockingÚlock_ttlNzReleasing lock {0})r   ÚLockr   Ú
lock_valueÚ_uuidr    r!   r.   ÚnameÚacquirer3   Úrelease)r'   r0   Úlockr,   r,   r-   Ú_queue_lockB   s   €ÿ
zChannel._queue_lockc              	   K   sš   || j |< |  |¡�9 z| jj|  |¡ddd�W W  d  ƒ S  tjyB   t d 	|¡¡ | jj
|  |¡d� Y W  d  ƒ S w 1 sFw   Y  dS )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 "{0}" already exists©r<   )r$   r;   r   Úwriter1   r   ÚEtcdNotFiler    r!   r.   Úread)r'   r0   Ú_r,   r,   r-   Ú
_new_queueY   s   
ÿþúüzChannel._new_queuec                 K   s0   z| j  |  |¡¡ W dS  tjy   Y dS w )z£Verify that queue exists.

        Returns:
            bool: Should return :const:`True` if the queue exists
                or :const:`False` otherwise.
        TF)r   rB   r1   r   ÚEtcdKeyNotFound)r'   r0   r)   r,   r,   r-   Ú
_has_queueh   s   ÿzChannel._has_queuec                 O   s   | j  |d¡ |  |¡ dS )z^Delete a `queue`.

        Arguments:
            queue (str): The name of the queue.
        N)r$   ÚpopÚ_purge)r'   r0   r(   rC   r,   r,   r-   Ú_deleteu   s   zChannel._deletec                 K   s^   |   |¡�  |  |¡}| jj|t|ƒdd�std |¡ƒ‚W d  ƒ dS 1 s(w   Y  dS )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<   r>   ÚappendzCannot add key {0!r} to etcdN)r;   r1   r   r@   r
   r   r.   )r'   r0   ÚpayloadrC   r<   r,   r,   r-   Ú_put~   s   	
ýü"þzChannel._putc                 C   sø   |   |¡�m |  |¡}t d|| j¡ z;| jj|d| j| jd�}|du r'tƒ ‚|j	d }t d 
|d ¡¡ t|d ƒ}| jj|d d	� |W W  d  ƒ S  tttjfyq } zt d
 
t|ƒ|¡¡ W Y d}~tƒ ‚d}~ww 1 suw   Y  dS )aI  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 {0}r<   r>   r?   z_get failed: {0}:{1})r;   r1   r    r!   rN   r   rB   r"   r   Ú	_childrenr.   r	   ÚdeleteÚ	TypeErrorÚ
IndexErrorr   ÚEtcdExceptionÚtype)r'   r0   r"   r<   ÚresultÚitemÚmsg_contentÚerrorr,   r,   r-   Ú_get�   s,   

þ
ï €ýîzChannel._getc                 C   sX   |   |¡� |  |¡}t d |¡¡ | jj|dd�W  d  ƒ S 1 s%w   Y  dS )zrRemove all `message`s from a `queue`.

        Arguments:
            queue (str): The name of the queue.
        zPurging queue at key {0}T)r<   rM   N)r;   r1   r    r!   r.   r   rQ   )r'   r0   r<   r,   r,   r-   rH   °   s
   
$ýzChannel._purgec              	   C   s˜   |   |¡�= d}z|  |¡}t d|| j¡ | jj|d| jd�}t|jƒ}W n	 t	y/   Y nw t d||| j¡ |W  d  ƒ S 1 sEw   Y  dS )zlReturn the size of the `queue`.

        Arguments:
            queue (str): The name of the queue.
        r   z)Fetching key recursively %s with index %sT)r<   rM   rN   z$Found %s keys under %s with index %sN)
r;   r1   r    r!   rN   r   rB   ÚlenrP   rR   )r'   r0   Úsizer<   rV   r,   r,   r-   Ú_size»   s(   
ÿþÿÿ$ñzChannel._sizec                 C   s   d  t ¡ t ¡ ¡S )Nz{0}.{1})r.   ÚsocketÚgethostnameÚosÚgetpid)r'   r,   r,   r-   r5   Ò   s   zChannel.lock_value)N)Ú__name__Ú
__module__Ú__qualname__Ú__doc__r/   rN   r"   Úsession_ttlr3   r   r1   r   r;   rD   rF   rI   rL   rZ   rH   r]   r   r5   Ú__classcell__r,   r,   r*   r-   r   "   s(    
	
!r   c                       sZ   e Zd ZdZeZeZdZdZdZ	e
jjjedgƒd�Z‡ fdd„Zd	d
„ Zdd„ Z‡  ZS )Ú	Transportz!Etcd storage Transport for Kombu.r   úpython-etcdé   Údirect)Úexchange_typec                    sN   t du rtdƒ‚tt| ƒj|i |¤Ž tjjt jf | _tjjt jf | _dS )z(Create a new instance of etcd.Transport.Nr   )	r   r   r   rh   r   r   Úconnection_errorsrT   Úchannel_errors)r'   r(   r)   r*   r,   r-   r   ä   s   ÿÿzTransport.__init__c                 C   sV   |j jp| j}|j jpt}t d||¡ ztj|t	|ƒd� W dS  t
y*   Y dS 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   r,   r,   r-   Úverify_connectionó   s   ýzTransport.verify_connectionc              	   C   sb   zddl }|jj ¡ D ]}| d¡r| d¡d   W S qW dS  ttfy0   t d¡ Y dS w )z„Return the version of the etcd library.

        .. note::
           python-etcd has no __version__. This is a workaround.
        r   Nri   z==r   z'Unable to find the python-etcd version.ÚUnknown)	Úpip.commands.freezeÚcommandsÚfreezeÚ
startswithÚsplitr   rS   r    Úwarning)r'   ÚpipÚxr,   r,   r-   Údriver_version  s   
ÿÿ
þzTransport.driver_version)rb   rc   rd   re   r   ÚDEFAULT_PORTr   Údriver_typeÚdriver_nameÚpolling_intervalr   rh   Ú
implementsÚextendÚ	frozensetr   rp   rz   rg   r,   r,   r*   r-   rh   ×   s    ÿrh   )re   Ú
__future__r   r   r`   r^   Úcollectionsr   Ú
contextlibr   Úkombu.exceptionsr   Ú
kombu.fiver   Ú	kombu.logr   Úkombu.utils.jsonr	   r
   Úkombu.utils.objectsr   Ú r   r   r   r    r{   r   r   rh   r,   r,   r,   r-   Ú<module>   s.    ÿ 6