o
    ›¨ÊhI  ã                   @   s  d Z ddlZddlZddl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 dd
lmZ ddlmZ dZddhZdZdZG dd„ dejƒZG dd„ dejƒZG dd„ dejƒZG dd„ dejƒZ G dd„ dejƒZ!G dd„ dejƒZ"dS )zWorker-level Bootsteps.é    N)ÚHub)Úget_event_loopÚset_event_loop)Ú	DummyLockÚLaxBoundedSemaphore)ÚTimer)Ú	bootsteps)Ú_set_task_join_will_block)ÚImproperlyConfigured)Ú
IS_WINDOWS)Úworker_logger)r   r   ÚPoolÚBeatÚStateDBÚConsumerÚeventletÚgeventzO-B option doesn't work with eventlet/gevent pools: use standalone beat instead.z©
The worker_pool setting shouldn't be used to select the eventlet/gevent
pools, instead you *must use the -P* argument so that patches are applied
as early as possible.
c                   @   s(   e Zd ZdZdd„ Zdd„ Zdd„ ZdS )	r   zTimer bootstep.c                 C   sF   |j rtdd�|_d S |js|jj|_| j|j|j| j| j	d�|_d S )Ng      $@)Úmax_interval)r   Úon_errorÚon_tick)
Úuse_eventloopÚ_TimerÚtimerÚ	timer_clsÚpool_clsr   ÚinstantiateÚtimer_precisionÚon_timer_errorÚon_timer_tick©ÚselfÚw© r"   úJ/var/www/html/env/lib/python3.10/site-packages/celery/worker/components.pyÚcreate#   s   
ýzTimer.createc                 C   s   t jd|dd� d S )NzTimer error: %rT)Úexc_info)ÚloggerÚerror)r    Úexcr"   r"   r#   r   1   s   zTimer.on_timer_errorc                 C   s   t  d|¡ d S )Nz Timer wake-up! Next ETA %s secs.)r&   Údebug)r    Údelayr"   r"   r#   r   4   ó   zTimer.on_timer_tickN)Ú__name__Ú
__module__Ú__qualname__Ú__doc__r$   r   r   r"   r"   r"   r#   r       s
    r   c                       sV   e Zd ZdZefZ‡ fdd„Zdd„ Zdd„ Zdd	„ Z	d
d„ Z
dd„ Zdd„ Z‡  ZS )r   zWorker starts the event loop.c                    s   d |_ tƒ j|fi |¤Ž d S ©N)ÚhubÚsuperÚ__init__©r    r!   Úkwargs©Ú	__class__r"   r#   r3   =   s   zHub.__init__c                 C   s   |j S r0   )r   r   r"   r"   r#   Ú
include_ifA   s   zHub.include_ifc                 C   sF   t ƒ |_|jd u rt|jdd ƒ}t|r|nt|jƒƒ|_|  |¡ | S )NÚrequires_hub)r   r1   ÚgetattrÚ	_conninfor   Ú_Hubr   Ú_patch_thread_primitives)r    r!   Úrequired_hubr"   r"   r#   r$   D   s   
ÿ
z
Hub.createc                 C   s   d S r0   r"   r   r"   r"   r#   ÚstartM   s   z	Hub.startc                 C   ó   |j  ¡  d S r0   ©r1   Úcloser   r"   r"   r#   ÚstopP   ó   zHub.stopc                 C   r@   r0   rA   r   r"   r"   r#   Ú	terminateS   rD   zHub.terminatec                 C   s<   t ƒ |jj_zddlm} W n
 ty   Y d S w t |_d S )Nr   )Úpool)r   ÚappÚclockÚmutexÚbilliardrF   ÚImportErrorÚLock)r    r!   rF   r"   r"   r#   r=   V   s   ÿ
zHub._patch_thread_primitives)r,   r-   r.   r/   r   Úrequiresr3   r8   r$   r?   rC   rE   r=   Ú__classcell__r"   r"   r6   r#   r   8   s    	r   c                       sP   e Zd ZdZefZd‡ fdd„	Zdd„ Zdd„ Zd	d
„ Z	dd„ Z
dd„ Z‡  ZS )r   a
  Bootstep managing the worker pool.

    Describes how to initialize the worker pool, and starts and stops
    the pool during worker start-up/shutdown.

    Adds attributes:

        * autoscale
        * pool
        * max_concurrency
        * min_concurrency
    Nc                    s€   d |_ d |_|j|_|j| _t|tƒr'| d¡\}}}t|ƒ|r$t|ƒp%dg}||_	|j	r4|j	\|_|_t
