o
    wvXj.  ã                   @   s   d Z ddlmZmZ ddlZddlmZmZ ddlm	Z	m
Z
 ddlmZmZ ddlmZ dd	lmZ zdd
lmZmZmZ W n% eyc   zdd
lmZmZmZ W n ey`   d Z ZZY nw Y nw dd„ ejD ƒZG dd„ dejƒZG dd„ dejƒZdS )aq  Azure Service Bus Message Queue transport.

The transport can be enabled by setting the CELERY_BROKER_URL to:

```
azureservicebus://{SAS policy name}:{SAS key}@{Service Bus Namespace}
```

Note that the Shared Access Policy used to connect to Azure Service Bus
requires Manage, Send and Listen claims since the broker will create new
queues and delete old queues as required.

Note that if the SAS key for the Service Bus account contains a slash, it will
have to be regenerated before it can be used in the connection URL.

More information about Azure Service Bus:
https://azure.microsoft.com/en-us/services/service-bus/

é    )Úabsolute_importÚunicode_literalsN)ÚEmptyÚtext_t)Úbytes_to_strÚsafe_str)ÚloadsÚdumps)Úcached_propertyé   )Úvirtual)ÚServiceBusServiceÚMessageÚQueuec                 C   s   i | ]}|d vrt |ƒd“qS )Ú_é_   )Úord)Ú.0Úc© r   ú\/var/www/html/myproject/venv/lib/python3.10/site-packages/kombu/transport/azureservicebus.pyÚ
<dictcomp>+   s    r   c                       sÎ   e Zd ZdZdZdZdZdZdZi Z	‡ fdd„Z
efd	d
„Zdd„ Z‡ fdd„Zdd„ Zd%dd„Zdd„ Zdd„ Zedd„ ƒZedd„ ƒZedd„ ƒZedd„ ƒZedd „ ƒZed!d"„ ƒZed#d$„ ƒZ‡  ZS )&ÚChannelzAzure Service Bus channel.i  é   Fzkombu%(vhost)sNc                    sD   t d u rtdƒ‚tt| ƒj|i |¤Ž | j ¡ D ]}|| j|< qd S )NzAAzure Service Bus transport requires the azure-servicebus library)r   ÚImportErrorÚsuperr   Ú__init__Úqueue_serviceÚlist_queuesÚ_queue_cache)ÚselfÚargsÚkwargsÚqueue©Ú	__class__r   r   r   :   s   ÿzChannel.__init__c                 C   s   t t|ƒƒ |¡S )z:Format AMQP queue name into a valid ServiceBus queue name.)r   r   Ú	translate)r    ÚnameÚtabler   r   r   Úentity_nameD   s   zChannel.entity_namec                 K   sZ   |   | j| ¡}z| j| W S  ty,   | jj|dd� | j |¡ }| j|< | Y S w )z$Ensure a queue exists in ServiceBus.F)Úfail_on_exist)r)   Úqueue_name_prefixr   ÚKeyErrorr   Úcreate_queueÚ	get_queue)r    r#   r"   Úqr   r   r   Ú
_new_queueH   s   ýzChannel._new_queuec                    s8   |   |¡}| j |d¡ | j |¡ tt| ƒ |¡ dS )zDelete queue by name.N)r)   r   Úpopr   Údelete_queuer   r   Ú_delete)r    r#   r!   r"   Ú
queue_namer$   r   r   r3   R   s   
zChannel._deletec                 K   s$   t t|ƒƒ}| j |  |¡|¡ dS )zPut message onto queue.N)r   r	   r   Úsend_queue_messager)   )r    r#   Úmessager"   Úmsgr   r   r   Ú_putY   s   zChannel._putc                 C   s>   | j j|  |¡|p| j| jd�}|jdu rtƒ ‚tt|jƒƒS )z/Try to retrieve a single message off ``queue``.)ÚtimeoutÚ	peek_lockN)	r   Úreceive_queue_messager)   Úwait_time_secondsr:   Úbodyr   r   r   )r    r#   r9   r6   r   r   r   Ú_get^   s   ý
zChannel._getc                 C   s   |   |¡jS )z)Return the number of messages in a queue.)r0   Úmessage_count)r    r#   r   r   r   Ú_sizek   s   zChannel._sizec                 C   s2   d}	 | j j|  |¡dd�}|js	 |S |d7 }q)z'Delete all current messages in a queue.r   Tgš™™™™™¹?)r9   r   )r   Úread_delete_queue_messager)   r=   )r    r#   Únr6   r   r   r   Ú_purgeo   s   
ÿþùzChannel._purgec                 C   s,   | j d u rt| jj| jj| jjd�| _ | j S )N)Úservice_namespaceÚshared_access_key_nameÚshared_access_key_value)Ú_queue_servicer   ÚconninfoÚhostnameÚuseridÚpassword©r    r   r   r   r   ~   s   
ýzChannel.queue_servicec                 C   s   | j jS ©N)Ú
connectionÚclientrL   r   r   r   rH   ˆ   s   zChannel.conninfoc                 C   s
   | j jjS rM   )rN   rO   Útransport_optionsrL   r   r   r   rP   Œ   s   
zChannel.transport_optionsc                 C   s   | j  d¡p| jS )NÚvisibility_timeout)rP   ÚgetÚdefault_visibility_timeoutrL   r   r   r   rQ   �   s   ÿzChannel.visibility_timeoutc                 C   s   | j  dd¡S )Nr+   Ú )rP   rR   rL   r   r   r   r+   •   s   zChannel.queue_name_prefixc                 C   ó   | j  d| j¡S )Nr<   )rP   rR   Údefault_wait_time_secondsrL   r   r   r   r<   ™   ó   ÿzChannel.wait_time_secondsc                 C   rU   )Nr:   )rP   rR   Údefault_peek_lockrL   r   r   r   r:   ž   rW   zChannel.peek_lockrM   )Ú__name__Ú
__module__Ú__qualname__Ú__doc__rS   rV   rX   Údomain_formatrG   r   r   ÚCHARS_REPLACE_TABLEr)   r0   r3   r8   r>   r@   rC   Úpropertyr   rH   rP   r
   rQ   r+   r<   r:   Ú__classcell__r   r   r$   r   r   0   s<    



	




r   c                   @   s   e Zd ZdZeZdZdZdS )Ú	TransportzAzure Service Bus transport.r   N)rY   rZ   r[   r\   r   Úpolling_intervalÚdefault_portr   r   r   r   ra   ¤   s
    ra   )r\   Ú
__future__r   r   ÚstringÚ
kombu.fiver   r   Úkombu.utils.encodingr   r   Úkombu.utils.jsonr   r	   Úkombu.utils.objectsr
   rT   r   Úazure.servicebusr   r   r   r   Úazure.servicebus.control_clientÚpunctuationr^   r   ra   r   r   r   r   Ú<module>   s.    ÿ€û	ÿt