
    AHjm.                       U d dl mZ d dlZd dlZd dlZd dlZd dlZd dlZd dl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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  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,m-Z-m.Z.m/Z/m0Z0m1Z1 de2d<   	  ed      j                  Z3dZ5dZ6 G d de7e      Z8 G d d      Z9d Z:y# e4$ r eZ3Y 'w xY w)    )annotationsN)Iterable)datetime)Enum)Processget_context)BaseProcess)uuid4)ConnectionPoolRedis)Pipeline   )parse_connection)DEFAULT_LOGGING_DATE_FORMATDEFAULT_LOGGING_FORMAT!DEFAULT_SCHEDULER_FALLBACK_PERIOD)SchedulerNotFound)Job)setup_loghandlers)Queue)ScheduledJobRegistry)resolve_serializer)current_timestampdecode_redis_hashnowparse_names	utcformatutcparseztype[BaseProcess]ForkProcessforkzrq:scheduler:%szrq:scheduler-lock:%sc                      e Zd ZdZdZdZy)SchedulerStatusstartedworkingstoppedN)__name__
__module____qualname__STARTEDWORKINGSTOPPED     D/root/tools/cai/cai_env/lib/python3.12/site-packages/rq/scheduler.pyr"   r"   '   s    GGGr-   r"   c                  (   e Zd ZeZdej                  eeddf	 	 	 	 	 ddZ	e
d        Ze
d        Ze
d        Ze
dd       Ze
d        Zdd	Zddd
Zedd       Zd dZd!dZdd"dZd#dZd$dZed%d       Zd Zd Zd&dZd Zd Zd Zd Z d Z!y)'RQSchedulerr   Nc	                B   t        t        |            | _        t               | _        g | _        d | _        t        |      \  | _        | _        | _	        t        |      | _        |xs t               j                  | _        t        j                          | _        t%               | _        d| _        d | _        d | _        || _        d| _        | j2                  j4                  | _        d | _        t;        j<                  t>              | _         tC        |t>        ||       y )Nr   F)levelname
log_formatdate_format)"setr   _queue_names_acquired_locks_scheduled_job_registrieslock_acquisition_timer   _connection_class_pool_class_pool_kwargsr   
serializerr
   hexr3   socketgethostnamehostnamer   
created_atpidlast_heartbeat_connectioninterval_stop_requestedStatusr+   _status_processlogging	getLoggerr&   logr   )	selfqueues
connectionrG   logging_levelr5   r4   r>   r3   s	            r.   __init__zRQScheduler.__init__4   s      F 34),EG&%)"FVWaFbC 0$2C,Z8 ,	#//1$'E/3 ${{**$$X.!#		
r-   c                    | j                   r| j                   S | j                  t        dd| j                  i| j                        | _         | j                   S )Nconnection_class)connection_poolr,   )rF   r;   r   r<   r=   rO   s    r.   rQ   zRQScheduler.connectionZ   sZ    ###11*bD<L<LbPTPaPab 2 
 r-   c                    | j                   S N)r8   rW   s    r.   acquired_lockszRQScheduler.acquired_locksc   s    ###r-   c                    | j                   S rY   )rJ   rW   s    r.   statuszRQScheduler.statusg   s    ||r-   c                (    t         | j                  z  S )z1Redis key holding this scheduler's metadata hash.)SCHEDULER_KEY_TEMPLATEr3   rW   s    r.   keyzRQScheduler.keyk   s     &		11r-   c                    | j                   | j                  k(  ry| j                  syt        j                         | j                  z
  j                         t        kD  S )zCReturns True if lock_acquisition_time is longer than 10 minutes agoFT)r7   rZ   r:   r   r   total_secondsr   rW   s    r.   should_reacquire_locksz"RQScheduler.should_reacquire_locksp   sM      3 33))!;!;;JJLOpppr-   c                   t               }| j                  j                  ddj                  | j                               | j                  D ]u  }| j
                  j                  | j                  |      | j                  d| j                  dz         sI| j                  j                  d|       |j                  |       w g | _        | j                  j                  |      | _        t        j                         | _        | j                  r8|r6| j"                  r| j"                  j%                         s| j'                          |S )z7Returns names of queue it successfully acquires lock onzAcquiring scheduler lock for %s, T<   )nxexzAcquired scheduler lock for %s)r6   rN   debugjoinr7   rQ   get_locking_keyr3   rG   infoaddr9   r8   unionr   r   r:   rK   is_alivestart)rO   
