
    AHj(                       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mZ d dlmZ d dlmZ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% ddl&m'Z' ddl(m)Z)m*Z* de+d<   	  ed      j                  Z,erd dlm.Z.  G d de      Z/ G d d      Z0e*ee!e%ddd f	 	 	 	 	 	 	 	 	 	 	 	 	 	 	 	 	 	 	 ddZ1y# e-$ r eZ,Y Fw xY w)    )annotationsN)Iterable)Enum)Processget_context)BaseProcess)TYPE_CHECKING
NamedTuple)uuid4)ConnectionPoolRedis)DefaultSerializer   )parse_connection)DEFAULT_LOGGING_DATE_FORMATDEFAULT_LOGGING_FORMAT)Job)setup_loghandlers)Queue)parse_names)
BaseWorkerWorkerztype[BaseProcess]ForkProcessfork)
Serializerc                  ,    e Zd ZU ded<   ded<   ded<   y)
WorkerDatastrnameintpidr   processN)__name__
__module____qualname____annotations__     F/root/tools/cai/cai_env/lib/python3.12/site-packages/rq/worker_pool.pyr   r   &   s    
I	Hr(   r   c                     e Zd Z G d de      Zdeeeef	 	 	 	 	 	 	 	 	 	 	 	 	 ddZ	e
dd       Ze
dd       Zd Zdd	Zdd
Zd ZddZdddZ	 	 d	 	 	 	 	 	 	 	 	 ddZ	 	 	 	 d	 	 	 	 	 	 	 ddZd d!dZej.                  fddZd Zd"d#dZy)$
WorkerPoolc                      e Zd ZdZdZdZy)WorkerPool.Statusr         N)r#   r$   r%   IDLESTARTEDSTOPPEDr'   r(   r)   Statusr-   -   s    r(   r3   r   c                   || _         g | _        t        dt        t        t
               t        j                  t
              | _        t        |      | _
        || _        t               j                  | _        d| _        d| _        | j"                  j$                  | _        || _        || _        || _        || _        i | _        t3        |      \  | _        | _        | _        y )NINFOr   Tr   )num_workers_workersr   r   r   r#   logging	getLoggerlogr   _queue_names
connectionr   hexr   _burst_sleepr3   r0   statusworker_class
serializer	job_classqueue_classworker_dictr   _connection_class_pool_class_pool_kwargs)
selfqueuesr=   r7   rB   rC   rD   rE   argskwargss
             r)   __init__zWorkerPool.__init__2   s     !,&(&"=?U\de#*#4#4X#>'26':$	 #';;#3#3.:&0$-(3 35FVWaFbC 0$2Cr(   c                v    | j                   D cg c]  }| j                  || j                        ! c}S c c}w )Returns a list of Queue objectsr=   )r<   rE   r=   )rJ   r   s     r)   rK   zWorkerPool.queuesR   s4     PTO`O`at  $// Baaas   $6c                ,    t        | j                        S )rP   )lenrF   rJ   s    r)   number_of_active_workersz#WorkerPool.number_of_active_workersW   s     4##$$r(   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SIGTERMrT   s    r)   _install_signal_handlersz#WorkerPool._install_signal_handlers\   s4     	fmmT%6%67fnnd&7&78r(   Nc                    | j                   j                  d       | j                  j                  | _        | j                          y)z8Toggle self._stop_requested that's checked on every loopz)Received SIGINT/SIGTERM, shutting down...N)r;   infor3   r2   rA   stop_workers)rJ   signumframes      r)   rY   zWorkerPool.request_stopc   s0    ABkk))r(   c                @    | j                          | j                  dk(  S )z)Returns True if all workers have stopped.r   )reap_workersrU   rT   s    r)   all_workers_have_stoppedz#WorkerPool.all_workers_have_stoppedi   s    ,,11r(   c                ~   | j                   j                  d       t        | j                  j	                               }|D ]z  }|j
                  j                  d       |j
                  j                         r2| j                   j                  d|j                  |j                         j| j                  |       | y)z%Removes dead workers from worker_dictzReaping dead workersg?zWorker %s with pid %d is aliveN)r;   debuglistrF   valuesr"   joinis_aliver   r!   handle_dead_worker)rJ   worker_datasdatas      r)   rb   zWorkerPool.reap_workerso   s    -.D,,3356  	DLLc"||$$&?DHHU''-	r(   c                   | j                   j                  d|j                  |j                         t	        j
                  t              5  | j                  j                  |j                         ddd       y# 1 sw Y   yxY w)z&
        Handle a dead worker
        zWorker %s with pid %d is deadN)	r;   r]   r   r!   
contextlibsuppressKeyErrorrF   pop)rJ   worker_datas     r)   rj   zWorkerPool.handle_dead_worker   s`     	5{7G7GY  * 	3  !1!12	3 	3 	3s   &A::Bc                `   | j                   j                  d       | j                          |r| j                  | j                  j
                  k7  r]| j                  t        | j                        z
  }|r8t        |      D ])  }| j                  | j                  | j                         + yyyy)z7
        Check whether workers are still alive
        zChecking worker processes)burstr@   N)r;   re   rb   rA   r3   r2   r7   rS   rF   rangestart_workerr?   r@   )rJ   respawndeltais       r)   check_workerszWorkerPool.check_workers   s     	23 t{{dkk&9&99$$s4+;+;'<<Eu MA%%DKK%LM  :7r(   c                    t        t        || j                  | j                  | j                  | j
                  f|||| j                  | j                  | j                  dd| d| j                   d      S )zReturns the worker process)r@   rt   logging_levelrB   rD   rC   zWorker z (WorkerPool ))targetrL   rM   r   )
