o
    ›¨Êh›,  ã                   @   sþ   d Z ddlmZmZ ddlmZ ddlmZ ddlmZ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yA   d
ZY nw erczddlmZ W n ey[   ddlmZ Y nw ddlmZ n
d
ZG dd„ deƒZdZeddgƒZG dd„ deƒZd
S )zMongoDB result store backend.é    )ÚdatetimeÚ	timedelta)ÚEncodeError)Úcached_property)Úmaybe_sanitize_urlÚurlparse)Ústates)ÚImproperlyConfiguredé   )ÚBaseBackendN)ÚBinary)ÚInvalidDocumentc                   @   s   e Zd ZdS )r   N)Ú__name__Ú
__module__Ú__qualname__© r   r   úI/var/www/html/env/lib/python3.10/site-packages/celery/backends/mongodb.pyr      s    r   )ÚMongoBackendÚpickleÚmsgpackc                       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	ZdZd3‡ fd
d„	Zedd„ ƒZdd„ Zdd„ Z‡ fdd„Z‡ fdd„Z	d4dd„Zdd„ Zdd„ Zdd„ Zdd„ Zd d!„ Zd"d#„ Zd5‡ fd%d&„	Zd'd(„ Ze d)d*„ ƒZ!e d+d,„ ƒZ"e d-d.„ ƒZ#e d/d0„ ƒZ$d6d1d2„Z%‡  Z&S )7r   z‘MongoDB result backend.

    Raises:
        celery.exceptions.ImproperlyConfigured:
            if module :pypi:`pymongo` is not available.
    NÚ	localhosti‰i  ÚceleryÚcelery_taskmetaÚcelery_groupmetaé
   Fc                    s¨  i | _ tƒ j|fi |¤Ž tstdƒ‚|  ¡  ¡ D ]\}}| j  ||¡ q| jr]|  	| j¡| _tj
 | j¡}dd„ |d D ƒ}|d | _|d | _|| _|d rU|d | _| j  |d ¡ | jj d	¡}|d urÒt|tƒsqtd
ƒ‚t|ƒ}d|v s}d|v r€d | _| d| j¡| _| d| j¡| _| d| j¡| _| d| j¡| _| d| j¡| _| d| j¡| _| d| j¡| _| d| j¡| _| j  | di ¡¡ | j  |¡ d S d S )NzCYou need to install the pymongo library to use the MongoDB backend.c                 S   s"   g | ]}|d  › d|d › �‘qS )r   ú:r
   r   )Ú.0Úxr   r   r   Ú
<listcomp>N   s    ÿz)MongoBackend.__init__.<locals>.<listcomp>ÚnodelistÚusernameÚpasswordÚdatabaseÚoptionsÚmongodb_backend_settingsz4MongoDB backend settings should be grouped in a dictÚhostÚportÚ
mongo_hostÚuserÚtaskmeta_collectionÚgroupmeta_collection)r#   ÚsuperÚ__init__Úpymongor	   Ú_prepare_client_optionsÚitemsÚ
setdefaultÚurlÚ_ensure_mongodb_uri_complianceÚ
uri_parserÚ	parse_urir(   r!   r'   Údatabase_nameÚupdateÚappÚconfÚgetÚ
isinstanceÚdictÚpopr%   r&   r)   r*   )Úselfr7   ÚkwargsÚkeyÚvalueÚuri_dataÚ	hostslistÚconfig©Ú	__class__r   r   r,   :   sX   ÿÿ



ÿÿÿèzMongoBackend.__init__c                 C   s2   t | ƒ}|j d¡sd| › �} | dkr| d7 } | S )NÚmongodbzmongodb+ú
mongodb://r   )r   ÚschemeÚ
startswith)r1   Ú
parsed_urlr   r   r   r2   v   s   
z+MongoBackend._ensure_mongodb_uri_compliancec                 C   s    t jdkr
d| jiS | jddœS )N)é   ÚmaxPoolSizeF)Úmax_pool_sizeÚauto_start_request)r-   Úversion_tuplerM   ©r=   r   r   r   r.   �   s
   

ÿz$MongoBackend._prepare_client_optionsc                 C   s”   | j du rGddlm} | j}|s&| j}t|tƒr&| d¡s&d|› d| j› �}t	| j
ƒ}||d< | jr7| j|d< | jr?| j|d< |d	i |¤Ž| _ | j S )
zConnect to the MongoDB server.Nr   )ÚMongoClientrG   r   r%   r    r!   r   )Ú_connectionr-   rQ   r'   r%   r:   ÚstrrI   r&   r;   r#   r(   r!   )r=   rQ   r%   r8   r   r   r   Ú_get_connectionˆ   s"   

ÿ


zMongoBackend._get_connectionc                    s0   | j dkr|S tƒ  |¡}| j tv rt|ƒ}|S ©NÚbson)Ú
serializerr+   ÚencodeÚBINARY_CODECSr   )r=   ÚdataÚpayloadrD   r   r   rX   ¥   s   

