o
    ›¨Êh\%  ã                   @   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W n ey7   d	ZY nw zdd	lZW n eyI   d	ZY nw d
ZdZG dd„ deƒZd	S )z#Elasticsearch result store backend.é    )Údatetime©Úbytes_to_str)Ú
_parse_url)Ústates)ÚImproperlyConfiguredé   )ÚKeyValueStoreBackendN)ÚElasticsearchBackendzVYou need to install the elasticsearch library to use the Elasticsearch result backend.c                       s¾   e Zd ZdZdZdZdZdZdZdZ	dZ
dZdZd	Zd&‡ 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‡ fdd„Z‡ fdd„Zdd„ Zd d!„ Zd"d#„ Zed$d%„ ƒZ‡  ZS )'r
   z–Elasticsearch Backend.

    Raises:
        celery.exceptions.ImproperlyConfigured:
            if module :pypi:`elasticsearch` is not available.
    ÚceleryNÚhttpÚ	localhostið#  Fé
   é   c                    s8  t ƒ j|i |¤Ž || _| jjj}td u rttƒ‚d  } } } } }	 }
}|rIt	|ƒ\}}}	}
}}}|dkr:d }|rI| 
d¡}| d¡\}}}|pM| j| _|pS| j| _|pY| j| _|p_| j| _|	pe| j| _|
pk| j| _|pq| j| _|dƒpy| j| _|dƒ}|d ur†|| _|dƒ}|d ur‘|| _|ddƒ| _d | _d S )NÚelasticsearchú/Úelasticsearch_retry_on_timeoutÚelasticsearch_timeoutÚelasticsearch_max_retriesÚelasticsearch_save_meta_as_textT)ÚsuperÚ__init__ÚurlÚappÚconfÚgetr   r   ÚE_LIB_MISSINGr   ÚstripÚ	partitionÚindexÚdoc_typeÚschemeÚhostÚportÚusernameÚpasswordÚes_retry_on_timeoutÚ
es_timeoutÚes_max_retriesÚes_save_meta_as_textÚ_server)Úselfr   ÚargsÚkwargsÚ_getr   r    r!   r"   r#   r$   r%   ÚpathÚ_r'   r(   ©Ú	__class__© úO/var/www/html/env/lib/python3.10/site-packages/celery/backends/elasticsearch.pyr   1   s<   

ÿ
zElasticsearchBackend.__init__c                 C   s2   t |tjjƒr|jdv rdS t |tjjƒrdS dS )N>   é‘  éô  éö  éø  úN/Aé™  TF)Ú
isinstancer   Ú
exceptionsÚApiErrorÚstatus_codeÚTransportError)r+   Úexcr3   r3   r4   Úexception_safe_to_retryZ   s   
z,ElasticsearchBackend.exception_safe_to_retryc              	   C   s`   z#|   |¡}z|d r|d d W W S W W d S  ttfy#   Y W d S w  tjjy/   Y d S w )NÚfoundÚ_sourceÚresult)r.   Ú	TypeErrorÚKeyErrorr   r<   ÚNotFoundError)r+   ÚkeyÚresr3   r3   r4   r   h   s   
ÿÿÿzElasticsearchBackend.getc                 C   s.   | j r| jj| j|| j d�S | jj| j|d�S ©N)r   Úidr    )r   rK   )r    Úserverr   r   ©r+   rH   r3   r3   r4   r.   s   s   ýþzElasticsearchBackend._getc                 C   s\   |d  t ¡  ¡ d d… ¡dœ}z
| j||d� W d S  tjjy-   |  |||¡ Y d S w )Nz{}Zéýÿÿÿ)rD   z
@timestamp)rK   Úbody)	Úformatr   ÚutcnowÚ	isoformatÚ_indexr   r<   ÚConflictErrorÚ_update)r+   rH   ÚvalueÚstaterO   r3   r3   r4   Ú_set_with_state€   s   ÿþþþz$ElasticsearchBackend._set_with_statec                 C   s   |   ||d ¡S ©N)rX   )r+   rH   rV   r3   r3   r4   Úset�   s   zElasticsearchBackend.setc                 K   sh   dd„ |  ¡ D ƒ}| jr!| jjdt|ƒ| j| j|ddidœ|¤ŽS | jjdt|ƒ| j|ddidœ|¤ŽS )Nc                 S   ó   i | ]	\}}t |ƒ|“qS r3   r   ©Ú.0ÚkÚvr3   r3   r4   Ú
<dictcomp>”   ó    z/ElasticsearchBackend._index.<locals>.<dictcomp>Úop_typeÚcreate©rK   r   r    rO   Úparams©rK   r   rO   re   r3   )Úitemsr    rL   r   r   )r+   rK   rO   r-   r3   r3   r4   rS   “   s&   ûú	üûzElasticsearchBackend._indexc           
      K   s�  dd„ |  ¡ D ƒ}z| j|d�}| d¡s | j||fi |¤ŽW S W n tjjy6   | j||fi |¤Ž Y S w z|  |d d ¡}W n tt	fyM   Y nw |d t
jkrYddiS |d t
jv ri|t
jv riddiS | d	d
¡}| dd
¡}| jr‘| jjdt|ƒ| j| jd|i||dœdœ|¤Ž}	n| jjdt|ƒ| jd|i||dœdœ|¤Ž}	|	d dkrÆtj dt ddt ¡ dt | j| j| j¡¡d¡‚|	S )au  Update state in a conflict free manner.

        If state is defined (not None), this will not update ES server if either:
        * existing state is success
        * existing state is a ready state and current state in not a ready state

        This way, a Retry state cannot override a Success or Failure, and chord_unlock
        will not retry indefinitely.
        c                 S   r[   r3   r   r\   r3   r3   r4   r`   ±   ra   z0ElasticsearchBackend._update.<locals>.<dictcomp>)rH   rB   rC   rD   ÚstatusÚnoopÚ_seq_nor   Ú_primary_termÚdoc)Úif_primary_termÚ	if_seq_nord   rf   z(conflicting update occurred concurrentlyr:   zHTTP/1.1r   Nr3   )rg   r.   r   rS   r   r<   rG   Údecode_resultrE   rF   r   ÚSUCCESSÚREADY_STATESÚUNREADY_STATESr    rL   Úupdater   r   rT   Úelastic_transportÚApiResponseMetaÚHttpHeadersÚ
