o
    wvXjU  ã                   @   sø  d Z ddlmZmZ ddl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mZ dd
lmZ ddlmZm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$m%Z% ddl&m'Z'm(Z(m)Z)m*Z* ddl+m,Z, ddl-m.Z. ddl/m0Z0 ddl1m2Z2 ddl3m4Z4 ddl5m6Z6m7Z7 ddl8m9Z9 ddl:m;Z;m<Z<m=Z= dZ>ej?Z?ej@Z@e?e@hZAe.eBƒZCeCjDeCjEeCjFeCjGeCjHf\ZDZEZIZGZJdZKdZLdZMdZNdZOd ZPd!ZQd"ZRd#ZSd$d%„ ZTe)G d&d'„ d'eUƒƒZVG d(d)„ d)ejWƒZXdS )*z¿Worker Consumer Blueprint.

This module contains the components responsible for consuming messages
from the broker, processing the messages and keeping the broker connections
up and running.
é    )Úabsolute_importÚunicode_literalsN)Údefaultdict)Úsleep)Úrestart_state)ÚRestartFreqExceeded)Ú	DummyLock)ÚContentDisallowedÚDecodeError)Ú_detect_environment)Úbytes_tÚ	safe_repr)ÚTokenBucket)ÚppartialÚpromise)Ú	bootstepsÚsignals)Úbuild_tracer)ÚInvalidTaskErrorÚNotRegistered)Úbuffer_tÚitemsÚpython_2_unicode_compatibleÚvalues)Únoop)Ú
get_logger)Úgethostname)ÚBunch)Útruncate)Úhumanize_secondsÚrate)Úloops)Úmaybe_shutdownÚreserved_requestsÚtask_reserved)ÚConsumerÚEvloopÚ	dump_bodyzMconsumer: Connection to broker lost. Trying to re-establish the connection...z0Trying again {when}... ({retries}/{max_retries})z'consumer: Cannot connect to %s: %s.
%s
zWill retry using next failover.zkReceived and deleted unknown message.  Wrong destination?!?

The full contents of the message body was: %s
aC  Received unregistered task of type %s.
The message has been ignored and discarded.

Did you remember to import the module containing this task?
Or maybe you're using relative imports?

Please see
http://docs.celeryq.org/en/latest/internals/protocol.html
for more information.

The full contents of the message body was:
%s
a  Received invalid task message: %s
The message has been ignored and discarded.

Please ensure your message conforms to the task
message protocol as described here:
http://docs.celeryq.org/en/latest/internals/protocol.html

The full contents of the message body was:
%s
zICan't decode message body: %r [type:%r encoding:%r headers:%s]

body: %s
zTbody: {0}
{{content_type:{1} content_encoding:{2}
  delivery_info:{3} headers={4}}}
c                 C   s@   |du r| j n|}t|tƒrt|ƒ}d tt|ƒdƒt| j ƒ¡S )z+Format message body for debugging purposes.Nz
{0} ({1}b)i   )ÚbodyÚ
isinstancer   r   Úformatr   r   Úlen)Úmr(   © r-   ú\/var/www/html/myproject/venv/lib/python3.10/site-packages/celery/worker/consumer/consumer.pyr'   r   s   
ÿr'   c                   @   sˆ  e Zd ZdZeZdZdZdZdZ	G dd„ de
jƒZedddddddddddfd	d
„Zdd„ Zdd„ Zdd„ Zdd„ ZdTd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d(d)„ Zd*d+„ Zd,d-„ Zd.d/„ Zd0d1„ Z d2d3„ Z!d4d5„ Z"dUd6d7„Z#dUd8d9„Z$d:d;„ Z%d<d=„ Z&d>d?„ Z'		dVd@dA„Z(dBdC„ Z)dDdE„ Z*dFdG„ Z+dHdI„ Z,dJdK„ Z-dLdM„ Z.dNdO„ Z/e0fdPdQ„Z1dRdS„ Z2dS )Wr%   úConsumer blueprint.Néÿÿÿÿc                   @   s$   e Zd ZdZdZg d¢Zdd„ ZdS )zConsumer.Blueprintr/   r%   )	z,celery.worker.consumer.connection:Connectionz$celery.worker.consumer.mingle:Minglez$celery.worker.consumer.events:Eventsz$celery.worker.consumer.gossip:Gossipz"celery.worker.consumer.heart:Heartz&celery.worker.consumer.control:Controlz"celery.worker.consumer.tasks:Tasksz&celery.worker.consumer.consumer:Evloopz"celery.worker.consumer.agent:Agentc                 C   s   |   |d¡ d S )NÚshutdown)Úsend_all)ÚselfÚparentr-   r-   r.   r1   Ÿ   ó   zConsumer.Blueprint.shutdownN)Ú__name__Ú
__module__Ú__qualname__Ú__doc__ÚnameÚdefault_stepsr1   r-   r-   r-   r.   Ú	Blueprint�   s
    r<   Fé   é   c                 K   s~  || _ || _|| _|ptƒ | _t ¡ | _|| _|| _	|  
