U
    dY                     @   s  d Z ddl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 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"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/m0Z0 ddl1m2Z2 ddl3m4Z4m5Z5m6Z6m7Z7 dZ8ej9Z9ej:Z:e9e:hZ;e'e<Z=e=j>e=j?e=j@e=jAe=jBf\Z>Z?ZCZAZDdZEdZFdZGdZHdZIdZJdZKd ZLd!ZMd"ZNd#ZOd$d% ZPG d&d' d'ZQG d(d) d)ejRZSdS )*zWorker 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.
    N)defaultdict)sleep)restart_state)RestartFreqExceeded)	DummyLock)ContentDisallowedDecodeError)_detect_environment)	safe_repr)TokenBucket)ppartialpromise)	bootstepssignals)build_tracer)CPendingDeprecationWarningInvalidTaskErrorNotRegistered)noop)
get_logger)gethostname)Bunch)truncate)humanize_secondsrate)loops)active_requests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
a  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

Thw full contents of the message headers:
%s

The delivery info for this task is:
%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}}}
z}Task %s cannot be acknowledged after a connection loss since late acknowledgement is enabled for it.
Terminating it instead.
a  
In Celery 5.1 we introduced an optional breaking change which
on connection loss cancels all currently executed tasks with late acknowledgement enabled.
These tasks cannot be acknowledged as the connection is gone, and the tasks are automatically redelivered back to the queue.
You can enable this behavior using the worker_cancel_long_running_tasks_on_connection_loss setting.
In Celery 5.1 it is set to False by default. The setting will be set to True by default in Celery 6.0.
c                 C   s.   |dkr| j n|}dtt|dt| j S )z+Format message body for debugging purposes.Nz{} ({}b)i   )bodyformatr   r
   len)mr#    r'   C/tmp/pip-unpacked-wheel-mu1yl971/celery/worker/consumer/consumer.pyr"      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#dVd8d9Z$d:d; Z%d<d= Z&d>d? Z'dWd@dAZ(dBdC Z)dDdE Z*dFdG Z+dHdI Z,dJdK Z-dLdM Z.dNdO Z/e0fdPdQZ1dRdS Z2dS )Xr    Consumer blueprint.Nc                	   @   s2   e Zd ZdZdZddddddd	d
dg	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(   r+      s    zConsumer.Blueprint.shutdownN)__name__
__module____qualname____doc__nameZdefault_stepsr+   r'   r'   r'   r(   	Blueprint   s   r4   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$st%| jddr|	| _&| j&d kr| j jj'| _&nd| _&t(| d	s |rt)j*nt)j+| _,t- d
kr6d | j j_.g | _/g | _0| j1| j j0d | j2d| _3| j3j4| ft5|
ppi f| d S )N   r6   )ZmaxRZmaxTr   c                   S   s   d S Nr'   r'   r'   r'   r(   <lambda>       z#Consumer.__init__.<locals>.<lambda>Zis_greenFloopZgeventZconsumer)stepson_close)6app
controllerinit_callbackr   hostnameosgetpidpidpooltimer
Strategies
strategiesconnection_for_readconninfoconnection_errorsZchannel_errorsr   _restart_stateloggerisEnabledForloggingINFOZ
_does_info_limit_orderon_task_requestseton_task_messageconfZbroker_heartbeat_checkrateamqheartbeat_ratedisable_rate_limitsinitial_prefetch_countprefetch_multiplierr   task_bucketsreset_rate_limitshubgetattramqheartbeatZbroker_heartbeathasattrr   ZasynloopZsynloopr;   r	   Zbroker_connection_timeout_pending_operationsr<   r4   r=   	blueprintapplydict)r-   rR   r@   rA   rE   r>   rF   r?   r\   r^   Zworker_optionsrW   rX   rY   kwargsr'   r'   r(   __init__   sN    






zConsumer.__init__c                 O   s2   t |f||}| jr"| j|S | j| |S r8   )r   r\   	call_soonr`   append)r-   pargsrd   r'   r'   r(   rf      s
    zConsumer.call_soonc              
   C   sR   | j sN| jrNz| j   W q tk