NodeConfigr!   r"   r#   )
r+   rK   rO   rW   r-   Úres_getÚmeta_present_on_backendÚseq_noÚ	prim_termrI   r3   r3   r4   rU   §   sb   

ÿÿÿûú	üû
ÿÿüzElasticsearchBackend._updatec                    sl   | j r	tƒ  |¡S t|tƒstƒ  |¡S | d¡r$|  |d ¡d |d< | d¡r4|  |d ¡d |d< |S )NrD   é   Ú	traceback)r)   r   Úencoder;   Údictr   Ú_encode)r+   Údatar1   r3   r4   r~   é   s   


zElasticsearchBackend.encodec                    sh   | j r	tƒ  |¡S t|tƒstƒ  |¡S | d¡r#tƒ  |d ¡|d< | d¡r2tƒ  |d ¡|d< |S )NrD   r}   )r)   r   Údecoder;   r   r   )r+   Úpayloadr1   r3   r4   r‚   õ   s   


zElasticsearchBackend.decodec                    s   ‡ fdd„|D ƒS )Nc                    s   g | ]}ˆ   |¡‘qS r3   )r   )r]   rH   ©r+   r3   r4   Ú
<listcomp>  s    z-ElasticsearchBackend.mget.<locals>.<listcomp>r3   )r+   Úkeysr3   r„   r4   Úmget  s   zElasticsearchBackend.mgetc                 C   s6   | j r| jj| j|| j d� d S | jj| j|d� d S rJ   )r    rL   Údeleter   rM   r3   r3   r4   rˆ     s   zElasticsearchBackend.deletec                 C   sL   d}| j r| jr| j | jf}tj| j› d| j› d| j› �| j| j| j	|d�S )z$Connect to the Elasticsearch server.Nz://ú:)Úretry_on_timeoutÚmax_retriesÚtimeoutÚ	http_auth)
r$   r%   r   ÚElasticsearchr!   r"   r#   r&   r(   r'   )r+   r�   r3   r3   r4   Ú_get_server
  s   ûz ElasticsearchBackend._get_serverc                 C   s   | j d u r
|  ¡ | _ | j S rY   )r*   r�   r„   r3   r3   r4   rL     s   

zElasticsearchBackend.serverrY   )Ú__name__Ú
__module__Ú__qualname__Ú__doc__r   r    r!   r"   r#   r$   r%   r&   r'   r(   r   rA   r   r.   rX   rZ   rS   rU   r~   r‚   r‡   rˆ   r�   ÚpropertyrL   Ú__classcell__r3   r3   r1   r4   r
      s6    )Br
   )r“   r   Úkombu.utils.encodingr   Úkombu.utils.urlr   r   r   Úcelery.exceptionsr   Úbaser	   r   ÚImportErrorrt   Ú__all__r   r
   r3   r3   r3   r4   Ú<module>   s(    ÿÿ