o
    wvXjÔ  ã                   @   sf  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	m
Z
 ddlmZmZ ddlmZ z]ddlZdd	lmZ dd
lmZ ejjejjejjejjejjejjejjejjejjf	Zejj ejj!ejj"ejjejjejjejj#ejj$ejjejj%ejj&ejj'ejjejj(ej)fZ*W n e+yš   dZd ZZ*Y nw dZ,dZ-G dd„ dej.ƒZ.G dd„ dej/ƒZ/dS )a¥  Zookeeper transport.

:copyright: (c) 2010 - 2013 by Mahendra M.
:license: BSD, see LICENSE for more details.

**Synopsis**

Connects to a zookeeper node as <server>:<port>/<vhost>
The <vhost> becomes the base for all the other znodes.  So we can use
it like a vhost.

This uses the built-in kazoo recipe for queues

**References**

- https://zookeeper.apache.org/doc/trunk/recipes.html#sc_recipes_Queues
- https://kazoo.readthedocs.io/en/latest/api/recipe/queue.html

**Limitations**
This queue does not offer reliable consumption.  An entry is removed from
the queue prior to being processed.  So if an error occurs, the consumer
has to re-queue the item or it will be lost.
é    )Úabsolute_importÚunicode_literalsN)ÚEmpty)Úbytes_to_strÚensure_bytes)ÚdumpsÚloadsé   )Úvirtual)ÚKazooClient)ÚQueue© i…  z!Mahendra M <mahendra.m@gmail.com>c                       s„   e Zd ZdZdZi 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dd„ Zdd„ Zdd„ Zedd„ ƒZ‡  ZS )ÚChannelzZookeeper Channel.Nc                    s8   t t| ƒj|fi |¤Ž | jjj}d | d¡¡| _d S )Nz/{}ú/)	Úsuperr   Ú__init__Ú
connectionÚclientÚvirtual_hostÚformatÚstripÚ_vhost)Úselfr   ÚkwargsÚvhost©Ú	__class__r   úV/var/www/html/myproject/venv/lib/python3.10/site-packages/kombu/transport/zookeeper.pyr   T   s   
zChannel.__init__c                 C   s   t j | j|¡S ©N)ÚosÚpathÚjoinr   )r   Ú
queue_namer   r   r   Ú	_get_pathY   s   zChannel._get_pathc                 C   s>   | j  |d ¡}|d u rt| j|  |¡ƒ}|| j |< t|ƒ |S r   )Ú_queuesÚgetr   r   r#   Úlen)r   r"   Úqueuer   r   r   Ú
_get_queue\   s   
zChannel._get_queuec                 K   s&   |   |¡jtt|ƒƒ| j|dd�d�S )NT)Úreverse)Úpriority)r(   Úputr   r   Ú_get_message_priority)r   r'   Úmessager   r   r   r   Ú_puth   s   

þzChannel._putc                 C   s,   |   |¡}| ¡ }|d u rtƒ ‚tt|ƒƒS r   )r(   r%   r   r   r   )r   r'   Úmsgr   r   r   Ú_getn   s
   
zChannel._getc                 C   s0   d}|   |¡}	 | ¡ }|d u r	 |S |d7 }q)Nr   Tr	   )r(   r%   )r   r'   Úcountr/   r   r   r   Ú_purgew   s   
þüzChannel._purgec                 O   s.   |   |¡r|  |¡ | j |  |¡¡ d S d S r   )Ú
_has_queuer2   r   Údeleter#   )r   r'   Úargsr   r   r   r   Ú_deleteƒ   s   

þzChannel._deletec                 C   s   |   |¡}t|ƒS r   )r(   r&   ©r   r'   r   r   r   Ú_sizeˆ   s   
zChannel._sizec                 K   s   |   |¡s|  |¡}d S d S r   )r3   r(   )r   r'   r   r   r   r   Ú
_new_queueŒ   s   
ÿzChannel._new_queuec                 C   s   | j  |  |¡¡d uS r   )r   Úexistsr#   r7   r   r   r   r3   �   s   zChannel._has_queuec              	   C   sê   | j j}g }|jrO|jD ]B}| d¡r|tdƒd … }|sqz| dd¡\}}|t|ƒf}W n tyH   ||jkrB||j	p?t
f}n|t
f}Y nw | |¡ q|j|j	pUt
f}||vra| d|¡ d dd„ |D ƒ¡}t|ƒ}| ¡  |S )Nzzookeeper://ú:r	   r   ú,c                 S   s   g | ]
\}}d ||f ‘qS )z%s:%sr   )Ú.0ÚhÚpr   r   r   Ú
<listcomp>¨   s    z!Channel._open.<locals>.<listcomp>)r   r   ÚaltÚ
startswithr&   ÚsplitÚintÚ
ValueErrorÚhostnameÚportÚDEFAULT_PORTÚappendÚinsertr!   r   Ústart)r   ÚconninfoÚhostsÚ	host_portÚhostrG   Úconn_strÚconnr   r   r   Ú_open“   s2   


€üzChannel._openc                 C   s   | j d u r
|  ¡ | _ | j S r   )Ú_clientrR   ©r   r   r   r   r   ­   s   

zChannel.client)Ú__name__Ú
__module__Ú__qualname__Ú__doc__rS   r$   r   r#   r(   r.   r0   r2   r6   r8   r9   r3   rR   Úpropertyr   Ú__classcell__r   r   r   r   r   N   s"    	r   c                       sT   e Zd ZdZeZdZeZej	j
e Z
ej	je ZdZdZ‡ fdd„Zdd„ Z‡  ZS )	Ú	TransportzZookeeper Transport.r	   Ú	zookeeperÚkazooc                    s*   t d u rtdƒ‚tt| ƒj|i |¤Ž d S )Nz"The kazoo library is not installed)r]   ÚImportErrorr   r[   r   )r   r5   r   r   r   r   r   Ã   s   zTransport.__init__c                 C   s   t jS r   )r]   Ú__version__rT   r   r   r   Údriver_versionÉ   s   zTransport.driver_version)rU   rV   rW   rX   r   Úpolling_intervalrH   Údefault_portr
   r[   Úconnection_errorsÚKZ_CONNECTION_ERRORSÚchannel_errorsÚKZ_CHANNEL_ERRORSÚdriver_typeÚdriver_namer   r`   rZ   r   r   r   r   r[   ´   s    
ÿ
ÿr[   )0rX   Ú
__future__r   r   r   ÚsocketÚ
kombu.fiver   Úkombu.utils.encodingr   r   Úkombu.utils.jsonr   r   Ú r
   r]   Úkazoo.clientr   Úkazoo.recipe.queuer   Ú
exceptionsÚSystemErrorExceptionÚConnectionLossExceptionÚMarshallingErrorExceptionÚUnimplementedExceptionÚOperationTimeoutExceptionÚNoAuthExceptionÚInvalidACLExceptionÚAuthFailedExceptionÚSessionExpiredExceptionrd   ÚRuntimeInconsistencyExceptionÚDataInconsistencyExceptionÚBadArgumentsExceptionÚApiErrorExceptionÚNoNodeExceptionÚNodeExistsExceptionÚ NoChildrenForEphemeralsExceptionÚNotEmptyExceptionÚInvalidCallbackExceptionÚerrorrf   r^   rH   Ú
__author__r   r[   r   r   r   r   Ú<module>   s\    ÷ñþf