o
    ›¨Êhñ  ã                   @   sÆ   d Z ddlZddl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eeƒZejejejZZZeej dd¡ƒZG dd„ de	jƒZG dd„ deƒZdS )zûPool Autoscaling.

This module implements the internal thread responsible
for growing and shrinking the pool according to the
current autoscale settings.

The autoscale thread is only enabled if
the :option:`celery worker --autoscale` option is used.
é    N)Ú	monotonicÚsleep)Ú	DummyLock)Ú	bootsteps)Ú
get_logger)ÚbgThreadé   )Ústate)ÚPool)Ú
AutoscalerÚWorkerComponentÚAUTOSCALE_KEEPALIVEé   c                   @   s>   e Zd ZdZdZdZefZdd„ Zdd„ Z	dd	„ Z
d
d„ ZdS )r   z?Bootstep that starts the autoscaler thread/timer in the worker.r   Tc                 K   s   |j | _d |_d S ©N)Ú	autoscaleÚenabledÚ
autoscaler)ÚselfÚwÚkwargs© r   úI/var/www/html/env/lib/python3.10/site-packages/celery/worker/autoscale.pyÚ__init__&   ó   
zWorkerComponent.__init__c                 C   s>   | j |j|j|j|j||jrtƒ nd d� }|_|js|S d S )N)ÚworkerÚmutex)ÚinstantiateÚautoscaler_clsÚpoolÚmax_concurrencyÚmin_concurrencyÚuse_eventloopr   r   )r   r   Úscalerr   r   r   Úcreate*   s   ýzWorkerComponent.createc                 C   s*   |j j |jj¡ | |jj|jj¡ d S r   )ÚconsumerÚon_task_messageÚaddr   Úmaybe_scaleÚcall_repeatedlyÚ	keepalive)r   r   Úhubr   r   r   Úregister_with_event_loop2   s   ÿz(WorkerComponent.register_with_event_loopc                 C   s   d|j  ¡ iS )zReturn `Autoscaler` info.r   )r   Úinfo)r   r   r   r   r   r,   8   s   zWorkerComponent.infoN)Ú__name__Ú
__module__Ú__qualname__Ú__doc__ÚlabelÚconditionalr
   Úrequiresr   r#   r+   r,   r   r   r   r   r      s    r   c                       s˜   e Zd ZdZddedf‡ fdd„	Zdd„ Zddd	„Zdd
d„Zddd„Z	dd„ Z
dd„ Zdd„ Zdd„ Zdd„ Zdd„ Zedd„ ƒZedd„ ƒZ‡  ZS ) r   z,Background thread to autoscale pool workers.r   Nc                    sN   t ƒ  ¡  || _|pt ¡ | _|| _|| _|| _d | _	|| _
| js%J dƒ‚d S )Nzcannot scale down too fast.)Úsuperr   r   Ú	threadingÚLockr   r   r    r)   Ú_last_scale_upr   )r   r   r   r    r   r)   r   ©Ú	__class__r   r   r   @   s   
zAutoscaler.__init__c                 C   s:   | j � |  ¡  W d   ƒ n1 sw   Y  tdƒ d S )Ng      ð?)r   r'   r   ©r   r   r   r   ÚbodyN   s   
ÿzAutoscaler.bodyc                 C   sZ   | j }t| j| jƒ}||kr|  || ¡ dS t| j| jƒ}||k r+|  || ¡ dS d S )NT)Ú	processesÚminÚqtyr   Úscale_upÚmaxr    Ú
scale_down)r   ÚreqÚprocsÚcurr   r   r   Ú_maybe_scaleS   s   þzAutoscaler._maybe_scalec                 C   s   |   |¡r| j ¡  d S d S r   )rE   r   Úmaintain_pool)r   rB   r   r   r   r'   ^   s   
ÿzAutoscaler.maybe_scalec                 C   s�   | j �; |d ur|| jk r|  | j| ¡ |  |¡ || _|d ur1|| jkr.|  || j ¡ || _| j| jfW  d   ƒ S 1 sAw   Y  d S r   )r   r<   Ú_shrinkÚ_update_consumer_prefetch_countr   Ú_growr    )r   r@   r=   r   r   r   Úupdateb   s   



$özAutoscaler.updatec                 C   s   t ƒ | _|  |¡S r   )r   r7   rI   ©r   Únr   r   r   r?   o   r   zAutoscaler.scale_upc                 C   s*   | j rtƒ | j  | jkr|  |¡S d S d S r   )r7   r   r)   rG   rK   r   r   r   rA   s   s
   
þzAutoscaler.scale_downc                 C   s   t d|ƒ | j |¡ d S )NzScaling up %s processes.)r,   r   ÚgrowrK   r   r   r   rI   x   s   
zAutoscaler._growc              
   C   sl   t d|ƒ z	| j |¡ W d S  ty   tdƒ Y d S  ty5 } ztd|dd� W Y d }~d S d }~ww )NzScaling down %s processes.z0Autoscaler won't scale down: all processes busy.zAutoscaler: scale_down: %rT)Úexc_info)r,   r   ÚshrinkÚ
ValueErrorÚdebugÚ	ExceptionÚerror)r   rL   Úexcr   r   r   rG   |   s   
€ÿzAutoscaler._shrinkc                 C   s$   || j  }|r| jj |¡ d S d S r   )r   r   r$   Ú_update_prefetch_count)r   Únew_maxÚdiffr   r   r   rH   …   s   
ÿÿz*Autoscaler._update_consumer_prefetch_countc                 C   s   | j | j| j| jdœS )N)r@   r=   Úcurrentr>   )r   r    r<   r>   r:   r   r   r   r,   Œ   s
   üzAutoscaler.infoc                 C   s
   t tjƒS r   )Úlenr	   Úreserved_requestsr:   r   r   r   r>   ”   s   
zAutoscaler.qtyc                 C   s   | j jS r   )r   Únum_processesr:   r   r   r   r<   ˜   s   zAutoscaler.processesr   )NN)r-   r.   r/   r0   r   r   r;   rE   r'   rJ   r?   rA   rI   rG   rH   r,   Úpropertyr>   r<   Ú__classcell__r   r   r8   r   r   =   s&    þ


	
r   )r0   Úosr5   Útimer   r   Úkombu.asynchronous.semaphorer   Úceleryr   Úcelery.utils.logr   Úcelery.utils.threadsr   Ú r	   Ú
componentsr
   Ú__all__r-   ÚloggerrQ   r,   rS   ÚfloatÚenvironÚgetr   ÚStartStopStepr   r   r   r   r   r   Ú<module>   s     	