o
    ›¨Êh|8  ã                   @   s.  d 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	 ddlm
Z ddlmZ dd	lmZmZ dd
lmZmZ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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'y…   dZ&Y nw dZ(dZ)dZ*dZ+G dd„ dƒZ,dS )aý  WorkController can be used to instantiate in-process workers.

The command-line interface for the worker is in :mod:`celery.bin.worker`,
while the worker program is in :mod:`celery.apps.worker`.

The worker program is responsible for adding signal handlers,
setting up logging, etc.  This is a bare-bones worker without
global side-effects (i.e., except for the global state stored in
:mod:`celery.worker.state`).

The worker consists of several components, all managed by bootsteps
(mod:`celery.bootsteps`).
é    N)Údatetime)Ú	cpu_count)Údetect_environment)Ú	bootsteps)Úconcurrency)Úsignals)ÚRUNÚ	TERMINATE)ÚImproperlyConfiguredÚTaskRevokedErrorÚWorkerTerminate)Ú
EX_FAILUREÚcreate_pidlock)Úreload_from_cwd)Úmlevel)Úworker_logger)Údefault_nodenameÚworker_direct)Ústr_to_list)Údefault_socket_timeouté   ©Ústate)ÚWorkControllerg      @zÑ
Trying to select queue subset of {0!r}, but queue {1} isn't
defined in the `task_queues` setting.

