o
    —¨Êh…G  ã                   @  sF  d Z ddlmZ ddlZddlmZ ddlmZmZm	Z	 ddl
ZddlZddlZddlmZmZmZmZmZ ddlmZ z
ddlmZmZ W n eyS   dZdZY nw dd	lmZmZ dd
lmZmZ ddl m!Z! ddl"m#Z# e$ej%ƒh d£ Z&e'dƒe'dƒidd„ e&D ƒ¥Z(G dd„ dƒZ)G dd„ de#j*ƒZ*G dd„ de#j+ƒZ+dS )aX  Azure Service Bus Message Queue transport module for kombu.

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.


Notes when using with Celery if you are experiencing issues with programs not
terminating properly. The Azure Service Bus SDK uses the Azure uAMQP library
which in turn creates some threads. If the AzureServiceBus Channel is closed,
said threads will be closed properly, but it seems there are times when Celery
does not do this so these threads will be left running. As the uAMQP threads
are not marked as Daemon threads, they will not be killed when the main thread
exits. Setting the ``uamqp_keep_alive_interval`` transport option to 0 will
prevent the keep_alive thread from starting


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

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 formats:

.. code-block::

    azureservicebus://SAS_POLICY_NAME:SAS_KEY@SERVICE_BUSNAMESPACE
    azureservicebus://DefaultAzureCredential@SERVICE_BUSNAMESPACE
    azureservicebus://ManagedIdentityCredential@SERVICE_BUSNAMESPACE

Transport Options
=================

* ``queue_name_prefix`` - String prefix to prepend to queue names in a
  service bus namespace.
* ``wait_time_seconds`` - Number of seconds to wait to receive messages.
  Default ``5``
* ``peek_lock_seconds`` - Number of seconds the message is visible for before
  it is requeued and sent to another consumer. Default ``60``
* ``uamqp_keep_alive_interval`` - Interval in seconds the Azure uAMQP library
  should send keepalive messages. Default ``30``
* ``retry_total`` - Azure SDK retry total. Default ``3``
* ``retry_backoff_factor`` - Azure SDK exponential backoff factor.
  Default ``0.8``
* ``retry_backoff_max`` - Azure SDK retry total time. Default ``120``
é    )ÚannotationsN)ÚEmpty)ÚAnyÚDictÚSet)ÚServiceBusClientÚServiceBusMessageÚServiceBusReceiveModeÚServiceBusReceiverÚServiceBusSender)ÚServiceBusAdministrationClient)ÚDefaultAzureCredentialÚManagedIdentityCredential)Úbytes_to_strÚsafe_str)ÚdumpsÚloads)Úcached_propertyé   )Úvirtual>   Ú.Ú_ú-r   r   c                 C  s   i | ]	}t |ƒt d ƒ“qS )r   )Úord)Ú.0Úc© r   úQ/var/www/html/env/lib/python3.10/site-packages/kombu/transport/azureservicebus.pyÚ
<dictcomp>Y   s    r   c                   @  s*   e Zd ZdZ		dddd„Zddd„ZdS )ÚSendReceivez"Container for Sender and Receiver.NÚreceiverúServiceBusReceiver | NoneÚsenderúServiceBusSender | Nonec                 C  s   || _ || _d S ©N)r    r"   )Úselfr    r"   r   r   r   Ú__init__`   s   
zSendReceive.__init__ÚreturnÚNonec                 C  s4   | j r| j  ¡  d | _ | jr| j ¡  d | _d S d S r$   )r    Úcloser"   ©r%   r   r   r   r)   f   s   


