o
    wvXjÿ  ã                   @   s(  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	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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.é    )Úabsolute_importÚunicode_literalsN)ÚHub)Úget_event_loopÚset_event_loop)Ú	DummyLockÚLaxBoundedSemaphore)ÚTimer)Ú	bootsteps)Ú_set_task_join_will_block)ÚImproperlyConfigured)Ústring_t)Ú
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%   úU/var/www/html/myproject/venv/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    5   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!   8   ó   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t| ƒj|fi |¤Ž d S ©N)ÚhubÚsuperr   Ú__init__©r#   r$   Úkwargs©Ú	__class__r%   r&   r6   A   s   zHub.__init__c                 C   s   |j S r3   )r   r"   r%   r%   r&   Ú
include_ifE   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   r4   ÚgetattrÚ	_conninfor   Ú_Hubr   Ú_patch_thread_primitives)r#   r$   Úrequired_hubr%   r%   r&   r'   H   s   
ÿ
z
Hub.createc                 C   s   d S r3   r%   r"   r%   r%   r&   ÚstartQ   s   z	Hub.startc                 C   ó   |j  ¡  d S r3   ©r4   Úcloser"   r%   r%   r&   ÚstopT   ó   zHub.stopc                 C   rC   r3   rD   r"   r%   r%   r&   Ú	terminateW   rG   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ÚbilliardrI   ÚImportErrorÚLock)r#   r$   rI   r%   r%   r&   r@   Z   s   ÿ
zHub._patch_thread_primitives)r/   r0   r1   r2   r	   Úrequiresr6   r;   r'   rB   rF   rH   r@   Ú__classcell__r%   r%   r9   r&   r   <   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
t| ƒj|fi |¤Ž d S )Nú,r   )rI   Úmax_concurrencyÚconcurrencyÚmin_concurrencyÚoptimizationÚ
isinstancer   Ú	partitionÚintÚ	autoscaler5   r   r6   )r#   r$   rZ   r8   Úmax_cÚ_Úmin_cr9   r%   r&   r6   v   s   
zPool.__init__c                 C   ó   |j r
|j  ¡  d S d S r3   )rI   rE   r"   r%   r%   r&   rE   ƒ   ó   ÿz
Pool.closec                 C   r^   r3   )rI   rH   r"   r%   r%   r&   rH   ‡   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_strategyrJ   )#rJ   ÚconfÚworker_poolÚGREEN_POOLSÚwarningsÚwarnÚUserWarningÚW_POOL_SETTINGr   r   rU   Ú_process_taskÚprocess_taskr   rl   ÚacquireÚ_quick_acquireÚreleaseÚ_quick_releaseÚpool_putlocksr   Úuses_semaphoreÚ_process_task_semÚpool_restartsr   ÚhostnameÚmax_tasks_per_childrc   Ú
time_limitÚsoft_time_limitÚworker_lost_waitrV   rI   r   Útask_join_will_block)r#   r$   rl   ri   ÚthreadedÚprocsrj   rI   r%   r%   r&   r'   ‹   sD   


ñ
zPool.createc                 C   s   d|j r	|j jiS diS )NrI   zN/A)rI   Úinfor"   r%   r%   r&   r‡   ¯   s   z	Pool.infoc                 C   s   |j  |¡ d S r3   )rI   Úregister_with_event_loop)r#   r$   r4   r%   r%   r&   rˆ   ²   r.   zPool.register_with_event_loopr3   )r/   r0   r1   r2   r   rP   r6   rE   rH   r'   r‡   rˆ   rQ   r%   r%   r9   r&   r   f   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                    s2   | | _ |_d |_tt| ƒj|fd|i|¤Ž d S )NÚbeat)Úenabledr‰   r5   r   r6   )r#   r$   r‰   r8   r9   r%   r&   r6   ¿   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   r0   Úendswithr   ÚERR_B_GREENrJ   rŒ   Ú	schedulerr‰   )r#   r$   r‹   Úbr%   r%   r&   r'   Ä   s   þzBeat.create)F)	r/   r0   r1   r2   ÚlabelÚconditionalr6   r'   rQ   r%   r%   r9   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t| ƒj|fi |¤Ž d S r3   )ÚstatedbrŠ   Ú_persistencer5   r   r6   r7   r9   r%   r&   r6   Ñ   s   zStateDB.__init__c                 C   s,   |j  |j |j|jj¡|_t |jj¡ d S r3   )	ÚstateÚ
Persistentr•   rJ   rK   r–   ÚatexitÚregisterÚsaver"   r%   r%   r&   r'   Ö   s   zStateDB.create)r/   r0   r1   r2   r6   r'   rQ   r%   r%   r9   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_countrI   r   rJ   Ú
controllerr4   Úworker_optionsÚdisable_rate_limitsÚprefetch_multiplier)rS   Úmaxr£   rT   r   Úconsumer_clsrv   r   r�   Úready_callbackrI   r   rJ   r4   Úoptionsr¢   Úconsumer)r#   r$   Úprefetch_countÚcr%   r%   r&   r'   à   s&   ózConsumer.createN)r/   r0   r1   r2   Úlastr'   r%   r%   r%   r&   r   Û   s    r   )(r2   Ú
__future__r   r   r™   rq   Úkombu.asynchronousr   r?   r   r   Úkombu.asynchronous.semaphorer   r   Úkombu.asynchronous.timerr	   r   Úceleryr
   Úcelery._stater   Úcelery.exceptionsr   Úcelery.fiver   Úcelery.platformsr   Úcelery.utils.logr   r)   Ú__all__rp   r�   rt   ÚStepÚStartStopStepr   r   r   r   r%   r%   r%   r&   Ú<module>   s0   *P