If you want to automatically declare unknown queues you can
enable the `task_create_missing_queues` setting.
ze
Trying to deselect queue subset of {0!r}, but queue {1} isn't
defined in the `task_queues` setting.
c                   @   s€  e Zd ZdZdZdZdZdZdZdZ	G dd„ de
jƒZdHdd„Z		dIdd„Zd	d
„ Zdd„ Zdd„ Zdd„ Zdd„ Zdd„ Zdd„ ZdJdd„Zdd„ Zdd„ Zdd„ Zdd „ Zd!d"„ Zd#d$„ Zd%d&„ Zd'd(„ Zd)d*„ ZdKd,d-„ZdLd.d/„Z dMd1d2„Z!dNd3d4„Z"dJd5d6„Z#dKd7d8„Z$d9d:„ Z%d;d<„ Z&d=d>„ Z'd?d@„ Z(dAdB„ Z)e*dCdD„ ƒZ+																					dOdFdG„Z,dS )Pr   zUnmanaged worker instance.Nc                   @   s   e Zd ZdZdZh d£ZdS )zWorkController.BlueprintzWorker bootstep blueprint.ÚWorker>   úcelery.worker.components:Hubúcelery.worker.components:Beatúcelery.worker.components:Poolúcelery.worker.components:Timerú celery.worker.components:StateDBú!celery.worker.components:Consumerú'celery.worker.autoscale:WorkerComponentN)Ú__name__Ú
__module__Ú__qualname__Ú__doc__ÚnameÚdefault_steps© r(   r(   úF/var/www/html/env/lib/python3.10/site-packages/celery/worker/worker.pyÚ	BlueprintK   s    r*   c                 K   s|   |p| j | _ t|ƒ| _t ¡ | _| j j ¡  | jdi |¤Ž | j	di |¤Ž | j
di |¤Ž | jdi | jdi |¤Ž¤Ž d S )Nr(   )Úappr   Úhostnamer   ÚutcnowÚstartup_timeÚloaderÚinit_workerÚon_before_initÚsetup_defaultsÚon_after_initÚsetup_instanceÚprepare_args)Úselfr+   r,   Úkwargsr(   r(   r)   Ú__init__Y   s   

 zWorkController.__init__c                 K   sð   || _ |  ||¡ |  t|ƒ¡ | js&ztƒ | _W n ty%   d| _Y nw t| jƒ| _|p0| j	| _
| j ¡ | _|d u r@|  ¡ n|| _|| _tjj| d� t | j¡| _g | _|  ¡  | j| jjd | j| j| jd�| _| jj| fi |¤Ž d S )Né   ©ÚsenderÚworker)ÚstepsÚon_startÚon_closeÚ
on_stopped)ÚpidfileÚsetup_queuesÚsetup_includesr   r   r   ÚNotImplementedErrorr   ÚloglevelÚon_consumer_readyÚready_callbackr+   Úconnection_for_readÚ	_conninfoÚshould_use_eventloopÚuse_eventloopÚoptionsr   Úworker_initÚsendÚ_concurrencyÚget_implementationÚpool_clsr=   Úon_init_blueprintr*   r>   r?   r@   Ú	blueprintÚapply)r6   ÚqueuesrG   rA   ÚincluderK   Úexclude_queuesr7   r(   r(   r)   r4   d   s6   
ÿþ
üzWorkController.setup_instancec                 C   ó   d S ©Nr(   ©r6   r(   r(   r)   rR   Œ   ó   z WorkController.on_init_blueprintc                 K   rX   rY   r(   ©r6   r7   r(   r(   r)   r1   �   r[   zWorkController.on_before_initc                 K   rX   rY   r(   r\   r(   r(   r)   r3   ’   r[   zWorkController.on_after_initc                 C   s   | j rt| j ƒ| _d S d S rY   )rA   r   ÚpidlockrZ   r(   r(   r)   r>   •   s   ÿzWorkController.on_startc                 C   rX   rY   r(   )r6   Úconsumerr(   r(   r)   rF   ™   r[   z WorkController.on_consumer_readyc                 C   s   | j j ¡  d S rY   )r+   r/   Úshutdown_workerrZ   r(   r(   r)   r?   œ   s   zWorkController.on_closec                 C   s,   | j  ¡  | j ¡  | jr| j ¡  d S d S rY   )ÚtimerÚstopr^   Úshutdownr]   ÚreleaserZ   r(   r(   r)   r@   Ÿ   s
   

ÿzWorkController.on_stoppedc              
   C   s¼   t |ƒ}t |ƒ}z
| jjj |¡ W n ty( } z
tt ¡  	||¡ƒ‚d }~ww z
| jjj 
|¡ W n tyI } z
tt ¡  	||¡ƒ‚d }~ww | jjjr\| jjj t| jƒ¡ d S d S rY   )r   r+   ÚamqprU   ÚselectÚKeyErrorr
   ÚSELECT_UNKNOWN_QUEUEÚstripÚformatÚdeselectÚDESELECT_UNKNOWN_QUEUEÚconfr   Ú
select_addr,   )r6   rV   ÚexcludeÚexcr(   r(   r)   rB   ¦   s*   ÿ€ÿÿ€ÿ
ÿzWorkController.setup_queuesc                    sf   t ˆ jjjƒ}|r|t |ƒ7 }‡ fdd„|D ƒ |ˆ _dd„ ˆ jj ¡ D ƒ}t t|ƒ|B ƒˆ jj_d S )Nc                    s   g | ]	}ˆ j j |¡‘qS r(   )r+   r/   Úimport_task_module©Ú.0ÚmrZ   r(   r)   Ú
<listcomp>¼   s    z1WorkController.setup_includes.<locals>.<listcomp>c                 S   s   h | ]}|j j’qS r(   )Ú	__class__r#   )rr   Útaskr(   r(   r)   Ú	<setcomp>¾   s    ÿz0WorkController.setup_includes.<locals>.<setcomp>)Útupler+   rl   rV   ÚtasksÚvaluesÚset)r6   ÚincludesÚprevÚtask_modulesr(   rZ   r)   rC   ¶   s   
ÿzWorkController.setup_includesc                 K   s   |S rY   r(   r\   r(   r(   r)   r5   Â   r[   zWorkController.prepare_argsc                 C   s   t jj| d� d S )Nr:   )r   Úworker_shutdownrN   rZ   r(   r(   r)   Ú_send_worker_shutdownÅ   s   z$WorkController._send_worker_shutdownc              
   C   sÀ   z	| j  | ¡ W d S  ty   |  ¡  Y d S  ty7 } ztjd|dd� | jtd� W Y d }~d S d }~w t	yP } z| j|j