ƒ j|fi |¤Ž d S )Nú,r   )rF   Úmax_concurrencyÚconcurrencyÚmin_concurrencyÚoptimizationÚ
isinstanceÚstrÚ	partitionÚintÚ	autoscaler2   r3   )r    r!   rX   r5   Úmax_cÚ_Úmin_cr6   r"   r#   r3   r   s   
zPool.__init__c                 C   ó   |j r
|j  ¡  d S d S r0   )rF   rB   r   r"   r"   r#   rB      ó   ÿz
Pool.closec                 C   r\   r0   )rF   rE   r   r"   r"   r#   rE   ƒ   r]   zPool.terminatec                 C   sâ   d }d }|j jjtv rt ttƒ¡ |j pt	}|j
}|j|_|s?t|ƒ }|_|jj|_|jj|_d}|jr?|jjr?|j|_|j}| j|j|j
|j |jf|j|j|j|j|joY||j|||d|| j|j d� }|_ t!|j"ƒ |S )Néd   T)ÚinitargsÚmaxtasksperchildÚmax_memory_per_childÚtimeoutÚsoft_timeoutÚputlocksÚlost_worker_timeoutÚthreadsÚmax_restartsÚallow_restartÚforking_enableÚ	semaphoreÚsched_strategyrG   )#rG   ÚconfÚworker_poolÚGREEN_POOLSÚwarningsÚwarnÚUserWarningÚW_POOL_SETTINGr   r   rR   Ú_process_taskÚprocess_taskr   rj   ÚacquireÚ_quick_acquireÚreleaseÚ_quick_releaseÚpool_putlocksr   Úuses_semaphoreÚ_process_task_semÚpool_restartsr   ÚhostnameÚmax_tasks_per_childra   Ú
time_limitÚsoft_time_limitÚworker_lost_waitrS   rF   r	   Útask_join_will_block)r    r!   rj   rg   ÚthreadedÚprocsrh   rF   r"   r"   r#   r$   ‡   sD   


ñ
zPool.createc                 C   s   d|j r	|j jiS diS )NrF   zN/A)rF   Úinfor   r"   r"   r#   r…   «   s   z	Pool.infoc                 C   s   |j  |¡ d S r0   )rF   Úregister_with_event_loop)r    r!   r1   r"   r"   r#   r†   ®   r+   zPool.register_with_event_loopr0   )r,   r-   r.   r/   r   rM   r3   rB   rE   r$   r…   r†   rN   r"   r"   r6   r#   r   b   s    $r   c                       s2   e Zd ZdZd ZdZd‡ fdd„	Zdd„ Z‡  ZS )	r   zWStep used to embed a beat process.

    Enabled when the ``beat`` argument is set.
    TFc                    s.   | | _ |_d |_tƒ j|fd|i|¤Ž d S )NÚbeat)Úenabledr‡   r2   r3   )r    r!   r‡   r5   r6   r"   r#   r3   »   s   zBeat.__init__c                 C   s@   ddl m} |jj d¡rttƒ‚||j|j|j	d� }|_
|S )Nr   )ÚEmbeddedService)r   r   )Úschedule_filenameÚscheduler_cls)Úcelery.beatr‰   r   r-   Úendswithr
   ÚERR_B_GREENrG   rŠ   Ú	schedulerr‡   )r    r!   r‰   Úbr"   r"   r#   r$   À   s   þzBeat.create)F)	r,   r-   r.   r/   ÚlabelÚconditionalr3   r$   rN   r"   r"   r6   r#   r   ²   s    r   c                       s(   e Zd ZdZ‡ fdd„Zdd„ Z‡  ZS )r   z:Bootstep that sets up between-restart state database file.c                    s&   |j | _d |_tƒ j|fi |¤Ž d S r0   )Ústatedbrˆ   Ú_persistencer2   r3   r4   r6   r"   r#   r3   Í   s   zStateDB.__init__c                 C   s,   |j  |j |j|jj¡|_t |jj¡ d S r0   )	ÚstateÚ
Persistentr“   rG   rH   r”   ÚatexitÚregisterÚsaver   r"   r"   r#   r$   Ò   s   zStateDB.create)r,   r-   r.   r/   r3   r$   rN   r"   r"   r6   r#   r   Ê   s    r   c                   @   s   e Zd ZdZdZdd„ ZdS )r   z)Bootstep starting the Consumer blueprint.Tc                 C   sn   |j rt|j dƒ|j }n|j|j }| j|j|j|j|j|j	||j
|j|j||j|j|j|jd� }|_|S )Né   )r}   Útask_eventsÚinit_callbackÚinitial_prefetch_countrF   r   rG   Ú
controllerr1   Úworker_optionsÚdisable_rate_limitsÚprefetch_multiplier)rP   Úmaxr¡   rQ   r   Úconsumer_clsrt   r}   r›   Úready_callbackrF   r   rG   r1   Úoptionsr    Úconsumer)r    r!   Úprefetch_countÚcr"   r"   r#   r$   Ü   s&   ózConsumer.createN)r,   r-   r.   r/   Úlastr$   r"   r"   r"   r#   r   ×   s    r   )#r/   r—   ro   Úkombu.asynchronousr   r<   r   r   Úkombu.asynchronous.semaphorer   r   Úkombu.asynchronous.timerr   r   Úceleryr   Úcelery._stater	   Úcelery.exceptionsr
   Úcelery.platformsr   Úcelery.utils.logr   r&   Ú__all__rn   rŽ   rr   ÚStepÚStartStopStepr   r   r   r   r"   r"   r"   r#   Ú<module>   s,    *P