o
    q~9S  ã                   @   sT  d dl mZmZ d dlZd dlmZ d dlmZ d dlm	Z	 z
d dlm
Z
mZ W n ey5   d Z
ZY nw 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mZmZ ddlmZ G dd„ deƒZd"dd„Zdd„ ZG dd„ deƒZG dd„ dejƒZ G dd„ de ƒZ!G dd„ de ƒZ"G dd„ de!ƒZ#G dd„ de!ƒZ$G d d!„ d!e ƒZ%dS )#é    )Úabsolute_importÚunicode_literalsN©Úwraps)Úcount)Ú
connection)ÚconnectionsÚrouter)Úmodels)ÚQuerySet)Úsettings)Úmaybe_timedeltaé   )Úcommit_on_successÚget_querysetÚrollback_unless_managed)Únowc                   @   s   e Zd ZdS )ÚTxIsolationWarningN)Ú__name__Ú
__module__Ú__qualname__© r   r   úC/var/www/html/env/lib/python3.10/site-packages/djcelery/managers.pyr      s    r   c                    s   ‡ fdd„}|S )z›Decorator for methods doing database operations.

    If the database operation fails, it will retry the operation
    at most ``max_retries`` times.

    c                    s   t ˆ ƒ‡ ‡fdd„ƒ}|S )Nc                     sl   |  dˆ¡}tdƒD ])}z
ˆ | i |¤ŽW   S  ty3   ||kr"‚ ztƒ  W n	 ty0   Y nw Y q
w d S )NÚexception_retry_countr   )Úpopr   Ú	Exceptionr   )ÚargsÚkwargsÚ_max_retriesÚretries)ÚfunÚmax_retriesr   r   Ú_inner%   s   
ÿ€öýz1transaction_retry.<locals>._outer.<locals>._innerr   )r    r"   ©r!   )r    r   Ú_outer#   s   z!transaction_retry.<locals>._outerr   )r!   r$   r   r#   r   Útransaction_retry   s   r%   c                    s"   ‡ fdd„|  ¡ D ƒ ˆ  ¡  ˆ S )Nc                    s   g | ]
\}}t ˆ ||ƒ‘qS r   )Úsetattr)Ú.0Ú	attr_nameÚ
attr_value©Úobjr   r   Ú
<listcomp>=   s    ÿz*update_model_with_dict.<locals>.<listcomp>)ÚitemsÚsave)r+   Úfieldsr   r*   r   Úupdate_model_with_dict<   s
   
ÿr0   c                   @   ó   e Zd Zdd„ ZdS )ÚExtendedQuerySetc                 K   s@   | j di |¤Ž\}}|st| di ¡ƒ}| |¡ t||ƒ |S )NÚdefaultsr   )Úget_or_createÚdictr   Úupdater0   )Úselfr   r+   Úcreatedr/   r   r   r   Úupdate_or_createE   s   

z!ExtendedQuerySet.update_or_createN)r   r   r   r9   r   r   r   r   r2   C   ó    r2   c                   @   s8   e Zd Zdd„ ZeZdd„ Zdd„ Zdd„ Zd	d
„ ZdS )ÚExtendedManagerc                 C   s
   t | jƒS ©N)r2   Úmodel©r7   r   r   r   r   R   s   
zExtendedManager.get_querysetc                 K   s   t | ƒjdi |¤ŽS )Nr   )r   r9   )r7   r   r   r   r   r9   V   s   z ExtendedManager.update_or_createc                 C   s   t r
t t | j¡ S tS r<   )r   r	   Údb_for_writer=   r   r>   r   r   r   Úconnection_for_writeY   s   z$ExtendedManager.connection_for_writec                 C   s   t rt | j S tS r<   )r   Údbr   r>   r   r   r   Úconnection_for_read^   s   
z#ExtendedManager.connection_for_readc                 C   s,   z	t j| j d W S  ty   t j Y S w )NÚENGINE)r   Ú	DATABASESrA   ÚAttributeErrorÚDATABASE_ENGINEr>   r   r   r   Úcurrent_enginec   s
   
ÿzExtendedManager.current_engineN)	r   r   r   r   Úget_query_setr9   r@   rB   rG   r   r   r   r   r;   P   s    r;   c                   @   s   e Zd Zdd„ Zdd„ ZdS )ÚResultManagerc                 C   s   | j tƒ t|ƒ d�S )zGet all expired task results.)Údate_done__lt)Úfilterr   r   )r7   Úexpiresr   r   r   Úget_all_expiredl   s   zResultManager.get_all_expiredc                 C   sd   | j j}tƒ �! |  |¡jdd� |  ¡  ¡ }| d |¡d¡ W d  ƒ dS 1 s+w   Y  dS )z#Delete all expired taskset results.T©Úhiddenú(DELETE FROM {0.db_table} WHERE hidden=%s©TN)	r=   Ú_metar   rM   r6   r@   ÚcursorÚexecuteÚformat)r7   rL   ÚmetarS   r   r   r   Údelete_expiredp   s   þ"ýzResultManager.delete_expiredN)r   r   r   rM   rW   r   r   r   r   rI   j   s    rI   c                   @   r1   )ÚPeriodicTaskManagerc                 C   ó   | j dd�S )NT)Úenabled©rK   r>   r   r   r   rZ   ~   ó   zPeriodicTaskManager.enabledN)r   r   r   rZ   r   r   r   r   rX   |   r:   rX   c                   @   s:   e Zd ZdZdZdd„ Zedd�	ddd„ƒZd	d
„ ZdS )ÚTaskManagerz/Manager for :class:`celery.models.Task` models.Nc                 C   sJ   z| j |d�W S  | jjy$   | j|kr|  ¡  || _| j|d� Y S w )aB  Get task meta for task by ``task_id``.

        :keyword exception_retry_count: How many times to retry by
            transaction rollback on exception. This could theoretically
            happen in a race condition if another worker is trying to
            create the same task. The default is to retry once.

        )Útask_id)Úgetr=   ÚDoesNotExistÚ_last_idÚwarn_if_repeatable_read)r7   r^   r   r   r   Úget_task†   s   	