¡ | _| j  ¡ | _| jj| _| jj| _tddd�| _t tj¡| _d| _|| _tƒ | _| j jj| _|| _|| _|| _ t!dd„ ƒ| _"|  #¡  || _$| j$snt%| jddƒr}|	| _&| j&d u r|| j jj'| _&nd| _&t(| d	ƒsŽ|rŠt)j*nt)j+| _,t-ƒ d
kr˜d | j j_.g | _/g | _0| j1| j j0d | j2d�| _3| j3j4| fi t5|
pµi fi |¤Ž¤Ž d S )Né   r>   )ÚmaxRÚmaxTr   c                   S   s   d S ©Nr-   r-   r-   r-   r.   Ú<lambda>À   s    z#Consumer.__init__.<locals>.<lambda>Úis_greenFÚloopÚgeventÚconsumer)ÚstepsÚon_close)6ÚappÚ
controllerÚinit_callbackr   ÚhostnameÚosÚgetpidÚpidÚpoolÚtimerÚ
StrategiesÚ
strategiesÚconnection_for_readÚconninfoÚconnection_errorsÚchannel_errorsr   Ú_restart_stateÚloggerÚisEnabledForÚloggingÚINFOÚ
_does_infoÚ_limit_orderÚon_task_requestÚsetÚon_task_messageÚconfÚbroker_heartbeat_checkrateÚamqheartbeat_rateÚdisable_rate_limitsÚinitial_prefetch_countÚprefetch_multiplierr   Útask_bucketsÚreset_rate_limitsÚhubÚgetattrÚamqheartbeatÚbroker_heartbeatÚhasattrr!   ÚasynloopÚsynlooprE   r   Úbroker_connection_timeoutÚ_pending_operationsrH   r<   rI   Ú	blueprintÚapplyÚdict)r3   r`   rL   rM   rQ   rJ   rR   rK   rk   rm   Úworker_optionsrf   rg   rh   Úkwargsr-   r-   r.   Ú__init__¢   sP   




€



þ(zConsumer.__init__c                 O   s8   t |g|¢R i |¤Ž}| jr| j |¡S | j |¡ |S rB   )r   rk   Ú	call_soonrs   Úappend)r3   ÚpÚargsrx   r-   r-   r.   rz   Ý   s
   zConsumer.call_soonc              
   C   s`   | j s,| jr.z| j ¡ ƒ  W n ty& } zt d|¡ W Y d }~nd }~ww | jsd S d S d S )NzPending callback raised: %r)rk   rs   ÚpopÚ	ExceptionrZ   Ú	exception©r3   Úexcr-   r-   r.   Úperform_pending_operationsä   s   €ÿ