auto_startsuccessful_locksr3   s       r.   acquire_lockszRQScheduler.acquire_locksy   s    58$))DDUDU:VW%% 	+D""4#7#7#=tyyTVZVcVcfhVh"i>E $$T*	+ *,&#3399:JK%-\\^" J==(>(>(@

r-   c                    g | _         |s| j                  }|D ]=  }| j                   j                  t        || j                  | j
                               ? y)z(Prepare scheduled job registries for userQ   r>   N)r9   r8   appendr   rQ   r>   )rO   queue_namesr3   s      r.   prepare_registrieszRQScheduler.prepare_registries   sR    )+&..K 	D**11$TdooRVRaRab	r-   c                    t         |z  S )z,Returns scheduler key for a given queue name)SCHEDULER_LOCKING_KEY_TEMPLATE)clsr3   s     r.   rj   zRQScheduler.get_locking_key   s     .44r-   c                    | j                   J | j                  | j                  t        | j                        dj                  | j                        t        | j                        t        | j                         dS )zBSerialize this scheduler's metadata for storage in its Redis hash.,)r3   rB   rD   rP   rC   rE   )	rE   r3   rB   strrD   ri   r7   r   rC   rW   s    r.   to_dictzRQScheduler.to_dict   sc    ""...IItxx=hht001#DOO4'(;(;<
 	
r-   c                d   t        |d      }|d   | _        |d   | _        t        |d         | _        |j                  d      rt        |d   j                  d            n	t               | _        t        |d         | _
        |j                  d	      rt        |d	         | _        y
d
| _        y
)z6Restore this scheduler's metadata from its Redis hash.T)decode_valuesr3   rB   rD   rP   r|   rC   rE   N)r   r3   rB   intrD   getr6   splitr7   r   rC   rE   )rO   raw_dataobjs      r.   restorezRQScheduler.restore   s    =K	Js5z?=@WWX=NCH 3 3C 89TWTY"3|#45ADIYAZhs+;'<=`dr-   c                   ||n| j                   j                         }|j                  | j                  | j	                                |j                  | j                  | j                  dz          ||j                          yy)z/Save this scheduler's metadata hash with a TTL.N)mappingre   )rQ   pipelinehsetr_   r~   expirerG   execute)rO   r   rQ   s      r.   savezRQScheduler.save   si    !)!5X4??;S;S;U
$,,.9$((DMMB$67  r-   c                    | j                   j                  d| j                         t        j                         | _        t               | _        | j                          y)zRegister this scheduler's birth by writing its metadata hash.

        Idempotent: re-registering the same name (e.g. when the worker restarts a crashed
        scheduler process) overwrites the existing hash rather than erroring.
        zScheduler %s: registering birthN)	rN   rh   r3   osgetpidrD   r   rE   r   rW   s    r.   register_birthzRQScheduler.register_birth   s;     	8$))D99;!e		r-   c                    | j                   j                  d| j                         t        | j                  j                  | j                              S )zRegister this scheduler's death by deleting its metadata hash.

        Returns:
            True if the scheduler metadata existed and was deleted, False if it was already absent.
        zScheduler %s: registering death)rN   rh   r3   boolrQ   deleter_   rW   s    r.   register_deathzRQScheduler.register_death   s9     	8$))DDOO**488455r-   c                    |j                  t        |z        }|st        d| d       | g ||      }|j                  |       |S )z<Fetch a scheduler by name, restoring it from its Redis hash.zScheduler with name 'z' not found)rQ   r3   )hgetallr^   r   r   )rz   r3   rQ   r   	schedulers        r.   fetchzRQScheduler.fetch   sT     %%&<t&CD#&;D6$MNNz=	(#r-   c           	        | j                   j                  | _        | j                  s| j                  r| j                          | j                  D ]  }t               }|j                  |      }|s!t        |j                  | j                  | j                        }| j                  j                         5 }t        j                  || j                  | j                        }|D ]'  }||j                  |||j!                                ) |D ]  }|j#                  ||        |j%                          ddd        | j                   j&                  | _        y# 1 sw Y   xY w)z+Enqueue jobs whose timestamp is in the pastrt   N)r   at_front)r   )rI   r*   rJ   r9   r8   rw   r   get_jobs_to_scheduler   r3   rQ   r>   r   r   
fetch_many_enqueue_jobshould_enqueue_at_frontremover   r)   )	rO   registry	timestampjob_idsqueuer   jobsjobjob_ids	            r.   enqueue_scheduled_jobsz"RQScheduler.enqueue_scheduled_jobs   s?   {{**--$2F2F##%66 	#H)+I 33I>G(--DOOPTP_P_`E))+ #x~~g$//VZVeVef kC**3CLgLgLi*jk & ?FOOFXO>?  "# #	#( {{**# #s   74E ,AE  E*	c                    t        j                   t         j                  | j                         t        j                   t         j                  | j                         y)zUInstalls signal handlers for handling SIGINT and SIGTERM
        gracefully.
        N)signalSIGINTrequest_stopSIGTERMrW   s    r.   _install_signal_handlersz$RQScheduler._install_signal_handlers   s4     	fmmT%6%67fnnd&7&78r-   c                    d| _         y)z8Toggle self._stop_requested that's checked on every loopTN)rH   )rO   signumframes      r.   r   zRQScheduler.request_stop   s
    #r-   c                :   | j                   j                  ddj                  | j                               t	               | _        | j                  j                         5 }|j                  | j                  dt        | j
                               |j                  | j                  | j                  dz          | j                  D ]0  }|j                  | j                  |      | j                  dz          2 |j                          ddd       y# 1 sw Y   yxY w)z?Refresh the TTL on the scheduler's metadata hash and its locks.z!Scheduler sending heartbeat to %srd   rE   re   N)rN   rh   ri   rZ   r   rE   rQ   r   r   r_   r   r   rG   r8   rj   r   )rO   r   r3   s      r.   	heartbeatzRQScheduler.heartbeat  s    :DIIdFYFY<Z[!e__%%' 	8MM$(($4i@S@S6TUOODHHdmmb&89,, P 4 4T :DMMB<NOP	 	 	s   B)DDc                    | j                   j                  ddj                  | j                               | j	                          | j
                  j                  | _        | j                          y )Nz-Scheduler stopping, releasing locks for %s...rd   )	rN   rk   ri   r8   release_locksrI   r+   rJ   r   rW   s    r.   stopzRQScheduler.stop  sN    EtyyQUQeQeGfg{{**r-   c                    | j                   D cg c]  }| j                  |       }} | j                  j                  |  t	               | _         yc c}w )zRelease acquired locksN)r8   rj   rQ   r   r6   )rO   r3   keyss      r.   r   zRQScheduler.release_locks  sJ    7;7K7KLt$$T*LL%"u Ms   Ac                    | j                   j                  | _        d | _        t	        t
        | fd      | _        | j                  j                          | j                  S )N	Scheduler)targetargsr3   )rI   r)   rJ   rF   r   runrK   ro   rW   s    r.   ro   zRQScheduler.start  sI    {{**  #3dW;O}}r-   c                6   | j                          | j                          	 | j                  r| j                          y | j                  r| j                          | j                          | j                          t        j                  | j                         yrY   )r   r   rH   r   rb   rr   r   r   timesleeprG   rW   s    r.   workzRQScheduler.work#  sr    %%'##		**""$'')NNJJt}}% r-   )rQ   r   rR   z	str | intr3   z
str | None)returnr}   )FrY   )rv   zIterable[str] | None)r3   r}   )r   dict)r   r   r   None)r   zPipeline | Noner   r   )r   r   )r   r   )r3   r}   rQ   r   r   r0   )NN)"r&   r'   r(   r"   rI   rL   INFOr   r   rS   propertyrQ   rZ   r\   r_   rb   rr   rw   classmethodrj   r~   r   r   r   r   r   r   r   r   r   r   r   ro   r   r,   r-   r.   r0   r0   -   s   
 F #*<</)$
 $

 !$
 $
L     $ $   2 2 q q , 5 5

e!	6  +:9$	%&r-   r0   c                   | j                   j                  ddj                  | j                        t	        j
                                	 | j                          | j                   j                  dt	        j
                                y #  | j                   j                  dt	        j
                         t        j                                 xY w)Nz$Scheduler for %s started with PID %srd   z*Scheduler [PID %s] raised an exception.
%sz!Scheduler with PID %d has stopped)
rN   rk   ri   r7   r   r   r   error	traceback
format_exc)r   s    r.   r   r   4  s    MM=tyyI_I_?`bdbkbkbmn MM:BIIKHI299;XaXlXlXnos   
B	 	AC);
__future__r   rL   r   r   r@   r   r   collections.abcr   r   enumr   multiprocessingr   r   multiprocessing.processr	   uuidr
   redisr   r   redis.clientr   connectionsr   defaultsr   r   r   
exceptionsr   r   r   logutilsr   r   r   r   r   serializersr   utilsr   r   r   r   r   r   __annotations__r   
ValueErrorr^   ry   r}   r"   r0   r   r,   r-   r.   <module>r      s    "  	     $   0 /  ' ! ) l l )  '  * + ^ ^ f%--K + !7 c4 D& D&NIi  Ks   C CC