þzSendReceive.close©NN)r    r!   r"   r#   ©r'   r(   )Ú__name__Ú
__module__Ú__qualname__Ú__doc__r&   r)   r   r   r   r   r   ]   s    þr   c                      sæ  e Zd ZU dZdZded< dZded< dZded< d	Zded
< dZ	ded< dZ
ded< dZded< i Zded< eƒ Zded< ‡ fdd„Zdjdd„Z‡ fdd„Z‡ fdd „Z	!	!dkdld(d)„Zdmd+d,„Zejd!fdnd1d2„Z	!dodpd5d6„Zdqd9d:„Zdmd;d<„Zdrd=d>„Zdrd?d@„Z	!dodsdDdE„Zdtdu‡ fdJdK„ZdvdLdM„ZdwdNdO„Z djdPdQ„Z!e"dxdSdT„ƒZ#e"dydVdW„ƒZ$e%dXdY„ ƒZ&e%dZd[„ ƒZ'e"dzd\d]„ƒZ(e"dwd^d_„ƒZ)e"dwd`da„ƒZ*e"dwdbdc„ƒZ+e"dwddde„ƒZ,e"d{dfdg„ƒZ-e"dwdhdi„ƒZ.‡  Z/S )|ÚChannelzAzure Service Bus channel.é   ÚintÚdefault_wait_time_secondsé<   Údefault_peek_lock_secondsé   Ú!default_uamqp_keep_alive_intervalé   Údefault_retry_totalgš™™™™™é?ÚfloatÚdefault_retry_backoff_factoréx   Údefault_retry_backoff_maxzkombu%(vhost)sÚstrÚdomain_formatzDict[str, SendReceive]Ú_queue_cachezSet[str]Ú_noack_queuesc                   s>   t ƒ j|i |¤Ž d | _d | _d | _d | _|  ¡  d| j_d S )NF)	Úsuperr&   Ú
_namespaceÚ_policyÚ_sas_keyÚ_connection_stringÚ_try_parse_connection_stringÚqosÚrestore_at_shutdown)r%   ÚargsÚkwargs©Ú	__class__r   r   r&   €   s   zChannel.__init__r'   r(   c                 C  s–   t  | jj¡\| _| _td urt| jtƒstd ur!t| jtƒr!d S d| jv r1| j 	dd¡\| _
| _d| j | j
| jdœ}d dd„ | ¡ D ƒ¡| _d S )Nú:r   zsb://)ÚEndpointÚSharedAccessKeyNameÚSharedAccessKeyú;c                 S  s   g | ]
\}}|d  | ‘qS )ú=r   )r   ÚkeyÚvaluer   r   r   Ú
<listcomp>¢   s    z8Channel._try_parse_connection_string.<locals>.<listcomp>)Ú	TransportÚ	parse_uriÚconninfoÚhostnamerD   Ú_credentialr   Ú
isinstancer   ÚsplitrE   rF   ÚjoinÚitemsrG   )r%   Ú	conn_dictr   r   r   rH   Œ   s&   ÿ
ÿ
ÿ
ý
ÿz$Channel._try_parse_connection_stringc                   s,   |r| j  |¡ tƒ j||g|¢R i |¤ŽS r$   )rB   ÚaddrC   Úbasic_consume)r%   ÚqueueÚno_ackrK   rL   rM   r   r   rc   ¤   s   ÿÿÿzChannel.basic_consumec                   s,   || j v r| j| }| j |¡ tƒ  |¡S r$   )Ú
_consumersÚ_tag_to_queuerB   ÚdiscardrC   Úbasic_cancel)r%   Úconsumer_tagrd   rM   r   r   ri   «   s   

zChannel.basic_cancelNÚnamer    r!   r"   r#   r   c                 C  sH   || j v r| j | }|jp||_|jp||_|S t||ƒ}|| j |< |S r$   )rA   r"   r    r   )r%   rk   r    r"   Úobjr   r   r   Ú_add_queue_to_cache±   s   