ýÿz#Consumer.perform_pending_operationsc                 C   s$   t t|dd ƒƒ}|rt|dd�S d S )NÚ
rate_limitr>   )Úcapacity)r    rl   r   )r3   ÚtypeÚlimitr-   r-   r.   Úbucket_for_taskì   s   zConsumer.bucket_for_taskc                    s&   ˆ j  ‡ fdd„tˆ jjƒD ƒ¡ d S )Nc                 3   s"   � | ]\}}|ˆ   |¡fV  qd S rB   )rˆ   )Ú.0ÚnÚt©r3   r-   r.   Ú	<genexpr>ñ   s   € 
ÿz-Consumer.reset_rate_limits.<locals>.<genexpr>)ri   Úupdater   rJ   ÚtasksrŒ   r-   rŒ   r.   rj   ð   s   
ÿzConsumer.reset_rate_limitsr   c                 C   s0   | j j}| jr	|sdS | j j| j | _|  |¡S )a•  Update prefetch count after pool/shrink grow operations.

        Index must be the change in number of processes as a positive
        (increasing) or negative (decreasing) number.

        Note:
            Currently pool grow operations will end up with an offset
            of +1 if the initial size of the pool was 0 (e.g.
            :option:`--autoscale=1,0 <celery worker --autoscale>`).
        N)rQ   Únum_processesrg   rh   Ú_update_qos_eventually)r3   Úindexr�   r-   r-   r.   Ú_update_prefetch_countõ   s   
ÿ
zConsumer._update_prefetch_countc                 C   s&   |dk r| j jn| j jt|ƒ| j ƒS )Nr   )ÚqosÚdecrement_eventuallyÚincrement_eventuallyÚabsrh   )r3   r’   r-   r-   r.   r‘     s   þzConsumer._update_qos_eventuallyc                 C   s   t |ƒ |  |¡ d S rB   )r$   r`   )r3   Úrequestr-   r-   r.   Ú_limit_move_to_pool  s   zConsumer._limit_move_to_poolc                 C   sˆ   	 z|  ¡ \}}W n
 ty   Y d S w | |¡r|  |¡ q |j ||f¡ | jd d  }| _| |¡}| jj	|| j
|f|d� d S )NTr>   é
   )Úpriority)r~   Ú
IndexErrorÚcan_consumer™   ÚcontentsÚ
appendleftr_   Úexpected_timerR   Ú
call_afterÚ_schedule_bucket_request)r3   Úbucketr˜   ÚtokensÚpriÚholdr-   r-   r.   r¢     s"   þ



þz!Consumer._schedule_bucket_requestc                 C   s   |  ||f¡ |  |¡S rB   )Úaddr¢   ©r3   r˜   r£   r¤   r-   r-   r.   Ú_limit_task)  s   
zConsumer._limit_taskc                 C   s"   | j  ¡  | ||f¡ |  |¡S rB   )r”   r•   r§   r¢   r¨   r-   r-   r.   Ú_limit_post_eta-  s   

zConsumer._limit_post_etac              
   C   s  | j }|jtvr�tƒ  | jr3z| j ¡  W n ty2 } ztd|dd� t	dƒ W Y d }~nd }~ww |  jd7  _z| 
| ¡ W nD | jy… } z7| jjjsP‚ t|tƒr\|jtjkr\‚ tƒ  |jtvr{| jrm|  |¡ n|  |¡ |  ¡  | | ¡ W Y d }~nd }~ww |jtvsd S d S )NzFrequent restarts detected: %rr>   ©Úexc_info)rt   ÚstateÚSTOP_CONDITIONSr"   Úrestart_countrY   Ústepr   Úcritr   ÚstartrW   rJ   rc   Úbroker_connection_retryr)   ÚOSErrorÚerrnoÚEMFILEÚ
connectionÚ#on_connection_error_after_connectedÚ$on_connection_error_before_connectedrI   Úrestart)r3   rt   r‚   r-   r-   r.   r²   2  s:   
€þ



€òõzConsumer.startc                 C   s   t t| j ¡ |dƒ d S )NzTrying to reconnect...)ÚerrorÚCONNECTION_ERRORrV   Úas_urir�   r-   r-   r.   r¹   O  s   ÿz-Consumer.on_connection_error_before_connectedc                 C   s2   t tdd� z| j ¡  W d S  ty   Y d S w )NTr«   )ÚwarnÚCONNECTION_RETRYr·   Úcollectr   r�   r-   r-   r.   r¸   S  s   ÿz,Consumer.on_connection_error_after_connectedc                 C   s   | j j| d|fdd� d S )NÚregister_with_event_loopzHub.register)r}   Údescription)rt   r2   )r3   rk   r-   r-   r.   rÁ   Z  s   
þz!Consumer.register_with_event_loopc                 C   ó   | j  | ¡ d S rB   )rt   r1   rŒ   r-   r-   r.   r1   `  r5   zConsumer.shutdownc                 C   rÃ   rB   )rt   ÚstoprŒ   r-   r-   r.   rÄ   c  r5   zConsumer.stopc                 C   s"   | j d }| _ |r|| ƒ d S d S rB   )rL   )r3   Úcallbackr-   r-   r.   Úon_readyf  s   ÿzConsumer.on_readyc              	   C   s(   | | j | j| j| j| j| j| jj| jf	S rB   )	r·   Útask_consumerrt   rk   r”   rm   rJ   Úclockre   rŒ   r-   r-   r.   Ú	loop_argsk  s   

