o
    wvXjk9  ã                   @   sT  d Z ddlmZm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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Z0dZ1eG dd„ de2ƒƒZ3d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`).
é    )Úabsolute_importÚunicode_literalsN)Údatetime)Ú	cpu_count)Údetect_environment)Ú	bootsteps)Úconcurrency)Úsignals)ÚRUNÚ	TERMINATE)ÚImproperlyConfiguredÚTaskRevokedErrorÚWorkerTerminate)Úpython_2_unicode_compatibleÚvalues)Ú
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,   úQ/var/www/html/myproject/venv/lib/python3.10/site-packages/celery/worker/worker.pyÚ	BlueprintQ   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/   r0   Úkwargsr,   r,   r-   Ú__init___   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_clsrA   Úon_init_blueprintr.   rB   rC   rD   Ú	blueprintÚapply)r:   ÚqueuesrK   rE   ÚincluderO   Úexclude_queuesr;   r,   r,   r-   r8   j   s6   
ÿþ
üzWorkController.setup_instancec                 C   ó   d S ©Nr,   ©r:   r,   r,   r-   rV   ’   ó   z WorkController.on_init_blueprintc                 K   r\   r]   r,   ©r:   r;   r,   r,   r-   r5   •   r_   zWorkController.on_before_initc                 K   r\   r]   r,   r`   r,   r,   r-   r7   ˜   r_   zWorkController.on_after_initc                 C   s   | j rt| j ƒ| _d S d S r]   )rE   r   Úpidlockr^   r,   r,   r-   rB   ›   s   ÿzWorkController.on_startc                 C   r\   r]   r,   )r:   Úconsumerr,   r,   r-   rJ   Ÿ   r_   z WorkController.on_consumer_readyc                 C   s   | j j ¡  d S r]   )r/   r3   Úshutdown_workerr^   r,   r,   r-   rC   ¢   s   zWorkController.on_closec                 C   s,   | j  ¡  | j ¡  | jr| j ¡  d S d S r]   )ÚtimerÚstoprb   Úshutdownra   Úreleaser^   r,   r,   r-   rD   ¥   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 r]   )r   r/   ÚamqprY   ÚselectÚKeyErrorr   ÚSELECT_UNKNOWN_QUEUEÚstripÚformatÚdeselectÚDESELECT_UNKNOWN_QUEUEÚconfr   Ú
select_addr0   )r:   rZ   ÚexcludeÚexcr,   r,   r-   rF   ¬   s*   ÿ€ÿÿ€ÿ
ÿzWorkController.setup_queuesc                    sf   t ˆ jjjƒ}|r|t |ƒ7 }‡ fdd„|D ƒ |ˆ _dd„ tˆ jjƒD ƒ}t t|ƒ|B ƒˆ jj_d S )Nc                    s   g | ]	}ˆ j j |¡‘qS r,   )r/   r3   Úimport_task_module©Ú.0Úmr^   r,   r-   Ú
<listcomp>Â   s    z1WorkController.setup_includes.<locals>.<listcomp>c                 S   s   h | ]}|j j’qS r,   )Ú	__class__r'   )rv   Útaskr,   r,   r-   Ú	<setcomp>Ä   s    ÿz0WorkController.setup_includes.<locals>.<setcomp>)Útupler/   rp   rZ   r   ÚtasksÚset)r:   ÚincludesÚprevÚtask_modulesr,   r^   r-   rG   ¼   s   
ÿzWorkController.setup_includesc                 K   s   |S r]   r,   r`   r,   r,   r-   r9   È   r_   zWorkController.prepare_argsc                 C   s   t jj| d� d S )Nr>   )r	   Úworker_shutdownrR   r^   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)rW   Ústartr   Ú	terminateÚ	ExceptionÚloggerÚcriticalre   r   Ú
SystemExitÚcodeÚKeyboardInterrupt)r:   rs   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)rW   Úsend_all)r:   Úhubr,   r,   r-   rŽ   Û   s   
þz'WorkController.register_with_event_loopc                 C   s   |   | j|¡S r]   )Ú_quick_acquireÚ_process_task©r:   Ú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 r]   )rb   Úcloser›   r^   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   rM   Ú	transportÚ
implementsÚasynchronousr/   Ú
IS_WINDOWSr^   r,   r,   r-   rN   ô   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…   rW   r   r
   r�   r™   Úsignal_safeÚ	_shutdownrƒ   )r:   Úin_sighandlerr…   r,   r,   r-   re   ù   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)rW   r   r   r�   r™   r¥   r¦   )r:   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‡   )rW   r   ÚSHUTDOWN_SOCKET_TIMEOUTre   Újoin)r:   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_modulesrb   Úupdate_strategiesÚreset_rate_limitsr™   ÚrestartrH   )r:   Ú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 r]   )Ú_maybe_reload_moduleru   ©r;   r:   r,   r-   Ú	<genexpr>  s
   € ÿ
ÿz1WorkController._reload_modules.<locals>.<genexpr>r,   )r~   r/   r3   r�   )r:   r±   r;   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/   r3   Úimport_from_cwdr   )r:   Ú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   r1   r2   r   Útotal_countÚosÚgetpidÚstrr/   r¼   ÚroundÚtotal_seconds)r:   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)ÚresourcerH   Ú	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)r:   Úsr,   r,   r-   Úrusage4  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Ä   ÚupdaterW   rb   ré   rH   )r:   rÄ   r,   r,   r-   ÚstatsK  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)r:   r   )rm   rW   Úhuman_stater^   r,   r,   r-   Ú__repr__U  s   þþzWorkController.__repr__c                 C   s   | j S )z#``str(worker) == worker.hostname``.)r0   r^   r,   r,   r-   Ú__str__\  s   zWorkController.__str__c                 C   s   t S r]   r   r^   r,   r,   r-   r   `  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/   ÚeitherrI   Úlogfiler   Útask_eventsrU   Ú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  )r:   r   rI   r  r  r™   r  r  r	  r  r  r  r
  ÚOr  r  r  r  rU   Ústate_dbrý   rþ   Úscheduler_clsr  r  r  r  r  r  Ú_kwr  r,   r,   r-   r6   d  sN   ÿ
ÿÿÿÿÿÿÿzWorkController.setup_defaults)NN)NNNNNNr]   )FN)F)T)NFN)Nrð   NNNNNNNNNNNNNNNNNNNNNNNNNN)-r&   r'   r(   r)   r/   ra   rW   r™   Ú	semaphorer…   r   r.   r<   r8   rV   r5   r7   rB   rJ   rC   rD   rF   rG   r9   rƒ   r†   rŽ   r—   r”   r�   rN   re   r‡   r¦   r²   r­   r³   rÄ   ré   rë   rî   rï   Úpropertyr   r6   r,   r,   r,   r-   r   C   s‚    

ÿ(










ìr   )4r)   Ú
__future__r   r   r¿   r¶   r   Úbilliardr   Úkombu.utils.compatr   Úceleryr   r   rS   r	   Úcelery.bootstepsr
   r   Úcelery.exceptionsr   r   r   Úcelery.fiver   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¨   rk   ro   Úobjectr   r,   r,   r,   r-   Ú<module>   s@   ÿ