þ
zChannel._add_queue_to_cacherd   c                 C  sD   | j  |d ¡}|d u s|jd u r | jj|| jd�}| j||d�}|S )N)Ú
keep_alive)r"   )rA   Úgetr"   Úqueue_serviceÚget_queue_senderÚuamqp_keep_alive_intervalrm   )r%   rd   Ú	queue_objr"   r   r   r   Ú_get_asb_sender¿   s   ÿzChannel._get_asb_senderÚ	recv_moder	   Úqueue_cache_keyú
str | Nonec                 C  sN   |p|}| j  |d ¡}|d u s|jd u r%| jj||| jd�}| j||d�}|S )N)Ú
queue_nameÚreceive_modern   )r    )rA   ro   r    rp   Úget_queue_receiverrr   rm   )r%   rd   ru   rv   Ú	cache_keyrs   r    r   r   r   Ú_get_asb_receiverÇ   s   þzChannel._get_asb_receiverÚtableúdict[int, int] | Nonec                 C  s   t t|ƒƒ |p	t¡S )z:Format AMQP queue name into a valid ServiceBus queue name.)r?   r   Ú	translateÚCHARS_REPLACE_TABLE)r%   rk   r}   r   r   r   Úentity_nameÔ   s   zChannel.entity_nameÚmessageúvirtual.base.Messagec                 C  s   d S r$   r   )r%   r‚   r   r   r   Ú_restoreÙ   s   zChannel._restorec                 K  s|   |   | j| ¡}z| j| W S  ty=   t tj| jd�¡}z
| jj	||d� W n t
jjjy5   Y nw |  |¡ Y S w )z$Ensure a queue exists in ServiceBus.)Úseconds)rx   Úlock_duration)r�   Úqueue_name_prefixrA   ÚKeyErrorÚisodateÚduration_isoformatÚDurationÚpeek_lock_secondsÚqueue_mgmt_serviceÚcreate_queueÚazureÚcoreÚ
exceptionsÚResourceExistsErrorrm   )r%   rd   rL   r†   r   r   r   Ú
_new_queueà   s    ÿ
ÿÿözChannel._new_queuec                 O  s>   |   | j| ¡}| j |¡ | j |d¡}|r| ¡  dS dS )zDelete queue by name.N)r�   r‡   r�   Údelete_queuerA   Úpopr)   )r%   rd   rK   rL   Úsend_receive_objr   r   r   Ú_deleteò   s   ÿzChannel._deletec                 K  s6   |   | j| ¡}tt|ƒƒ}|  |¡}|j |¡ dS )zPut message onto queue.N)r�   r‡   r   r   rt   r"   Úsend_messages)r%   rd   r‚   rL   Úmsgrs   r   r   r   Ú_putû   s   
zChannel._putÚtimeoutúfloat | int | Noneúdict[str, Any]c           	      C  sª   || j v rtjntj}|  | j| ¡}|  ||¡}|jjd|p!| j	d�}|s)t
ƒ ‚|d }t|jtƒs:d |j¡}n|j}tt|ƒƒ}||d d d< ||d d d< |S )	z/Try to retrieve a single message off ``queue``.r   ©Úmax_message_countÚmax_wait_timer   ó    Ú
propertiesÚdelivery_infoÚazure_messageÚazure_queue_name)rB   r	   ÚRECEIVE_AND_DELETEÚ	PEEK_LOCKr�   r‡   r|   r    Úreceive_messagesÚwait_time_secondsr   r]   ÚbodyÚbytesr_   r   r   )	r%   rd   r›   ru   rs   Úmessagesr‚   rª   r™   r   r   r   Ú_get  s(   
ÿÿþzChannel._getFÚdelivery_tagÚmultipleÚboolc                   s°   z	| j  |¡j}W n ty   tƒ  |¡ Y d S w |d }|  |¡}z
|j |d ¡ W n" t	j
jjy@   tƒ  |¡ Y d S  tyO   tƒ  |¡ Y d S w tƒ  |¡ d S )Nr¥   r¤   )rI   ro   r£   rˆ   rC   Ú	basic_ackr|   r    Úcomplete_messager�   Ú
servicebusr‘   ÚMessageAlreadySettledÚ	ExceptionÚbasic_reject)r%   r®   r¯   r£   rd   rs   rM   r   r   r±   #  s"   ÿ
ÿÿzChannel.basic_ackc                 C  s"   |   | j| ¡}| j |¡}|jS )z)Return the number of messages in a queue.)r�   r‡   r�   Úget_queue_runtime_propertiesÚtotal_message_count)r%   rd   Úpropsr   r   r   Ú_size7  s   zChannel._sizec                 C  sˆ   d}d}|   | j| ¡}| j |d¡}|| jvs!|du s!|jdu r+|  |tjd| ¡}	 |jj	|dd�}|t
|ƒ7 }t
|ƒ|k rC	 |S q,)z'Delete all current messages in a queue.r   é
   NÚpurge_Tgš™™™™™É?rž   )r�   r‡   rA   ro   rB   r    r|   r	   r¦   r¨   Úlen)r%   rd   ÚnÚmax_purge_countrs   r¬   r   r   r   Ú_purge>  s(   

þþözChannel._purgec                 C  sP   | j s$d| _ | j ¡ D ]}| ¡  q| j ¡  | jd ur&| j | ¡ d S d S d S )NT)ÚclosedrA   Úvaluesr)   ÚclearÚ
connectionÚclose_channel)r%   rs   r   r   r   r)   Z  s   


ùzChannel.closer   c                 C  s<   | j rtj| j | j| j| jd�S t| j| j| j| j| jd�S )N)Úretry_totalÚretry_backoff_factorÚretry_backoff_max)rG   r   Úfrom_connection_stringrÆ   rÇ   rÈ   rD   r\   r*   r   r   r   rp   e  s   üûzChannel.queue_servicer   c                 C  s    | j r	t | j ¡S t| j| jƒS r$   )rG   r   rÉ   rD   r\   r*   r   r   r   r�   w  s   ÿÿzChannel.queue_mgmt_servicec                 C  s   | j jS r$   )rÄ   Úclientr*   r   r   r   rZ   ‚  s   zChannel.conninfoc                 C  s
   | j jjS r$   )rÄ   rÊ   Útransport_optionsr*   r   r   r   rË   †  s   
