
    AHj                    x    d dl mZ d dlmZmZmZ d dlmZ d dlmZ d dl	m
Z
 erddlmZ ddlmZ  G d	 d
      Zy)    )annotations)datetime	timedeltatimezone)TYPE_CHECKING)Redis)now   )Queue)
BaseWorkerc                  ^    e Zd ZddZedd       ZddZddZddZddZ	ddZ
ddZdd	Zy
)IntermediateQueuec                L    || _         | j                  |      | _        || _        y )N)	queue_keyget_intermediate_queue_keykey
connection)selfr   r   s      M/root/tools/cai/cai_env/lib/python3.12/site-packages/rq/intermediate_queue.py__init__zIntermediateQueue.__init__   s"    "229=$    c                    | dS )zReturns the intermediate queue key for a given queue key.

        Args:
            key (str): The queue key

        Returns:
            str: The intermediate queue key
        z:intermediate )clsr   s     r   r   z,IntermediateQueue.get_intermediate_queue_key   s     M**r   c                $    | j                    d| S )zReturns the first seen key for a given job ID.

        Args:
            job_id (str): The job ID

        Returns:
            str: The first seen key
        z:first_seen:)r   r   job_ids     r   get_first_seen_keyz$IntermediateQueue.get_first_seen_key!   s     ((<x00r   c                    t        | j                  j                  | j                  |      t	               j                         dd            S )zSets the first seen timestamp for a job.

        Args:
            job_id (str): The job ID
            timestamp (float): The timestamp
        TiQ )nxex)boolr   setr   r	   	timestampr   s     r   set_first_seenz IntermediateQueue.set_first_seen,   s=     DOO''(?(?(GIZ_chq'rssr   c                    | j                   j                  | j                  |            }|r.t        j                  t        |      t        j                        S y)zReturns the first seen timestamp for a job.

        Args:
            job_id (str): The job ID

        Returns:
            Optional[datetime]: The timestamp
        )tzN)r   getr   r   fromtimestampfloatr   utc)r   r   r$   s      r   get_first_seenz IntermediateQueue.get_first_seen6   sE     OO''(?(?(GH	))%	*:x||LLr   c                ^    | j                  |      }|syt               |z
  t        d      kD  S )a  Returns whether a job should be cleaned up.
        A job in intermediate queue should be cleaned up if it has been there for more than 1 minute.

        Args:
            job_id (str): The job ID

        Returns:
            bool: Whether the job should be cleaned up
        Fr
   )minutes)r,   r	   r   )r   r   
first_seens      r   should_be_cleaned_upz&IntermediateQueue.should_be_cleaned_upD   s1     ((0
uz!Ia$888r   c                    | j                   j                  | j                  dd      D cg c]  }|j                          c}S c c}w )zlReturns the job IDs in the intermediate queue.

        Returns:
            List[str]: The job IDs
        r   )r   lranger   decoder   s     r   get_job_idszIntermediateQueue.get_job_idsT   s6     /3oo.D.DTXXqRT.UVFVVVs   Ac                R    | j                   j                  | j                  d|       y)zgRemoves a job from the intermediate queue.

        Args:
            job_id (str): The job ID
        r
   N)r   lremr   r   s     r   removezIntermediateQueue.remove\   s     	TXXq&1r   c                ,   | j                         }|D ]  }|j                  |      }||j                  vs#|s| j                  |       7| j	                  |      rI| j                  |      s[|j                  ||d       | j                  |        y )Nz$Job was stuck in intermediate queue.)
exc_string)r5   	fetch_jobstarted_job_registryr8   r%   r0   handle_job_failure)r   workerqueuejob_idsr   jobs         r   cleanupzIntermediateQueue.cleanupd   s    ""$ 	(F//&)CU777KK' &&v.,,V4--c5Ek-lKK'!	(r   N)r   strr   r   )r   rC   returnrC   )r   rC   rD   rC   )r   rC   rD   r"   )r   rC   rD   zdatetime | None)rD   z	list[str])r   rC   rD   None)r>   r   r?   r   rD   rE   )__name__
__module____qualname__r   classmethodr   r   r%   r,   r0   r5   r8   rB   r   r   r   r   r      s@    %
 	+ 	+	1t9 W2(r   r   N)
__future__r   r   r   r   typingr   redisr   rq.utilsr	   r?   r   r>   r   r   r   r   r   <module>rN      s,    " 2 2    "h( h(r   