o
    ›¨Êhy  ã                   @   sè   d 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 z"dd	lZdd
lmZ ddlmZmZmZ ddlmZ ddlmZ W n ey_   d	 Z Z Z Z Z ZZY nw dZdZdZe
eƒZG dd„ deƒZd	S )z3The CosmosDB/SQL backend for Celery (experimental).é    )Úcached_property)Úbytes_to_str)Ú
_parse_url)ÚImproperlyConfigured)Ú
get_loggeré   )ÚKeyValueStoreBackendN)ÚDocumentClient)ÚConnectionPolicyÚConsistencyLevelÚPartitionKind)ÚHTTPFailure)ÚRetryOptions)ÚCosmosDBSQLBackendi”  i™  c                       s¢   e Zd ZdZ						d‡ fdd„	Zedd„ ƒZedd„ ƒZd	d
„ Z	dd„ Z
edd„ ƒZedd„ ƒZdd„ Zedd„ ƒZdd„ Zdd„ Zdd„ Zdd„ Z‡  ZS )r   z CosmosDB/SQL backend for Celery.Nc           
         s¨   t ƒ j|i |¤Ž td u rtdƒ‚| jj}	|  |¡\| _| _|p#|	d | _	|p*|	d | _
ztt|p4|	d ƒ| _W n tyC   tdƒ‚w |pI|	d | _|pP|	d | _d S )NzIYou need to install the pydocumentdb library to use the CosmosDB backend.Úcosmosdbsql_database_nameÚcosmosdbsql_collection_nameÚcosmosdbsql_consistency_levelz"Unknown CosmosDB consistency levelÚcosmosdbsql_max_retry_attemptsÚcosmosdbsql_max_retry_wait_time)ÚsuperÚ__init__Úpydocumentdbr   ÚappÚconfr   Ú	_endpointÚ_keyÚ_database_nameÚ_collection_nameÚgetattrr   Ú_consistency_levelÚAttributeErrorÚ_max_retry_attemptsÚ_max_retry_wait_time)
ÚselfÚurlÚdatabase_nameÚcollection_nameÚconsistency_levelÚmax_retry_attemptsÚmax_retry_wait_timeÚargsÚkwargsr   ©Ú	__class__© úM/var/www/html/env/lib/python3.10/site-packages/celery/backends/cosmosdbsql.pyr   !   s8   	ÿþþ
ýÿþþzCosmosDBSQLBackend.__init__c                 C   sZ   t |ƒ\}}}}}}}|r|stdƒ‚|sd}|dkrdnd}|› d|› d|› �}||fS )NzInvalid URLi»  ÚhttpsÚhttpz://ú:)r   r   )Úclsr$   Ú_ÚhostÚportÚpasswordÚschemeÚendpointr.   r.   r/   r   M   s   zCosmosDBSQLBackend._parse_urlc                 C   sJ   t ƒ }t| j| jd�|_t| jd| ji|| jd�}|  |¡ |  	|¡ |S )zÄReturn the CosmosDB/SQL client.

        If this is the first call to the property, the client is created and
        the database and collection are initialized if they don't yet exist.

        )Úmax_retry_attempt_countÚmax_wait_time_in_secondsÚ	masterKey)Úconnection_policyr'   )
r
   r   r!   r"   r	   r   r   r   Ú_create_database_if_not_existsÚ _create_collection_if_not_exists)r#   r=   Úclientr.   r.   r/   Ú_client[   s   þü

zCosmosDBSQLBackend._clientc              
   C   sZ   z
|  d| ji¡ W n ty# } z|jtkr‚ W Y d }~d S d }~ww t d| j¡ d S )NÚidzCreated CosmosDB database %s)ÚCreateDatabaser   r   Ústatus_codeÚERROR_EXISTSÚLOGGERÚinfo©r#   r@   Úexr.   r.   r/   r>   s   s   
ÿ€ÿÿz1CosmosDBSQLBackend._create_database_if_not_existsc              
   C   sn   z|  | j| jdgtjdœdœ¡ W n ty+ } z|jtkr ‚ W Y d }~d S d }~ww t 	d| j