rJ } ztd| W 5 d }~X Y qX qd S )NzPending callback raised: %r)r\   r`   pop	ExceptionrM   	exceptionr-   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 )NZ
rate_limitr6   )capacity)r   r]   r   )r-   typelimitr'   r'   r(   bucket_for_task   s    zConsumer.bucket_for_taskc                    s&    j  fdd jj D  d S )Nc                 3   s    | ]\}}|  |fV  qd S r8   )rs   ).0ntr-   r'   r(   	<genexpr>   s    z-Consumer.reset_rate_limits.<locals>.<genexpr>)rZ   updater>   tasksitemsrw   r'   rw   r(   r[      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)rE   num_processesrX   rY   _update_qos_eventually)r-   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_eventuallyZincrement_eventuallyabsrY   )r-   r~   r'   r'   r(   r}     s    zConsumer._update_qos_eventuallyc                 C   s   t | | | d S r8   )r   rR   )r-   requestr'   r'   r(   _limit_move_to_pool  s    zConsumer._limit_move_to_poolc                 C   s   z|  \}}W n tk
r(   Y qY nX ||rB| | q q |j||f | jd d  }| _||}| jj	|| j
|f|d qq d S )Nr6   
   )priority)rj   
IndexErrorZcan_consumer   contents
appendleftrQ   Zexpected_timerF   Z
call_after_schedule_bucket_request)r-   bucketr   tokenspriZholdr'   r'   r(   r     s"    



  z!Consumer._schedule_bucket_requestc                 C   s   | ||f | |S r8   )addr   r-   r   r   r   r'   r'   r(   _limit_task7  s    zConsumer._limit_taskc                 C   s"   | j   |||f | |S r8   )r   r   r   r   r   r'   r'   r(   _limit_post_eta;  s    
zConsumer._limit_post_etac              
   C   s  | j }|jtkrt  | jrfz| j  W n8 tk
rd } ztd|dd t	d W 5 d }~X Y nX |  jd7  _z|
|  W q | jk
r
 } zf| jjjs t|tr|jtjkr t  |jtkr| jr| | n
| | |   ||  W 5 d }~X Y qX qd S )NzFrequent restarts detected: %rr6   exc_info)ra   stateSTOP_CONDITIONSr   restart_countrL   stepr   critr   startrK   r>   rU   broker_connection_retry
isinstanceOSErrorerrnoZEMFILE
connection#on_connection_error_after_connected$on_connection_error_before_connectedr=   Zrestart)r-   ra   rn   r'   r'   r(   r   @  s0    


zConsumer.startc                 C   s   t t| j |d d S )NzTrying to reconnect...)errorCONNECTION_ERRORrJ   as_urirm   r'   r'   r(   r   ]  s    z-Consumer.on_connection_error_before_connectedc                 C   s~   t tdd z| j  W n tk
r.   Y nX | jjjrntt	D ](}|j
jrB|jsBt t| || j qBnt tt d S )NTr   )warnCONNECTION_RETRYr   Zcollectrk   r>   rU   Z3worker_cancel_long_running_tasks_on_connection_losstupler   taskZ	acks_lateZacknowledged3TERMINATING_TASK_ON_RESTART_AFTER_A_CONNECTION_LOSScancelrE   warningsCANCEL_TASKS_BY_DEFAULTr   )r-   rn   r   r'   r'   r(   r   a  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)ri   description)ra   r,   )r-   r\   r'   r'   r(   r   q  s      z!Consumer.register_with_event_loopc                 C   s   | j |  d S r8   )ra   r+   rw   r'   r'   r(   r+   w  s    zConsumer.shutdownc                 C   s   | j |  d S r8   )ra   stoprw   r'   r'   r(   r   z  s    zConsumer.stopc                 C   s   | j d  }| _ |r||  d S r8   )r@   )r-   callbackr'   r'   r(   on_ready}  s    zConsumer.on_readyc              	   C   s(   | | j | j| j| j| j| j| jj| jf	S r8   )	r   task_consumerra   r\   r   r^   r>   ZclockrV   rw   r'   r'   r(   	loop_args  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.
        r6   r   N)	r   MESSAGE_DECODE_ERRORcontent_typecontent_encodingr
   headersr"   r#   Zack)r-   messagern   r'   r'   r(   on_decode_error  s    
   
zConsumer.on_decode_errorc                 C   sj   | j r| j jr| j j  | jr*| j  | j D ]}|r4|  q4t  | jrf| jj	rf| j	  d S r8   )
r?   Z	semaphoreclearrF   rZ   valuesZclear_pendingr   rE   flush)r-   r   r'   r'   r(   r=     s    

zConsumer.on_closec                 C   s*   | j | jd}| jr&|j|j| j |S )zEstablish the broker connection used for consuming tasks.

        Retries establishing the connection if the
        :setting:`broker_connection_retry` setting is enabled
        	heartbeat)rI   r^   r\   	transportr   r   )r-   connr'   r'   r(   connect  s    zConsumer.connectc                 C   s   |  | jj|dS Nr   )ensure_connectedr>   rI   r-   r   r'   r'   r(   rI     s    zConsumer.connection_for_readc                 C   s   |  | jj|dS r   )r   r>   connection_for_writer   r'   r'   r(   r     s    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 )NZaltr   in r5   )whenretriesmax_retries)r]   CONNECTION_FAILOVERr$   r   intr>   rU   broker_connection_max_retriesr   r   r   )rn   intervalZ	next_stepr   r-   r'   r(   _error_handler  s    

z1Consumer.ensure_connected.<locals>._error_handler)r   )CONNECTION_RETRY_STEPr>   rU   r   r   Zensure_connectionr   r   )r-   r   r   r'   r   r(   r     s    
 zConsumer.ensure_connectedc                 C   s   | j r| j   d S r8   )event_dispatcherr   rw   r'   r'   r(   _flush_events  s    zConsumer._flush_eventsc                 C   s   | j r| j j| j d S r8   )r\   Z_readyr   r   rw   r'   r'   r(   on_send_event_buffered  s    zConsumer.on_send_event_bufferedc           	      K   s   | j }| jjj}||kr"|| }n:|d kr.|n|}|d kr>dn|}|j|f|||d|}||s|| |  td| d S )Ndirect)exchangeexchange_typerouting_keyzStarted consuming from %s)	r   r>   amqpqueuesZ
select_addZconsuming_fromZ	add_queueconsumeinfo)	r-   queuer   r   r   optionsZ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   r>   r   r   Zdeselectr   Zcancel_by_queue)r-   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   rR   r   r   )r-   r   r'   r'   r(   apply_eta_task  s    
zConsumer.apply_eta_taskc                 C   s0   t t||t|jt|jt|jt|jS r8   )MESSAGE_REPORTr$   r"   r
   r   r   delivery_infor   r-   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   rn   )	r   UNKNOWN_FORMATr   reject_log_errorrM   rK   r   task_rejectedsendr   r'   r'   r(   on_unknown_message  s    zConsumer.on_unknown_messagec           	      C   s   t t|t|||j|jdd z&|jd |jd  }}|jd}W n0 tk
rt   |j}|d |d  }}d }Y nX t|d ||j	d|j	dd d}|
t| j | jjj|t||d	 | jr| jjd
|d|dd tjj| ||||d d S )NTr   idr   root_idcorrelation_idreply_to)r3   Zchordr   r   r   Zerrbacks)r   ztask-failedzNotRegistered())uuidrl   )r   r   rn   r3   r   )r   UNKNOWN_TASK_ERRORr"   r   r   getKeyErrorpayloadr   Z
propertiesr   rM   rK   r>   backendZmark_as_failurer   r   r   r   Ztask_unknown)	r-   r#   r   rn   Zid_r3   r   r   r   r'   r'   r(   on_unknown_task  sR    
  

   
    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   rM   rK   r   r   r   )r-   r#   r   rn   r'   r'   r(   on_invalid_task(  s
    zConsumer.on_invalid_taskc                 C   sN   | j j}| j j D ]4\}}|| j | | j|< t|||| j| j d|_qd S )N)r>   )	r>   loaderrz   r{   Zstart_strategyrH   r   rA   Z	__trace__)r-   r   r3   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 n tk
r0   d |  Y S  tk
r   z|  }W n6 tk
r } z| | W Y  Y S d }~X Y nX z|d | }}W n& ttfk
r   ||  Y  Y S X Y nX z| }W n4 tk
r } zd | | W Y S d }~X Y nX z(|| | | jf | jf W nj tt	fk
rd } z|| | W Y S d }~X Y n4 t
k
r } z| | W Y S d }~X Y nX d S )Nr   )r   	TypeErrorr   decoderk   r   Zack_log_errorr   r   r   r   )r   r   type_rn   Zstrategyrf   	callbacksr   r   r   r   r-   rH   r'   r(   on_task_received=  s<    &"  z6Consumer.create_task_handler.<locals>.on_task_received)rH   r   r   r   rT   rf   )r-   r   r  r'   r   r(   create_task_handler5  s    "zConsumer.create_task_handlerc                 C   s   dj | | j dS )z``repr(self)``.z%<Consumer: {self.hostname} ({state})>)r-   r   )r$   ra   Zhuman_staterw   r'   r'   r(   __repr__a  s     zConsumer.__repr__)r   )N)N)NNN)3r/   r0   r1   r2   rc   rG   r@   rE   rF   r   r   r4   r   re   rf   ro   rs   r[   r   r}   r   r   r   r   r   r   r   r   r+   r   r   r   r   r=   r   rI   r   r   r   r   r   r   r   r   r   r   r   r   r   r  r  r'   r'   r'   r(   r       st          
;


  
!,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 r8   )	patch_allr;   r   r-   cr'   r'   r(   r   r  s    
zEvloop.startc                 C   s   t  |j_d S r8   )r   r   Z_mutexr  r'   r'   r(   r  v  s    zEvloop.patch_allN)r/   r0   r1   r2   labellastr   r  r'   r'   r'   r(   r!   h  s
   r!   )Tr2   r   rO   rB   r   collectionsr   timer   Zbilliard.commonr   Zbilliard.exceptionsr   Zkombu.asynchronous.semaphorer   Zkombu.exceptionsr   r   Zkombu.utils.compatr	   Zkombu.utils.encodingr
   Zkombu.utils.limitsr   Zviner   r   Zceleryr   r   Zcelery.app.tracer   Zcelery.exceptionsr   r   r   Zcelery.utils.functionalr   Zcelery.utils.logr   Zcelery.utils.nodenamesr   Zcelery.utils.objectsr   Zcelery.utils.textr   Zcelery.utils.timer   r   Zcelery.workerr   Zcelery.worker.stater   r   r   r   __all__ZCLOSEZ	TERMINATEr   r/   rM   debugr   warningr   criticalr   r   r   r   r   r   r   r   r   r   r   r   r   r"   r    ZStartStopStepr!   r'   r'   r'   r(   <module>   sf    	   `