þzConsumer.loop_argsc              	   C   s4   t t||j|jt|jƒt||jƒdd� | ¡  dS )a.  Callback called if an error occurs while decoding a message.

        Simply logs the error and acknowledges the message so it
        doesn't enter a loop.

        Arguments:
            message (kombu.Message): The message received.
            exc (Exception): The exception being handled.
        r>   r«   N)	r±   ÚMESSAGE_DECODE_ERRORÚcontent_typeÚcontent_encodingr   Úheadersr'   r(   Úack)r3   Úmessager‚   r-   r-   r.   Úon_decode_errorp  s   

ýzConsumer.on_decode_errorc                 C   sr   | j r| j jr| j j ¡  | jr| j ¡  t| jƒD ]}|r"| ¡  qt ¡  | jr5| jj	r7| j 	¡  d S d S d S rB   )
rK   Ú	semaphoreÚclearrR   r   ri   Úclear_pendingr#   rQ   Úflush)r3   r£   r-   r-   r.   rI   €  s   
€ÿzConsumer.on_closec                 C   s*   | j | jd�}| jr|j |j| j¡ |S )z´Establish the broker connection used for consuming tasks.

        Retries establishing the connection if the
        :setting:`broker_connection_retry` setting is enabled
        ©Ú	heartbeat)rU   rm   rk   Ú	transportrÁ   r·   )r3   Úconnr-   r-   r.   Úconnect�  s   zConsumer.connectc                 C   ó   |   | jj|d�¡S ©NrÕ   )Úensure_connectedrJ   rU   ©r3   rÖ   r-   r-   r.   rU   š  ó   ÿzConsumer.connection_for_readc                 C   rÚ   rÛ   )rÜ   rJ   Úconnection_for_writerÝ   r-   r-   r.   rß   ž  rÞ   zConsumer.connection_for_writec                    sB   t f‡ ‡fdd„	}ˆjjjsˆ  ¡  ˆ S ˆ j|ˆjjjtd�‰ ˆ S )Nc                    sT   t ˆ dd ƒr|dkrt}|jt|ddƒt|d ƒˆjjjd�}tt	ˆ  
¡ | |ƒ d S )NÚaltr   Úinú r=   )ÚwhenÚretriesÚmax_retries)rl   ÚCONNECTION_FAILOVERr*   r   ÚintrJ   rc   Úbroker_connection_max_retriesr»   r¼   r½   )r‚   ÚintervalÚ	next_step©rØ   r3   r-   r.   Ú_error_handler¥  s   

ýz1Consumer.ensure_connected.<locals>._error_handler)rÅ   )ÚCONNECTION_RETRY_STEPrJ   rc   r³   rÙ   Úensure_connectionrè   r"   )r3   rØ   rì   r-   rë   r.   rÜ   ¢  s   