üzTaskManager.get_taské   r#   c                 C   s   | j ||||d|idœd�S )a+  Store the result and status of a task.

        :param task_id: task id

        :param result: The return value of the task, or an exception
            instance raised by the task.

        :param status: Task status. See
            :meth:`celery.result.AsyncResult.get_status` for a list of
            possible status values.

        :keyword traceback: The traceback at the point of exception (if the
            task failed).

        :keyword children: List of serialized results of subtasks
            of this task.

        :keyword exception_retry_count: How many times to retry by
            transaction rollback on exception. This could theoretically
            happen in a race condition if another worker is trying to
            create the same task. The default is to retry twice.

        Úchildren)ÚstatusÚresultÚ	tracebackrV   )r^   r3   ©r9   )r7   r^   rg   rf   rh   re   r   r   r   Ústore_result—   s   ýÿzTaskManager.store_resultc                 C   sX   d|   ¡  ¡ v r&|  ¡  ¡ }| d¡r(| ¡ d }|dkr*t tdƒ¡ d S d S d S d S )NÚmysqlzSELECT @@tx_isolationr   zREPEATABLE-READz²Polling results with transaction isolation level repeatable-read within the same transaction may give outdated results. Be sure to commit the transaction for each poll iteration.)	rG   ÚlowerrB   rS   rT   ÚfetchoneÚwarningsÚwarnr   )r7   rS   Ú	isolationr   r   r   rb   ·   s   

ÿûz#TaskManager.warn_if_repeatable_read)NN)	r   r   r   Ú__doc__ra   rc   r%   rj   rb   r   r   r   r   r]   ‚   s    ÿr]   c                   @   s2   e Zd ZdZdd„ Zdd„ Zedd�dd	„ ƒZd
S )ÚTaskSetManagerz2Manager for :class:`celery.models.TaskSet` models.c                 C   s(   z| j |d�W S  | jjy   Y dS w )z,Get the async result instance by taskset id.)Ú
taskset_idN)r_   r=   r`   )r7   rs   r   r   r   Úrestore_tasksetÇ   s
   ÿzTaskSetManager.restore_tasksetc                 C   s   |   |¡}|r| ¡  dS dS )zDelete a saved taskset result.N)rt   Údelete)r7   rs   Úsr   r   r   Údelete_tasksetÎ   s   
ÿzTaskSetManager.delete_tasksetrd   r#   c                 C   s   | j |d|id�S )z—Store the async result instance of a taskset.

        :param taskset_id: task set id

        :param result: The return value of the taskset

        rg   )rs   r3   ri   )r7   rs   rg   r   r   r   rj   Ô   s   	ÿzTaskSetManager.store_resultN)r   r   r   rq   rt   rw   r%   rj   r   r   r   r   rr   Ä   s    rr   c                   @   s0   e Zd Zdd„ Zefdd„Zdd„ Zdd„ Zd	S )
ÚTaskStateManagerc                 C   rY   )NFrN   r[   r>   r   r   r   Úactiveã   r\   zTaskStateManager.activec                 C   s   | j ||ƒ t|ƒ d�S )N)Ú	state__inÚtstamp__lte)rK   r   )r7   ÚstatesrL   Únowfunr   r   r   Úexpiredæ   s   ÿzTaskStateManager.expiredc                 C   s    |d ur|   ||¡jdd�S d S )NTrN   )r~   r6   )r7   r|   rL   r   r   r   Úexpire_by_statesê   s   ÿz!TaskStateManager.expire_by_statesc                 C   sR   t ƒ � | jj}|  ¡  ¡ }| d |¡d¡ W d   ƒ d S 1 s"w   Y  d S )NrP   rQ   )r   r=   rR   r@   rS   rT   rU   )r7   rV   rS   r   r   r   Úpurgeî   s   þ"ýzTaskStateManager.purgeN)r   r   r   ry   r   r~   r   r€   r   r   r   r   rx   á   s
    rx   )r   )&Ú
__future__r   r   rn   Ú	functoolsr   Ú	itertoolsr   Ú	django.dbr   r   r	   ÚImportErrorr
   Údjango.db.models.queryr   Údjango.confr   Úcelery.utils.timeutilsr   rA   r   r   r   Úutilsr   ÚUserWarningr   r%   r0   r2   ÚManagerr;   rI   rX   r]   rr   rx   r   r   r   r   Ú<module>   s4    ÿ
 B