r   
run_workerr<   rG   rH   rI   rB   rD   rC   r   )rJ   r   rt   r@   r|   s        r)   get_worker_processzWorkerPool.get_worker_process   sx     ))4+A+A4CSCSUYUfUfg !. $ 1 1!^^"oo 4&dii[:
 	
r(   c                   t               j                  }| j                  ||||      }|j                          t	        ||j
                  |      }|| j                  |<   | j                  j                  d||j
                         y)z
        Starts a worker and adds the data to worker_datas.
        * sleep: waits for X seconds before creating worker, for testing purposes
        rt   r@   r|   )r   r!   r"   zSpawned worker: %s with PID %dN)	r   r>   r   startr   r!   rF   r;   re   )rJ   countrt   r@   r|   r   r"   rr   s           r)   rv   zWorkerPool.start_worker   sm     w{{))$eFZg)h dWM!,7w{{Kr(   c                    | j                   j                  d| j                   d       t        | j                        D ]  }| j	                  |dz   |||        y)zx
        Run the workers
        * sleep: waits for X seconds before creating worker, only for testing purposes
        z	Spawning z workersr   r   N)r;   re   r7   ru   rv   )rJ   rt   r@   r|   ry   s        r)   start_workerszWorkerPool.start_workers   s[    
 	4#3#3"4H=>t''( 	^Aa!e5}]	^r(   c                2   	 t        j                  |j                  |       | j                  j	                  d|j                         y# t
        $ rD}|j                  t        j                  k(  r| j                  j                  d       n Y d}~yd}~ww xY w)zm
        Send stop signal to worker and catch "No such process" error if the worker is already dead.
        z'Sent shutdown command to worker with %szHorse already deadN)	oskillr!   r;   r]   OSErrorerrnoESRCHre   )rJ   rr   siges       r)   stop_workerzWorkerPool.stop_worker   si    	GGKOOS)HHMMC[__U 	ww%++%34 5	s   AA	 		B:BBc                    | j                   j                  dt        | j                               t	        | j                  j                               }|D ]  }| j                  |        y)zSend SIGINT to all workersz!Sending stop signal to %s workersN)r;   r]   rS   rF   rf   rg   r   )rJ   rk   rr   s      r)   r^   zWorkerPool.stop_workers   sV    93t?O?O;PQD,,3356' 	*K[)	*r(   c                    || _         | }t        |t        t        t               | j
                  j                  d| j                   dt        j                                | j                  j                  | _        | j                  | j                   |       | j                          	 | j                  | j                  j                  k(  r]| j!                         r| j
                  j                  d       y | j
                  j                  d       t#        j$                  d       | j'                  |       |r+| j(                  d	k(  r| j
                  j                  d       y t#        j$                  d       )
Nr6   zStarting worker pool z with pid %d...)rt   r|   zAll workers stopped, exiting...z"Waiting for workers to shutdown...r   )rw   r   )r?   r   r   r   r#   r;   r]   r   r   getpidr3   r1   rA   r   r[   r2   rc   timesleeprz   rU   )rJ   rt   r|   rw   s       r)   r   zWorkerPool.start   s   )-)DF\ckl-dii[H"))+Vkk))MJ%%'{{dkk111002HHMM"CDHHMM"FGJJqM""7"3T::a?HHMM"CD

1 r(   )rK   zIterable[str | Queue]r=   r   r7   r    rB   type[BaseWorker]rC   r   rD   	type[Job]rE   type[Queue])returnzlist[Queue])r   r    )NN)r   bool)rr   r   )T)rw   r   r   None)r   r5   )
r   r   rt   r   r@   floatr|   r   r   r   )NTr   r5   )r   z
int | Nonert   r   r@   r   r|   r   )Tr   r5   )rt   r   r@   r   r|   r   )Fr5   )rt   r   r|   r   )r#   r$   r%   r   r3   r   r   r   r   rN   propertyrK   rU   r[   rY   rc   rb   rj   rz   r   rv   r   rW   rX   r   r^   r   r'   r(   r)   r+   r+   ,   sS     )/!2"#(c%c c 	c
 'c c c !c@ b b % %9243M$ #

 
 	

 
 

0 !#LL L 	L
 L$^ 8>}} *r(   r+   Tr5   c                .    |t        dd|i|      }|D cg c]  } |||       }} ||| ||||      }|j                  j                  dt        j                                t        j                  |       |j                  |	d|
       y c c}w )	Nconnection_class)connection_poolrQ   )r   r=   rC   rD   rE   z#Starting worker started with PID %sT)rt   with_schedulerr|   r'   )r   r;   r]   r   r   r   r   work)worker_namequeue_namesr   connection_pool_classconnection_pool_kwargsrB   rC   rD   rE   rt   r|   r@   r=   r   rK   workers                   r)   r   r      s     "&h8MhQghJ DOO4k$:6OFOF JJOO9299;GJJv
KKeDKN Ps   B)r   r   r   zIterable[str]r   dictrB   r   rC   r   rD   r   rE   r   rt   r   r|   r   r@   r    )2
__future__r   rn   r   r9   r   rW   r   collections.abcr   enumr   multiprocessingr   r   multiprocessing.processr   typingr	   r
   uuidr   redisr   r   rq.serializersr   connectionsr   defaultsr   r   jobr   logutilsr   queuer   utilsr   r   r   r   r&   r   
ValueErrorr   r   r+   r   r'   r(   r)   <module>r      s!   "    	   $  0 / ,  ' , ) I  '   & f%--K ) P Pr &,.$OOO
 !O #O O O O O O OA  Ks   
C C#"C#