zChannel.transport_optionsc                 C  s   | j  dd¡S )Nr‡   Ú )rË   ro   r*   r   r   r   r‡   Š  s   zChannel.queue_name_prefixc                 C  ó   | j  d| j¡S )Nr©   )rË   ro   r4   r*   r   r   r   r©   Ž  s   ÿzChannel.wait_time_secondsc                 C  s   t | j d| j¡dƒS )NrŒ   i,  )ÚminrË   ro   r6   r*   r   r   r   rŒ   “  s
   
ÿþzChannel.peek_lock_secondsc                 C  rÍ   )Nrr   )rË   ro   r8   r*   r   r   r   rr   ™  s   þz!Channel.uamqp_keep_alive_intervalc                 C  rÍ   )NrÆ   )rË   ro   r:   r*   r   r   r   rÆ      ó   ÿzChannel.retry_totalc                 C  rÍ   )NrÇ   )rË   ro   r<   r*   r   r   r   rÇ   ¥  rÏ   zChannel.retry_backoff_factorc                 C  rÍ   )NrÈ   )rË   ro   r>   r*   r   r   r   rÈ   ª  rÏ   zChannel.retry_backoff_maxr,   r+   )rk   r?   r    r!   r"   r#   r'   r   )rd   r?   r'   r   )rd   r?   ru   r	   rv   rw   r'   r   r$   )rk   r?   r}   r~   r'   r?   )r‚   rƒ   r'   r(   )rd   r?   r'   r(   )rd   r?   r›   rœ   r'   r�   )F)r®   r?   r¯   r°   r'   r(   )rd   r?   r'   r3   )r'   r3   )r'   r   )r'   r   )r'   r?   )r'   r;   )0r-   r.   r/   r0   r4   Ú__annotations__r6   r8   r:   r<   r>   r@   rA   ÚsetrB   r&   rH   rc   ri   rm   rt   r	   r§   r|   r�   r„   r“   r—   rš   r­   r±   rº   rÀ   r)   r   rp   r�   ÚpropertyrZ   rË   r‡   r©   rŒ   rr   rÆ   rÇ   rÈ   Ú__classcell__r   r   rM   r   r1   o   sp   
 
ý

ýÿ



	
þ 





r1   c                   @  s>   e Zd ZdZeZdZdZdZedd	d
„ƒZ	e
dddd„ƒZdS )rX   zAzure Service Bus transport.r   NTÚurir?   r'   úDtuple[str, str | DefaultAzureCredential | ManagedIdentityCredential]c                 C  s¸   |   dd¡} |  dd¡\}}| d¡s|d7 }d ¡ | ¡ kr+td u r'tdƒ‚tƒ }n#d	 ¡ | ¡ kr?td u r;td
ƒ‚tƒ }n| dd¡\}}|› d|› �}t||gƒsXt	dƒ‚||fS )Nzazureservicebus://rÌ   ú@r   z.netz.servicebus.windows.netr   z]Azure Service Bus transport with a DefaultAzureCredential requires the azure-identity libraryr   z`Azure Service Bus transport with a ManagedIdentityCredential requires the azure-identity libraryrO   z|Need a URI like azureservicebus://{SAS policy name}:{SAS key}@{ServiceBus Namespace} or the azure Endpoint connection string)
ÚreplaceÚrsplitÚendswithÚlowerr   ÚImportErrorr   r^   ÚallÚ
ValueError)rÔ   Ú
credentialÚ	namespaceÚpolicyÚsas_keyr   r   r   rY   ¹  s&   	
ÿzTransport.parse_uriFú**c                 C  sZ   |   |¡\}}t|tƒr%d|v r%| dd¡\}}d ||r!||¡S ||¡S d |jj|¡S )NrO   r   zazureservicebus://{}:{}@{}zazureservicebus://{}@{})rY   r]   r?   r^   ÚformatrN   r-   )ÚclsrÔ   Úinclude_passwordÚmaskrß   rÞ   rà   rá   r   r   r   Úas_uriä  s   ýýþzTransport.as_uri)rÔ   r?   r'   rÕ   )Frâ   )rÔ   r?   r'   r?   )r-   r.   r/   r0   r1   Úpolling_intervalÚdefault_portÚcan_parse_urlÚstaticmethodrY   Úclassmethodrç   r   r   r   r   rX   °  s    *rX   ),r0   Ú
__future__r   Ústringrd   r   Útypingr   r   r   Úazure.core.exceptionsr�   Úazure.servicebus.exceptionsr‰   Úazure.servicebusr   r   r	   r
   r   Úazure.servicebus.managementr   Úazure.identityr   r   rÛ   Úkombu.utils.encodingr   r   Úkombu.utils.jsonr   r   Úkombu.utils.objectsr   rÌ   r   rÑ   ÚpunctuationÚPUNCTUATIONS_TO_REPLACEr   r€   r   r1   rX   r   r   r   r   Ú<module>   s<    9þÿþ  C