þzConsumer.ensure_connectedc                 C   s   | j r
| j  ¡  d S d S rB   )Úevent_dispatcherrÔ   rŒ   r-   r-   r.   Ú_flush_events»  s   ÿzConsumer._flush_eventsc                 C   s   | j r| j j | j¡ d S d S rB   )rk   Ú_readyr§   rð   rŒ   r-   r-   r.   Úon_send_event_buffered¿  s   ÿzConsumer.on_send_event_bufferedc           	      K   sŠ   | j }| jjj}||v r|| }n|d u r|n|}|d u rdn|}|j|f|||dœ|¤Ž}| |¡sC| |¡ | ¡  td|ƒ d S d S )NÚdirect)ÚexchangeÚexchange_typeÚrouting_keyzStarted consuming from %s)	rÇ   rJ   ÚamqpÚqueuesÚ
select_addÚconsuming_fromÚ	add_queueÚconsumeÚinfo)	r3   Úqueuerô   rõ   rö   ÚoptionsÚcsetrø   Úqr-   r-   r.   Úadd_task_queueÃ  s(   

ÿýý

ýzConsumer.add_task_queuec                 C   s*   t d|ƒ | jjj |¡ | j |¡ d S )NzCanceling queue %s)rý   rJ   r÷   rø   ÚdeselectrÇ   Úcancel_by_queue)r3   rþ   r-   r-   r.   Úcancel_task_queueÙ  s   
zConsumer.cancel_task_queuec                 C   s    t |ƒ |  |¡ | j ¡  dS )zAMethod called by the timer to apply a task with an ETA/countdown.N)r$   r`   r”   r•   )r3   Útaskr-   r-   r.   Úapply_eta_taskÞ  s   
zConsumer.apply_eta_taskc                 C   s0   t  t||ƒt|jƒt|jƒt|jƒt|jƒ¡S rB   )ÚMESSAGE_REPORTr*   r'   r   rË   rÌ   Údelivery_inforÍ   ©r3   r(   rÏ   r-   r-   r.   Ú_message_reportä  s   üzConsumer._message_reportc                 C   s6   t t|  ||¡ƒ | t| j¡ tjj| |d d� d S )N©ÚsenderrÏ   r‚   )	r¾   ÚUNKNOWN_FORMATr  Úreject_log_errorrZ   rW   r   Útask_rejectedÚsendr
  r-   r-   r.   Úon_unknown_messageë  s   zConsumer.on_unknown_messagec           	      C   sî   t t|t||ƒdd� z|jd |jd }}|j d¡}W n ty5   |j}|d |d }}d }Y nw t|d ||j d¡|j d¡d d�}| 	t
| j¡ | jjj|t|ƒ|d	� | jrj| jjd
|d |¡d� tjj| ||||d� d S )NTr«   Úidr  Úroot_idÚcorrelation_idÚreply_to)r:   Úchordr  r  r  Úerrbacks)r˜   ztask-failedzNotRegistered({0!r}))Úuuidr€   )r  rÏ   r‚   r:   r  )r»   ÚUNKNOWN_TASK_ERRORr'   rÍ   ÚgetÚKeyErrorÚpayloadr   Ú
propertiesr  rZ   rW   rJ   ÚbackendÚmark_as_failurer   rï   r  r*   r   Útask_unknown)	r3   r(   rÏ   r‚   Úid_r:   r  r  r˜   r-   r-   r.   Úon_unknown_taskð  s6   ý

ü
ÿþ

ÿzConsumer.on_unknown_taskc                 C   s:   t t|t||ƒdd� | t| j¡ tjj| ||d� d S )NTr«   r  )	r»   ÚINVALID_TASK_ERRORr'   r  rZ   rW   r   r  r  )r3   r(   rÏ   r‚   r-   r-   r.   Úon_invalid_task  s   zConsumer.on_invalid_taskc                 C   sN   | j j}t| j jƒD ]\}}| | j | ¡| j|< t|||| j| j d�|_q
d S )N)rJ   )	rJ   Úloaderr   r�   Ústart_strategyrT   r   rM   Ú	__trace__)r3   r&  r:   r  r-   r-   r.   Úupdate_strategies  s   
ÿþzConsumer.update_strategiesc                    sB   ˆj ‰ˆj‰ˆj‰ˆj‰ˆj‰ˆj‰ ‡ ‡‡‡‡‡‡‡fdd„}|S )Nc                    s†  d }z| j d }W nS ty   ˆd | ƒ Y S  ty\   z|  ¡ }W n ty= } zˆ | |¡W  Y d }~ Y S d }~ww z	|d |}}W n ttfyY   ˆ|| ƒ Y  Y S w Y nw zˆ| }W n ty{ } zˆd | |ƒW  Y d }~S d }~ww z|| |ˆˆ | jfƒˆˆ | jfƒˆƒ W d S  tt	fy« } zˆ|| |ƒW  Y d }~S d }~w t
yÂ } zˆ | |¡W  Y d }~S d }~ww )Nr  )rÍ   Ú	TypeErrorr  Údecoder   rÐ   Úack_log_errorr  r   r	   r
   )rÏ   r  Útype_r‚   Ústrategy©rz   Ú	callbacksr%  r  r#  r   r3   rT   r-   r.   Úon_task_received   sN   €ÿÿÿú	€ÿ
ü€€ÿz6Consumer.create_task_handler.<locals>.on_task_received)rT   r  r#  r%  rb   rz   )r3   r   r1  r-   r/  r.   Úcreate_task_handler  s   "zConsumer.create_task_handlerc                 C   s   dj | | j ¡ d�S )z``repr(self)``.z%<Consumer: {self.hostname} ({state})>)r3   r­   )r*   rt   Úhuman_staterŒ   r-   r-   r.   Ú__repr__D  s   
ÿzConsumer.__repr__)r   rB   )NNN)3r6   r7   r8   r9   rv   rS   rL   rQ   rR   r¯   r   r<   r   ry   rz   rƒ   rˆ   rj   r“   r‘   r™   r¢   r©   rª   r²   r¹   r¸   rÁ   r1   rÄ   rÆ   rÉ   rÐ   rI   rÙ   rU   rß   rÜ   rð   rò   r  r  r  r  r  r#  r%  r)  r   r2  r4  r-   r-   r-   r.   r%   |   sh    
û;



ÿ,r%   c                   @   s(   e Zd ZdZdZdZdd„ Zdd„ ZdS )	r&   zHEvent loop service.

    Note:
        This is always started last.
    z
event loopTc                 C   s   |   |¡ |j| ¡ Ž  d S rB   )Ú	patch_allrE   rÉ   ©r3   Úcr-   r-   r.   r²   U  s   
zEvloop.startc                 C   s   t ƒ |j_d S rB   )r   r”   Ú_mutexr6  r-   r-   r.   r5  Y  s   zEvloop.patch_allN)r6   r7   r8   r9   ÚlabelÚlastr²   r5  r-   r-   r-   r.   r&   K  s    r&   )Yr9   Ú
__future__r   r   rµ   r\   rN   Úcollectionsr   Útimer   Úbilliard.commonr   Úbilliard.exceptionsr   Úkombu.asynchronous.semaphorer   Úkombu.exceptionsr	   r
   Úkombu.utils.compatr   Úkombu.utils.encodingr   r   Úkombu.utils.limitsr   Úviner   r   Úceleryr   r   Úcelery.app.tracer   Úcelery.exceptionsr   r   Úcelery.fiver   r   r   r   Úcelery.utils.functionalr   Úcelery.utils.logr   Úcelery.utils.nodenamesr   Úcelery.utils.objectsr   Úcelery.utils.textr   Úcelery.utils.timer   r    Úcelery.workerr!   Úcelery.worker.stater"   r#   r$   Ú__all__ÚCLOSEÚ	TERMINATEr®   r6   rZ   Údebugrý   Úwarningr»   Úcriticalr¾   r±   r¿   rí   r¼   ræ   r  r  r$  rÊ   r  r'   Úobjectr%   ÚStartStopStepr&   r-   r-   r-   r.   Ú<module>   sf   ÿ
   Q