zMongoBackend.encodec                    s   | j dkr|S tƒ  |¡S rU   )rW   r+   Údecode)r=   rZ   rD   r   r   r\   °   s   
zMongoBackend.decodec           	   
   K   s`   | j |  |¡|||dd�}||d< z| jjd|i|dd� W |S  ty/ } zt|ƒ‚d}~ww )z1Store return value and state of an executed task.F)ÚresultÚstateÚ	tracebackÚrequestÚformat_dateÚ_idT©ÚupsertN)Ú_get_result_metarX   Ú
collectionÚreplace_oner   r   )	r=   Útask_idr]   r^   r_   r`   r>   ÚmetaÚexcr   r   r   Ú_store_resultµ   s   þý€ÿzMongoBackend._store_resultc                 C   sÀ   | j  d|i¡}|rZ| jj dd¡r?|  |d |d |d |d |d |d |d	 |d
 |d |d |d |  |d ¡dœ¡S |  |d |d |  |d ¡|d |d |d dœ¡S tjddœS )z$Get task meta-data for a task by id.rb   Úextendedr]   ÚnameÚargsÚqueuer>   ÚstatusÚworkerÚretriesÚchildrenÚ	date_doner_   )rm   rn   rh   ro   r>   rp   rq   rr   rs   rt   r_   r]   )rh   rp   r]   rt   r_   rs   N)rp   r]   )	rf   Úfind_oner7   r8   Úfind_value_for_keyÚmeta_from_decodedr\   r   ÚPENDING)r=   rh   Úobjr   r   r   Ú_get_task_meta_forÅ   s4   ôúzMongoBackend._get_task_meta_forc                 C   s:   ||   dd„ |D ƒ¡t ¡ dœ}| jjd|i|dd� |S )zSave the group result.c                 S   s   g | ]}|j ‘qS r   )Úid)r   Úir   r   r   r   æ   s    z,MongoBackend._save_group.<locals>.<listcomp>)rb   r]   rt   rb   Trc   )rX   r   ÚutcnowÚgroup_collectionrg   )r=   Úgroup_idr]   ri   r   r   r   Ú_save_groupâ   s   ýzMongoBackend._save_groupc                    sD   ˆ j  d|i¡}|r |d |d ‡ fdd„ˆ  |d ¡D ƒdœS dS )z!Get the result for a group by id.rb   rt   c                    s   g | ]}ˆ j  |¡‘qS r   )r7   ÚAsyncResult)r   ÚtaskrP   r   r   r   ó   s    
ÿÿz/MongoBackend._restore_group.<locals>.<listcomp>r]   )rh   rt   r]   N)r~   ru   r\   )r=   r   ry   r   rP   r   Ú_restore_groupì   s   
þýÿzMongoBackend._restore_groupc                 C   ó   | j  d|i¡ dS )zDelete a group by id.rb   N)r~   Ú
delete_one)r=   r   r   r   r   Ú_delete_groupù   s   zMongoBackend._delete_groupc                 C   r„   )zšRemove result from MongoDB.

        Raises:
            pymongo.exceptions.OperationsError:
                if the task_id could not be removed.
        rb   N)rf   r…   )r=   rh   r   r   r   Ú_forgetý   s   
zMongoBackend._forgetc                 C   sN   | j sdS | j dd| j ¡ | j ii¡ | j dd| j ¡ | j ii¡ dS )zDelete expired meta-data.Nrt   z$lt)Úexpiresrf   Údelete_manyr7   ÚnowÚexpires_deltar~   rP   r   r   r   Úcleanup	  s   ÿÿzMongoBackend.cleanupr   c                    s(   |si n|}t ƒ  |t|| j| jd�¡S )N)rˆ   r1   )r+   Ú
__reduce__r;   rˆ   r1   )r=   rn   r>   rD   r   r   r�     s   ÿzMongoBackend.__reduce__c                 C   s   |   ¡ }|| j S ©N)rT   r5   )r=   Úconnr   r   r   Ú_get_database  s   
zMongoBackend._get_databasec                 C   s   |   ¡ S )z]Get database from MongoDB connection.

        performs authentication if necessary.
        )r�   rP   r   r   r   r"     s   zMongoBackend.databasec                 C   ó   | j | j }|jddd� |S ©z"Get the meta-data task collection.rt   T)Ú
background)r"   r)   Úcreate_index©r=   rf   r   r   r   rf   &  ó   zMongoBackend.collectionc                 C   r‘   r’   )r"   r*   r”   r•   r   r   r   r~   0  r–   zMongoBackend.group_collectionc                 C   s   t | jd�S )N)Úseconds)r   rˆ   rP   r   r   r   r‹   :  s   zMongoBackend.expires_deltac                 C   sL   | j sdS |r
| j S d| j vrt| j ƒS | j  dd¡\}}d t|ƒ|g¡S )z~Return the backend as an URI.

        Arguments:
            include_password (bool): Password censored if disabled.
        rG   ú,r
   )r1   r   ÚsplitÚjoin)r=   Úinclude_passwordÚuri1Ú	remainderr   r   r   Úas_uri>  s   

zMongoBackend.as_urirŽ   )NN)r   N)F)'r   r   r   Ú__doc__r'   r%   r&   r(   r!   r5   r)   r*   rM   r#   Úsupports_autoexpirerR   r,   Ústaticmethodr2   r.   rT   rX   r\   rk   rz   r€   rƒ   r†   r‡   rŒ   r�   r�   r   r"   rf   r~   r‹   rž   Ú__classcell__r   r   rD   r   r   #   sP    <


ÿ


	
	
r   )rŸ   r   r   Úkombu.exceptionsr   Úkombu.utils.objectsr   Úkombu.utils.urlr   r   r   r   Úcelery.exceptionsr	   Úbaser   r-   ÚImportErrorÚbson.binaryr   Úpymongo.binaryÚpymongo.errorsr   Ú	ExceptionÚ__all__Ú	frozensetrY   r   r   r   r   r   Ú<module>   s2    ÿÿ