d� W Y d }~d S d }~w ty_   | jtd� Y d S w )NzUnrecoverable error: %rT)Úexc_info)Úexitcode)rS   Ústartr   Ú	terminateÚ	ExceptionÚloggerÚcriticalra   r   Ú
SystemExitÚcodeÚKeyboardInterrupt)r6   ro   r(   r(   r)   rƒ   È   s   €€ÿzWorkController.startc                 C   s   | j j| d|fdd� d S )NÚregister_with_event_loopzhub.register)ÚargsÚdescription)rS   Úsend_all)r6   Úhubr(   r(   r)   r‹   Õ   s   
þz'WorkController.register_with_event_loopc                 C   s   |   | j|¡S rY   )Ú_quick_acquireÚ_process_task©r6   Úreqr(   r(   r)   Ú_process_task_semÛ   s   z WorkController._process_task_semc                 C   sJ   z	|  | j¡ W dS  ty$   z|  ¡  W Y dS  ty#   Y Y dS w w )z2Process task by sending it to the pool of workers.N)Úexecute_using_poolÚpoolr   Ú_quick_releaseÚAttributeErrorr’   r(   r(   r)   r‘   Þ   s   ÿýzWorkController._process_taskc                 C   s&   z| j  ¡  W d S  ty   Y d S w rY   )r^   Úcloser˜   rZ   r(   r(   r)   Úsignal_consumer_closeè   s
   ÿz$WorkController.signal_consumer_closec                 C   s    t ƒ dko| jjjjo| jj S )NÚdefault)r   rI   Ú	transportÚ
implementsÚasynchronousr+   Ú
IS_WINDOWSrZ   r(   r(   r)   rJ   î   s
   

ÿþz#WorkController.should_use_eventloopFc                 C   sF   |dur|| _ | jjtkr|  ¡  |r| jjr| jdd� |  ¡  dS )z'Graceful shutdown of the worker server.NT©Úwarm)	r‚   rS   r   r   rš   r–   Úsignal_safeÚ	_shutdownr€   )r6   Úin_sighandlerr‚   r(   r(   r)   ra   ó   s   zWorkController.stopc                 C   s8   | j jtkr|  ¡  |r| jjr| jdd� dS dS dS )z.Not so graceful shutdown of the worker server.Fr    N)rS   r   r	   rš   r–   r¢   r£   )r6   r¤   r(   r(   r)   r„   ý   s   ýzWorkController.terminateTc                 C   sX   | j d ur*ttƒ� | j j| | d� | j  ¡  W d   ƒ d S 1 s#w   Y  d S d S )N)r„   )rS   r   ÚSHUTDOWN_SOCKET_TIMEOUTra   Újoin)r6   r¡   r(   r(   r)   r£     s   

"þÿzWorkController._shutdownc                 C   sT   t | j|||d�ƒ | jr| j ¡  | j ¡  z| j ¡  W d S  ty)   Y d S w )N)Úforce_reloadÚreloader)ÚlistÚ_reload_modulesr^   Úupdate_strategiesÚreset_rate_limitsr–   ÚrestartrD   )r6   ÚmodulesÚreloadr¨   r(   r(   r)   r¯     s   ÿ

ÿzWorkController.reloadc                    s4   ‡ ‡fdd„t |d u rˆjjjƒD ƒS |pdƒD ƒS )Nc                 3   s"   � | ]}ˆj |fi ˆ ¤ŽV  qd S rY   )Ú_maybe_reload_modulerq   ©r7   r6   r(   r)   Ú	<genexpr>  s
   € ÿ