| j¡ d S )Nz/id)ÚpathsÚkind)rB   ÚpartitionKeyz!Created CosmosDB collection %s/%s)ÚCreateCollectionÚ_database_linkr   r   ÚHashr   rD   rE   rF   rG   r   rH   r.   r.   r/   r?   }   s$   ÿÿþ
ÿ€ÿÿz3CosmosDBSQLBackend._create_collection_if_not_existsc                 C   s
   d| j  S )Nzdbs/)r   ©r#   r.   r.   r/   rN   ‹   s   
z!CosmosDBSQLBackend._database_linkc                 C   s   | j d | j S )Nz/colls/)rN   r   rP   r.   r.   r/   Ú_collection_link�   s   z#CosmosDBSQLBackend._collection_linkc                 C   s   | j d | S )Nz/docs/)rQ   ©r#   Úkeyr.   r.   r/   Ú_get_document_link“   s   z%CosmosDBSQLBackend._get_document_linkc                 C   s   |r|  ¡ r
tdƒ‚d|iS )Nz(Key cannot be none, empty or whitespace.rL   )ÚisspaceÚ
ValueError)r3   rS   r.   r.   r/   Ú_get_partition_key–   s   z%CosmosDBSQLBackend._get_partition_keyc              
   C   sx   t |ƒ}t d| j| j|¡ z| j |  |¡|  |¡¡}W n t	y6 } z|j
tkr+‚ W Y d}~dS d}~ww | d¡S )zxRead the value stored at the given key.

        Args:
              key: The key for which to read the value.

        z"Getting CosmosDB document %s/%s/%sNÚvalue)r   rF   Údebugr   r   rA   ÚReadDocumentrT   rW   r   rD   ÚERROR_NOT_FOUNDÚget)r#   rS   ÚdocumentrI   r.   r.   r/   r\   �   s    
ÿþ
€ý
zCosmosDBSQLBackend.getc                 C   s>   t |ƒ}t d| j| j|¡ | j | j||dœ|  |¡¡ dS )z˜Store a value for a given key.

        Args:
              key: The key at which to store the value.
              value: The value to store.

        z#Creating CosmosDB document %s/%s/%s)rB   rX   N)	r   rF   rY   r   r   rA   ÚCreateDocumentrQ   rW   )r#   rS   rX   r.   r.   r/   Úset³   s   
ÿýzCosmosDBSQLBackend.setc                    s   ‡ fdd„|D ƒS )zqRead all the values for the provided keys.

        Args:
              keys: The list of keys to read.

        c                    s   g | ]}ˆ   |¡‘qS r.   )r\   )Ú.0rS   rP   r.   r/   Ú
<listcomp>Ë   s    z+CosmosDBSQLBackend.mget.<locals>.<listcomp>r.   )r#   Úkeysr.   rP   r/   ÚmgetÄ   s   zCosmosDBSQLBackend.mgetc                 C   s:   t |ƒ}t d| j| j|¡ | j |  |¡|  |¡¡ dS )zlDelete the value at a given key.

        Args:
              key: The key of the value to delete.

        z#Deleting CosmosDB document %s/%s/%sN)	r   rF   rY   r   r   rA   ÚDeleteDocumentrT   rW   rR   r.   r.   r/   ÚdeleteÍ   s   
ÿþzCosmosDBSQLBackend.delete)NNNNNN)Ú__name__Ú
__module__Ú__qualname__Ú__doc__r   Úclassmethodr   r   rA   r>   r?   rN   rQ   rT   rW   r\   r_   rc   re   Ú__classcell__r.   r.   r,   r/   r      s4    ú,





	r   )ri   Úkombu.utilsr   Úkombu.utils.encodingr   Úkombu.utils.urlr   Úcelery.exceptionsr   Úcelery.utils.logr   Úbaser   r   Úpydocumentdb.document_clientr	   Úpydocumentdb.documentsr
   r   r   Úpydocumentdb.errorsr   Úpydocumentdb.retry_optionsr   ÚImportErrorÚ__all__r[   rE   rf   rF   r   r.   r.   r.   r/   Ú<module>   s2    ÿÿþ