ÿz1WorkController._reload_modules.<locals>.<genexpr>r(   )r{   r+   r/   r~   )r6   r®   r7   r(   r±   r)   rª     s   
ÿþÿþzWorkController._reload_modulesc                 C   sH   |t jvrt d|¡ | jj |¡S |r"t d|¡ tt j| |ƒS d S )Nzimporting module %szreloading module %s)Úsysr®   r†   Údebugr+   r/   Úimport_from_cwdr   )r6   Úmoduler§   r¨   r(   r(   r)   r°     s   
þz#WorkController._maybe_reload_modulec                 C   s4   t  ¡ | j }| jjt ¡ t| jj	ƒt
| ¡ ƒdœS )N)ÚtotalÚpidÚclockÚuptime)r   r-   r.   r   Útotal_countÚosÚgetpidÚstrr+   r¹   ÚroundÚtotal_seconds)r6   rº   r(   r(   r)   Úinfo'  s   

ýzWorkController.infoc                 C   s    t d u rtdƒ‚t  t j¡}i d|j“d|j“d|j“d|j“d|j“d|j	“d|j
“d	|j“d
|j“d|j“d|j“d|j“d|j“d|j“d|j“d|j“S )Nz%rusage not supported by this platformÚutimeÚstimeÚmaxrssÚixrssÚidrssÚisrssÚminfltÚmajfltÚnswapÚinblockÚoublockÚmsgsndÚmsgrcvÚnsignalsÚnvcswÚnivcsw)ÚresourcerD   Ú	getrusageÚRUSAGE_SELFÚru_utimeÚru_stimeÚ	ru_maxrssÚru_ixrssÚru_idrssÚru_isrssÚ	ru_minfltÚ	ru_majfltÚru_nswapÚ
ru_inblockÚ
ru_oublockÚ	ru_msgsndÚ	ru_msgrcvÚru_nsignalsÚru_nvcswÚ	ru_nivcsw)r6   Úsr(   r(   r)   Úrusage.  sH   ÿþýüûúùø	÷
öõôóòñðzWorkController.rusagec                 C   s`   |   ¡ }| | j  | ¡¡ | | jj  | j¡¡ z	|  ¡ |d< W |S  ty/   d|d< Y |S w )Nræ   zN/A)rÁ   ÚupdaterS   r^   ræ   rD   )r6   rÁ   r(   r(   r)   ÚstatsE  s   þ
þzWorkController.statsc                 C   s"   dj | | jr| j ¡ d�S dd�S )z``repr(worker)``.z#<Worker: {self.hostname} ({state})>ÚINIT)r6   r   )ri   rS   Úhuman_staterZ   r(   r(   r)   Ú__repr__O  s   þþzWorkController.__repr__c                 C   s   | j S )z#``str(worker) == worker.hostname``.)r,   rZ   r(   r(   r)   Ú__str__V  s   zWorkController.__str__c                 C   s   t S rY   r   rZ   r(   r(   r)   r   Z  s   zWorkController.stateÚWARNc                 K   s  | j j}|| _|| _|d|ƒ| _|d|ƒ| _|d||ƒ| _|d|ƒ| _|d|ƒ| _|d|ƒ| _	|p2|| _
|d|	ƒ| _|d|
ƒ| _|d	|ƒ| _|d
||ƒ| _|d|ƒ| _|d||ƒ| _|d||ƒ| _|d||ƒ| _|d|ƒ| _|d|ƒ| _t|d|ƒƒ| _|d|ƒ| _|d|ƒ| _d S )NÚworker_concurrencyÚworker_send_task_eventsÚworker_poolÚworker_consumerÚworker_timerÚworker_timer_precisionÚworker_autoscalerÚworker_pool_putlocksÚworker_pool_restartsÚworker_state_dbÚbeat_schedule_filenameÚbeat_schedulerÚtask_time_limitÚtask_soft_time_limitÚworker_max_tasks_per_childÚworker_max_memory_per_childÚworker_prefetch_multiplierÚworker_disable_rate_limitsÚworker_lost_wait)r+   ÚeitherrE   Úlogfiler   Útask_eventsrQ   Úconsumer_clsÚ	timer_clsÚtimer_precisionÚoptimizationÚautoscaler_clsÚpool_putlocksÚpool_restartsÚstatedbÚschedule_filenameÚ	schedulerÚ
time_limitÚsoft_time_limitÚmax_tasks_per_childÚmax_memory_per_childÚintÚprefetch_multiplierÚdisable_rate_limitsr   )r6   r   rE   r  r  r–   r  r  r  r  r	  r
  r  ÚOr  r  r  r  rQ   Ústate_dbrú   rû   Úscheduler_clsr  r  r  r  r   r  Ú_kwr  r(   r(   r)   r2   ^  sN   ÿ
ÿÿÿÿÿÿÿzWorkController.setup_defaults)NN)NNNNNNrY   )FN)F)T)NFN)Nrí   NNNNNNNNNNNNNNNNNNNNNNNNNN)-r"   r#   r$   r%   r+   r]   rS   r–   Ú	semaphorer‚   r   r*   r8   r4   rR   r1   r3   r>   rF   r?   r@   rB   rC   r5   r€   rƒ   r‹   r”   r‘   rš   rJ   ra   r„   r£   r¯   rª   r°   rÁ   ræ   rè   rë   rì   Úpropertyr   r2   r(   r(   r(   r)   r   >   s‚    

ÿ(










ìr   )-r%   r¼   r³   r   Úbilliardr   Úkombu.utils.compatr   Úceleryr   r   rO   r   Úcelery.bootstepsr   r	   Úcelery.exceptionsr
   r   r   Úcelery.platformsr   r   Úcelery.utils.importsr   Úcelery.utils.logr   r   r†   Úcelery.utils.nodenamesr   r   Úcelery.utils.textr   Úcelery.utils.threadsr   Ú r   rÒ   ÚImportErrorÚ__all__r¥   rg   rk   r   r(   r(   r(   r)   Ú<module>   s:    ÿ